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

如何将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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 14:10:49