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

.NET 7 BackgroundService结合RabbitMQ优雅停止异常排查

问题描述

基于.NET 7实现了一个BackgroundService,通过RabbitMQ消费消息,采用微软官方的IBackgroundTaskQueue、BackgroundTaskQueue和QueuedHostedService组件实现任务队列。消息在Consumer_ReceivedAsync事件中被加入队列,由BackgroundService取出并执行Process方法。

应用运行正常,但调用StopAsync停止服务时出现两个问题:

  • 有时Process方法抛出AlreadyClosedException
  • 有时CleanRabbit方法抛出「通道已关闭」的异常

不确定取消令牌的处理是否正确,需要实现服务的优雅停止并正确释放RabbitMQ资源。


问题根源分析

  1. 令牌逻辑混乱:QueuedBackgroundService自行维护_tokenSource,未与宿主stoppingToken关联,导致停止逻辑不一致;CleanRabbit被多分支调用,重复释放资源。
  2. 资源释放时机错误:停止服务时未等待队列任务完成就关闭RabbitMQ资源,导致正在执行的Process访问已关闭资源;重复调用CleanRabbit引发重复关闭异常。
  3. 消费停止不及时:收到停止信号后仍继续接收新消息,导致资源释放时还有任务在处理。

解决方案

1. 统一取消令牌,关联宿主停止信号

删除QueuedBackgroundService自定义_tokenSource,直接使用宿主提供的stoppingToken,保证停止逻辑一致:

protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
    await BackgroundProcessing(stoppingToken);
}
    
private async Task BackgroundProcessing(CancellationToken cancellationToken)
{
    while (!cancellationToken.IsCancellationRequested)
    {
        try
        {
            var workItem = await TaskQueue.DequeueAsync(cancellationToken);
            await workItem(cancellationToken).ConfigureAwait(false);
        }
        catch (OperationCanceledException)
        {
            break;
        }
        catch (Exception ex)
        {
            _logger.LogError(ex, "执行任务时发生错误");
        }
    }
}

public override async Task StopAsync(CancellationToken stoppingToken)
{
    await base.StopAsync(stoppingToken);
}

2. 控制RabbitMQ资源释放的原子性

添加标志位和锁,确保CleanRabbit只执行一次:

private bool _isDisposed = false;
private readonly object _disposeLock = new object();

private void CleanRabbit()
{
    lock (_disposeLock)
    {
        if (_isDisposed) return;
        _isDisposed = true;
    }

    try
    {
        _basicConsumer.Received -= Consumer_Received;
        if (_consumerChannel != null)
        {
            if (_consumerChannel.IsOpen)
            {
                try { _consumerChannel.Close(); }
                catch (Exception ex) { _logger.LogError(ex, "关闭通道出错"); }
            }
            _consumerChannel.Dispose();
        }

        if (_persistentConnection != null)
        {
            if (_persistentConnection.IsConnected)
            {
                try { _persistentConnection.Close(); }
                catch (Exception ex) { _logger.LogError(ex, "关闭连接出错"); }
            }
            _persistentConnection.Dispose();
        }
    }
    catch (Exception ex)
    {
        _logger.LogError(ex, "清理RabbitMQ资源出错");
    }
}

3. 调整优雅停止流程

先停止接收新消息,等待队列任务完成后再释放资源:

public EventBusRabbitMQ(IHostApplicationLifetime hostApplicationLifetime, IBackgroundTaskQueue taskQueue, /*其他参数*/)
{
    _cancellationToken = hostApplicationLifetime.ApplicationStopping;
    _taskQueue = taskQueue;

    hostApplicationLifetime.ApplicationStopping.Register(() =>
    {
        _logger.LogInformation("开始优雅停止,停止接收新消息");
        _basicConsumer.Received -= Consumer_Received;
        _ = WaitForQueueEmptyAndCleanupAsync();
    });
}

private async Task WaitForQueueEmptyAndCleanupAsync()
{
    while (!_taskQueue.IsEmpty && !_cancellationToken.IsCancellationRequested)
    {
        await Task.Delay(100, _cancellationToken);
    }

    _logger.LogInformation("队列任务处理完成,清理RabbitMQ资源");
    CleanRabbit();
}

private async Task Consumer_ReceivedAsync(object sender, DataReceived @event)
{
    try
    {
        if (_cancellationToken.IsCancellationRequested)
        {
            _logger.LogInformation("已收到停止信号,拒绝处理新消息");
            return;
        }
        await ProcessQueueEvent(@event, @event.KeyName);
    }
    catch(Exception ex)
    {
        _logger.LogError("Consumer_ReceivedAsync出错: {Message}", ex.Message);
    }
}

4. 修改Process方法,移除重复清理逻辑

private async ValueTask Process(DataReceived dataReceived, string routingKey, CancellationToken token)
{
    try
    {
        token.ThrowIfCancellationRequested();
        await DoWork(dataReceived, routingKey);
    }
    catch (OperationCanceledException) when (token.IsCancellationRequested)
    {
        _logger.LogInformation("任务已被取消");
    }
    catch (AlreadyClosedException ac)
    {
        _logger.LogError("通道已关闭: {Message}", ac.Message);
    }
    catch (Exception ex)
    {
        HandleException(dataReceived, ex);
    }
}

5. 完善队列空状态判断

在BackgroundTaskQueue中添加IsEmpty属性:

public bool IsEmpty => _queue.Reader.Count == 0;

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 13:22:04