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
相关产品推荐
相关产品推荐

