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

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被错误分组。正确做法是:

  1. 以TickTime的整秒值作为分组依据
  2. 等待足够长的时间(覆盖Broker的最大延迟)再关闭分组,确保所有属于该秒的Tick都被收集
  3. 空分组自动输出空列表,满足无报价时的需求

修正后的代码

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 10:05:56