Chuyển đến nội dung chính

第 15 課:訊息佇列和後台作業

RabbitMQ 與 lapin,Kafka 與 rdkafka。 NATS 訊息傳遞。後台作業處理。事件驅動的架構模式。 Redis 發佈/訂閱。

💻 程式設計 — 第 15 課 第 15 課:訊息佇列和後台作業

Rust:從基礎到高級

第 4 部分:進階後端

亞洲開發網

1.RabbitMQ 與 Lapin

use lapin::{Connection, ConnectionProperties, options::*, types::FieldTable, BasicProperties};

async fn publish(channel: &Channel, message: &str) -> Result<(), lapin::Error> {
    channel.basic_publish(
        "",           // exchange
        "task_queue", // routing key
        BasicPublishOptions::default(),
        message.as_bytes(),
        BasicProperties::default().with_delivery_mode(2), // persistent
    ).await?.await?;
    Ok(())
}

async fn consume(channel: &Channel) -> Result<(), lapin::Error> {
    channel.basic_qos(1, BasicQosOptions::default()).await?;

    let mut consumer = channel.basic_consume(
        "task_queue",
        "worker",
        BasicConsumeOptions::default(),
        FieldTable::default(),
    ).await?;

    while let Some(delivery) = consumer.next().await {
        let delivery = delivery?;
        let message = String::from_utf8_lossy(&delivery.data);
        println!("Received: {}", message);

        // Process...

        delivery.ack(BasicAckOptions::default()).await?;
    }
    Ok(())
}

2.Kafka 與 rdkafka

use rdkafka::producer::{FutureProducer, FutureRecord};
use rdkafka::consumer::{StreamConsumer, Consumer};
use rdkafka::ClientConfig;

// Producer
async fn produce(producer: &FutureProducer, topic: &str, key: &str, payload: &str) {
    producer.send(
        FutureRecord::to(topic).key(key).payload(payload),
        Duration::from_secs(5),
    ).await.unwrap();
}

// Consumer
async fn consume_kafka(brokers: &str, group_id: &str, topic: &str) {
    let consumer: StreamConsumer = ClientConfig::new()
        .set("bootstrap.servers", brokers)
        .set("group.id", group_id)
        .set("auto.offset.reset", "earliest")
        .create()
        .unwrap();

    consumer.subscribe(&[topic]).unwrap();

    let mut stream = consumer.stream();
    while let Some(result) = stream.next().await {
        match result {
            Ok(msg) => {
                let payload = msg.payload_view::<str>().unwrap().unwrap();
                println!("Key: {:?}, Payload: {}", msg.key(), payload);
            }
            Err(e) => eprintln!("Error: {}", e),
        }
    }
}

3.Redis 發布/訂閱

use redis::AsyncCommands;

async fn publisher(client: &redis::Client) -> redis::RedisResult<()> {
    let mut conn = client.get_multiplexed_async_connection().await?;
    conn.publish("events", "order:created:123").await?;
    Ok(())
}

async fn subscriber(client: &redis::Client) -> redis::RedisResult<()> {
    let mut pubsub = client.get_async_pubsub().await?;
    pubsub.subscribe("events").await?;
    let mut stream = pubsub.on_message();
    while let Some(msg) = stream.next().await {
        let payload: String = msg.get_payload()?;
        println!("Event: {}", payload);
    }
    Ok(())
}

4. 後台作業模式

use tokio::sync::mpsc;

enum Job {
    SendEmail { to: String, subject: String, body: String },
    ProcessImage { path: String },
    GenerateReport { user_id: String },
}

async fn job_worker(mut rx: mpsc::Receiver<Job>) {
    while let Some(job) = rx.recv().await {
        match job {
            Job::SendEmail { to, subject, .. } => {
                println!("Sending email to {} - {}", to, subject);
            }
            Job::ProcessImage { path } => {
                tokio::task::spawn_blocking(move || {
                    // CPU-intensive image processing
                }).await.unwrap();
            }
            Job::GenerateReport { user_id } => {
                println!("Generating report for {}", user_id);
            }
        }
    }
}

下一篇: 快取、CLI 工具和宏。