如何创建可异步输出海量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
相关产品推荐
相关产品推荐

