.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
相关产品推荐
相关产品推荐

