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

如何订阅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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 17:52:59