如何将IAsyncEnumerable转换函数适配为Dataflow的IPropagatorBlock?
适配异步枚举转换函数到Dataflow双向块的最优方案
直接自定义一个同时实现ITargetBlock<TIn>和ISourceBlock<TOut>的Dataflow块,利用Channel做轻量输入缓冲,直接将转换后的异步枚举元素推送给下游,全程不需要额外的BufferBlock和后台Task,性能更优。
核心实现代码
using System.Threading.Channels; using System.Threading.Tasks.Dataflow; public class AsyncEnumerableTransformBlock<TIn, TOut> : ITargetBlock<TIn>, ISourceBlock<TOut> { private readonly Channel<TIn> _inputChannel; private readonly Func<IAsyncEnumerable<TIn>, IAsyncEnumerable<TOut>> _transform; private readonly DataflowBlockOptions _options; private readonly Task _processingTask; private readonly PropagatorBlock<TOut, TOut> _outputPropagator; public AsyncEnumerableTransformBlock(Func<IAsyncEnumerable<TIn>, IAsyncEnumerable<TOut>> transform, DataflowBlockOptions options = null) { _options = options ?? new DataflowBlockOptions(); _inputChannel = Channel.CreateBounded<TIn>(new BoundedChannelOptions(_options.BoundedCapacity) { FullMode = _options.BoundedCapacity > 0 ? BoundedChannelFullMode.Wait : BoundedChannelFullMode.None, SingleReader = true }); _transform = transform; _outputPropagator = new PropagatorBlock<TOut, TOut>(_options); _processingTask = ProcessAsync(); } private async Task ProcessAsync() { try { // 将输入通道转为异步枚举,传入转换函数 var inputEnumerable = _inputChannel.Reader.ReadAllAsync(); var outputEnumerable = _transform(inputEnumerable); // 流式遍历转换结果,直接推送给下游块 await foreach (var item in outputEnumerable.ConfigureAwait(false)) { await _outputPropagator.SendAsync(item, _options.CancellationToken).ConfigureAwait(false); } } catch (Exception ex) { ((IDataflowBlock)_outputPropagator).Fault(ex); } finally { _inputChannel.Writer.Complete(); ((IDataflowBlock)_outputPropagator).Complete(); } } #region ITargetBlock<TIn> 实现 DataflowMessageStatus ITargetBlock<TIn>.OfferMessage(DataflowMessageHeader messageHeader, TIn messageValue, ISourceBlock<TIn> source, bool consumeToAccept) { if (_inputChannel.Writer.TryWrite(messageValue)) { return DataflowMessageStatus.Accepted; } // 通道满时按配置等待或拒绝 if (_options.BoundedCapacity > 0 && _options.EnsureOrdered) { try { _inputChannel.Writer.WriteAsync(messageValue, _options.CancellationToken).Wait(_options.CancellationToken); return DataflowMessageStatus.Accepted; } catch (OperationCanceledException) { return DataflowMessageStatus.Declined; } } return DataflowMessageStatus.Declined; } #endregion #region ISourceBlock<TOut> 实现 TOut ISourceBlock<TOut>.ConsumeMessage(DataflowMessageHeader messageHeader, ITargetBlock<TOut> target, out bool messageConsumed) { return _outputPropagator.ConsumeMessage(messageHeader, target, out messageConsumed); } bool ISourceBlock<TOut>.ReserveMessage(DataflowMessageHeader messageHeader, ITargetBlock<TOut> target) { return _outputPropagator.ReserveMessage(messageHeader, target); } void ISourceBlock<TOut>.ReleaseReservation(DataflowMessageHeader messageHeader, ITargetBlock<TOut> target) { _outputPropagator.ReleaseReservation(messageHeader, target); } #endregion #region IDataflowBlock 实现 Task IDataflowBlock.Completion => _processingTask.ContinueWith(_ => _outputPropagator.Completion).Unwrap(); void IDataflowBlock.Fault(Exception exception) { _inputChannel.Writer.TryComplete(exception); ((IDataflowBlock)_outputPropagator).Fault(exception); _processingTask.Wait(_options.CancellationToken); } void IDataflowBlock.Complete() { _inputChannel.Writer.TryComplete(); ((IDataflowBlock)_outputPropagator).Complete(); } #endregion }
方案优势
- 轻量缓冲:用
Channel<TIn>替代BufferBlock<TIn>,Channel是.NET原生的异步队列,比Dataflow块的开销更低,且原生支持IAsyncEnumerable - 无额外开销:输出侧直接通过
PropagatorBlock的SendAsync推送转换后的元素,不需要额外的BufferBlock<TOut>和后台Task,减少线程切换和块调度开销 - 流式兼容:完全支持
BoundedCapacity配置,能处理百万级数据的流式传输,避免内存过载 - 生命周期整合:严格遵循Dataflow块的
Complete、Fault生命周期规范,和现有管道无缝集成
使用示例
// 定义你的异步转换函数 Func<IAsyncEnumerable<int>, IAsyncEnumerable<string>> transform = input => input.SelectAsync(async num => $"Processed: {num}"); // 创建自定义转换块 var transformBlock = new AsyncEnumerableTransformBlock<int, string>(transform, new DataflowBlockOptions { BoundedCapacity = 1000, // 限制缓冲大小,避免内存溢出 CancellationToken = CancellationToken.None }); // 连接到下游处理块 var printBlock = new ActionBlock<string>(Console.WriteLine); transformBlock.LinkTo(printBlock, new DataflowLinkOptions { PropagateCompletion = true }); // 批量发送输入数据 for (int i = 0; i < 1_000_000; i++) { await transformBlock.SendAsync(i); } // 标记输入完成 ((IDataflowBlock)transformBlock).Complete(); // 等待整个管道处理完成 await printBlock.Completion;
内容的提问来源于stack exchange,提问作者Bouke
相关产品推荐
相关产品推荐

