合并带取消支持的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; } }
问题分析与解决方案
问题根源
ValueTask重复等待:原代码将MoveNextAsync()返回的ValueTask<bool>转换为Task<bool>后,在catch块中再次等待该任务。由于ValueTask设计为不支持多次等待,这会直接触发NotSupportedException。- 异常处理逻辑错误:原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(); } } }
修正说明
- 避免
ValueTask重复等待:仅在需要使用Task.WhenAny时才转换ValueTask为Task,且每个ValueTask仅被等待一次。 - 正确传递异常:catch块中不再吞掉异常,直接重新抛出,保证上层能捕获到主源的原始异常。
- 完善清理逻辑:添加finally块,确保主源无论正常结束还是异常终止,都会触发副源的清理操作。
内容的提问来源于stack exchange,提问作者Arokh
相关产品推荐
相关产品推荐

