1分钟聚合问题:缺少交易量数据
场景是将大约5200只股票的L1 Tick数据聚合成1 分钟的K 线数据。
当前的数据源在关键时间点之后继续推送数据(如9:25:00、11:30:00、15:00:00),直到该时刻对应的所有数据都已完全传输为止。例如,贵州茅台在集合竞价阶段的最后一笔Tick会在9:25:04推送,其中包含交易量数据;而在此时间点之前推送的Tick,处于集合竞价阶段,交易量为0。
为了生成标准化的时间戳(例如将9:25:07推送的数据的1 分钟时间戳设为9:25:00),源端数据将执行以下处理:
-
数据缓存:在即将到来的关键时间点之前的几秒钟开始(如从9:24:55起),所有接收到的Tick数据会缓存在Python内存中,暂时不向下游转发。
-
统一处理:在该关键时间点之后的几秒钟(如9:25:20),对缓存数据进行处理:
-
将9:25:00之后接收的所有数据的时间戳强制重置为9:25:00。
-
基于tradetime进行去重,只保留关键时间点之后最后接收的数据。例如,如果对贵州茅台在9:25:04和 9:25:07收到两笔Tick,则只保留9:25:07的那笔(其中包含交易量数据),并将其时间戳重置为9:25:00。
-
数据写入:处理完成后,按时间戳顺序写入流表。
在Tick级别,看起来实际数据已成功接收(已持久化到本地数据库,Tick流表仍在写入)。

然而,问题在于1 分钟聚合结果中的交易量字段stock_1min_amount显示为无数据。

我的引擎创建代码如下,其中stock_1min_amount通过对stock_tick_delta_amount求和得到,而stock_tick_delta_amount是在reactiveStateEngine中通过amount - prev(amount) 计算得到的。
Engine_DTS_FTkStk_FB1mStk_RT_20260303 = createDailyTimeSeriesEngine(
name="Engine_DTS_FTkStk_FB1mStk_RT_20260303",
windowSize=60000,
step=60000,
metrics=<[first(stock_tick_current), last(stock_tick_current), max(stock_tick_current), min(stock_tick_current), sum(stock_tick_delta_volume), sum(stock_tick_delta_amount), last(stock_tick_amount)]>,
dummyTable=FactorStreamTickStock,
outputTable=FactorStreamBase1MinStock,
keyColumn=`code,
timeColumn=`tradetime,
sessionBegin=[time(09:15:00),time(09:30:00),time(13:00:00)],
sessionEnd=[time(09:25:00),time(11:30:00),time(15:00:00)],
garbageSize=50000,
useSystemTime=false,
useWindowStartTime=false,
mergeSessionEnd=true,
mergeLastWindow=true,
roundTime=true,
forceTriggerSessionEndTime=6000,
closed="right",
fill=["null", "null", "null", "null", "ffill", "ffill", "ffill"],
parallelism=4
)
所以,我想请教这是我的引擎配置问题,还是性能方面的问题?
解决方案
从你的描述来看,分钟级的最后一笔Tick似乎是在聚合窗口已经输出之后才到达。
即便你将时间戳重设为9:25:00,大多数流处理引擎仍然使用事件到达时间或水印来决定何时关闭窗口。如果窗口在晚到Tick处理完成前就已经关闭,该Tick的交易量将不会被包含在1 分钟蜡烛中。
你可能需要检查引擎是否支持晚到事件处理或聚合窗口的水印延迟。允许几秒钟的晚到(例如5–10秒)通常可以解决最后一笔交易刚好在边界之外到达的这种情况。
还要再核实,窗口触发是否基于到达时间而不是修改后的时间戳。
如果你愿意分享9:25附近的一段Tick流样本,可能更容易复现该行为。