无需等待结果的任务队列实现问题求助
背景任务队列优化方案:清理已完成任务与优雅停机处理
针对你提出的两个问题,下面给出具体的代码修改方案和关键说明:
一、移除队列中已完成的任务,避免队列膨胀
现有代码直接将Task实例存入队列,若任务已完成仍会占据队列空间。我们可以在入队/出队操作时定期清理已完成的任务,同时建议调整队列存储逻辑,从存储已启动的Task改为存储任务执行委托Func<CancellationToken, Task>,这样能更好地控制任务执行时机,避免任务在入队前就启动,从根源减少无效任务堆积。
二、应用关闭时优雅处理任务队列
应用停止时需要完成三个核心动作:
- 禁止新任务入队
- 处理队列中所有待执行任务
- 确保已启动的任务能响应停止信号并安全收尾
修改后的完整代码
public interface IBackgroundTaskQueue { // 改为接受任务委托,而非已启动的Task bool QueueBackgroundWorkItem(Func<CancellationToken, Task> workItem); Task<Func<CancellationToken, Task>?> DequeueAsync(CancellationToken cancellationToken); } public class BackgroundTaskQueue : IBackgroundTaskQueue { private readonly ConcurrentQueue<Func<CancellationToken, Task>> _workItems = new(); private readonly ILogger<BackgroundTaskQueue> _logger; private readonly SemaphoreSlim _signal = new(0); private readonly CancellationTokenSource _shutdownTokenSource = new(); private bool _isShuttingDown; public BackgroundTaskQueue(ILogger<BackgroundTaskQueue> logger, IHostApplicationLifetime lifetime) { _logger = logger; // 注册应用停止事件 lifetime.ApplicationStopping.Register(OnStopping); } public bool QueueBackgroundWorkItem(Func<CancellationToken, Task> workItem) { if (_isShuttingDown || workItem == null) return false; // 入队前清理无效任务(如果是存储Task实例的场景,此处改为清理已完成的Task) CleanupStaleItems(); _workItems.Enqueue(workItem); _signal.Release(); return true; } private void CleanupStaleItems() { // 若坚持存储Task实例,替换为以下逻辑: // while (_workItems.TryPeek(out var task) && task?.IsCompleted == true) // { // _workItems.TryDequeue(out _); // } } private void OnStopping() { _isShuttingDown = true; _shutdownTokenSource.Cancel(); try { // 释放信号量,让DequeueAsync退出等待 _signal.Release(); // 处理队列中剩余的所有任务 ProcessRemainingTasksAsync().GetAwaiter().GetResult(); } catch (Exception ex) { _logger.LogWarning(ex, "清理任务队列时发生错误"); } } private async Task ProcessRemainingTasksAsync() { while (_workItems.TryDequeue(out var workItem)) { try { await workItem(_shutdownTokenSource.Token); } catch (Exception ex) { _logger.LogError(ex, "执行剩余任务时失败"); } } } public async Task<Func<CancellationToken, Task>?> DequeueAsync(CancellationToken cancellationToken) { try { // 合并停止令牌和业务令牌,确保任务能响应停止信号 using var linkedTokenSource = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, _shutdownTokenSource.Token); await _signal.WaitAsync(linkedTokenSource.Token); if (_isShuttingDown) return null; _workItems.TryDequeue(out var workItem); return workItem; } catch (OperationCanceledException) { // 停止时取消等待,无需记录日志 } catch (Exception ex) { _logger.LogWarning(ex, "出队任务时发生错误"); } return null; } }
关键修改说明
任务存储逻辑调整
- 将队列存储从
Task改为Func<CancellationToken, Task>,确保任务由队列调度执行,而非入队前就启动,从根源避免队列中存在已完成的任务。 - 若必须存储
Task实例,可在CleanupStaleItems方法中添加遍历清理逻辑,移除IsCompleted为true的任务。
- 将队列存储从
已完成任务清理
- 在每次入队时触发清理操作,定期移除队列中的无效任务。高并发场景下,也可单独启动后台定时任务执行清理,避免入队时的性能开销。
优雅停机实现
- 添加
_isShuttingDown标记,禁止停止期间新任务入队。 - 在
OnStopping中,先触发停止令牌,再逐个处理队列中剩余任务,确保应用关闭前完成所有待执行任务。 - 通过令牌合并,让所有任务能响应停止信号,安全收尾。
- 添加
后台Worker配合示例
任务队列需要配合后台Worker服务使用,示例Worker代码如下:
public class BackgroundTaskWorker : BackgroundService { private readonly IBackgroundTaskQueue _taskQueue; private readonly ILogger<BackgroundTaskWorker> _logger; public BackgroundTaskWorker(IBackgroundTaskQueue taskQueue, ILogger<BackgroundTaskWorker> logger) { _taskQueue = taskQueue; _logger = logger; } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { _logger.LogInformation("后台任务Worker已启动"); while (!stoppingToken.IsCancellationRequested) { var workItem = await _taskQueue.DequeueAsync(stoppingToken); if (workItem == null) continue; try { await workItem(stoppingToken); } catch (Exception ex) { _logger.LogError(ex, "执行后台任务失败"); } } _logger.LogInformation("后台任务Worker已停止"); } }
最后在DI容器中注册服务:
services.AddSingleton<IBackgroundTaskQueue, BackgroundTaskQueue>(); services.AddHostedService<BackgroundTaskWorker>();
内容的提问来源于stack exchange,提问作者Jimbo Jones
相关产品推荐
相关产品推荐

