创建sessionWindowEngine报错‘表不能含重复列名’的排查与解决
创建sessionWindowEngine时服务器报错:A table can't contain duplicate column names,但输入表与输出表均无重复列名,求报错原因及脚本修改方案。
输入表定义
ofpanel_colnames = [ "TimeStamp_1", 'ID_1', 'UpperLimit_1', 'LowerLimit_1', 'PreOpenInt_1', 'OpenPrice_1', 'HighPrice_1', 'LowPrice_1', 'LastPrice_1', 'BidPrice_1', 'BidVol_1', 'AskPrice_1', 'AskVol_1', 'Volume_1', 'Amount_1', 'OpenInt_1', 'ClosePrice_1', 'SettlePrice_1', 'ReceiveTime_1', 'UnderlyingID', "TimeStamp_2", 'ID_2', 'UpperLimit_2', 'LowerLimit_2', 'BidPrice_2', 'BidVol_2', 'AskPrice_2', 'AskVol_2', 'ReceiveTime_2' ] ofpanel_colTypes = [ TIMESTAMP, SYMBOL, DOUBLE, DOUBLE, DOUBLE, DOUBLE, DOUBLE, DOUBLE, DOUBLE, DOUBLE, INT, DOUBLE, INT, INT, DOUBLE, DOUBLE, DOUBLE, DOUBLE, NANOTIMESTAMP, STRING, TIMESTAMP, SYMBOL, DOUBLE, DOUBLE, DOUBLE, INT, DOUBLE, INT, NANOTIMESTAMP ] share streamTable(1:0, ofpanel_colnames, ofpanel_colTypes ) as ofPanel_stream
输出表定义
try{dropStreamTable("SampleStream")}catch(ex){print(ex)} share streamTable( 1:0, ["LogTime","InstrumentID","TimeStamp_1","ID_1","UpperLimit_1","LowerLimit_1","PreOpenInt_1"], [TIMESTAMP,STRING,TIMESTAMP,STRING,DOUBLE, DOUBLE, INT] ) as SampleStream
引擎创建脚本
try{dropStreamEngine(name="DataProcessor")}catch(ex){print(ex)} engine_ = createSessionWindowEngine( name = 'DataProcessor', sessionGap = 30*1000, // 30s response window metrics=<[ last(TimeStamp_1), last(ID_1), last(UpperLimit_1), last(LowerLimit_1), last(PreOpenInt_1) ]>, dummyTable = ofPanel_stream, outputTable = SampleStream, useSystemTime=true, keyColumn = `ID_1, useSessionStartTime=true )
报错原因
key列重复输出
指定keyColumn =ID_1后,引擎会自动将ID_1作为输出结果的第一列;但你在metrics中又添加了last(ID_1),这会在metrics结果中再生成一个ID_1列,最终输出结果中出现两个ID_1`列,触发重复列名错误。列顺序不匹配
当useSessionStartTime=true时,引擎默认会添加名为LogTime的会话开始时间列,紧随key列之后。你的输出表列顺序是LogTime在前、ID_1在后,与引擎实际输出顺序(ID_1→LogTime→ metrics列)完全不符,进一步加剧了列名冲突。未匹配的自定义列
输出表中的InstrumentID列在metrics中没有对应计算逻辑,引擎无法自动填充该列,后续也会引发数据写入错误。
脚本修改方案
步骤1:调整metrics,移除重复的key列计算
删除metrics中的last(ID_1),因为key列会由引擎自动输出。
步骤2:调整输出表列顺序,匹配引擎输出逻辑
引擎输出顺序为:key列(ID_1) → 会话开始时间列(LogTime) → metrics列 → 自定义列(需在metrics中添加对应计算)。这里假设InstrumentID对应输入表的UnderlyingID,在metrics中添加last(UnderlyingID)来填充。
修改后的完整脚本
1. 重新定义输出表
try{dropStreamTable("SampleStream")}catch(ex){print(ex)} share streamTable( 1:0, ["ID_1","LogTime","TimeStamp_1","UpperLimit_1","LowerLimit_1","PreOpenInt_1","InstrumentID"], [STRING,TIMESTAMP,TIMESTAMP,DOUBLE, DOUBLE, INT,STRING] ) as SampleStream
2. 重新创建sessionWindowEngine
try{dropStreamEngine(name="DataProcessor")}catch(ex){print(ex)} engine_ = createSessionWindowEngine( name = 'DataProcessor', sessionGap = 30*1000, // 30s response window metrics=<[ last(TimeStamp_1), last(UpperLimit_1), last(LowerLimit_1), last(PreOpenInt_1), last(UnderlyingID) // 对应输出表的InstrumentID列 ]>, dummyTable = ofPanel_stream, outputTable = SampleStream, useSystemTime=true, keyColumn = `ID_1, useSessionStartTime=true )
可选方案:自定义会话开始时间列名
如果希望保留原输出表的列顺序,可以指定sessionStartTimeColumn参数自定义会话开始时间列名,同时调整metrics和输出表列对应关系:
输出表修改
try{dropStreamTable("SampleStream")}catch(ex){print(ex)} share streamTable( 1:0, ["LogTime","InstrumentID","ID_1","TimeStamp_1","UpperLimit_1","LowerLimit_1","PreOpenInt_1"], [TIMESTAMP,STRING,STRING,TIMESTAMP,DOUBLE, DOUBLE, INT] ) as SampleStream
引擎创建修改
try{dropStreamEngine(name="DataProcessor")}catch(ex){print(ex)} engine_ = createSessionWindowEngine( name = 'DataProcessor', sessionGap = 30*1000, // 30s response window metrics=<[ last(ID_1), last(TimeStamp_1), last(UpperLimit_1), last(LowerLimit_1), last(PreOpenInt_1), last(UnderlyingID) ]>, dummyTable = ofPanel_stream, outputTable = SampleStream, useSystemTime=true, keyColumn = `ID_1, useSessionStartTime=true, sessionStartTimeColumn = `LogTime // 指定会话开始时间列名与输出表匹配 )
内容的提问来源于stack exchange,提问作者user31077149

