.NET 7 BackgroundService结合RabbitMQ优雅停止异常排查
问题描述
基于.NET 7实现了一个BackgroundService,通过RabbitMQ消费消息,采用微软官方的IBackgroundTaskQueue、BackgroundTaskQueue和QueuedHostedService组件实现任务队列。消息在Consumer_ReceivedAsync事件中被加入队列,由BackgroundService取出并执行Process方法。
应用运行正常,但调用StopAsync停止服务时出现两个问题:
- 有时
Process方法抛出AlreadyClosedException - 有时
CleanRabbit方法抛出「通道已关闭」的异常
不确定取消令牌的处理是否正确,需要实现服务的优雅停止并正确释放RabbitMQ资源。
问题根源分析
- 令牌逻辑混乱:
QueuedBackgroundService自行维护_tokenSource,未与宿主stoppingToken关联,导致停止逻辑不一致;CleanRabbit被多分支调用,重复释放资源。 - 资源释放时机错误:停止服务时未等待队列任务完成就关闭RabbitMQ资源,导致正在执行的
Process访问已关闭资源;重复调用CleanRabbit引发重复关闭异常。 - 消费停止不及时:收到停止信号后仍继续接收新消息,导致资源释放时还有任务在处理。
解决方案
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
相关产品推荐
相关产品推荐

