中介式IAsyncEnumerable管道的异常处理与重新抛出实现咨询
解决方案
核心问题在于C#编译器禁止yield return出现在带有catch子句的try块中,我们可以通过两种方式实现需求:
方案一:利用await using自动管理资源(最简实现)
如果仅需要完成捕获异常、释放生产者、传播异常这三个基础需求,无需额外异常处理逻辑,await using是最优选择。它会自动在枚举结束(无论正常完成还是异常终止)时调用生产者的DisposeAsync,异常会直接传播给消费者。
public class Operator<TItem> : IOperator<TItem> { private IProducer<TItem>? _producer; public async IAsyncEnumerable<TItem> GetItemsAsync([EnumeratorCancellation] CancellationToken cancellationToken = default) { // await using 自动管理生产者生命周期 await using var producer = CreateProducer(); _producer = producer; // 流式传输数据到消费者 await foreach (var item in producer.GetItemsAsync().WithCancellation(cancellationToken).ConfigureAwait(false)) { yield return item; } } // 创建具体生产者实例的逻辑 private IProducer<TItem> CreateProducer() { return new ConcreteProducer<TItem>(); } // 实现IOperator的IAsyncDisposable,确保提前释放时清理生产者 public async ValueTask DisposeAsync() { if (_producer != null) { await _producer.DisposeAsync().ConfigureAwait(false); _producer = null; } } }
说明
当生产者在枚举过程中抛出异常时,await using会自动触发生产者的DisposeAsync,异常会直接传递给消费者,完全满足你的三个要求。
方案二:手动实现IAsyncEnumerator(支持自定义异常处理)
如果需要在捕获异常时执行额外操作(比如日志记录、错误统计),可以手动实现异步枚举器,完全控制枚举过程中的异常处理和资源释放。
第一步:实现自定义异步枚举器
internal sealed class ProducerEnumerator<T> : IAsyncEnumerator<T> { private readonly IProducer<T> _producer; private IAsyncEnumerator<T>? _innerEnumerator; private bool _isDisposed; public ProducerEnumerator(IProducer<T> producer, CancellationToken cancellationToken) { _producer = producer; _innerEnumerator = producer.GetItemsAsync().GetAsyncEnumerator(cancellationToken); } public T Current => _innerEnumerator!.Current; public async ValueTask<bool> MoveNextAsync() { try { return await _innerEnumerator!.MoveNextAsync().ConfigureAwait(false); } catch (Exception ex) { // 这里添加自定义异常处理逻辑,比如日志 // Logger.Error("生产者枚举失败", ex); // 异常发生时优先释放生产者资源 await DisposeProducerAsync().ConfigureAwait(false); // 重新抛出异常给消费者 throw; } } public async ValueTask DisposeAsync() { if (_isDisposed) return; // 释放内部枚举器 if (_innerEnumerator != null) { await _innerEnumerator.DisposeAsync().ConfigureAwait(false); _innerEnumerator = null; } // 释放生产者 await DisposeProducerAsync().ConfigureAwait(false); _isDisposed = true; } private async ValueTask DisposeProducerAsync() { await _producer.DisposeAsync().ConfigureAwait(false); } }
第二步:在操作器中使用自定义枚举器
public class Operator<TItem> : IOperator<TItem> { private IProducer<TItem>? _producer; public IAsyncEnumerable<TItem> GetItemsAsync([EnumeratorCancellation] CancellationToken cancellationToken = default) { _producer = CreateProducer(); // 创建包含自定义枚举器的异步可枚举对象 return AsyncEnumerable.CreateEnumerable(() => new ProducerEnumerator<TItem>(_producer, cancellationToken)); } private IProducer<TItem> CreateProducer() { return new ConcreteProducer<TItem>(); } public async ValueTask DisposeAsync() { if (_producer != null) { await _producer.DisposeAsync().ConfigureAwait(false); _producer = null; } } }
说明
这种方式完全掌控了枚举流程,异常触发时可先执行自定义逻辑,再释放生产者并重新抛出异常,同时保证所有资源都能被正确清理。
内容的提问来源于stack exchange,提问作者Taras
相关产品推荐
相关产品推荐

