如何在.NET中实现不丢消息的RabbitMQ异步Consume方法?
问题根源
你的代码出现消息丢失的核心原因是用单个变量存储消息,当RabbitMQ高并发推送消息时,后续触发的Received事件会直接覆盖_message变量中未被消费的前置消息;同时_lock和信号量的组合逻辑存在竞态条件,导致部分消息被覆盖后无法被消费。
可行解决方案
方案一:ConcurrentQueue + SemaphoreSlim(线程安全队列+信号量)
用线程安全的队列缓存所有收到的消息,避免单变量覆盖问题,同时通过信号量协调消费节奏:
private readonly ConcurrentQueue<string> _messageQueue = new ConcurrentQueue<string>(); private readonly SemaphoreSlim _signal = new SemaphoreSlim(0); private readonly IModel _channel; private readonly RabbitConfig _rabbitConfig; public void Configure() { var consumer = new EventingBasicConsumer(_channel); consumer.Received += (sender, ea) => { // 将消息存入线程安全队列 var message = Encoding.UTF8.GetString(ea.Body.ToArray()); _messageQueue.Enqueue(message); // 释放信号量通知消费者 _signal.Release(); }; _channel.BasicConsume(queue: _rabbitConfig.Queue, autoAck: true, consumer: consumer); } public async Task<string> Consume(CancellationToken cancellationToken) { // 等待消息可用 await _signal.WaitAsync(cancellationToken); // 从队列取出消息(理论上不会失败,因为信号量和入队操作一一对应) if (_messageQueue.TryDequeue(out var message)) { return message; } throw new InvalidOperationException("信号量触发但队列无消息"); }
优化点(可选,保证消息不丢失)
如果需要避免消费失败导致的消息丢失,建议关闭autoAck,手动确认消息:
// 定义消息模型存储内容和DeliveryTag public class RabbitMessage { public string Content { get; set; } public ulong DeliveryTag { get; set; } } // 修改队列类型 private readonly ConcurrentQueue<RabbitMessage> _messageQueue = new ConcurrentQueue<RabbitMessage>(); // Configure方法中修改Received事件 consumer.Received += (sender, ea) => { var rabbitMsg = new RabbitMessage { Content = Encoding.UTF8.GetString(ea.Body.ToArray()), DeliveryTag = ea.DeliveryTag }; _messageQueue.Enqueue(rabbitMsg); _signal.Release(); }; // 关闭autoAck _channel.BasicConsume(queue: _rabbitConfig.Queue, autoAck: false, consumer: consumer); // Consume方法修改为: public async Task<string> Consume(CancellationToken cancellationToken) { await _signal.WaitAsync(cancellationToken); if (_messageQueue.TryDequeue(out var rabbitMsg)) { try { // 这里可以添加消息处理逻辑 _channel.BasicAck(rabbitMsg.DeliveryTag, false); return rabbitMsg.Content; } catch (Exception ex) { // 处理失败,重新入队或发送到死信队列 _channel.BasicNack(rabbitMsg.DeliveryTag, false, true); throw; } } throw new InvalidOperationException("信号量触发但队列无消息"); }
方案二:使用System.Threading.Channels(推荐,异步生产者消费者专用)
.NET Core 2.1+提供的System.Threading.Channels是专门为异步生产者消费者场景设计的组件,内部已封装线程安全和异步等待逻辑,代码更简洁高效:
private readonly Channel<string> _messageChannel; private readonly IModel _channel; private readonly RabbitConfig _rabbitConfig; // 构造函数初始化通道 public YourConsumerClass() { // 创建无界通道,也可根据需求设置有界限制(防止内存溢出) _messageChannel = Channel.CreateUnbounded<string>(); } public void Configure() { var consumer = new EventingBasicConsumer(_channel); consumer.Received += async (sender, ea) => { var message = Encoding.UTF8.GetString(ea.Body.ToArray()); // 写入通道,异步等待队列可用 await _messageChannel.Writer.WriteAsync(message); }; _channel.BasicConsume(queue: _rabbitConfig.Queue, autoAck: true, consumer: consumer); } public async Task<string> Consume(CancellationToken cancellationToken) { // 从通道读取消息,自动等待新消息 return await _messageChannel.Reader.ReadAsync(cancellationToken); }
同样,如果需要手动确认消息,只需将通道类型改为Channel<RabbitMessage>,在消费成功后调用BasicAck即可。
内容的提问来源于stack exchange,提问作者TGuerin
相关产品推荐
相关产品推荐

