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

TPL Dataflow结合Rx:Observer丢失消息问题咨询

TPL Dataflow与Rx混合使用的消息接收差异解析

在下面的示例代码中,source(BroadcastBlock<int>实例)同时扮演两个角色:它是TPL Dataflow管道的起始节点,又通过AsObservable()转换为IObservable<int>供Rx订阅。运行后会出现特殊现象:管道末端的sink能接收所有消息,但订阅Observable的Observer仅能收到第一条和最后一条消息。

代码示例

public class Test
{
    public static async Task Run()
    {
        var source = new BroadcastBlock<int>(x => x);
        var transform = new TransformBlock<int, int>(x => x * 2);
        var sink = new ActionBlock<int>(x => Console.WriteLine($"from ActionBlock: {x}"));

        source.LinkTo(transform);
        transform.LinkTo(sink);

        var observable = source.AsObservable();

        IDisposable subscription = observable.Subscribe(x => Console.WriteLine($"from observer: {x}"));

        await source.SendAsync(1);
        await source.SendAsync(2);
        await source.SendAsync(3);
        await source.SendAsync(4);
        await source.SendAsync(5);
        await source.SendAsync(6);

        subscription.Dispose();
    }
}

运行输出

from observer: 1
from ActionBlock: 2
from ActionBlock: 4
from ActionBlock: 6
from ActionBlock: 8
from ActionBlock: 10
from ActionBlock: 12
from observer: 6

原因说明

  • BroadcastBlock的默认行为是保留最后一条处理过的消息,当新订阅建立时会立即推送这条留存消息,所以订阅一开始就收到了第一条消息1。
  • 后续发送的2-5号消息会不断覆盖BroadcastBlock留存的内容,但Rx订阅不会实时接收这些中间消息——BroadcastBlock作为Observable时,仅会在订阅终止(调用Dispose)时,推送最后留存的消息6。
  • TPL Dataflow内部的LinkTo机制是主动转发每一条消息:所有进入source的消息都会被转发给transform,因此sink能收到所有经转换后的消息。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 20:10:36