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

RabbitMQ消费者托管服务仅消费首批消息,需重启才能继续的问题排查

RabbitMQ BackgroundService仅能消费首批消息的问题解决

问题场景

在Linux环境下使用C#实现的BackgroundService作为RabbitMQ消费者,启动服务后可正常接收队列中的首批消息,但后续发送到队列的新消息无法被自动接收,必须重启服务才能处理。

相关代码

ConsumerWorker类

public class ConsumerWorker : BackgroundService
{
    private ConnectionFactory _connectionFactory;
    private IConnection _connection;
    private IModel _channel;

    public ConsumerWorker()
    {
    }

    public override Task StartAsync(CancellationToken cancellationToken)
    {
        _connectionFactory = new ConnectionFactory
        {
            HostName = "localhost",
            Port = 5672,
            UserName = "User",
            Password = "password",
            DispatchConsumersAsync = true,
        };
        _connection = _connectionFactory.CreateConnection();
        _channel = _connection.CreateModel();

        _channel.QueueDeclare(queue: "myqueue",
                                durable: true,
                                exclusive: false,
                                autoDelete: false,
                                arguments: null);
        _channel.BasicQos(0, 1, false);

        return base.StartAsync(cancellationToken);
    }

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        stoppingToken.ThrowIfCancellationRequested();

        var consumer = new AsyncEventingBasicConsumer(_channel);
        consumer.Received += async (bc, ea) =>
        {
            var bodyBytes = ea.Body.ToArray();

            _channel.BasicAck(ea.DeliveryTag, false);
        };

        _channel.BasicConsume(queue: "myqueue", autoAck: true, consumer: consumer);

        await Task.CompletedTask;
    }

    public override async Task StopAsync(CancellationToken cancellationToken)
    {
        await base.StopAsync(cancellationToken);
        _connection.Close();
    }
}

Program.cs

IHost host = Host.CreateDefaultBuilder(args)
    .ConfigureServices(services =>
    {
        services.AddHostedService<ConsumerWorker>();
    })
    .Build();

await host.RunAsync();

问题原因与解决方案

问题出在BasicConsume方法的autoAck参数设置上:

_channel.BasicConsume(queue: "myqueue", autoAck: true, consumer: consumer);

当autoAck设为true时,RabbitMQ会在消息推送给消费者后自动确认消息已处理;但代码中又手动调用了_channel.BasicAck(ea.DeliveryTag, false),这导致确认逻辑冲突。消费者处理完首批消息后,RabbitMQ因自动确认机制认为消费流程已完成,不再继续推送新消息。

将autoAck改为false即可解决问题:

_channel.BasicConsume(queue: "myqueue", autoAck: false, consumer: consumer);

修改后,消息确认完全由手动调用的BasicAck控制,结合之前配置的BasicQos(0,1,false)(每次仅推送1条未确认消息),RabbitMQ会在收到手动确认后,继续推送队列中的下一条新消息,服务就能持续消费后续消息。

内容的提问来源于stack exchange,提问作者Pea Kay See Es

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 09:03:14