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

使用crossSectionalEngine时,如何让keyedStreamTable保留最新记录?

如何让KeyedStreamTable保留最新聚合数据而非首次插入记录?

针对你遇到的问题——crossSectionalEngine分批触发时,KeyedStreamTable因默认保留首次插入记录导致聚合结果偏小,有两种可行方案实现覆盖更新:

方案一:使用KeyedStreamTable的updateWhenDuplicate参数(推荐,需DolphinDB 2.00.9+)

DolphinDB 2.00.9及以后版本为keyedStreamTable新增了updateWhenDuplicate参数,设置为true后,同主键的新插入数据会直接覆盖旧记录,无需额外逻辑。

修改你的键值流表创建代码:

// 创建键值流表: FactorStreamAggBase1MinSector_shsz (主键: code, tradetime),开启重复主键覆盖更新
colNames = [`tradetime, `code, `sector_agg_1min_shsz_amount]
colTypes = [TIMESTAMP, SYMBOL, DOUBLE]
keyColumns = [`code, `tradetime]
share(table=keyedStreamTable(keyColumns, 1:0, colNames, colTypes, updateWhenDuplicate=true),
  sharedName = `FactorStreamAggBase1MinSector_shsz)

保持crossSectionalEngine的配置不变即可,后续同一(code, tradetime)的新聚合结果会自动替换旧值。

方案二:普通流表+触发器+Upsert(兼容低版本)

如果你的DolphinDB版本不支持updateWhenDuplicate,可以通过中间流表配合upsert!函数实现覆盖更新:

  1. 创建普通流表作为crossSectionalEngine的输出目标
  2. 注册触发器,将中间流表的数据upsert!到键值表(重复主键自动覆盖)

示例代码:

// 1. 创建中间普通流表
colNames = [`tradetime, `code, `sector_agg_1min_shsz_amount]
colTypes = [TIMESTAMP, SYMBOL, DOUBLE]
share(table=streamTable(1:0, colNames, colTypes), sharedName=`TempStreamTable)

// 2. 创建最终存储用的键值表
finalKeyedTable = keyedTable(keyColumns=`code`tradetime, colNames=colNames, colTypes=colTypes)
share(finalKeyedTable, `FactorStreamAggBase1MinSector_shsz)

// 3. 定义触发器处理函数:用upsert!实现覆盖更新
def handleUpsert(data) {
    upsert!(finalKeyedTable, data, keyColumns=`code`tradetime, update=true)
}

// 4. 给中间流表注册触发器
subscribeTable(tableName=`TempStreamTable, actionName="upsert_to_keyed_table", handler=handleUpsert, msgAsTable=true)

// 5. 修改crossSectionalEngine的输出表为中间流表
Engine_CS_IT001_FAggB1mSecshsz_RT_20251219 = createCrossSectionalEngine(
   name="Engine_CS_IT001_FAggB1mSecshsz_RT_20251219",
   metrics=<[sum(stock_1min_amount)]>,
   dummyTable=IT001,
   outputTable=TempStreamTable,
   timeColumn=`tradetime,
   useSystemTime=false,
   keyColumn=`code,
   triggeringInterval=100,
   contextByColumn=`stock_daily_shsz,
   keyFilter=<stock_daily_shsz != '' and code != ''>
)

问题根源说明

你之前的聚合结果偏小,是因为crossSectionalEngine在数据分批到达时,会对同一时间戳多次触发计算(第一次可能是部分数据的聚合,后续是更完整数据的聚合),但默认的KeyedStreamTable会忽略后续同主键的插入,只保留第一次的不完整结果。开启覆盖更新后,每次新的聚合结果会替换旧值,最终得到的是最新、最准确的计算结果。

内容的提问来源于stack exchange,提问作者Jane

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.11 20:12:37