Rx.NET 报价聚合时Buffer延迟问题求助
问题:基于TickTime的1秒K线分组聚合异常
我正在使用响应式库进行外汇和CFD品种报价的1秒K线聚合,此前尝试IAsyncEnumerable扩展未成功,改用推送策略引入响应式操作符,但因对响应式编程不熟悉遇到问题。
需求是按TickTime的秒数对品种报价进行分组,但Broker端的TickTime与通道接收的ReceivedTime存在10-30ms的延迟。尝试过Delay操作符与Buffer(TimeSpan.FromSeconds(1))结合,但效果不佳,知道应该使用Window操作符但不清楚具体用法。
注意:若1秒内无报价,Observable需输出空列表,且不支持后续更新。
当前代码
await Task.Run(() => { var asset = symbol; var observable = streamQuotes.ToObservable() .Delay(TimeSpan.FromMilliseconds(15)) .Buffer(TimeSpan.FromSeconds(1), Scheduler.Default) .Subscribe(quotes => { var topic = Topic.ToObject($"{provider}-{asset}-{Topic.SECOND}-1"); if (DebugMode) Logger.LogInformation( $"{topic} {JsonConvert.SerializeObject(quotes, Formatting.Indented)}"); //_ = AggregateAsync(topic, quotes, epochBegin, epochStep++), cancellationToken); }); cancellationToken.Register(() => observable.Dispose()); Subscriptions.Add(observable); }, cancellationToken);
当前错误输出
[ { "Id": 142, "Bid": 45282.72, "Ask": 45341.63, "Symbol": "BTCUSD", "Time": 1707431630102, "TickTime": "2024-02-08T22:33:50.102+00:00", "ReceivedTime": "2024-02-08T22:33:50.1187799+00:00" }, { "Id": 142, "Bid": 45283.05, "Ask": 45341.96, "Symbol": "BTCUSD", "Time": 1707431630552, "TickTime": "2024-02-08T22:33:50.552+00:00", "ReceivedTime": "2024-02-08T22:33:50.5654973+00:00" }, { "Id": 142, "Bid": 45283.86, "Ask": 45342.77, "Symbol": "BTCUSD", "Time": 1707431630591, "TickTime": "2024-02-08T22:33:50.591+00:00", "ReceivedTime": "2024-02-08T22:33:50.6046997+00:00" }, { "Id": 142, "Bid": 45285.07, "Ask": 45343.98, "Symbol": "BTCUSD", "Time": 1707431630702, "TickTime": "2024-02-08T22:33:50.702+00:00", "ReceivedTime": "2024-02-08T22:33:50.7155601+00:00" }, { "Id": 142, "Bid": 45284.73, "Ask": 45343.64, "Symbol": "BTCUSD", "Time": 1707431630753, "TickTime": "2024-02-08T22:33:50.753+00:00", "ReceivedTime": "2024-02-08T22:33:50.7658716+00:00" }, { "Id": 142, "Bid": 45284.45, "Ask": 45343.36, "Symbol": "BTCUSD", "Time": 1707431630853, "TickTime": "2024-02-08T22:33:50.853+00:00", "ReceivedTime": "2024-02-08T22:33:50.866821+00:00" }, { "Id": 142, "Bid": 45284.54, "Ask": 45343.46, "Symbol": "BTCUSD", "Time": 1707431631003, "TickTime": "2024-02-08T22:33:50.999+00:00", "ReceivedTime": "2024-02-08T22:33:51.005728+00:00" }, { "Id": 142, "Bid": 45284.64, "Ask": 45343.55, "Symbol": "BTCUSD", "Time": 1707431631040, "TickTime": "2024-02-08T22:33:51.03+00:00", "ReceivedTime": "2024-02-08T22:33:51.0530009+00:00" }, { "Id": 142, "Bid": 45284.64, "Ask": 45343.55, "Symbol": "BTCUSD", "Time": 1707431631040, "TickTime": "2024-02-08T22:33:51.05+00:00", "ReceivedTime": "2024-02-08T22:33:51.056112+00:00" } ]
当前输出错误,最后两个对象(TickTime为22:33:51.03和22:33:51.05)属于下一秒的分组,却被混入上一秒的结果中。
解决方案
核心思路
默认的Buffer/Window是按接收时间而非业务时间(TickTime)切分窗口,加上延迟的不可控性,导致跨秒Tick被错误分组。正确做法是:
- 以TickTime的整秒值作为分组依据
- 等待足够长的时间(覆盖Broker的最大延迟)再关闭分组,确保所有属于该秒的Tick都被收集
- 空分组自动输出空列表,满足无报价时的需求
修正后的代码
await Task.Run(() => { var asset = symbol; const int maxDelayMs = 30; // 覆盖Broker的最大延迟,可根据实际情况调整 var observable = streamQuotes.ToObservable() // 按TickTime的整秒分组,直到最大延迟时间过后关闭分组 .GroupByUntil( // 分组键:TickTime的整秒时间点 quote => new DateTime( quote.TickTime.Year, quote.TickTime.Month, quote.TickTime.Day, quote.TickTime.Hour, quote.TickTime.Minute, quote.TickTime.Second ), // 分组关闭条件:该秒的最大延迟时间过后,触发分组关闭 group => Observable.Timer(TimeSpan.FromMilliseconds(maxDelayMs), Scheduler.Default) ) // 将每个分组转换为列表,空分组会生成空列表 .SelectMany(group => group.ToList()) .Subscribe(quotes => { var topic = Topic.ToObject($"{provider}-{asset}-{Topic.SECOND}-1"); if (DebugMode) Logger.LogInformation( $"{topic} {JsonConvert.SerializeObject(quotes, Formatting.Indented)}"); //_ = AggregateAsync(topic, quotes, epochBegin, epochStep++), cancellationToken); }); cancellationToken.Register(() => observable.Dispose()); Subscriptions.Add(observable); }, cancellationToken);
关键操作符说明
- GroupByUntil:
- 第一个参数:生成分组键,确保同一整秒的Tick被分到同一组
- 第二个参数:定义分组的生命周期,等待
maxDelayMs后关闭分组,确保所有延迟到达的Tick都被纳入对应分组
- SelectMany + ToList:将每个分组的Observable序列转换为列表,空分组会输出空列表,完美满足无报价时的需求
优化建议
- 可根据实际运行中的延迟数据调整
maxDelayMs,平衡数据完整性和输出及时性 - 如果需要严格按整秒边界输出(比如每到整秒就输出前一秒结果,允许少量延迟Tick丢失),可结合
Observable.Interval作为时间触发源,但需评估业务可接受的丢数风险
内容的提问来源于stack exchange,提问作者Strausa
相关产品推荐
相关产品推荐

