如何订阅Observable并获取事件及全部历史值(System.Reactive)
解决方案
核心思路是利用Rx的可重播序列缓存csvAppendMock的所有历史值,无需自行维护本地列表。每次定时器或文件事件触发时,直接从重播序列中获取全部历史数据。
完整实现代码
using System; using System.Reactive.Linq; using System.Threading.Tasks; class Program { static void Main() { // 1. 模拟文件系统监视器(键盘输入),去重仅触发新名称 var csvAppendMock = Observable.Create<string>(async observer => { Console.WriteLine("请输入文件名(输入exit退出):"); while (true) { var input = await Task.Run(() => Console.ReadLine()); if (input?.Trim().ToLower() == "exit") { observer.OnCompleted(); break; } observer.OnNext(input?.Trim() ?? string.Empty); } }) .DistinctUntilChanged(); // 仅新名称触发 // 2. 创建重播序列,自动缓存所有历史值 var replayCsv = csvAppendMock.Replay(); replayCsv.Connect(); // 启动连接,开始缓存数据 // 3. 3秒定时器序列 var timer = Observable.Interval(TimeSpan.FromSeconds(3)); // 4. 合并触发源:定时器/文件事件都作为触发信号 var triggerSignals = timer .Select(_ => Unit.Default) .Merge(csvAppendMock.Select(_ => Unit.Default)); // 5. 订阅合并后的信号,每次触发时获取所有历史值 triggerSignals .SelectMany(_ => replayCsv.ToList()) // 从重播序列获取全部历史值 .Subscribe(allHistory => { Console.WriteLine("\n当前所有历史文件名:"); foreach (var fileName in allHistory) { Console.WriteLine($" - {fileName}"); } Console.WriteLine("------------------------"); }); Console.ReadLine(); } }
关键说明
- 重播序列的作用:
Replay()操作符会创建一个可重播的Observable,自动缓存所有已推送的值。后续任何订阅都能立即获取到之前的所有历史数据,无需手动维护列表。 - 合并触发逻辑:将定时器和文件事件都转换为统一的触发信号,确保两者任一触发时都能触发历史值的获取。
- 避免ToList()阻塞问题:之前直接用
csvAppendMock.ToList()会无限等待序列完成(因为输入序列是无限的),而重播序列是热序列,缓存的历史值会立即推送给新订阅,所以replayCsv.ToList()会快速返回当前所有缓存数据。
内容的提问来源于stack exchange,提问作者Jochen van Wylick
相关产品推荐
相关产品推荐

