多级流计算链路时延优化最佳实践
在实时计算场景中,端到端时延是衡量流计算系统性能的核心指标之一。随着业务逻辑复杂度提升,流计算任务往往由多个算子和多级处理链路组成,链路层级的叠加可能导致整体时延显著增加。
在实际项目中,用户经常会遇到这样的现象:
- 单个计算环节并不慢,但端到端时延偏高
- 对局部代码进行优化后收益有限
- 性能瓶颈隐藏在复杂链路的某一段自定义逻辑中
本文结合一个典型的多级流计算链路案例,总结一套可复用的优化方法,帮助用户系统性降低流计算链路的端到端时延。
1. 一个典型的多级流计算链路
在实际应用中,流计算任务往往并非单一计算环节即可完成,而是由多个处理阶段串联形成完整链路。例如,数据在接入后可能需要经过实时过滤、窗口聚合、关联计算以及自定义逻辑处理等环节,最终输出到下游系统。随着链路层级增加,端到端延时通常会呈现叠加效应,性能瓶颈也更难通过直觉判断定位。
为了说明多级流计算链路中延时问题的典型来源,以及如何系统性地完成优化,本文将以下面的流计算链路为例展开分析:
该链路包含四个主要阶段:
- 数据接入与快照合成(步骤 ①)订阅上游 Level 2 逐笔合并数据,输入到快照合成引擎生成盘口快照和自定义因子;
- 多级响应式状态引擎计算(步骤 ②、步骤 ④)将快照合成结果依次注入两个响应式状态引擎(rse1、rse2),计算累计窗口、滑动窗口及 zscore 等指标;
- 中间结果处理(步骤 ③)在两个响应式状态引擎中间引入流数据表 st1,完成空值填充和表连接等中间结果处理;
- 结果输出(步骤 ⑤)将最终计算结果写入结果表 stRes,供下游消费。
基于该链路,下面将逐步展示从时延拆解、瓶颈定位到算子与结构优化的完整过程。
2. 多级流计算链路的时延拆解与统计
在多级流计算链路中,端到端时延通常由多个环节共同构成。为了准确定位瓶颈,首先对链路时延进行分段拆解。
在 DolphinDB 中,不同类型的链路组件对应不同的时延统计方式,主要包括以下三类:
2.1 流订阅链路时延
对于流订阅场景,推荐在数据的发布、接收以及计算结果输出等关键节点记录时间戳,例如通过 now(true) 函数标记数据到达各节点的时刻,或使用 setStreamTableTimestamp 函数自动记录写入流表的时间戳。通过对不同节点时间戳作差,可以统计订阅链路各环节的传输与处理时延。
2.2 流计算引擎处理时延
对于响应式状态引擎、时序聚合引擎、横截面引擎等大部分流计算引擎,DolphinDB 提供了 outputElapsedMicroseconds 参数,可在结果表中输出每批数据从注入引擎到计算输出的总耗时。该指标能够直接反映引擎内部计算阶段的延时表现,是定位流计算性能瓶颈的重要依据。
2.3 快照合成引擎的时延
对于快照合成引擎,可通过输出结果中的 UpdateTime1 和 UpdateTime2 衡量快照合成的响应时延。其中,UpdateTime1 是触发窗口关闭的那条输入数据的 receiveTime,UpdateTime2 是触发窗口关闭后计算完成的系统时刻,即输出计算结果时的系统时间。UpdateTime2 - UpdateTime1,即为从逐笔数据写入引擎到引擎完成本次快照盘口合成并输出的整体耗时。
需要特别说明的是,用户自定义因子计算(userDefinedMetrics)的耗时不包含在该指标中。若需统计自定义逻辑的开销,可在自定义函数中返回一列 now(true) 时间戳,并与 UpdateTime2 作差计算,从而得到自定义因子计算耗时。
2.4 时延的分级统计
根据上面不同的统计方式,对第一节的多级流计算链路进行时延的分级统计。
在步骤 ① 中,需要在数据写入快照合成引擎前,记录数据的接收时间 step1_receiveTime:
def step1_handler(mutable msg){
update msg set step1_receiveTime = now(true)
getStreamEngine("orderbookEngine").append!(msg)
}
在快照合成引擎中,将该列指定为输入表的 ReceiveTime 字段,并输出 UpdateTime1 和 UpdateTime2:
// 创建引擎参数 inputColMap,即指定输入表各字段的含义,其中将 step1_receiveTime 指定为输入表的 ReceiveTime
inputColMap = dict(`codeColumn`timeColumn`msgTypeColumn`typeColumn`priceColumn`qtyColumn`sideColumn`buyOrderColumn`sellOrderColumn`seqColumn`receiveTime, `SecurityID`MDTime`SourceType`Type`Price`Qty`BSFlag`BuyNo`SellNo`ApplSeqNum`step1_receiveTime)
// 创建引擎参数 outputColMap,即需要输出的字段名称。time = true 即可输出 UpdateTime1 和 UpdateTime2
outputColMap = genOutputColumnsForOBSnapshotEngine(basic=true, time=true, depth=(10, true), tradeDetail=true, orderDetail=true, withdrawDetail=true, orderBookDetailDepth=0, prevDetail=false, seqDetail=true, residualDetail=true)[0]
在快照合成引擎的 userDefinedMetrics 参数指定的函数中,记录自定义因子计算完成的时间 calcTime:
def userDefinedFunc(t){
... // 自定义因子计算过程
calcTime = now(true)
return (..., calcTime)
}
由此,可以分别统计出快照合成引擎合成并输出盘口的耗时(latency1)和自定义因子计算耗时(latency2):
// 合成并输出盘口的耗时(avg)
latency1 = exec avg((UpdateTime2 - UpdateTime1) \ 1000000) from outputTable
// 自定义因子计算的耗时(avg)
latency2 = exec avg((calcTime - UpdateTime2) \ 1000000) from outputTable
在步骤 ② 和步骤 ④ 中,需要设置响应式状态引擎的 outputElapsedMicroseconds 参数为 true,统计每批数据从注入引擎到计算输出的总耗时:
rse1 = createReactiveStateEngine(name="rse1", metrics=rseMetrics1, dummyTable=dummyOrderOutput, outputTable=st1, keyColumn="code", outputElapsedMicroseconds = true)
rse2 = createReactiveStateEngine(name="rse2", metrics=rseMetrics2, dummyTable=dummySt1Output, outputTable=stRes, keyColumn="code", outputElapsedMicroseconds = true)
由此,可以统计出响应式状态引擎计算和输出的耗时(latency3、latency5):
// rse1 计算和输出的耗时(avg)
latency3 = exec avg(rse1ElapsedTime) \ 1000 from st1
// rse2 计算和输出的耗时(avg)
latency5 = exec avg(rse2ElapsedTime) \ 1000 from stRes
在步骤 ③ 中,需要在订阅 st1 的 handler 里分别记录计算前和计算后的时间戳:
def step4_handler(mutable msg){
beforeCalcTime = now(true)
... // 计算逻辑
afterCalcTime = now(true)
update res set beforeStep4Time = beforeCalcTime, afterStep4Time = afterCalcTime
getStreamEngine("rse2").append!(res)
}
subscribeTable(tableName = "st1", handler = step4_handler, msgAsTable = true)
由此,可以统计出 st1 处理中间结果的耗时(latency4):
// st1 处理中间结果的耗时(avg)
latency4 = exec avg((afterStep4Time - beforeStep4Time) \ 1000000) from stRes
3. 多级流计算链路的性能优化方法
本节介绍一些流计算的性能优化方法,并结合前文的多级流计算链路进行讲解。
3.1 善用内存结构优化计算
在极高频的流数据回调计算中,频繁对下游流表进行 SQL 查询会成为严重的性能瓶颈。对共享字典的访问性能远优于对流表的 SQL 查询。因此,可以用一个共享字典来存储当前状态数据,并在下一次计算时直接查询共享字典,避免全表扫描与复杂查询。
例如,在前文的快照合成引擎中,计算了自定义因子 rqfir,它需要基于当条快照和上一条快照的剩余委托订单量来计算。
原始代码中的计算方法如下:
def userDefinedFunc(t){
...
secID = t.code[0]
if (t.mdTime[0] == 09:30:00.000) {
rqfir = 0.0
}
else {
curResidualBidQty = long(sum(t.ResidualBidQtyList[0]))
curResidualAskQty = long(sum(t.ResidualAskQtyList[0]))
prevResidualBidQty = exec long(sum(ResidualBidQtyList[0])) from st1 where code = secID and mdTime = t.mdTime[0] - 500
prevResidualAskQty = exec long(sum(ResidualAskQtyList[0])) from st1 where code = secID and mdTime = t.mdTime[0] - 500
rqfir = calcRQFIR(prevResidualBidQty, prevResidualAskQty, curResidualBidQty, curResidualAskQty)
}
...
}
每次计算时,都需要去查询上一条快照的剩余委托订单量进行计算,这样的频繁查询非常耗时。对共享字典的访问性能优于对流表的 SQL 查询,因此,可以用一个共享字典来存储当前的剩余委托订单量,并在下一次计算时直接查询共享字典,而无需使用 SQL 对下游流表进行查询:
def userDefinedFunc(t){
...
secID = t.code[0]
curResidualBidQty = long(sum(t.ResidualBidQtyList[0]))
curResidualAskQty = long(sum(t.ResidualAskQtyList[0]))
prevOBCache_dict = objByName(`prevOBCache)
if (t.mdTime[0] == 09:30:00.000){
rqfir = 0.0
}
else{
// 从共享字典读取上一条快照的剩余委托量,避免对流表进行 SQL 查询
prev = prevOBCache_dict[secID]
prevResidualBidQty = prev[`ResidualBidQty]
prevResidualAskQty = prev[`ResidualAskQty]
rqfir = calcRQFIR(prevResidualBidQty, prevResidualAskQty, curResidualBidQty, curResidualAskQty)
}
// 更新共享字典,供下一条快照使用
prevOBCache_dict[secID] = dict(`ResidualBidQty`ResidualAskQty, [curResidualBidQty, curResidualAskQty])
...
}
3.2 合并引擎,避免使用中间表
在处理多级逻辑时,中间表和多级发布订阅会引入额外的内存与时间开销。为了提升性能,需要尽可能把中间表的处理逻辑合并到上游环节中。
判断一段中间处理逻辑能否上移,核心看它是否依赖本级引擎之外的、无法在上游获得的动态状态。如果一段逻辑的输入在上游就已经确定,那么它通常可以被消除或前移。反过来,真正依赖下游多级状态、无法在单级内闭合的逻辑则不宜强行上移,例如需要跨多个引擎、多个 keyColumn 维度反复聚合的计算,强行合并反而会让单级引擎的逻辑变得臃肿、难以维护。合并的目标是消除冗余的中转和查询,而不是把所有逻辑堆到一个引擎里。
例如,在前文的步骤 ② 和步骤 ④ 之间,由一张中间表 st1 来承接中间的计算结果,并且由中间表的发布订阅来串联两个响应式状态引擎,存在一定的内存和耗时的开销。为了提升性能,需要把中间表的处理逻辑合并到上游的环节中,然后将两个响应式状态引擎合并为一个。
中间表的处理逻辑包含两个部分:
- 对 rse1 计算的滑动窗口因子进行空值填充:
def step4_handler(mutable msg){
...
update msg set priceMomentum_m60 = nullFill(priceMomentum_m60, midPrice - first(midPrice)),
priceRange_m60 = nullFill(priceRange_m60, cummaxPrice - cumminPrice),
avgSpread_m60 = nullFill(avgSpread_m60, cumavgSpread),
spreadVol_m60 = nullFill(spreadVol_m60, cumSpreadVol)
...
}
- 和外部表 zscoreTable 进行表关联,以供 rse2 计算 zscore:
def step4_handler(mutable msg){
...
zscoreTable = loadText("./zscoreTable.csv")
res = lj(msg, zscoreTable, `code)
getStreamEngine("rse2").append!(res)
}
为了将这两部分处理逻辑合并到上游,分别做出如下改动:
- 在 rse1 的 m 系列算子中设置 minPeriod,即滑动窗口中最少包含的观测值数据,直接通过计算得到相同结果,避免填充空值:
rseMetrics = sqlCol(orderOutputColNames[1:]).join([
<cummax(midPrice)>,
<cummin(midPrice)>,
<cumavg(asksPrice1 - bidsPrice1)>,
<cumstd(asksPrice1 - bidsPrice1)>,
<midPrice - mfirst(midPrice, 60, 1)>,
<mmax(midPrice, 60, 1) - mmin(midPrice, 60, 1)>,
<mavg(asksPrice1 - bidsPrice1, 60, 1)>,
<mstd(asksPrice1 - bidsPrice1, 60, 1)>
]).join(parseExpr(zscoreMetrics)).flatten()
当 minPeriod = 1 时,窗口未满时直接用已有数据参与计算,结果在业务上等价于原先的空值填充逻辑。
- 提前把 zscoreTable 的内容初始化为一个共享字典,在快照合成引擎的 userDefinedMetrics 参数指定的函数中和其他字段一起返回,避免后续的表连接:
syncDict(STRING, ANY, `zscore_sync_dict)
zscoreTable = loadText("./zscoreTable.csv")
zscore_sync_dict[zscoreTable.colNames()[2:]] = each(first, zscoreTable.values()[2:])
def userDefinedFunc(t){
...
zscore_dict = objByName(`zscore_sync_dict)
...
return (..., zscore_dict.priceMomentum_Mean_1, zscore_dict.priceMomentum_Std_1, ...)
}
随后,把两个响应式状态引擎合并为一个:
rse1 = createReactiveStateEngine(name="rse1", metrics=rseMetrics, dummyTable=dummyOrderOutput, outputTable=stRes, keyColumn="code", outputElapsedMicroseconds = true)
3.3 开启引擎的低延时模式
为了满足用户对高频数据逐条处理的微秒级低延时的需求,DolphinDB 对多个流计算引擎内部的存储结构、函数执行逻辑进行了低延时优化。在使用方面,仅在原有的引擎接口上新增了参数 lowLatency,默认为 false。当设置为 false 时,系统采用传统的流批一体处理模式,支持接入实时数据,但以列式批量处理为主,具备更高的吞吐量。而当设置为 true 时,则启用低延时优化,切换为逐条数据处理模式,响应速度更快,更适合对单条数据进行实时处理的场景。
在使用低延时模式时,需要注意:
- 在低延时模式下,每个引擎只支持一个数据流的写入,不支持多个线程同时向同一个引擎写数据。如果需要同时处理多个数据流,建议为不同的数据流启动多个引擎。
- 在低延时模式下,引擎只支持单条数据写入。如果数据是以微批形式产生,可以先将微批数据拆分成单条,再依次写入引擎。
在前文的例子中,快照合成引擎不建议开启低延时模式,因为上游的 Level 2 逐笔合并数据频率高、数据量大,要拆分为单条写入需要引入额外的复杂处理。响应式状态引擎适合开启低延时模式,因为快照合成引擎每固定频率输出一条对应标的快照,在多标的合成场景中也可以拆分后写入响应式状态引擎。
由于在 3.2 节中,已经把两个响应式状态引擎合并为了一个,只需要在该引擎的参数设置中开启低延时模式即可:
rse1 = createReactiveStateEngine(name="rse1", metrics=rseMetrics, dummyTable=dummyOrderOutput, outputTable=stRes, keyColumn="code", outputElapsedMicroseconds = true, lowLatency = true)
4. 性能测试
4.1 硬件配置
本次测试的服务器配置信息如下:
主机:PowerEdge R750,双路服务器架构(2 Socket,2 NUMA 节点)
CPU:Intel(R) Xeon(R) Silver 4314 CPU @ 2.40GHz 64cores 2线程
内存:16GB*32 DDR4 内存,实际运行频率 2666 MT/s
硬盘:HDD 1.8T*1
网络:万兆以太网
OS:Rocky Linux 9.7 (Blue Onyx)
4.2 测试数据
本文使用的测试数据为 2023.12.12 的上交所 Level2 通道 2 全量逐笔合并数据,约 1500 万条,大小约 1.2G。在快照合成引擎中设置只输出单只标的 600189 的快照,共 28800 条。
4.3 测试环境
本次测试使用的是 DolphinDB 单节点环境,版本为 3.00.4.2。与本次测试相关的关键配置参数如下:
maxMemSize=800
workerNum=16
maxPubConnections=64
subExecutors=16
localSubscriberNum=16
subThrottle=1
4.4 优化前的测试结果
按照 2.4 节的时延分级统计方法,基于附录 streamTest.dos 的性能测试结果如下:
表4-1 优化前的流计算链路性能测试结果
| 阶段 | latency1 | latency2 | latency3 | latency4 | latency5 | 总计 |
|---|---|---|---|---|---|---|
| 时延 / ms | 1.887 | 0.892 | 0.190 | 1.006 | 0.318 | 4.293 |
4.5 优化后的链路结构与测试结果
按照第3小节优化后的多级流计算链路,结构如下图所示:
合并引擎后,中间流表 st1 和 step4_handler 订阅均被移除,latency4(中间结果处理耗时)和 latency5(rse2 计算耗时)不再适用。latency3 改为从 stRes 统计合并后的 rse1 耗时。
latency3 = exec avg(rse1ElapsedTime) \ 1000 from stRes
基于附录 streamTestOptimized.dos 的性能测试结果如下:
表4-2 优化后的流计算链路性能测试结果
| 阶段 | latency1 | latency2 | latency3 | 总计 |
|---|---|---|---|---|
| 时延 / ms | 1.839 | 0.421 | 0.138 | 2.398 |
可以看到:
- latency1 优化前后几乎没有差异。这一步的时延来自于内存密集型的合成盘口输出操作,可以考虑配置单路、高主频、低内存延迟的平台,从硬件层面进行优化;
- latency2 优化后性能提升了1倍左右。这一步的性能提升来自于用读写共享字典代替了频繁查询共享内存表的操作;
- latency3 优化后,相较于优化前的 latency3 + latency4 + latency5 性能提升了10倍以上。该部分避免了中间表的数据处理操作,还开启了响应式状态引擎的低延时模式优化计算性能。
综合来看,不同优化手段的收益差异明显,遇到流计算链路的时延优化问题时,建议按以下优先级依次尝试:
- 优先审视链路结构,消除中间表和多级发布订阅。这类改造收益最大(本文中该环节时延下降超过 10 倍),每减少一级中间表,就省去一次序列化、发布、订阅和内存拷贝的开销。
- 其次优化高频的数据访问方式,用共享字典替代对流表的 SQL 查询和与静态表的关联,收益次之。
- 再对适合的引擎开启低延时模式,进一步压低逐条数据的处理时延。
- 最后才考虑硬件层面的优化,用于应对已无代码优化空间的内存密集型环节(如本文的 latency1)。
5. 总结
本文拆解了一个典型的多级流计算链路,说明了各个环节的时延统计方法,并介绍了一些性能优化的技巧。用户可以基于本文,对自己的流计算链路进行性能统计和优化,以实现更低延迟的实时流计算效果。
6. 附录
逐笔合并数据:testData.zip
scodeTable:zscoreTable.csv
优化前的多级流计算链路脚本:streamTest.dos
优化后的多级流计算链路脚本:streamTestOptimized.dos
