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

无需等待结果的任务队列实现问题求助

背景任务队列优化方案:清理已完成任务与优雅停机处理

针对你提出的两个问题,下面给出具体的代码修改方案和关键说明:

一、移除队列中已完成的任务,避免队列膨胀

现有代码直接将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;
    }
}

关键修改说明

  1. 任务存储逻辑调整

    • 将队列存储从Task改为Func<CancellationToken, Task>,确保任务由队列调度执行,而非入队前就启动,从根源避免队列中存在已完成的任务。
    • 若必须存储Task实例,可在CleanupStaleItems方法中添加遍历清理逻辑,移除IsCompleted为true的任务。
  2. 已完成任务清理

    • 在每次入队时触发清理操作,定期移除队列中的无效任务。高并发场景下,也可单独启动后台定时任务执行清理,避免入队时的性能开销。
  3. 优雅停机实现

    • 添加_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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 17:15:47