You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用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>
  • 不要直接将路由键作为Target名称,否则RabbitMQ会尝试创建同名队列,这就是你之前触发SenderAttachError的原因

内容的提问来源于stack exchange,提问作者V.Lorz

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.31 13:15:24