.NET Kestrel:如何自定义控制服务Shutdown流程?
.NET RabbitMQ桥接服务自定义Shutdown流程解决方案
核心需求回顾
你需要实现自定义的服务关闭流程:
- 收到关闭信号后,先停止RabbitMQ事件订阅
- 处理完内部缓冲队列的所有事件(包括处理过程中触发的新回复消息)
- 最后执行服务关闭
但Kestrel默认会在IHostApplicationLifetime.ApplicationStopping触发时停止接收新请求,阻碍了队列清空的过程。
解决思路与实现步骤
1. 手动接管关闭信号,绕过默认触发逻辑
不要依赖Kestrel默认的停止触发,自己监听系统关闭信号(比如控制台Ctrl+C、Windows服务停止事件),然后按自定义流程执行。可以通过注册HostedService来封装关闭逻辑:
public class ShutdownManager : IHostedService { private readonly IHostApplicationLifetime _appLifetime; private readonly RabbitMQSubscriber _rabbitSubscriber; private readonly InternalQueueProcessor _queueProcessor; private CancellationTokenSource _shutdownCts = new(); public ShutdownManager(IHostApplicationLifetime appLifetime, RabbitMQSubscriber rabbitSubscriber, InternalQueueProcessor queueProcessor) { _appLifetime = appLifetime; _rabbitSubscriber = rabbitSubscriber; _queueProcessor = queueProcessor; } public Task StartAsync(CancellationToken cancellationToken) { // 监听控制台关闭信号 Console.CancelKeyPress += OnCancelKeyPress; // 监听容器/系统发起的停止请求 _appLifetime.ApplicationStopping.Register(OnApplicationStopping); return Task.CompletedTask; } private void OnCancelKeyPress(object? sender, ConsoleCancelEventArgs e) { e.Cancel = true; // 阻止系统立即终止进程 InitiateCustomShutdown().Wait(); } private void OnApplicationStopping() { InitiateCustomShutdown().Wait(); } private async Task InitiateCustomShutdown() { if (_shutdownCts.IsCancellationRequested) return; _shutdownCts.Cancel(); // 步骤1:停止RabbitMQ订阅,不再接收新消息 _rabbitSubscriber.StopSubscription(); // 步骤2:等待内部队列所有消息处理完成(含触发的回复) await _queueProcessor.WaitForAllProcessingComplete(_shutdownCts.Token); // 步骤3:通知Kestrel和应用正式停止 _appLifetime.StopApplication(); } public Task StopAsync(CancellationToken cancellationToken) { Console.CancelKeyPress -= OnCancelKeyPress; _shutdownCts.Dispose(); return Task.CompletedTask; } }
2. 控制Kestrel停止时机
Kestrel的停止完全由IHostApplicationLifetime.StopApplication()触发,只要在自定义流程(队列清空)完成后再调用这个方法,就能让Kestrel延迟到队列处理完毕后再停止,无需修改Kestrel的默认配置。
3. 确保内部队列处理完所有消息(含新产生的回复)
你的队列处理器需要跟踪活跃任务数,确保所有消息(包括处理中生成的回复)都完成投递:
public class InternalQueueProcessor { private readonly ConcurrentQueue<Message> _inboxQueue; private readonly ConcurrentQueue<Message> _outboxQueue; private int _activeTaskCount; public async Task WaitForAllProcessingComplete(CancellationToken token) { // 等待活跃任务数归零 while (_activeTaskCount > 0 && !token.IsCancellationRequested) { await Task.Delay(100, token); } // 确保outbox剩余消息全部投递完成 await ProcessOutboxUntilEmpty(token); } private async Task ProcessOutboxUntilEmpty(CancellationToken token) { while (_outboxQueue.TryDequeue(out var message) && !token.IsCancellationRequested) { Interlocked.Increment(ref _activeTaskCount); try { await DeliverToRabbitMQ(message); } finally { Interlocked.Decrement(ref _activeTaskCount); } } } // 处理inbox消息的示例方法 public async Task ProcessInboxMessage(Message message) { Interlocked.Increment(ref _activeTaskCount); try { // 转发消息到其他服务并获取回复 var reply = await ForwardToTargetService(message); if (reply != null) { _outboxQueue.Enqueue(reply); // 立即处理回复消息 await ProcessOutboxUntilEmpty(CancellationToken.None); } } finally { Interlocked.Decrement(ref _activeTaskCount); } } }
4. 注册服务到DI容器
在Program.cs中完成服务注册:
var builder = WebApplication.CreateBuilder(args); // 注册业务组件 builder.Services.AddSingleton<RabbitMQSubscriber>(); builder.Services.AddSingleton<InternalQueueProcessor>(); // 注册自定义关闭管理器 builder.Services.AddHostedService<ShutdownManager>(); var app = builder.Build(); // 中间件配置... app.Run();
关键注意点
- 必须阻止系统默认的立即终止行为(如
Console.CancelKeyPress中设置e.Cancel=true),否则进程会被强行杀死,导致队列处理中断。 - 跟踪活跃任务数时要保证线程安全,优先使用
Interlocked类或SemaphoreSlim。 - 回复消息的投递操作必须完成后再递减任务数,避免遗漏未完成的投递。
内容的提问来源于stack exchange,提问作者JesperGJensen
相关产品推荐
相关产品推荐

