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

使用lapin消费RabbitMQ消息时Delivery无ack方法问题求助

问题解决:lapin库消费RabbitMQ时ack方法不存在报错

报错根因

该报错由lapin版本API不兼容导致:

  • 官方仓库master分支的basic_consume返回的迭代器产出Delivery类型,该类型自带ack方法
  • crates.io发布的2.x稳定版迭代器产出的是(Channel, Delivery)元组,元组本身没有ack方法,需要通过Channel调用ack接口

解决办法

方案1:适配2.x稳定版(推荐,稳定性更高)

修改消费循环代码即可,同时补充消息载荷获取逻辑:

while let Some(delivery) = consumer.next().await {
    info!(message=?delivery, "received message");
    if let Ok((channel, delivery)) = delivery {
        // delivery.data即为消息载荷的字节数组,此处转为字符串打印
        let payload = String::from_utf8_lossy(&delivery.data);
        info!("消息载荷内容:{}", payload);
        // 通过channel调用ack方法,传入当前消息的delivery_tag
        channel
            .basic_ack(delivery.delivery_tag, BasicAckOptions::default())
            .await
            .expect("basic_ack");
    }
}

对应Cargo.toml依赖写法参考:

lapin = "2.3.1"
futures-lite = "1.13.0"
tracing = "0.1.37"
tracing-subscriber = "0.3.17"
async-global-executor = "2.3.1"

方案2:使用master分支的lapin版本

将lapin依赖指向官方仓库master分支,即可直接运行官方示例无需修改代码。

完整可运行的稳定版消费示例

use futures_lite::StreamExt;
use lapin::{options::*, types::FieldTable, Connection, ConnectionProperties};
use tracing::info;

fn main() {
    if std::env::var("RUST_LOG").is_err() {
        std::env::set_var("RUST_LOG", "info");
    }

    tracing_subscriber::fmt::init();

    let addr = std::env::var("AMQP_ADDR").unwrap_or_else(|_| "amqp://127.0.0.1:5672/%2f".into());

    async_global_executor::block_on(async {
        let conn = Connection::connect(&addr, ConnectionProperties::default())
            .await
            .expect("connection error");

        info!("CONNECTED");

        // 创建消费通道
        let channel = conn.create_channel().await.expect("create_channel");
        info!(state=?conn.status().state());

        let queue = channel
            .queue_declare(
                "hello",
                QueueDeclareOptions::default(),
                FieldTable::default(),
            )
            .await
            .expect("queue_declare");
        info!(state=?conn.status().state());
        info!(?queue, "声明队列完成");

        info!("开始消费");
        let mut consumer = channel
            .basic_consume(
                "hello",
                "my_consumer",
                BasicConsumeOptions::default(),
                FieldTable::default(),
            )
            .await
            .expect("basic_consume");
        info!(state=?conn.status().state());

        while let Some(delivery) = consumer.next().await {
            info!(message=?delivery, "收到消息");
            if let Ok((channel, delivery)) = delivery {
                // 读取并打印消息载荷
                let payload = String::from_utf8_lossy(&delivery.data);
                info!("消息载荷:{}", payload);
                // 手动ack确认消息
                channel
                    .basic_ack(delivery.delivery_tag, BasicAckOptions::default())
                    .await
                    .expect("basic_ack");
            }
        }
    })
}

内容的提问来源于stack exchange,提问作者Allan K

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 04:06:03