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

C#如何并行接收多个IAsyncEnumerable流数据并处理多源异常

并行消费多个IAsyncEnumerable的异常安全实现

原有实现的核心问题

原有代码在多数据源同时抛出异常时触发崩溃,本质是几个设计缺陷导致的:

  • 调用SafeFireAndForget启动并行消费逻辑时没有做完整的异常兜底:TPL Dataflow块在多任务失败场景下,仅会将第一个异常挂载到Completion任务上,其余未被观察的异常会触发.NET的未观察任务异常回调,直接造成进程崩溃
  • Channel容量配置逻辑错误:传入的outputQueueCapacity参数未被实际使用,创建有界通道时错误使用数据源数量作为队列容量
  • 缺少异常/终止场景的资源回收逻辑:任意数据源抛出异常、或者消费方提前终止枚举时,没有给其他正在运行的数据源传递取消信号,会造成协程、连接等资源泄漏
  • 写入Channel前重复调用WaitToWriteAsync属于冗余操作,WriteAsync本身已经内置了队列空位等待逻辑

异常安全的实现代码

以下实现不需要依赖AsyncAwaitBestPractices包,原生支持多异常聚合、取消传播、正确的队列容量配置:

public static async IAsyncEnumerable<T> ExecuteSimultaneouslyAsync<T>(
    this IEnumerable<IAsyncEnumerable<T>> sources,
    int outputQueueCapacity = 0,
    [EnumeratorCancellation] CancellationToken cancellationToken = default)
{
    var sourceList = sources.ToList();
    var linkedCts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken);
    
    // 按传入参数初始化Channel
    var channel = outputQueueCapacity > 0
        ? Channel.CreateBounded<T>(outputQueueCapacity)
        : Channel.CreateUnbounded<T>();

    // 启动所有数据源的并行枚举任务
    var enumerateTasks = sourceList.Select(async source =>
    {
        await foreach (var item in source.WithCancellation(linkedCts.Token).ConfigureAwait(false))
        {
            await channel.Writer.WriteAsync(item, linkedCts.Token).ConfigureAwait(false);
        }
    }).ToList();

    // 等待所有枚举任务完成(无论成功失败),之后关闭Channel写入端
    _ = Task.WhenAll(enumerateTasks)
        .ContinueWith(task =>
        {
            channel.Writer.Complete(task.Exception);
            linkedCts.Dispose();
        }, CancellationToken.None, TaskContinuationOptions.ExecuteSynchronously, TaskScheduler.Default);

    // 消费Channel中的数据,消费方提前终止时自动触发取消
    try
    {
        await foreach (var item in channel.Reader.ReadAllAsync(cancellationToken).ConfigureAwait(false))
        {
            yield return item;
        }
    }
    finally
    {
        // 无论消费是正常结束、异常、还是被取消,都通知所有数据源停止枚举
        linkedCts.Cancel();
    }
}

实现特性

  • 所有数据源抛出的异常会被聚合后传递给消费端,不会出现未观察异常导致的进程崩溃
  • 全链路取消支持:消费方提前终止枚举、任意数据源抛出异常时,都会第一时间给所有运行中的数据源传递取消信号,无资源泄漏
  • 队列容量参数正确生效,支持自定义有界/无界通道配置
  • 移除冗余逻辑,通过ConfigureAwait(false)减少不必要的线程上下文切换,性能更优
  • 不需要额外依赖第三方NuGet包,仅依赖.NET原生提供的Channel和IAsyncEnumerable能力

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 19:33:35