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

ASP.NET Core中如何将任务加入后台工作队列并等待处理结果

实现可返回处理结果的顺序执行单实例后台工作队列

官方提供的即发即忘队列实现只做了单向的任务投递,没有给任务预留结果回传的通道,要实现调用方等待返回的顺序队列,核心是给每个入队任务绑定一个TaskCompletionSource作为结果回传的载体,后台服务单线程串行消费队列即可,完全适配Cv2 Dnn不可并发调用的场景。

第一步:定义队列模型与线程安全队列实现

每个队列项需要包含待处理的图像数据,以及一个TaskCompletionSource<T>实例负责传递处理结果/异常,是实现调用方异步等待返回的核心组件,队列采用Channel实现原生支持异步生产消费,比手写ConcurrentQueue加信号量的方案更轻量稳定。

// 图像处理结果模型,可根据实际业务调整字段
public class ImageProcessResult
{
    public bool Success { get; set; }
    public byte[]? ProcessedImage { get; set; }
    public string? ErrorMsg { get; set; }
}

// 队列内部存储的任务项
internal class ImageProcessQueueItem
{
    public byte[] ImageData { get; init; }
    public TaskCompletionSource<ImageProcessResult> TaskCompletionSource { get; init; }
}

// 队列服务抽象
public interface IImageProcessQueue
{
    // 入队方法直接返回Task,调用方await即可等待处理结果
    Task<ImageProcessResult> EnqueueAsync(byte[] imageData, CancellationToken cancellationToken = default);
}

public class ImageProcessQueue : IImageProcessQueue
{
    private readonly Channel<ImageProcessQueueItem> _queue;

    public ImageProcessQueue()
    {
        // 可根据业务承载能力调整队列容量上限
        var options = new BoundedChannelOptions(1000)
        {
            FullMode = BoundedChannelFullMode.Wait
        };
        _queue = Channel.CreateBounded<ImageProcessQueueItem>(options);
    }

    public async Task<ImageProcessResult> EnqueueAsync(byte[] imageData, CancellationToken cancellationToken = default)
    {
        var tcs = new TaskCompletionSource<ImageProcessResult>(TaskCreationOptions.RunContinuationsAsynchronously);
        var item = new ImageProcessQueueItem
        {
            ImageData = imageData,
            TaskCompletionSource = tcs
        };

        await _queue.Writer.WriteAsync(item, cancellationToken);
        // 直接返回tcs的Task,后续处理完成后会自动将结果写入该Task
        return await tcs.Task;
    }

    // 供后台服务调用的出队方法
    internal ValueTask<ImageProcessQueueItem> DequeueAsync(CancellationToken cancellationToken = default)
    {
        return _queue.Reader.ReadAsync(cancellationToken);
    }
}

第二步:实现单线程串行消费的后台工作服务

后台服务的执行逻辑为单线程循环出队,同一时间永远只会执行一个Cv2 Dnn处理任务,从根源上避免并发导致的Dnn崩溃、CPU抢占损耗问题。

public class ImageProcessBackgroundService : BackgroundService
{
    private readonly IImageProcessQueue _queue;
    private readonly ILogger<ImageProcessBackgroundService> _logger;
    // 注入你已实现的Cv2 Dnn处理封装类
    private readonly ICv2DnnProcessor _dnnProcessor;

    public ImageProcessBackgroundService(
        IImageProcessQueue queue,
        ICv2DnnProcessor dnnProcessor,
        ILogger<ImageProcessBackgroundService> logger)
    {
        _queue = queue;
        _dnnProcessor = dnnProcessor;
        _logger = logger;
    }

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        _logger.LogInformation("单实例图像处理后台服务启动");
        // 单线程循环,严格按入队顺序处理任务
        while (!stoppingToken.IsCancellationRequested)
        {
            ImageProcessQueueItem? queueItem = null;
            try
            {
                queueItem = await _queue.DequeueAsync(stoppingToken);
                // 执行Cv2 Dnn图像处理逻辑,此处始终为串行执行
                var processResult = _dnnProcessor.Process(queueItem.ImageData);
                // 将结果写入tcs,等待的请求端会立刻收到结果
                queueItem.TaskCompletionSource.SetResult(processResult);
            }
            catch (Exception ex)
            {
                _logger.LogError(ex, "图像处理任务执行失败");
                if (queueItem != null)
                {
                    // 异常必须回传给调用方,不可吞掉
                    queueItem.TaskCompletionSource.SetException(ex);
                }
            }
        }
    }

    public override async Task StopAsync(CancellationToken cancellationToken)
    {
        _logger.LogInformation("图像处理后台服务停止,清空待处理任务");
        // 服务停止时为所有待处理任务返回停止异常,避免请求永久等待
        while (_queue is ImageProcessQueue concreteQueue && concreteQueue.TryDequeue(out var item))
        {
            item.TaskCompletionSource.SetException(new InvalidOperationException("图像处理服务已停止"));
        }
        await base.StopAsync(cancellationToken);
    }
}

第三步:服务注册与接口调用

在Program.cs中注册相关服务,注意队列和后台服务均需注册为单例:

// 注册队列与Dnn处理器
builder.Services.AddSingleton<IImageProcessQueue, ImageProcessQueue>();
builder.Services.AddSingleton<ICv2DnnProcessor, Cv2DnnProcessor>();
// 注册后台工作服务
builder.Services.AddHostedService<ImageProcessBackgroundService>();

在API控制器中注入队列服务,直接调用入队方法等待结果即可:

[ApiController]
[Route("api/image")]
public class ImageProcessController : ControllerBase
{
    private readonly IImageProcessQueue _processQueue;

    public ImageProcessController(IImageProcessQueue processQueue)
    {
        _processQueue = processQueue;
    }

    [HttpPost("process")]
    public async Task<IActionResult> ProcessImage(IFormFile imageFile, CancellationToken cancellationToken)
    {
        if (imageFile == null || imageFile.Length == 0)
        {
            return BadRequest("请上传有效的图像文件");
        }
        using var ms = new MemoryStream();
        await imageFile.CopyToAsync(ms, cancellationToken);
        var imageBytes = ms.ToArray();

        // 入队并异步等待处理结果,不会阻塞请求线程
        var result = await _processQueue.EnqueueAsync(imageBytes, cancellationToken);

        if (!result.Success)
        {
            return StatusCode(500, result.ErrorMsg);
        }
        return File(result.ProcessedImage!, "image/png");
    }
}

关键注意点

  • 整个后台处理逻辑是单线程串行执行的,完全符合全CPU负载优化、避免Dnn并发崩溃的要求
  • 没有使用即发即忘模式,每个请求都能拿到对应的处理结果或异常,和直接调用本地处理方法的体验完全一致
  • 初始化TaskCompletionSource时必须指定TaskCreationOptions.RunContinuationsAsynchronously参数,避免请求响应的续跑逻辑占用后台工作线程,拖慢队列处理效率
  • 如果需要支持任务取消,可在入队时传入请求的取消令牌,注册令牌回调时将对应任务的tcs设置为取消状态即可

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.01 20:18:29