在ASP.NET Core中,RabbitMQ的消费者未被调用

编程语言 2026-07-11

我有一个RabbitMQ发布者,代码如下:

var factory = new ConnectionFactory { HostName = "localhost" };
using var connection = await factory.CreateConnectionAsync();
using var channel = await connection.CreateChannelAsync();

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

const string message = "Hello World!";
var body = Encoding.UTF8.GetBytes(message);

await channel.BasicPublishAsync(exchange: string.Empty, routingKey: "hello", body: body);
Console.WriteLine($" [x] Sent {message}");

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

我创建了一个ASP.NET Core MVC应用作为消费者。下面我添加了一个将从RabbitMQ接收消息的类:

public class RMQConfig
{
    public async void AddConsumer()
    {
        var factory = new ConnectionFactory { HostName = "localhost" };
        using var connection = await factory.CreateConnectionAsync();
        using var channel = await connection.CreateChannelAsync();
        await channel.QueueDeclareAsync(queue: "hello", durable: false, exclusive: false, autoDelete: false, arguments: null);
        var consumer = new AsyncEventingBasicConsumer(channel);

        consumer.ReceivedAsync += (model, ea) =>
        {
            var body = ea.Body.ToArray();
            var message = Encoding.UTF8.GetString(body);
            Console.WriteLine($" [x] Received {message}");
            return Task.CompletedTask;
        };

        await channel.BasicConsumeAsync("hello", autoAck: true, consumer: consumer);
    }
}

在Program类中,我把该类的方法注册为单例:

builder.Services.AddSingleton(new RMQConfig().AddConsumer);

问题在于消息没有被消费到。我在 consumer.ReceivedAsync() 上设置了断点,但断点从未被命中。

我已经验证过,如果把 AddConsumer() 方法的全部代码放到 Program 类中,当RabbitMQ发送消息时,它就会被调用。问题出在哪里?

解决方案

问题在于连接和通道对象的生命周期。

在你的方法中,你有:

using var connection = await factory.CreateConnectionAsync();
using var channel = await connection.CreateChannelAsync();

AddConsumer() 执行完成后,方法返回,using 里的dispose会关闭连接和通道。这会立即停止消费者。

重新启动应用时,方法会再次执行,短暂地创建消费者,处理排队中的消息,然后再次释放。

为了解决这个问题,你需要把 connectionchannel 存储在字段中,例如:

public class RMQConfig
{
    private IConnection? _connection;
    private IChannel? _channel;

    public async Task AddConsumer()
    {
        var factory = new ConnectionFactory { HostName = "localhost" };

        _connection = await factory.CreateConnectionAsync();
        _channel = await _connection.CreateChannelAsync();

        // ...
    }
}

你可能需要实现 IDisposable,以正确关闭 connectionchannel,然后把 RMQConfig 注册为单例。

站内所有文章版权归属LeftHeroAI导航站,无授权禁止任何主体转载、抄袭、复制内容,亦不得私自架设镜像站点。一经侵权,本站将通过法律途径追责。

相关文章