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

如何创建可异步输出海量CSV数据的IObservable<T>实例并解决现有实现无响应问题

问题分析与解决方案

你遇到的核心问题有几个:误用Publish()导致Observable未激活、日期递增的语法bug、资源管理不规范,还有订阅后立即释放导致流被中断。我们一步步来解决这些问题:

1. 最关键问题:Publish()的误用

你在Observable.Create后添加了.Publish(),这个操作会把冷Observable转换成热Observable。热Observable需要手动调用Connect()才会开始执行内部逻辑,而你的场景是每次订阅都应该重新读取文件(冷Observable的默认行为),所以直接去掉Publish()就能让内部代码正常执行,这是你看不到任何输出的核心原因。

2. 修复日期递增的bug

原代码里date.AddDays(1);是无效的:DateTime是值类型,AddDays()会返回新的实例,不会修改原变量。必须改成date = date.AddDays(1);,否则循环会无限停留在初始日期,永远无法推进。

3. 修正文件操作与资源管理

  • 原代码使用FileMode.Create打开文件,这会覆盖原有文件,读取文件应该用FileMode.Open;
  • 用using语句自动管理FileStream和StreamReader的生命周期,避免手动Close()可能带来的资源泄漏;
  • 把取消令牌token传递给ReadLineAsync(),并在循环中检查token.IsCancellationRequested,及时响应订阅取消的请求。

修正后的完整代码

public static IObservable<Data> GetHistoricalData(DateTime startDate, DateTime endDate)
{
    return Observable.Create<Data>(async (observer, token) =>
    {
        try
        {
            DateTime date = startDate;
            // 循环时检查取消令牌,响应订阅取消
            while (date < endDate && !token.IsCancellationRequested)
            {
                string path = string.Format(@"C:\MYPATH\{0}-{1}-{2}.csv", 
                    date.Year, 
                    date.Month.ToString("00"), 
                    date.Day.ToString("00"));

                if (!File.Exists(path))
                {
                    date = date.AddDays(1);
                    continue;
                }

                // 使用using自动释放文件资源
                using (var stream = File.Open(path, FileMode.Open, FileAccess.Read))
                using (var sr = new StreamReader(stream))
                {
                    while (!sr.EndOfStream && !token.IsCancellationRequested)
                    {
                        // 传递取消令牌给异步读取方法
                        string line = await sr.ReadLineAsync(token);
                        Data d = ParseData(line);
                        observer.OnNext(d);
                    }
                }

                date = date.AddDays(1);
            }

            // 所有数据处理完成后发送完成信号
            observer.OnCompleted();
        }
        catch (Exception ex)
        {
            observer.OnError(ex);
        }
    });
}

4. 修正订阅代码的问题

你当前的调用代码在订阅后立即执行subscription.Dispose(),这会立刻取消订阅,导致Observable还没来得及处理任何数据就停止了。如果需要等待所有数据处理完成再释放订阅,可以用以下两种方式:

方式一:异步等待(推荐,非阻塞)

using (var subscription = signals.Subscribe(ShowSignal))
{
    // 等待Observable发送完成信号(仅当Observable会调用OnCompleted时有效)
    await signals.LastOrDefaultAsync();
}

方式二:阻塞等待(不推荐在UI线程使用)

using (var subscription = signals.Subscribe(ShowSignal))
{
    // 阻塞直到Observable完成
    signals.Wait();
}

入门资源推荐

  • Microsoft Docs Reactive Extensions 入门:官方文档清晰讲解了IObservable/IObserver的核心概念、冷热Observable的区别、常用操作符用法,适合新手建立基础认知;
  • 《Intro to Rx》:经典的Rx入门书籍,从基础概念到实际应用都有细致讲解,重点帮助你从传统同步/Task思维切换到响应式流思维。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.30 22:24:09