使用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
相关产品推荐
相关产品推荐

