1分钟聚合数据异常:成交额字段缺失问题排查
问题分析:1分钟K线成交额字段无数据排查
场景与数据处理逻辑
- 处理规模:约5200只股票的L1逐笔Tick数据聚合为1分钟K线
- 数据源特性:关键时间点(9:25:00、11:30:00、15:00:00)后持续推送数据,直至该时段数据传输完成(如贵州茅台集合竞价最后一笔Tick在9:25:04推送,含成交数据,此前Tick成交额为0)
- 数据预处理步骤:
- 数据缓存:关键时间点前几秒(如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数据已成功持久化至本地数据库,但1分钟K线的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数据,且统一重置了时间戳。对于
reactiveStateEngine来说,单只股票的单个1分钟窗口内只有一条数据,prev(amount)无法获取到前序值,导致stock_tick_delta_amount计算结果为null,求和后stock_1min_amount自然无数据。 - 聚合指标依赖的前提不成立:
sum(stock_tick_delta_amount)需要连续的Tick数据来累计增量,但预处理阶段丢弃了同窗口内的前置Tick,直接切断了增量计算的数据源。
解决方向:
- 调整预处理逻辑:保留同窗口内的所有Tick数据,确保
prev(amount)能获取到前序值; - 替换聚合指标:将
sum(stock_tick_delta_amount)改为last(stock_tick_amount) - first(stock_tick_amount),直接用窗口首尾的成交额差值计算时段总成交额,无需依赖增量字段。
内容的提问来源于stack exchange,提问作者Jane
相关产品推荐
相关产品推荐

