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

.NET C#环境下RabbitMQ Topic Exchange消费消息失败求助

问题排查与解决方案

你的代码核心问题在于执行顺序和资源释放时机,导致消息无法被消费者接收,具体原因和修复方案如下:

1. 执行顺序错误

你在Main方法里先调用CreatePublisher发送消息,再调用CreateConsumer启动消费者。此时消息发送时,消费者还未完成队列绑定,RabbitMQ的Topic Exchange会直接丢弃找不到匹配队列的消息,所以消费者启动后根本看不到这条消息。

2. 发布者资源提前释放

CreatePublisher里的connection和channel使用了using语句,方法执行完毕后会立即释放这些资源。RabbitMQ的消息发送是异步过程,可能资源释放时消息还没完全发送到服务器,导致消息丢失。

修复后的代码示例

调整执行顺序,让消费者先启动

修改Main方法,通过线程启动消费者避免阻塞主线程,确保消费者完成队列绑定后再发送消息:

public static void Main(string[] args)
{
    // 用线程启动消费者,避免阻塞主线程
    var consumerThread = new Thread(CreateConsumer);
    consumerThread.Start();
    
    // 给消费者预留队列绑定的时间
    Thread.Sleep(1000);
    
    CreatePublisher();
    
    // 等待消费者线程结束
    consumerThread.Join();
}

public static void CreatePublisher() 
{
    var factory = new ConnectionFactory { HostName = "localhost", 
        UserName = "myUser", Password = "myPass" };
    using var connection = factory.CreateConnection();
    using var channel = connection.CreateModel();

    channel.ExchangeDeclare(exchange: "logs", type: 
        ExchangeType.Topic);

    var message = "testeee";
    var body = Encoding.UTF8.GetBytes(message);
    channel.BasicPublish(exchange: "logs",
        routingKey: "logs",
        basicProperties: null,
        body: body);
    Console.WriteLine($" [x] Sent {message}");

    // 延迟释放资源,确保消息完成发送
    Thread.Sleep(500);
}

public static void CreateConsumer()
{
    var factory = new ConnectionFactory { HostName = "localhost", 
        UserName = "myUser", Password = "myPass" };
    using var connection = factory.CreateConnection();
    using var channel = connection.CreateModel();

    channel.ExchangeDeclare(exchange: "logs", type: 
        ExchangeType.Topic);

    var queueName = channel.QueueDeclare().QueueName;
    channel.QueueBind(queue: queueName,
                        exchange: "logs",
                        routingKey: "logs");

    Console.WriteLine(" [*] Waiting for logs.");

    var consumer = new EventingBasicConsumer(channel);
    consumer.Received += (model, ea) =>
    {
        byte[] body = ea.Body.ToArray();
        var message = Encoding.UTF8.GetString(body);
        Console.WriteLine($" [x] {message}");
    };
    channel.BasicConsume(queue: queueName,
                         autoAck: true,
                         consumer: consumer);

    Console.ReadLine();
}

额外注意事项

  • 生产环境中不要依赖Thread.Sleep做同步,建议用信号量等可靠的同步机制。
  • 若需要消息持久化,可以在声明队列和发送消息时配置持久化参数。
  • 确认RabbitMQ服务正常运行,且配置的用户拥有访问logs交换机的权限。

内容的提问来源于stack exchange,提问作者Gabriel

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 01:17:28