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

如何高效合并同源冷热Observable并避免重复与内存膨胀?

冷热Observable合并的无内存增长去重实现

我有一个数据源,既能通过查询获取全量数据(对应cold Observable),也会在新元素新增时发布事件(对应hot Observable)。目前有三种合并冷热Observable的方案,但各有问题:

  • 方案1:var o = cold.Merge(hot);
    逻辑可行,但会丢失订阅间隙中hot发出的元素——也就是在订阅并处理cold全量数据的过程中,hot可能已经推送了新元素,这部分元素会被遗漏。

  • 方案2:var o = hot.Merge(cold);
    不会丢失元素,但会产生重复元素——cold的查询结果中可能包含hot在订阅前已经推送的元素,合并后这些元素会重复出现。

  • 方案3:var o = hot.Merge(cold).Distinct();
    解决了重复问题,但Distinct()会维护一个持续增长的元素列表,长期运行会导致内存占用不断增加,存在性能隐患。

需求目标

实现类似Distinct()的去重效果,但避免内存持续增长。核心思路是:通过Subject缓存hot事件,先完成cold全量数据的订阅与去重记录,再重放缓存的hot事件并去重,之后正常转发后续的hot事件。

现有方案行为演示代码

using System;
using System.Reactive.Linq;
using System.Reactive.Subjects;
using System.Threading.Tasks;

public class Program
{
    public static async Task Main()
    {
        // 模拟cold Observable:查询现有全量数据
        var cold = Observable.Return(new[] { "Item1", "Item2" }).SelectMany(x => x);
        // 模拟hot Observable:新增元素事件流
        var hot = new Subject<string>();

        // 提前触发hot事件,模拟订阅前已推送的元素
        hot.OnNext("Item2");
        hot.OnNext("Item3");

        Console.WriteLine("=== 方案1:cold.Merge(hot) ===");
        var o1 = cold.Merge(hot);
        o1.Subscribe(x => Console.WriteLine($"Received: {x}"));
        hot.OnNext("Item4");
        await Task.Delay(100);

        Console.WriteLine("\n=== 方案2:hot.Merge(cold) ===");
        var o2 = hot.Merge(cold);
        o2.Subscribe(x => Console.WriteLine($"Received: {x}"));
        hot.OnNext("Item5");
        await Task.Delay(100);

        Console.WriteLine("\n=== 方案3:hot.Merge(cold).Distinct() ===");
        var o3 = hot.Merge(cold).Distinct();
        o3.Subscribe(x => Console.WriteLine($"Received: {x}"));
        hot.OnNext("Item6");
        await Task.Delay(100);
    }
}

无内存增长的去重合并实现

using System;
using System.Collections.Generic;
using System.Reactive;
using System.Reactive.Linq;
using System.Reactive.Subjects;
using System.Threading.Tasks;

public static class ObservableExtensions
{
    public static IObservable<T> MergeColdAndHotWithoutDuplicates<T>(
        this IObservable<T> cold, 
        IObservable<T> hot, 
        IEqualityComparer<T> comparer = null)
    {
        comparer ??= EqualityComparer<T>.Default;
        var hotBuffer = new Subject<T>();
        var coldCompletionSignal = new TaskCompletionSource<bool>();
        var seenItems = new HashSet<T>(comparer);

        // 先订阅hot,将所有事件缓存到Subject中,避免丢失
        var hotSubscription = hot.Subscribe(
            item => hotBuffer.OnNext(item),
            ex => hotBuffer.OnError(ex),
            () => hotBuffer.OnCompleted()
        );

        return Observable.Create<T>(observer =>
        {
            // 第一步:处理cold全量数据,记录所有已出现的元素
            var coldSubscription = cold.Subscribe(
                item =>
                {
                    if (seenItems.Add(item))
                    {
                        observer.OnNext(item);
                    }
                },
                ex =>
                {
                    observer.OnError(ex);
                    coldCompletionSignal.TrySetResult(true);
                },
                () => coldCompletionSignal.TrySetResult(true)
            );

            // 第二步:cold处理完成后,重放缓存的hot事件并去重,之后转发新的hot事件
            var replayTask = coldCompletionSignal.Task.ContinueWith(_ =>
            {
                hotBuffer.Subscribe(
                    item =>
                    {
                        if (seenItems.Add(item))
                        {
                            observer.OnNext(item);
                        }
                    },
                    observer.OnError,
                    observer.OnCompleted
                );
            }, TaskScheduler.Default);

            // 清理所有订阅资源
            return new CompositeDisposable(coldSubscription, hotSubscription, replayTask);
        });
    }
}

方案优势

  1. 无元素丢失:通过Subject提前缓存hot事件,避免订阅cold过程中遗漏新元素;
  2. 无重复元素:先通过cold全量数据初始化去重集合,再处理缓存的hot事件和后续hot事件,确保每个元素只推送一次;
  3. 内存可控:去重集合仅存储cold全量数据+订阅cold期间缓存的hot事件,不会持续无限制增长(除非数据源重复推送旧元素,这种场景本身不符合业务逻辑)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 03:15:17