RabbitMQ如何在.NET Framework 4.8中创建持续监听队列的消费者
.NET Framework 4.8 环境RabbitMQ持久监听消费者实现
现有代码问题说明
你给出的核心代码逻辑本身是通顺的,但存在一个致命问题:控制台程序的Main方法在执行完BasicConsume注册逻辑后会直接运行结束,导致进程退出,无法持续监听队列消息。同时还有部分可靠性优化点可以补充。
修正后的完整实现代码
using System; using System.Text; using System.Threading; using System.Threading.Tasks; using RabbitMQ.Client; using RabbitMQ.Client.Events; class Program { // 用于阻塞进程保持运行的信号量 private static readonly ManualResetEvent _exitEvent = new ManualResetEvent(false); static void Main(string[] args) { // 初始化RabbitMQ连接(替换为你实际的连接参数) RabbitHelper.Init("host", 123, "user", "password", "Provider", "virtual"); const string queueName = "MYQUEUE"; var consumer = new AsyncEventingBasicConsumer(RabbitHelper.channel); consumer.Received += Consumer_Received; // 注册消费者,第二个参数为true是自动ack,改为false可以开启手动ack保证消息可靠性 RabbitHelper.channel.BasicConsume(queueName, true, consumer); Console.WriteLine("消费者已启动,持续监听队列中..."); // 阻塞进程,直到收到退出信号 _exitEvent.WaitOne(); } private static async Task Consumer_Received(object sender, BasicDeliverEventArgs eevent) { try { var message = Encoding.UTF8.GetString(eevent.Body.ToArray()); Console.WriteLine($"开始处理消息:{message}"); // 替换为你实际的业务处理逻辑 await Task.Delay(250); Console.WriteLine($"消息处理完成:{message}"); // 如果上面BasicConsume关了自动ack,这里需要加手动确认逻辑 // RabbitHelper.channel.BasicAck(eevent.DeliveryTag, false); } catch (Exception ex) { Console.WriteLine($"消息处理异常:{ex.Message}"); // 异常场景下可以选择重新入队或者丢弃 // RabbitHelper.channel.BasicNack(eevent.DeliveryTag, false, true); } } } // RabbitHelper参考实现,你也可以保留你自己的封装逻辑 public static class RabbitHelper { public static IModel channel; private static IConnection _connection; public static void Init(string host, int port, string userName, string password, string provider, string virtualHost) { var factory = new ConnectionFactory { HostName = host, Port = port, UserName = userName, Password = password, VirtualHost = virtualHost, // 开启自动重连,避免网络波动导致连接断开后无法恢复 AutomaticRecoveryEnabled = true, NetworkRecoveryInterval = TimeSpan.FromSeconds(5) }; _connection = factory.CreateConnection(); channel = _connection.CreateModel(); } }
注意事项
- 若需要保证消息不丢失,可将
BasicConsume的第二个参数改为false,处理完业务逻辑后手动调用BasicAck确认消息 - 代码中已经开启了RabbitMQ客户端的自动重连能力,网络波动后会自动恢复监听,不需要额外写重连逻辑
- 若需要部署为后台服务,建议把控制台程序改成Windows服务,避免需要保持控制台窗口开启
内容的提问来源于stack exchange,提问作者Max Bertoli
相关产品推荐
相关产品推荐

