在ASP.NET Core中,RabbitMQ的消费者未被调用
我有一个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会关闭连接和通道。这会立即停止消费者。
重新启动应用时,方法会再次执行,短暂地创建消费者,处理排队中的消息,然后再次释放。
为了解决这个问题,你需要把 connection 和 channel 存储在字段中,例如:
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,以正确关闭 connection 和 channel,然后把 RMQConfig 注册为单例。
站内所有文章版权归属LeftHeroAI导航站,无授权禁止任何主体转载、抄袭、复制内容,亦不得私自架设镜像站点。一经侵权,本站将通过法律途径追责。