C#实现可无标签合流并还原源的IObservable组合算子
Rx.NET 无标签可还原双源流合并实现
现有实现参考
以下是当前基于Zip算子的双源组合代码:
public class TransformScript { public IObservable<Tuple<bool,bool>> Process(IObservable<bool> source1, IObservable<bool> source2) { return source1.Zip(source2,(s1,s2) => Tuple.Create(s1,s2)); } }
现有常用组合算子的行为差异:
Zip:必须等待两个流各产出一个配对元素才输出Tuple,会缓存先到的元素等待另一流CombineLatest:首次凑齐两个流的最新值后触发输出,之后任意流更新都输出两个流的最新值组合Merge:直接转发所有流的元素,但完全丢失元素来源信息,接收端无法拆分还原原始双源
需求明确
需要实现的组合器满足以下要求:
- 任意源流推送新元素时立刻转发,不做等待缓存
- 输出为和输入元素类型完全一致的单流,适配单通道传输瓶颈
- 禁止给元素附加显式来源标签、禁止用包装类携带来源元数据
- 接收端可从合并后的单流无损拆分出两个原始源流,不丢失任何元素和来源归属
弹珠行为对比如下:
Merge算子输出(无法拆分):
source 1 -----1-----1-----1----- source 2 ---2----2---------2---- merge ---m-m--m--m-----mm----
预期输出(可无损拆分):
source 1 -----1-----1-----1----- source 2 ---2----2---------2---- output ---2-1--2--1-----12----
实现代码
核心逻辑是利用收发两端对齐的计数隐式携带来源信息,不需要修改元素本身,完整实现如下:
using System.Reactive; using System.Reactive.Disposables; using System.Reactive.Linq; using System.Threading; public class TransformScript { /// <summary> /// 合并两个同类型流为单流,无额外标签,接收端可无损拆分 /// </summary> public IObservable<T> Process<T>(IObservable<T> source1, IObservable<T> source2) { return Observable.Create<T>(observer => { long sendCount1 = 0; long sendCount2 = 0; var syncLock = new object(); var sub1 = source1.Synchronize(syncLock).Subscribe( value => { Interlocked.Increment(ref sendCount1); observer.OnNext(value); }, observer.OnError, () => { if (Volatile.Read(ref sendCount1) == Volatile.Read(ref sendCount2)) observer.OnCompleted(); }); var sub2 = source2.Synchronize(syncLock).Subscribe( value => { Interlocked.Increment(ref sendCount2); observer.OnNext(value); }, observer.OnError, () => { if (Volatile.Read(ref sendCount1) == Volatile.Read(ref sendCount2)) observer.OnCompleted(); }); return new CompositeDisposable(sub1, sub2); }); } /// <summary> /// 接收端拆分方法,从合并流无损还原两个原始源流 /// </summary> public (IObservable<T> source1, IObservable<T> source2) Split<T>(IObservable<T> mergedStream) { var published = mergedStream.Publish().RefCount(); var syncLock = new object(); long recvCount1 = 0; long recvCount2 = 0; var restored1 = published .Synchronize(syncLock) .Where(_ => { if (recvCount1 <= recvCount2) { recvCount1++; return true; } recvCount2++; return false; }); var restored2 = published .Synchronize(syncLock) .Where(_ => { if (recvCount2 < recvCount1) { recvCount2++; return true; } recvCount1++; return false; }); return (restored1, restored2); } }
逻辑说明
- 合并端没有对元素做任何包装、附加标签,输出流元素类型和输入完全一致
- 合并端通过同步锁严格保证元素推送顺序和源流实际到达顺序完全一致,同时维护两个流的发送计数
- 拆分端不需要任何额外传输的元数据,仅通过本地维护的接收计数做匹配:每次收到新元素时,哪个流的已接收计数更小,当前元素就归属该流,和合并端的发送计数严格对齐,可100%无损还原原始流
- 元素到达即转发,不会被缓存等待另一流的元素,满足低延迟要求
- 实现基于Rx原生的线程同步机制,可处理多线程并发推送的场景
该实现依赖传输通道的顺序保证:即接收端收到的元素顺序和合并端发送顺序完全一致,这是Rx序列的默认契约,常规同进程Rx管线、支持顺序保证的消息传输通道都满足该要求。
内容的提问来源于stack exchange,提问作者Pablo
相关产品推荐
相关产品推荐

