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

RabbitMQ实现生产者多消费者单向通信及消息溯源方案咨询

问题解答

这个场景在RabbitMQ中完全可行,核心是拆分两个方向的消息流,用不同的交换器/队列分别处理广播和单向通信需求。

原代码问题分析

你当前使用单一的Fanout交换器处理所有消息发送,而Fanout交换器的特性是把收到的消息路由到所有绑定的队列,所以消费者发送的消息会被所有绑定到该交换器的消费者队列接收,这就导致了不符合需求的情况。

解决方案设计

  • 生产者→所有消费者(广播):保留Fanout交换器,每个消费者创建独立的临时队列并绑定到该交换器,确保生产者的消息能被所有消费者接收。
  • 消费者→生产者(单向):新增一个Direct交换器,生产者声明一个专属队列并绑定到该交换器;消费者发送消息时,指定这个Direct交换器和对应的路由键,同时在消息内容中携带自身标识,让生产者能识别消息来源。

修改后示例代码

生产者代码

Console.WriteLine("Producer");
var factory = new ConnectionFactory() { HostName = "localhost" };
using (var connection = factory.CreateConnection())
using (var channel = connection.CreateModel())
{
    // 1. 声明生产者广播用的Fanout交换器
    const string broadcastExchange = "producer-broadcast-exchange";
    channel.ExchangeDeclare(broadcastExchange, ExchangeType.Fanout);

    // 2. 声明消费者给生产者发消息用的Direct交换器和专属队列
    const string directExchange = "consumer-to-producer-exchange";
    const string producerQueue = "producer专属队列";
    channel.ExchangeDeclare(directExchange, ExchangeType.Direct);
    channel.QueueDeclare(producerQueue, durable: false, exclusive: false, autoDelete: false, arguments: null);
    // 绑定专属队列到Direct交换器,使用固定路由键
    channel.QueueBind(producerQueue, directExchange, "producer-route-key");

    // 监听消费者发来的消息
    var consumer = new EventingBasicConsumer(channel);
    consumer.Received += (model, ea) =>
    {
        var body = ea.Body.ToArray();
        var message = System.Text.Encoding.UTF8.GetString(body);
        Console.WriteLine($" [x] 收到来自消费者的消息: {message}");
    };
    channel.BasicConsume(producerQueue, autoAck: true, consumer);

    // 发送广播消息
    while (true)
    {
        var message = Console.ReadLine();
        if (message == "exit") break;

        var body = System.Text.Encoding.UTF8.GetBytes(message);
        channel.BasicPublish(broadcastExchange, "", null, body);
        Console.WriteLine(" [x] 发送广播消息: {0}", message);
    }
}

消费者1代码

Console.WriteLine("Consumer 1");
var factory = new ConnectionFactory() { HostName = "localhost" };
using (var connection = factory.CreateConnection())
using (var channel = connection.CreateModel())
{
    // 绑定到生产者的广播交换器
    const string broadcastExchange = "producer-broadcast-exchange";
    channel.ExchangeDeclare(broadcastExchange, ExchangeType.Fanout);
    var consumerQueue = channel.QueueDeclare().QueueName;
    channel.QueueBind(consumerQueue, broadcastExchange, "");

    // 监听生产者的广播消息
    var broadcastConsumer = new EventingBasicConsumer(channel);
    broadcastConsumer.Received += (model, ea) =>
    {
        var body = ea.Body.ToArray();
        var message = System.Text.Encoding.UTF8.GetString(body);
        Console.WriteLine($" [x] 收到广播消息: {message}");
    };
    channel.BasicConsume(consumerQueue, autoAck: true, broadcastConsumer);

    // 发送消息给生产者
    const string directExchange = "consumer-to-producer-exchange";
    while (true)
    {
        var input = Console.ReadLine();
        if (input == "exit") break;

        // 消息中携带消费者标识
        var message = $"Consumer 1: {input}";
        var body = System.Text.Encoding.UTF8.GetBytes(message);
        channel.BasicPublish(directExchange, "producer-route-key", null, body);
        Console.WriteLine(" [x] 发送消息给生产者: {0}", message);
    }
}

消费者2代码

Console.WriteLine("Consumer 2");
var factory = new ConnectionFactory() { HostName = "localhost" };
using (var connection = factory.CreateConnection())
using (var channel = connection.CreateModel())
{
    // 绑定到生产者的广播交换器
    const string broadcastExchange = "producer-broadcast-exchange";
    channel.ExchangeDeclare(broadcastExchange, ExchangeType.Fanout);
    var consumerQueue = channel.QueueDeclare().QueueName;
    channel.QueueBind(consumerQueue, broadcastExchange, "");

    // 监听生产者的广播消息
    var broadcastConsumer = new EventingBasicConsumer(channel);
    broadcastConsumer.Received += (model, ea) =>
    {
        var body = ea.Body.ToArray();
        var message = System.Text.Encoding.UTF8.GetString(body);
        Console.WriteLine($" [x] 收到广播消息: {message}");
    };
    channel.BasicConsume(consumerQueue, autoAck: true, broadcastConsumer);

    // 发送消息给生产者
    const string directExchange = "consumer-to-producer-exchange";
    while (true)
    {
        var input = Console.ReadLine();
        if (input == "exit") break;

        // 消息中携带消费者标识
        var message = $"Consumer 2: {input}";
        var body = System.Text.Encoding.UTF8.GetBytes(message);
        channel.BasicPublish(directExchange, "producer-route-key", null, body);
        Console.WriteLine(" [x] 发送消息给生产者: {0}", message);
    }
}

代码说明

  • 生产者的专属队列仅自身绑定到consumer-to-producer-exchange交换器,确保只有生产者能接收消费者发送的消息。
  • 消费者发送消息时在内容中加入自身标识,生产者可直接识别消息来源。
  • 广播部分继续使用Fanout交换器,保证生产者的消息能被所有消费者接收。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 07:45:35