关于subscribeTable中batchSize的理解及流处理逻辑的疑问
针对DolphinDB流处理疑问的解答
问题1:自定义函数的输入是否始终是长度为1的向量?
不是。自定义函数的输入是当前处理批次的整列向量,向量长度等于批次内的数据行数。你看到输入长度始终为1,是因为当前流数据是逐条到达的,每批只有1条数据,所以向量长度固定为1。
问题2:是否可以确认类的append方法按行处理,而自定义函数按批处理且批大小恒为1?
不能。在ReactiveStateEngine中,自定义函数和类的append方法的输入逻辑完全一致,都是接收当前批次的列向量(而非单条行数据)。你观察到结果一致,是因为当前每批只有1条数据,两种方法处理的都是长度为1的向量,输出自然没有差异。
问题3:是否是subscribeTable配置错误导致每批仅处理一行?如何修正?
核心问题不在subscribeTable,而在于replay的参数配置:
你当前的replay未指定replayRate,默认会按照原始数据的时间间隔(每条间隔1秒)逐条发送数据,导致subscribeTable无法攒够批量。即使设置了batchSize,也因为数据到达间隔太长,每批只能处理1条。
修正方案:
- 修改
replay参数,添加replayRate并设置为极大值(比如1000000),让数据批量快速发送,而非按真实时间间隔回放; - 调整
subscribeTable的batchSize和throttle参数,控制批量大小和等待时间; - 保持
msgAsTable=true(已设置),让subscribeTable以表的形式传递批量数据给引擎。
修正后的关键代码片段:
// 调整subscribeTable,启用批量处理参数 subscribeTable( tableName="input_stream", actionName="tes", handler=getStreamEngine("tes"), batchSize=200000, // 设置期望的批处理大小 throttle=0.1, // 设置等待超时时间(秒),超时后即使未达batchSize也处理 reconnect=true, msgAsTable=true ) // 修改replay,添加replayRate实现批量发送 timing = now() replay(inputTables=t, outputTables = input_stream, timeColumn=`datetime, replayRate=1000000)
验证方式:
修改后,在自定义函数中加入print(size(volume_)),就能看到向量长度等于你设置的batchSize(或最后一批的剩余行数),此时取volume_[0]会得到批次第一条数据的volume值,不再是严格递增的连续序列。
内容的提问来源于stack exchange,提问作者user32013486
相关产品推荐
相关产品推荐

