如何实现System.Reactive中两个Observable的全排列组合操作?
实现Rx的全排列组合算子
要实现你需要的CreatePermutations方法,核心思路是缓存两个Observable序列的所有历史元素,当任意一个序列发射新元素时,将其与另一个序列的全部历史元素进行组合并推送结果。下面是具体的实现方案:
完整代码实现
using System; using System.Collections.Generic; using System.Reactive.Disposables; using System.Reactive.Linq; public static class ObservableExtensions { public static IObservable<TResult> CreatePermutations<T1, T2, TResult>( this IObservable<T1> source1, IObservable<T2> source2, Func<T1, T2, TResult> selector) { // 参数校验 if (source1 == null) throw new ArgumentNullException(nameof(source1)); if (source2 == null) throw new ArgumentNullException(nameof(source2)); if (selector == null) throw new ArgumentNullException(nameof(selector)); // 缓存两个序列的所有历史元素 var cachedItems1 = new List<T1>(); var cachedItems2 = new List<T2>(); // 线程安全锁:Rx的OnNext可能在多线程环境触发,避免并发修改集合 var syncLock = new object(); return Observable.Create<TResult>(observer => { bool isSource1Completed = false; bool isSource2Completed = false; // 检查是否两个源都已完成,完成则结束结果序列 void CheckCompletion() { if (isSource1Completed && isSource2Completed) { observer.OnCompleted(); } } // 订阅第一个序列 var sub1 = source1.Subscribe( item1 => { lock (syncLock) { cachedItems1.Add(item1); // 用新元素和第二个序列的所有历史元素组合 foreach (var item2 in cachedItems2) { observer.OnNext(selector(item1, item2)); } } }, observer.OnError, () => { lock (syncLock) { isSource1Completed = true; CheckCompletion(); } }); // 订阅第二个序列 var sub2 = source2.Subscribe( item2 => { lock (syncLock) { cachedItems2.Add(item2); // 用新元素和第一个序列的所有历史元素组合 foreach (var item1 in cachedItems1) { observer.OnNext(selector(item1, item2)); } } }, observer.OnError, () => { lock (syncLock) { isSource2Completed = true; CheckCompletion(); } }); // 组合两个订阅,取消订阅时自动释放资源 return new CompositeDisposable(sub1, sub2); }); } }
实现逻辑说明
- 缓存历史元素:用两个
List分别存储source1和source2发射过的所有元素,确保新元素进来时能和全部历史元素组合。 - 线程安全保障:通过
lock关键字保护缓存集合的读写操作,避免Rx多线程发射元素导致的并发修改问题。 - 新元素处理逻辑:
- 当
source1发射新元素时,先将其加入缓存,再遍历source2的所有缓存元素,调用选择器生成结果并推送。 - 当
source2发射新元素时,同理,遍历source1的所有缓存元素生成组合结果。
- 当
- 完成事件处理:当两个源序列都完成时,结果序列也会触发
OnCompleted,符合Rx的序列生命周期规范。
验证你的测试场景
用你给出的测试代码验证:
var o1 = new Subject<string>(); var o2 = new Subject<string>(); IObservable<string> result = o1.CreatePermutations(o2, (item1, item2) => $"{item1}{item2}"); // 订阅结果序列,打印推送内容 result.Subscribe(Console.WriteLine); o1.OnNext("A"); // 无输出(o2还没有元素) o2.OnNext("1"); // 输出:A1 o1.OnNext("B"); // 输出:B1 o2.OnNext("2"); // 输出:A2、B2 o1.OnNext("C"); // 输出:C1、C2
完全符合你的预期结果。
注意事项
- 内存占用:如果序列会发射大量元素,缓存集合会持续增长,可能导致内存压力。如果不需要永久保留全部历史,可以根据需求添加缓存清理逻辑(比如限制缓存大小)。
- 错误处理:任意一个源序列触发
OnError,结果序列会立即转发错误并终止,这是Rx的标准行为。
内容的提问来源于stack exchange,提问作者Bogey
相关产品推荐
相关产品推荐

