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

合并带取消支持的IAsyncEnumerable实例时遇NotSupportedException问题

问题描述

需要合并两个同类型IAsyncEnumerable实例,满足以下要求:

  • 主数据源结束时,副数据源同步停止枚举
  • Merge函数调用方能在主源结束时执行副源的清理操作

当前实现存在问题:当主数据源抛出异常时,会触发NotSupportedException,打乱上层的异常处理逻辑。


重现代码
// 需要安装 "System.Interactive.Async" NuGet包
using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
using System.Runtime.CompilerServices;
using System.Threading.Channels;

while (true) {
    try {
        using var cts = new CancellationTokenSource();

        var secondaryDataSource = Channel.CreateUnbounded<object?>();

        var reader = PrimaryDataSource(cts.Token).Do(x => {
            Console.WriteLine("---");
        });

        var combinedSource = Merge(
            reader,
            secondaryDataSource.Reader.ReadAllAsync(cts.Token),
            () => secondaryDataSource.Writer.TryComplete(),
            cts.Token
        );

        await foreach (var item in combinedSource) {
            Console.WriteLine(item ?? "<Null>");
            throw new Exception();
        }

    } catch (NotSupportedException ex) {
        throw;
    } catch (Exception) {
    }
}

static async IAsyncEnumerable<object?> PrimaryDataSource([EnumeratorCancellation] CancellationToken ct) {
    while (true) {
        await Task.Delay(Random.Shared.Next(100), ct).ConfigureAwait(false);
        yield return "PrimaryDataSource";

        throw new Exception();
        yield break;
    }
}

static async IAsyncEnumerable<object?> Merge(IAsyncEnumerable<object?> primary, IAsyncEnumerable<object?> secondary, Action primaryFinished, [EnumeratorCancellation] CancellationToken ct) {
    var pIt = primary.GetAsyncEnumerator(ct);
    var sIt = secondary.GetAsyncEnumerator(ct);

    await using var _pIt = pIt.ConfigureAwait(false);
    await using var _sIt = sIt.ConfigureAwait(false);

    Task<bool>? pItTask = null;
    Task<bool>? sItTask = null;

    IAsyncEnumerator<object?>? it = null;
    while (true) {
        Task<bool> task;

        try {

            pItTask ??= pIt.MoveNextAsync().AsTask();
            sItTask ??= sIt.MoveNextAsync().AsTask();


            task = await Task.WhenAny(pItTask, sItTask).ConfigureAwait(false);
            if (pItTask == task) {
                pItTask = null;
                it = pIt;
                if (!task.Result) {
                    primaryFinished();
                    yield break;
                }

            } else {

                sItTask = null;
                it = sIt;
            }
        } catch (Exception) {
            primaryFinished();
            if (pItTask != null) await pItTask;
            if (sItTask != null) await sItTask;
            yield break;
        }

        if (task.Result) yield return it.Current;
    }
}

问题分析与解决方案

问题根源

  1. ValueTask重复等待:原代码将MoveNextAsync()返回的ValueTask<bool>转换为Task<bool>后,在catch块中再次等待该任务。由于ValueTask设计为不支持多次等待,这会直接触发NotSupportedException。
  2. 异常处理逻辑错误:原catch块吞掉了主源抛出的原始异常,仅执行清理后就yield break,导致枚举器内部状态异常,最终抛出非预期的NotSupportedException。

修正后的Merge函数

static async IAsyncEnumerable<object?> Merge(IAsyncEnumerable<object?> primary, IAsyncEnumerable<object?> secondary, Action primaryFinished, [EnumeratorCancellation] CancellationToken ct) {
    await using var pIt = primary.GetAsyncEnumerator(ct);
    await using var sIt = secondary.GetAsyncEnumerator(ct);

    bool primaryCompleted = false;
    bool secondaryCompleted = false;

    try {
        while (!primaryCompleted && !secondaryCompleted && !ct.IsCancellationRequested) {
            var pMoveNext = pIt.MoveNextAsync();
            var sMoveNext = sIt.MoveNextAsync();

            // 优先处理已完成的操作,避免不必要的Task转换
            if (pMoveNext.IsCompleted) {
                primaryCompleted = await pMoveNext.ConfigureAwait(false);
                if (primaryCompleted) {
                    yield return pIt.Current;
                } else {
                    primaryFinished();
                    yield break;
                }
            } else if (sMoveNext.IsCompleted) {
                secondaryCompleted = await sMoveNext.ConfigureAwait(false);
                if (secondaryCompleted) {
                    yield return sIt.Current;
                }
            } else {
                // 仅在需要时将ValueTask转为Task,且只等待一次
                var pTask = pMoveNext.AsTask();
                var sTask = sMoveNext.AsTask();
                var completedTask = await Task.WhenAny(pTask, sTask).ConfigureAwait(false);

                if (completedTask == pTask) {
                    primaryCompleted = await pTask.ConfigureAwait(false);
                    if (primaryCompleted) {
                        yield return pIt.Current;
                    } else {
                        primaryFinished();
                        yield break;
                    }
                } else {
                    secondaryCompleted = await sTask.ConfigureAwait(false);
                    if (secondaryCompleted) {
                        yield return sIt.Current;
                    }
                }
            }
        }
    } catch {
        primaryFinished();
        // 重新抛出原始异常,交由上层处理
        throw;
    } finally {
        // 确保主源无论正常/异常结束,都执行副源清理
        if (!primaryCompleted) {
            primaryFinished();
        }
    }
}

修正说明

  1. 避免ValueTask重复等待:仅在需要使用Task.WhenAny时才转换ValueTask为Task,且每个ValueTask仅被等待一次。
  2. 正确传递异常:catch块中不再吞掉异常,直接重新抛出,保证上层能捕获到主源的原始异常。
  3. 完善清理逻辑:添加finally块,确保主源无论正常结束还是异常终止,都会触发副源的清理操作。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 04:09:52