为何BroadcastBlock搭配AsObservable使用时无法接收全部值?
老哥我太懂你踩的这个坑了!你是不是碰到了这种情况:用BroadcastBlock链接两个ActionBlock的时候,0到4的所有值都能被正常接收打印,但换成AsObservable()订阅观察者之后,就收不全所有值了?先把你那段不完整的代码补全(我猜你后面的代码大概是这样的):
var broadCastBlock = new BroadcastBlock<int>(i => i); Func<string, Action<int>> actionGenerator = (s) => (int i) => Console.WriteLine(s + " " + i); Action<int> actionFoo = actionGenerator("FOO"); Action<int> actionBar = actionGenerator("BAR"); // 第一段:用ActionBlock链接,正常接收所有值 var actionBlockFoo = new ActionBlock<int>(actionFoo); var actionBlockBar = new ActionBlock<int>(actionBar); broadCastBlock.LinkTo(actionBlockFoo); broadCastBlock.LinkTo(actionBlockBar); for (int i = 0; i < 5; i++) { broadCastBlock.Post(i); } Thread.Sleep(1000); Console.WriteLine("-Observable-"); // 第二段:换成AsObservable订阅,发现收不全值 broadCastBlock = new BroadcastBlock<int>(i => i); var observerFoo = Observer.Create<int>(actionFoo); var observerBar = Observer.Create<int>(actionBar); var observable = broadCastBlock.AsObservable(); // 这里如果是先Post再订阅,就会出问题! for (int i = 0; i < 5; i++) { broadCastBlock.Post(i); } observable.Subscribe(observerFoo); observable.Subscribe(observerBar); Thread.Sleep(1000);
咱们先拆解一下为什么会出现这种差异:
1. ActionBlock链接的情况为什么正常?
当你调用broadCastBlock.LinkTo(actionBlockFoo)的时候,是在Post任何值之前就把两个ActionBlock和BroadcastBlock绑定好了。BroadcastBlock的逻辑是:每收到一个Post的值,就会把这个值发送给所有已经链接好的消费者,所以0到4的每个值都会同时发给两个ActionBlock,自然能打印出全部内容。
2. AsObservable订阅为什么收不全?
这就要说到BroadcastBlock的核心特性了:它默认只保留最新的一条消息,旧的消息会被新的消息直接覆盖掉。而且用AsObservable()的时候,你如果是先Post值,再订阅观察者,那订阅发生的时候,之前的0到3早就被新的4覆盖丢弃了,观察者只能拿到当前保留的最新值4;就算你是先订阅再Post,但如果你的订阅时机晚了(比如Post了几个值之后才订阅),那之前的旧值也拿不到。
另外还有个隐藏点:AsObservable()的每个订阅者,其实都是BroadcastBlock的一个新的隐式消费者,但BroadcastBlock只会把订阅之后收到的新值和当前保留的最新值发给这个新消费者,之前的历史消息一概没有。
怎么解决这个问题?
给你几个实用的解决方案,根据你的需求选:
方案一:先订阅,再Post值
这是最直接的,确保所有观察者都在Post任何值之前完成订阅,这样每个Post的值都会被所有订阅者收到,和ActionBlock的效果一致:Console.WriteLine("-Observable-"); broadCastBlock = new BroadcastBlock<int>(i => i); var observable = broadCastBlock.AsObservable(); // 先完成所有订阅! observable.Subscribe(Observer.Create<int>(actionFoo)); observable.Subscribe(Observer.Create<int>(actionBar)); // 再开始Post值 for (int i = 0; i < 5; i++) { broadCastBlock.Post(i); } Thread.Sleep(1000);方案二:用BufferBlock替代BroadcastBlock(如果需要广播所有历史消息)
如果你希望不管订阅时机,晚订阅的观察者也能拿到之前的所有消息,那BroadcastBlock就不合适了——它的设计就是广播最新值,不是全量值。换成BufferBlock就可以,它会缓存所有收到的消息,每个新订阅者都能拿到完整的历史:var bufferBlock = new BufferBlock<int>(); var observable = bufferBlock.AsObservable(); // 先Post值 for (int i = 0; i <5; i++) { bufferBlock.Post(i); } // 后订阅依然能拿到所有0-4 observable.Subscribe(Observer.Create<int>(i => Console.WriteLine("Late FOO " + i))); Thread.Sleep(1000);方案三:用Rx的ReplaySubject做全量广播
如果你本来就想走Rx的路子,直接用ReplaySubject就完事了,它会缓存所有发送的值,不管什么时候订阅都能拿到全量:var replaySubject = new ReplaySubject<int>(); replaySubject.Subscribe(Observer.Create<int>(actionFoo)); replaySubject.Subscribe(Observer.Create<int>(actionBar)); for (int i = 0; i <5; i++) { replaySubject.OnNext(i); } Thread.Sleep(1000);
最后再划个重点
BroadcastBlock是TPL Dataflow里专门用来广播最新值的组件,不是用来全量广播所有消息的。如果你的需求是把所有产生的消息都发给订阅者,那它就不是最佳选择,要么调整订阅时机,要么换用其他组件。
备注:内容来源于stack exchange,提问作者Ivan Petrov

