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

.NET中支持等待已提交任务/作业完成的生产者/消费者模式

嘿,这个场景其实是生产者/消费者模式的常见变种,我来给你梳理几个可行的实现方案,刚好能满足你说的“部分生产者要等待任务完成”的需求,而且完全基于BlockingCollection来扩展:

核心思路拆解

1. 先搭好基础的单消费者多生产者框架

BlockingCollection本身就完美适配多生产者单消费者的场景,我们只需要给每个任务加一个完成信号,让生产者能追踪任务的执行状态。先定义一个包含任务逻辑和完成信号的结构体/记录:

// 封装带完成信号的队列任务
public record QueuedTask(Func<Task> WorkLogic, TaskCompletionSource<bool> CompletionSignal);

// 全局任务队列
private readonly BlockingCollection<QueuedTask> _taskQueue = new BlockingCollection<QueuedTask>();

// 消费者工作循环
private async Task WorkerLoop()
{
    // GetConsumingEnumerable会自动阻塞等待新任务,直到队列被标记为完成添加
    foreach (var taskItem in _taskQueue.GetConsumingEnumerable())
    {
        try
        {
            // 执行任务逻辑
            await taskItem.WorkLogic();
            // 任务完成,给生产者发信号
            taskItem.CompletionSignal.SetResult(true);
        }
        catch (Exception ex)
        {
            // 任务出错,把异常传递给等待的生产者
            taskItem.CompletionSignal.SetException(ex);
        }
    }
}

2. 单个任务提交后等待完成

对于需要等待任务执行结果的生产者,只需要在提交任务时创建一个TaskCompletionSource,把它和任务一起放进队列,然后await这个TCS的Task即可:

// 生产者方法:提交任务并等待完成
public async Task SubmitTaskAndWait(Func<Task> taskWork)
{
    var completionSignal = new TaskCompletionSource<bool>();
    _taskQueue.Add(new QueuedTask(taskWork, completionSignal));
    // 这里会异步等待直到消费者执行完该任务
    await completionSignal.Task;
}

3. 批量提交任务后等待全部完成

如果生产者需要提交一批任务,然后等所有任务都执行完毕,就收集每个任务的CompletionSignal.Task,最后用Task.WhenAll等待全部完成:

// 生产者方法:批量提交并等待所有任务完成
public async Task SubmitBatchAndWaitAll(IEnumerable<Func<Task>> batchTasks)
{
    var completionTasks = new List<Task>();
    foreach (var taskWork in batchTasks)
    {
        var completionSignal = new TaskCompletionSource<bool>();
        _taskQueue.Add(new QueuedTask(taskWork, completionSignal));
        completionTasks.Add(completionSignal.Task);
    }
    // 等待批量里的所有任务都执行完成
    await Task.WhenAll(completionTasks);
}

4. 额外的细节提醒

  • 异常处理:上面的代码已经把任务执行时的异常传递给了生产者,所以生产者可以用try/catch包裹await来捕获错误
  • 队列关闭:当系统要停止接收新任务时,记得调用_taskQueue.CompleteAdding(),这样消费者的GetConsumingEnumerable会在处理完剩余任务后自动退出循环
  • 返回值支持:如果任务需要返回结果,只需要把TaskCompletionSource<bool>改成TaskCompletionSource<TResult>,同时把Func<Task>改成Func<Task<TResult>>,生产者就能拿到任务的返回值
  • 线程安全:BlockingCollection本身是线程安全的,多个生产者同时调用Add完全没问题,不需要额外加锁

这样既保留了BlockingCollection带来的简洁、高效的生产者/消费者模型,又完美满足了部分生产者需要等待任务完成的需求~

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 06:35:44