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

RabbitMQ单消费者高效处理耗时事件的生产环境可行性问询

RabbitMQ单消费者处理耗时事件的实验与生产环境疑问

需求:使用单个消费者尽可能快速处理耗时事件。
测试场景:5个事件,每个事件处理耗时2秒。

第一次尝试:官方Hello World同步方案

采用RabbitMQ官方Hello World示例,代码如下:

static void Main(string[] args)
{
    // Sender

    var factory = new ConnectionFactory { HostName = "localhost"};
    using var connection = factory.CreateConnection();
    using var channel = connection.CreateModel();

    channel.QueueDeclare(queue: "hello", durable: false, exclusive: false, autoDelete: false, arguments: null);

    const string message = "Hello World!";

    for (int i = 1; i <= 5; i++)
    {
        var body = Encoding.UTF8.GetBytes($"{message} {i}");
        channel.BasicPublish(exchange: string.Empty,
                                routingKey: "hello",
                                basicProperties: null,
                                body: body);

        Console.WriteLine($"Sent {message} {i}");
    }

    // Consumer

    var consumer = new EventingBasicConsumer(channel);
    consumer.Received += (model, ea) =>
    {
        var body = ea.Body.ToArray();
        var message = Encoding.UTF8.GetString(body);
        Thread.Sleep(2000);
        Console.WriteLine($"Received {message}");
    };

    channel.BasicConsume(queue: "hello",
                            autoAck: true,
                            consumer: consumer);


    Console.WriteLine("Press [enter] to exit.");
    Console.ReadLine();
}

结论:处理速度过慢,因事件为同步处理。

第二次尝试:AsyncEventingBasicConsumer异步方案

代码修改点:

  • 在ConnectionFactory中设置DispatchConsumersAsync属性
  • 使用AsyncEventingBasicConsumer
  • 为consumer.Received绑定异步方法

代码如下:

static void Main(string[] args)
{
    // Sender

    var factory = new ConnectionFactory { HostName = "localhost", DispatchConsumersAsync = true };
    using var connection = factory.CreateConnection();
    using var channel = connection.CreateModel();

    channel.QueueDeclare(queue: "hello", durable: false, exclusive: false, autoDelete: false, arguments: null);

    const string message = "Hello World!";

    for (int i = 1; i <= 5; i++)
    {
        var body = Encoding.UTF8.GetBytes($"{message} {i}");
        channel.BasicPublish(exchange: string.Empty,
                                routingKey: "hello",
                                basicProperties: null,
                                body: body);

        Console.WriteLine($"Sent {message} {i}");
    }

    // Consumer

    var consumer = new AsyncEventingBasicConsumer(channel);
    consumer.Received += async (model, ea) =>
    {
        var body = ea.Body.ToArray();
        var message = Encoding.UTF8.GetString(body);
        await Task.Delay(2000);
        Console.WriteLine($"Received {message}");
    };

    channel.BasicConsume(queue: "hello",
                            autoAck: true,
                            consumer: consumer);


    Console.WriteLine("Press [enter] to exit.");
    Console.ReadLine();
}

结论:处理速度仍过慢,事件依旧为同步处理。

第三次尝试:EventingBasicConsumer绑定异步处理器

代码修改点:

  • 移除ConnectionFactory中的DispatchConsumersAsync属性
  • 重新使用EventingBasicConsumer

代码如下:

static void Main(string[] args)
{
    // Sender

    var factory = new ConnectionFactory { HostName = "localhost" };
    using var connection = factory.CreateConnection();
    using var channel = connection.CreateModel();

    channel.QueueDeclare(queue: "hello", durable: false, exclusive: false, autoDelete: false, arguments: null);

    const string message = "Hello World!";

    for (int i = 1; i <= 101; i++)
    {
        var body = Encoding.UTF8.GetBytes($"{message} {i}");
        channel.BasicPublish(exchange: string.Empty,
                                routingKey: "hello",
                                basicProperties: null,
                                body: body);

        Console.WriteLine($"Sent {message} {i}");
    }

    // Consumer

    var consumer = new EventingBasicConsumer(channel);
    consumer.Received += async (model, ea) =>
    {
        var body = ea.Body.ToArray();
        var message = Encoding.UTF8.GetString(body);
        await Task.Delay(2000);
        Console.WriteLine($"Received {message}");
    };

    channel.BasicConsume(queue: "hello",
                            autoAck: true,
                            consumer: consumer);


    Console.WriteLine("Press [enter] to exit.");
    Console.ReadLine();
}

结论:达到预期效果,即使发布1000个事件,消息处理速度也非常快。

疑问

目前第三次方案在测试环境表现良好,但不确定该方案是否适合生产环境使用,希望有人解答该问题,或分享生产环境中快速处理事件的相关经验。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 01:00:59