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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 21:45:07