Cross-sectional Engine中enlist时正常时异常的原因及相关技术疑问
我在使用DolphinDB的Cross-sectional Engine时,需要将多行字段存储为单个数组向量,过程中遇到以下场景问题:
场景1:无聚合函数的标量字段正常执行
当sym和price不使用聚合函数时,引擎可直接执行,返回聚合部分的展开结果。
代码示例:
unsubscribeTable(tableName="trades1", actionName="tradesStats") dropAggregator("csEngineDemo1") share streamTable(10:0,`time`sym`price`volume,[TIMESTAMP,SYMBOL,DOUBLE,INT]) as trades1 share table(1:0, `time`sym`avgPrice`volume`dollarVolume`count, [TIMESTAMP,SYMBOL,DOUBLE,INT,DOUBLE,INT]) as outputTable csEngine1=createCrossSectionalEngine(name="csEngineDemo1", metrics=<[sym, price, avg(volume), sum(price*volume), count(price)]>, dummyTable=trades1, outputTable=outputTable, keyColumn=`sym, triggeringPattern="perRow", useSystemTime=false, timeColumn=`time) subscribeTable(tableName="trades1", actionName="tradesStats", offset=-1, handler=append!{csEngine1}, msgAsTable=true) insert into trades1 values(2020.08.12T09:30:00.000 + 123 234 456 678 890 901, `A`B`A`B`B`A, 10 20 10.1 20.1 20.2 10.2, 20 10 20 30 40 20);
执行结果:
场景2:非聚合标量字段+enlist数组字段执行失败
当sym不使用聚合函数,volume用enlist(输出字段定义为INT[])时,引擎计算失败,报错:Output results with error: Failed to append data to column 'volume'。
代码示例:
unsubscribeTable(tableName="trades1", actionName="tradesStats") dropAggregator("csEngineDemo1") share streamTable(10:0,`time`sym`price`volume,[TIMESTAMP,SYMBOL,DOUBLE,INT]) as trades1 share table(1:0, `time`sym`avgPrice`volume`dollarVolume`count, [TIMESTAMP,SYMBOL,DOUBLE,INT[],DOUBLE,INT]) as outputTable csEngine1=createCrossSectionalEngine(name="csEngineDemo1", metrics=<[sym, avg(price), enlist(volume), sum(price*volume), count(price)]>, dummyTable=trades1, outputTable=outputTable, keyColumn=`sym, triggeringPattern="perRow", useSystemTime=false, timeColumn=`time) subscribeTable(tableName="trades1", actionName="tradesStats", offset=-1, handler=append!{csEngine1}, msgAsTable=true) insert into trades1 values(2020.08.12T09:30:00.000 + 123 234 456 678 890 901, `A`B`A`B`B`A, 10 20 10.1 20.1 20.2 10.2, 20 10 20 30 40 20); getStreamEngineStat()
场景3:聚合标量字段+enlist数组字段执行正常
当sym使用聚合函数(first(sym)),volume用enlist(输出字段定义为INT[])时,引擎计算正常。
代码示例:
unsubscribeTable(tableName="trades1", actionName="tradesStats") dropAggregator("csEngineDemo1") share streamTable(10:0,`time`sym`price`volume,[TIMESTAMP,SYMBOL,DOUBLE,INT]) as trades1 share table(1:0, `time`sym`avgPrice`volume`dollarVolume`count, [TIMESTAMP,SYMBOL,DOUBLE,INT[],DOUBLE,INT]) as outputTable csEngine1=createCrossSectionalEngine(name="csEngineDemo1", metrics=<[first(sym), avg(price), enlist(volume), sum(price*volume), count(price)]>, dummyTable=trades1, outputTable=outputTable, keyColumn=`sym, triggeringPattern="perRow", useSystemTime=false, timeColumn=`time) subscribeTable(tableName="trades1", actionName="tradesStats", offset=-1, handler=append!{csEngine1}, msgAsTable=true) insert into trades1 values(2020.08.12T09:30:00.000 + 123 234 456 678 890 901, `A`B`A`B`B`A, 10 20 10.1 20.1 20.2 10.2, 20 10 20 30 40 20); getStreamEngineStat() select * from outputTable
执行结果:
技术疑问
- 场景2的问题是否源于数组向量不支持像传统聚合标量那样展开为多行?
- 尽管
enlist原则上返回元组,但在此场景下可正常使用,是否预期其写入数组向量? - 尝试用
toArray替代enlist未被支持,是否因Cross-sectional Engine暂不支持该用法?
解答
场景2失败原因:是的。当
metrics中存在未使用聚合函数的字段(如sym)时,Cross-sectional Engine会按该字段的唯一值展开为多行结果;但数组向量类型的字段(enlist(volume))无法被拆分到多行,导致类型不匹配,写入输出表时失败。只有当所有metrics字段都为聚合结果(标量或数组)时,引擎才会按key聚合为单行结果,数组向量才能正常写入。enlist写入数组向量的合理性:是预期行为。
enlist在聚合场景下,会将同一key对应的多行volume值打包为一个元组,而DolphinDB的数组向量(INT[])与元组兼容,引擎会自动将元组转换为数组向量存入输出表,这是设计时支持的用法。toArray不支持的原因:Cross-sectional Engine的
metrics参数仅支持聚合函数或返回标量/可转换为数组的聚合表达式。toArray并非聚合函数,无法在引擎的聚合逻辑中完成多行到数组的转换;而enlist作为特殊的聚合用法,能被引擎识别为将多行值打包为单一组的操作,因此仅支持enlist实现该需求。
内容的提问来源于stack exchange,提问作者xiao feng

