使用createTimeSeriesAggregator生成的时序聚合结果存在乱序问题
使用
createTimeSeriesAggregator时序聚合器输出时间顺序异常问题 场景说明
- 源表
MD_PTB中的数据按时间排序 - 创建了窗口大小120000毫秒、步长120000毫秒的时序聚合器
tsAggr - 聚合结果输出到
OHLC_PTB表
实现代码
OHLC_PTB=enableTableShareAndPersistence(table=keyedStreamTable(`datetime`symbol`price,1:0,colname_OHLC,coltype_OHLC), tableName="OHLC_PTB", cacheSize=5000, preCache = 1000) tsAggr = createTimeSeriesAggregator( name="tsAggr", windowSize=120000, step=120000, metrics=<[last(TradingDay), last(tradeTime), sum(askvol)]>, dummyTable=MD_PTB, outputTable=OHLC_PTB, timeColumn=`dateTime, keyColumn=`symbol`price )
问题现象
输入表MD_PTB时间有序,但输出表OHLC_PTB时间戳出现跳变,未严格按时间顺序排列:
| dateTime | symbol | |
|---|---|---|
| 0 | 2025.07.01 09:00:00.000 | eg2509 |
| 1 | 2025.07.01 09:02:00.000 | eg2509 |
| 2 | 2025.07.01 09:02:00.000 | eg2509 |
| 3 | 2025.07.01 09:04:00.000 | eg2509 |
| 14 | 2025.07.01 09:08:00.000 | eg2509 |
| 27 | 2025.07.01 09:14:00.000 | eg2509 |
| 33 | 2025.07.01 09:14:00.000 | eg2509 |
解决方案
1. 修正分组键(核心优化)
你的keyColumn包含price字段,这会导致每个不同价格都成为独立分组,触发大量并行窗口计算,不同分组的窗口输出时间相互交错,直接打乱整体时间顺序。时序聚合通常仅按标的(symbol)分组,修改keyColumn为:
keyColumn=`symbol
2. 让输出表自动维护有序
如果必须保留多分组键,可在创建持久化表时指定sortColumns参数,让输出表自动按时间列排序:
OHLC_PTB=enableTableShareAndPersistence( table=keyedStreamTable(`datetime`symbol`price,1:0,colname_OHLC,coltype_OHLC), tableName="OHLC_PTB", cacheSize=5000, preCache = 1000, sortColumns=`dateTime -- 指定按时间列排序 )
3. 处理迟到数据(可选)
如果存在少量迟到数据,可在createTimeSeriesAggregator中设置lateOrderTime参数,让窗口等待指定时长后再输出,避免提前输出旧窗口结果:
tsAggr = createTimeSeriesAggregator( name="tsAggr", windowSize=120000, step=120000, metrics=<[last(TradingDay), last(tradeTime), sum(askvol)]>, dummyTable=MD_PTB, outputTable=OHLC_PTB, timeColumn=`dateTime, keyColumn=`symbol, lateOrderTime=1000 -- 允许1秒的迟到时间,可按需调整 )
4. 查询时排序(最终兜底)
如果对实时性要求不高,可在查询输出表时手动排序:
select * from OHLC_PTB order by dateTime, symbol
内容的提问来源于stack exchange,提问作者haru
相关产品推荐
相关产品推荐

