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

如何实现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);
        });
    }
}

实现逻辑说明

  1. 缓存历史元素:用两个List分别存储source1和source2发射过的所有元素,确保新元素进来时能和全部历史元素组合。
  2. 线程安全保障:通过lock关键字保护缓存集合的读写操作,避免Rx多线程发射元素导致的并发修改问题。
  3. 新元素处理逻辑:
    • 当source1发射新元素时,先将其加入缓存,再遍历source2的所有缓存元素,调用选择器生成结果并推送。
    • 当source2发射新元素时,同理,遍历source1的所有缓存元素生成组合结果。
  4. 完成事件处理:当两个源序列都完成时,结果序列也会触发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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 04:28:14