如何高效合并同源冷热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); }); } }
方案优势
- 无元素丢失:通过Subject提前缓存hot事件,避免订阅cold过程中遗漏新元素;
- 无重复元素:先通过cold全量数据初始化去重集合,再处理缓存的hot事件和后续hot事件,确保每个元素只推送一次;
- 内存可控:去重集合仅存储cold全量数据+订阅cold期间缓存的hot事件,不会持续无限制增长(除非数据源重复推送旧元素,这种场景本身不符合业务逻辑)。
内容的提问来源于stack exchange,提问作者Nick Strupat
相关产品推荐
相关产品推荐

