RabbitMQ消息Unacked问题:IHostedService消费未正常确认消息
问题分析与解决
消息进入RabbitMQ队列后处于Unacked状态,核心原因是消费逻辑中同步调用异步方法引发死锁,阻塞了线程,导致BasicAck无法执行,RabbitMQ始终认为消息未处理完成。
你的代码里,HandleMessage方法用.Result同步等待_mediator.Send(command)的结果,这种操作在异步场景下极易卡死线程,使得BasicAck步骤永远无法触发。
修正方案
将消费逻辑改为异步处理,避免同步阻塞:
- 将消息处理方法改为异步,用
await替代同步等待 - 调整事件处理逻辑以支持异步执行
- 增加异常捕获,确保处理失败时能正确反馈给RabbitMQ(可选,防止消息无限阻塞)
修改后的完整代码
using EMS.Models; using RabbitMQ.Client; using RabbitMQ.Client.Events; using System.Text; using Newtonsoft.Json; using EMS.Models.Command; using MediatR; using Microsoft.Extensions.Hosting; namespace EMS.RabbitMQ { public class RabbitMqEmployeeConsumer : IHostedService, IDisposable { private readonly IMediator _mediator; private readonly IConnection _connection; private readonly IModel _channel; private readonly string _queueName = "employeeRequests"; private readonly CancellationTokenSource _cancellationTokenSource; public RabbitMqEmployeeConsumer(IMediator mediator) { _mediator = mediator; _cancellationTokenSource = new CancellationTokenSource(); var factory = new ConnectionFactory() { HostName = "localhost" }; _connection = factory.CreateConnection(); _channel = _connection.CreateModel(); _channel.QueueDeclare(queue: _queueName, durable: false, exclusive: false, autoDelete: false, arguments: null); var consumer = new EventingBasicConsumer(_channel); // 改为异步事件处理 consumer.Received += async (model, ea) => { var body = ea.Body.ToArray(); var message = Encoding.UTF8.GetString(body); await HandleMessageAsync(ea, message); }; _channel.BasicConsume(queue: _queueName, autoAck: false, consumer: consumer); } // 改为异步方法 private async Task HandleMessageAsync(BasicDeliverEventArgs ea, string message) { try { var command = JsonConvert.DeserializeObject<AddEmployeeCommand>(message); await _mediator.Send(command); // 用await替代.Result _channel.BasicAck(deliveryTag: ea.DeliveryTag, multiple: false); Console.WriteLine("消息处理完成,数据已插入数据库。"); } catch (Exception ex) { // 处理失败时Nack,requeue设为true会重新入队,false则丢弃 _channel.BasicNack(deliveryTag: ea.DeliveryTag, multiple: false, requeue: true); Console.WriteLine($"消息处理失败:{ex.Message}"); } } public Task StartAsync(CancellationToken cancellationToken) { return Task.CompletedTask; } public Task StopAsync(CancellationToken cancellationToken) { _cancellationTokenSource.Cancel(); _channel.Dispose(); _connection.Dispose(); return Task.CompletedTask; } public void Dispose() { _cancellationTokenSource.Cancel(); _channel.Dispose(); _connection.Dispose(); _cancellationTokenSource.Dispose(); } } }
额外注意事项
- 确保
AddEmployeeCommand的MediatR处理器是异步实现(返回Task),否则await无法发挥作用 - 生产环境建议给RabbitMQ连接添加重连逻辑,避免连接断开后消费终止
- 若消息处理耗时较长,可通过
BasicQos调整消费者预取数,防止同时获取过多未处理消息
内容的提问来源于stack exchange,提问作者Yubin Budhathoki
相关产品推荐
相关产品推荐

