如何用fe2o3-amqp与RabbitMQ通过Topic Exchange交换AMQP 1.0消息
使用fe2o3-amqp与RabbitMQ Topic Exchange收发消息(AMQP 1.0)
核心差异说明
你用的pika基于AMQP 0-9-1协议,直接指定Exchange和Routing Key即可发送消息;但AMQP 1.0采用Link+Target/Source的模型,RabbitMQ对AMQP 1.0提供了特定的地址格式来兼容原有交换器模型,这是你之前配置出错的核心原因。
配置Sender(发送到Topic Exchange)
要将消息发送到MyExchange并使用路由键my.topic,需将Sender的Target地址设置为RabbitMQ专属格式:/exchange/<交换器名称>/<路由键>,以此明确指定消息的投递目标。
示例代码:
use fe2o3_amqp::{Connection, Session, Sender, Target}; use fe2o3_amqp::types::messaging::Message; use fe2o3_amqp::types::primitives::Value; #[tokio::main] async fn main() -> Result<(), Box<dyn std::error::Error>> { // 1. 创建连接,替换为你的RabbitMQ地址、用户名、密码 let mut connection = Connection::builder() .container_id("rust-sender-container") .scheme("amqp") .host("localhost") .port(5672) .username("guest") .password("guest") .open() .await?; // 2. 启动Session let mut session = Session::begin(&mut connection).await?; // 3. 构建Sender,关键是指定Exchange+路由键的Target地址 let sender = Sender::builder() .name("rust-sender-link-1") .target(Target::address("/exchange/MyExchange/my.topic")) .attach(&mut session) .await?; // 4. 发送消息 let message = Message::builder() .body(Value::String("Hello topic!".to_string())) .build(); sender.send(message).await?; // 优雅清理资源 sender.close().await?; session.end().await?; connection.close().await?; Ok(()) }
配置Receiver(订阅绑定的队列)
要从MyQueue接收消息,需将Receiver的Source地址设置为/queue/<队列名称>,直接绑定到已创建的队列。
示例代码:
use fe2o3_amqp::{Connection, Session, Receiver, Source}; use fe2o3_amqp::types::messaging::Delivery; #[tokio::main] async fn main() -> Result<(), Box<dyn std::error::Error>> { // 1. 创建连接 let mut connection = Connection::builder() .container_id("rust-receiver-container") .scheme("amqp") .host("localhost") .port(5672) .username("guest") .password("guest") .open() .await?; // 2. 启动Session let mut session = Session::begin(&mut connection).await?; // 3. 构建Receiver,指定要订阅的队列地址 let mut receiver = Receiver::builder() .name("rust-receiver-link-1") .source(Source::address("/queue/MyQueue")) .attach(&mut session) .await?; // 4. 循环接收并确认消息 while let Some(delivery) = receiver.recv().await? { let body = delivery.body::<String>()?; println!("Received message: {}", body); receiver.accept(&delivery).await?; // 向RabbitMQ确认消息已处理 } // 优雅清理资源 receiver.close().await?; session.end().await?; connection.close().await?; Ok(()) }
关键注意事项
- 确保RabbitMQ已启用AMQP 1.0插件:执行命令
rabbitmq-plugins enable rabbitmq_amqp1_0 - 严格遵循RabbitMQ的AMQP 1.0地址格式:
- 发送到Exchange:
/exchange/<exchange-name>/<routing-key> - 订阅Queue:
/queue/<queue-name>
- 发送到Exchange:
- 不要直接将路由键作为Target名称,否则RabbitMQ会尝试创建同名队列,这就是你之前触发
SenderAttachError的原因
内容的提问来源于stack exchange,提问作者V.Lorz
相关产品推荐
相关产品推荐

