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 工具和宏。