Orca 实时计算平台技术原理
1. 为什么需要 Orca
DolphinDB 已经具备比较完整的流计算能力:流表负责接收数据,发布订阅机制负责传递数据,时间序列引擎(time-series engine)、响应式状态引擎(reactive state engine)、连接引擎(join engine)等负责计算,持久化和高可用(HA)机制负责可靠性。业务比较简单时,开发者直接创建这些对象即可完成任务。问题出现在业务变复杂之后。
以实时行情处理为例,逐笔成交进入系统后,可能要同时生成 1 分钟、5 分钟、日内多周期 K 线;每条 K 线后面还会接移动窗口指标、规则判断、风险监控、落库和下游订阅。从业务角度看,这是一条清晰的数据处理链路;落到脚本里,往往会变成多张输入表、中间表、输出表,多个引擎,多个订阅关系,再加上一组节点部署和重启恢复逻辑。
更难维护的是这些对象之间的关系。哪些表和引擎属于同一条业务链路,哪些对象必须先创建,停止任务时应该先停输入还是先删中间对象,某个节点故障后订阅位点(offset,即从哪条记录接着消费)和引擎状态从哪里恢复,这些问题都需要仔细考虑清楚,才能真正部署到生产环境。
如果这些关系全部写在用户脚本里,系统规模越大,越容易出现这样的情况:计算逻辑本身没有错误,但部署顺序、恢复顺序或对象清理出了问题。轻则排查困难,重则出现重复消费、部分链路仍在写入、下游看到不一致结果等事故。
Orca 解决的就是这类问题。用户只需要提交一张流图(Stream Graph,即一条实时计算任务的完整描述),说明数据从哪里来、经过哪些计算、最后写到哪里。Orca 根据这张图生成实际运行所需的流表、引擎、订阅、任务、节点部署和恢复计划。这样,开发者只需要维护业务数据流,而由 Orca 负责维护分布式执行过程中容易出错的那部分细节。
2. 流图是什么
流图是 Orca 对一条实时计算任务的完整描述。它把输入表、计算引擎、中间结果、输出表和它们之间的数据流关系放在同一张图里管理。
编写流图的时候,用户主要使用以下几类接口:
- 输入:
source(以及keyedSource、haSource等变体),对应一张可被外部写入的公共流表。 - 计算:
timeSeriesEngine、reactiveStateEngine等引擎接口。底层仍然是 DolphinDB 已有的流引擎。 - 中间可见结果:
buffer,把某一步的中间结果落成公共流表。(默认情况下,引擎之间的中间结果由 Orca 内部对象承接,外部不可见;使用buffer显式声明后,这一步的结果才会成为可以被外部查询和订阅的公共流表)后面也可以继续接计算 - 输出:
sink,把当前这一跳接到公共流表、分布式存储 (DFS) 表或函数。 - 分支 / 并行:
fork把同一份数据复制成多路;parallelize按指定列拆分并行计算,再用sync汇合。
编写流图时,开发者关心的是业务链路:成交数据进入系统,按股票代码并行聚合,生成 K 线,再计算移动指标,最后写入输出表。Orca 接收这张图之后,可以在提交前检查结构,在提交时拆分任务,在运行时维护状态,在故障后按同一份元数据进行恢复。
把流计算链路抽象成流图后,相较于普通的脚本封装,Orca
可以做更多的事情。提交阶段检查表结构(schema,即列名和类型)是否衔接、source 和
sink
的使用是否合法;调度阶段可以看到整张图的并行关系和节点约束;停止或销毁图时,可以按照依赖关系逐步关闭订阅、表和引擎;出现故障时,也可以从元数据和检查点(checkpoint,即整张图的统一恢复点)信息恢复出原来的运行结构。脚本封装只是把创建表、创建引擎、创建订阅这些动作隐藏起来;流图则把这些动作背后的业务关系保存下来,让系统后续还能继续调度、恢复和排查。
从实现上看,Orca 位于流计算子系统(Streaming)之上。数据处理仍然由 DolphinDB 已有的流引擎完成;Orca 增加的是围绕流图的管理能力,包括建图、调度、分发、权限、恢复和检查点。这样既可以复用已经成熟的流计算能力,又把分布式编排放到单独的一层处理。
运行时可以分成两个部分看。控制台(Stream Master),负责保存元数据、生成物理图(Physical Graph)、调度流任务(Stream Task,即真正放到节点上执行的最小单位,第 4.3 节详述)、向节点下发构建动作,并协调检查点,本身不参与数据计算。工作节点(Worker)负责真正创建流表、引擎、订阅,以及任务之间传递数据用的数据通道(Orca Channel)。Orca Channel 是系统自动插入的内部对象,用户脚本中不会出现它:当两个任务不在同一段连续计算里、数据要从一个任务流向另一个任务时,系统会在中间插入这个通道,既传递普通数据,也携带检查点的进度标记。数据主要在工作节点之间流动;控制台不参与数据计算,负责的是整体的调度和权限控制。
3. 逻辑与执行分离
开发者编写流图时,描述的是业务逻辑。下面的例子表达的是:输入成交数据,生成 K 线,再计算几个移动指标,最后写入输出表。每一步用到的函数标注在注释中。
if (!existsCatalog("Orca")) {
createCatalog("Orca") // catalog:给流图、表、引擎一个命名空间
}
go
use catalog orca
// createStreamGraph:创建一张名为 graph1 的流图
g = createStreamGraph(`graph1)
// source:定义输入表 Trade;表已存在时会校验 schema
// timeSeriesEngine:按 1 分钟窗口聚合 K 线
// reactiveStateEngine:在 K 线上计算滚动指标
// sink:把结果写入输出表 output
g.source("Trade", `symbol`datetime`price`volume, [SYMBOL, TIMESTAMP, DOUBLE, INT])
.timeSeriesEngine(
60*1000, 60*1000,
<[first(price), max(price), min(price), last(price), sum(volume)]>,
"datetime", false, "symbol")
.reactiveStateEngine(
<[datetime, first_price, max_price, min_price, last_price, sum_volume,
mmax(max_price, 5), mavg(sum_volume, 5)]>,
`symbol)
.sink("output")
// submit:把逻辑图交给 Orca,去集群里部署和启动
g.submit()
这段脚本没有说明中间是否需要创建私有流表,也没有说明跨节点时如何订阅、并行任务如何编号、source 订阅应该何时启动。对开发者来说,这些细节不应该混在业务逻辑里;但在生产环境中,它们又必须被处理清楚。
调用 submit() 之后,Orca
会把用户编写的逻辑图转换成物理图。逻辑图回答「数据应该怎样计算」,物理图回答「这条数据链路在集群里怎样运行」。转换顺序大致是:先检查图是否合法(例如是否有环、引擎后面是否有输出),再补充必要的中间表、删除多余的中间表,把图切分成若干子图(subgraph,即整张流图中切出来的一段,详见
4.3 节),再按并行度拆分成流任务,最后在需要跨子图传递数据的位置加上 Orca Channel。
这种分离与数据库中的 SQL 类似。用户编写 SQL 时只描述要查询什么,优化器负责生成执行计划;在 Orca 中,用户链式写出的数据流(DStream)负责描述业务数据流,物理图负责生成可部署、可调度、可恢复的流计算执行计划。
这样分开以后,业务表达可以保持稳定,执行方式也可以持续演进。以后并行度、节点分布、检查点策略或调度策略发生变化时,不需要把这些细节反向暴露给业务脚本。
4. 提交时完成的转换
一张图真正运行之前,Orca 需要先把“业务逻辑”翻译成“集群里的运行对象”。这一步需要回答以下几个问题:
-
每个
source、引擎、sink最终对应哪些运行对象。 -
哪些对象可以放在同一个节点内一起执行,哪些需要拆分到不同节点或任务中。
-
哪些位置需要中间流表承接数据。
-
哪些跨节点的数据流需要 Orca Channel。
-
哪些对象先创建,哪些对象后创建。
这些问题如果留给用户手写脚本处理,不仅琐碎,而且出错概率很高。Orca 在提交阶段统一完成这部分转换。
4.1 预先检查
提交流图时,Orca 会先检查流图结构,比如表结构是否衔接,source 和 sink
的使用是否符合规则,整张图是否完整。这类错误需要在调度之前就被彻底检查和发现。
如果等到远程节点已经创建了一部分表和引擎之后才失败,问题会更难排查:有的对象已经存在,有的订阅尚未建立,有的输入可能已经开始写入。严重时,后续的数据清理和问题排查都会消耗大量时间。
4.2 私有流表:连接上下游的中间表
用户编写流图时,经常会直接写成 engine -> engine 或 engine -> sink,但真实运行时,上下游之间不一定能简单地直接相连。比如一个引擎后面接两个下游分支。逻辑上只是一个输出分给两个分支;运行时两个分支可能在不同流任务上,也可能在不同节点上,各自的订阅批次、过滤条件和检查点进度都不一样。用户侧对应的写法是:先写出该引擎,再 fork 成两路。
use catalog orca
g = createStreamGraph("engine_fork_demo")
// 同一个 timeSeriesEngine 后面接两条下游
eng = g.source(
"Trade",
`symbol`datetime`price`volume,
[SYMBOL, TIMESTAMP, DOUBLE, INT]
).timeSeriesEngine(
60*1000, 60*1000,
<[first(price), max(price), min(price), last(price), sum(volume)]>,
"datetime", false, "symbol")
// fork(2):把当前引擎的输出复制成两路,count 必须大于 1
branches = eng.fork(2)
branches[0]
.sink("output_kline") // 一路直接落 K 线
branches[1]
.reactiveStateEngine( // 另一路继续算滚动指标
<[datetime, first_price, max_price, min_price, last_price, sum_volume,
mmax(max_price, 5), mavg(sum_volume, 5)]>,
`symbol)
.sink("output_indicator")
g.submit()
再比如上游只有一路,下游按 symbol 拆成 4 路并行计算;或者上游 4 路计算完成后,下游需要通过 sync() 汇合。这些场景都需要一个稳定的中间位置来承接数据。
私有流表就是为解决这类情况而引入的内部流表。用户通常不直接读写它;它属于 Orca 生成的中间对象,用来连接上下游,名字一般以
private_stream_table_
开头。实际上,sink、fork、map、parallelize
如果接在引擎后面,逻辑图阶段就会先插入这样一张表,并不一定要等到物理图转换。
有了私有流表,上下游之间就多了一个明确的中间承接点。下游订阅创建失败时,系统可以只处理这条订阅和对应的中间表,而不必把上游引擎一起回滚;上下游并行度不同时,可以通过中间表和订阅关系完成数据重分布;检查点传播到跨任务链路时,也有明确的位置记录数据和控制消息的流向。
当然,并非每一条逻辑边都需要保留私有流表。当前实现会先按较保守的方式生成中间对象,保证复杂场景下不缺少必要的连接点,随后在优化(optimize)阶段删除不需要的私有表。为 sink / fork 插入的那张表通常会保留,因为它后面还要承担订阅。
4.3 流任务:最小执行单位
流图提交后,Orca 不会直接按单个引擎调度,也不会把整张流图当成一个任务调度。
如果按单个引擎调度,任务数量会膨胀,很多本来可以在同一节点内连续执行的计算也会被拆散,订阅和网络开销都会增加。底层流计算已经支持本地级联执行,Orca 应当利用这种能力。
如果整张图只作为一个任务,又会走向另一个极端:不同分支不能分布到不同节点,parallelize 的并行度无法发挥;公共流表和计算引擎也必须跟随同一个任务放置,调度器很难分别满足表位置和计算资源这两类需求。
流任务是 Orca 选定的执行单位。物理图先被划分成子图,再根据并行度拆分成多个流任务。一个流任务内部尽量保留本地连续执行;流任务之间通过流表、订阅和 Orca Channel 连接。
子图不是用户需要编写的对象。用户编写的是整张流图;提交之后,系统不会把这张图当成一个整体去调度,而是先把图切分成若干段,每一段就是一个子图。
parallelize("symbol", 4) 的含义也在这里:这段流计算可以按 symbol 拆成 4 份独立任务。后续调度器再根据资源和约束,决定这些流任务分别放置到哪些节点上。
4.4 Orca Channel:任务之间的数据通道
当数据从一个流任务流向另一个流任务时,这条连接除了传递数据,还要保证上下游的处理进度能够对齐。因为 Orca 需要在故障后恢复整张流图,而恢复不能只看某个任务自身的状态:上游已经处理到哪一批数据,下游已经接收到哪一批数据,中间链路里是否还有未处理的数据,这些信息都需要有办法记录。
检查点屏障(checkpoint barrier)可以理解为数据流里的一个进度标记。它和普通数据走同一条路径向下游传播。下游看到屏障时,就知道在它之前的数据已经进入本次检查点,在它之后的数据属于后续处理。这样,多个任务才能围绕同一个检查点保存状态。
Orca Channel 就位于这种跨任务的链路上。它负责把普通数据送到下游,也负责识别和转发屏障;如果一条链路有多个输入,它还需要在必要时等待其他输入也到达同一个检查点,再继续向后推进。这样恢复时,系统得到的是一组进度一致的状态,而不是各个任务各自保存的零散状态。
5. 调度:把任务放到合适节点上
调度要解决的问题是:一张图被拆成多个执行任务以后,每个任务应该放到哪个节点上。这个决定不能只看哪个节点空闲,还要看任务本身的要求。
有些要求来自表的位置。公共流表通常应该放在数据节点(data node)上,因为它需要被外部读写,也会作为元数据长期存在;如果任务依赖某张已经存在的公共流表,调度器会尽量让任务靠近这张表,减少不必要的跨节点访问。
有些要求来自计算资源。计算任务可以放在数据节点或计算节点(compute node)上;如果用户指定了计算组(compute group),就要遵守计算组的隔离约束,不能随意放到其他计算资源组里。
用户自定义函数(UDF)中的共享变量也是需要特别处理的约束。UDF 可以读写共享表(shared table)、共享字典(shared dict)、共享键值表(shared keyed table)这类共享状态。如果访问同一份共享状态的任务分散到多个节点,就会引入跨节点一致性、锁、复制和恢复问题。当前 Orca 在调度阶段把这些任务放到同一个约束组里,尽量让它们落到同一节点。这里主要考虑的是执行语义清晰,而不只是优化性能。
在满足这些约束之后,调度器会考察节点当前负载、剩余内存、数据节点磁盘空间、网络 IO 和优先节点(preferred node)权重。当前实现采用启发式评分,不追求全局最优解。对在线提交的流图来说,调度结果需要足够快、足够稳定,也要便于解释;严格最优调度需要对未来流量、状态大小、网络成本和故障概率建模,代价很高,实际收益也不一定稳定。
调度实现中还会根据任务复杂度(task complexity)扣减节点得分。一个节点被分配任务后,后续得分会下降,避免所有任务都集中到初始得分最高的节点上。
6. 部署和启动:按顺序创建运行对象
调度结束后,系统只是确定了每个任务应该放在哪个节点,流图还没有真正运行。接下来要做的是,把这些任务对应的流表、引擎和订阅关系创建到目标节点上。
Orca 不会一次性启动所有对象,而是分几步执行:
-
先创建本地流表和引擎,把每个任务内部的计算链路准备好。
-
再创建非
source订阅,让中间和下游链路先连接起来。 -
最后创建
source订阅,让输入数据进入整张图。
这个顺序主要是为了避免半启动状态。实时计算链路一旦打开入口,数据就会持续进入。如果 source 订阅先启动,而下游表或引擎尚未建好,系统就可能出现一部分数据已经进入、一部分链路尚未接通的情况。后续排查时,用户看到的就不是一次清晰的启动失败,而是一条状态不完整的链路。
销毁时同样要注意顺序。系统会先停止订阅,再删除表和引擎,避免对象已经删除、上游还在继续推送数据的情况。
这些创建和销毁动作通过延迟执行接口(Lazy API)表达。控制台不需要直接操作远程节点上的 C++ 对象,它只生成一组「该节点要执行哪些动作」的描述,再交给对应的工作节点执行。工作节点执行完成后,把任务状态和错误原因返回给控制台。这样无论是提交、恢复还是销毁,系统都可以复用同一套执行方式;失败时,也能明确问题出在哪个阶段、哪个任务上。
7. 生命周期:提交、停止和销毁
流图从提交到销毁,会经历多个阶段:创建元数据、部署任务、进入运行状态、停止、重新启动、重新提交,最后销毁。每一步都会影响元数据,也会影响多个节点上的流表、引擎和订阅关系,所以不能只当成一次普通的函数调用。
以提交为例,系统要保存图的序列化数据(graph blob),更新图的元数据,创建或更新公共流表元数据,记录每个任务被放置到哪个节点,以及任务当前运行成功还是失败。如果开启了检查点,还要把检查点配置与这张图关联起来。只要其中一步失败,系统就要知道失败发生在哪个阶段。
这些状态如果记录不清楚,就容易出现半完成状态:比如表已经创建,但元数据尚未更新;某些任务已经运行,但图仍然显示在构建(building)状态;销毁图时订阅没有停止干净,后续仍有数据写入已删除的对象。
控制台会把这些步骤串联起来,并记录图和任务当前所处的阶段。这样用户调用 startStreamGraph、stopStreamGraph、resubmitStreamGraph 这类函数时,不需要自行判断哪些订阅要先停止、哪些表和引擎可以保留、失败后应该从哪一步重新处理,系统会根据当前状态完成对应的启动、停止或重新提交。
8. Checkpoint 和故障恢复
故障恢复时,系统需要保证各个任务回到同一个处理进度。比如上游已经处理到第 10000 条数据,下游只处理到第 9800 条,中间通道里还留有一批未处理完的数据。如果每个节点各自保存状态,恢复后可能漏数据,也可能重复处理。
检查点要解决的就是这个问题:为整张流图确定一个统一的恢复点。Orca 使用屏障(barrier)做这个标记。负责检查点的组件会从
source 端插入屏障;屏障和普通数据走同一条链路向下游传播。
当某个任务收到屏障时,说明在这条输入上,屏障之前的数据已经处理到这里。任务此时会保存自身状态,并返回保存结果。等相关任务都完成状态保存并返回结果后,这次检查点才算成功。
如果一个任务有多个上游输入,情况会更复杂一些。它不能只看到其中一路屏障,就认为自己可以保存状态,因为另一路可能还停留在更早的位置。遇到这种情况,Orca Channel 会让先到的屏障暂时等待,等其他输入也收到同一次检查点的屏障后,再让任务保存状态并继续处理。这样恢复时,各路输入对应的是同一批处理进度。
恢复时,Orca 可以基于最近一次成功的检查点重建图内状态和处理进度。因为恢复点已经由系统统一记录,用户不需要逐个检查引擎、订阅和中间通道的进度,也不需要手工拼接各个节点上零散的状态。
9. 查询和管理
实时计算任务通常会长期运行。提交成功以后,用户还需要持续查看它的状态:图是否仍在运行,输出表由哪些计算产生,节点失败后影响了哪些任务。需要停止或重启流图时,也应该有统一的入口,而不是让用户自行处理每个表、引擎和订阅。
如果没有统一的元数据,排查时就要回头翻脚本、查订阅配置、看各个节点的状态和日志。Orca 的做法是把图、表、引擎都登记到目录(catalog)中,让它们成为可以查询和管理的对象。这样用户看到的不只是某个会话里创建过的变量,而是一组有名字、有元数据、有生命周期的实时计算资源。
有了这些元数据,排查可以从流图开始。查图,可以看到整张图处于运行(running)、失败(failed)还是停止(stopped)状态;查表,可以看到某张 Orca 流表在哪个节点、被哪些图引用;查检查点,可以看到恢复过程卡在哪个任务上;查数据血缘,可以知道一张输出表来自哪张图、经过了哪些计算。
权限也可以按对象控制。计算组主要限制任务能使用哪些计算资源;Orca 的图、表、引擎权限则控制谁可以创建图、停止图、读写 Orca 流表、管理引擎。这样既能限制任务运行在哪些资源上,也能限制用户能操作哪些实时计算对象。
10. 以多周期 K 线为例
下面用多周期 K 线的例子把前面的过程串起来。业务目标是:输入逐笔成交 Trade,同时生成 1 分钟和 5 分钟 K
线,再计算移动指标,最终写入两个 Orca 输出表。
if (!existsCatalog("orca")) {
createCatalog("orca")
}
go
use catalog orca
g = createStreamGraph(`kline_graph)
sourceStreams = g.source(
"Trade",
`symbol`datetime`price`volume,
[SYMBOL, TIMESTAMP, DOUBLE, INT]
).fork(2)
stream_1min = sourceStreams[0]
.parallelize("symbol", 4)
.timeSeriesEngine(
60*1000, 60*1000,
<[first(price), max(price), min(price), last(price), sum(volume)]>,
"datetime", false, "symbol")
.reactiveStateEngine(
<[datetime, first_price, max_price, min_price, last_price, sum_volume,
mmax(max_price, 5), mavg(sum_volume, 5)]>,
`symbol)
.sync()
.sink("output_1min")
stream_5min = sourceStreams[1]
.parallelize("symbol", 4)
.timeSeriesEngine(
5*60*1000, 5*60*1000,
<[first(price), max(price), min(price), last(price), sum(volume)]>,
"datetime", false, "symbol")
.reactiveStateEngine(
<[datetime, first_price, max_price, min_price, last_price, sum_volume,
mmax(max_price, 5), mavg(sum_volume, 5)]>,
`symbol)
.sync()
.sink("output_5min")
g.submit()
这段脚本在提交前,只是在描述业务链路:Trade 作为输入,数据被分成两条分支,分别计算 1 分钟和 5 分钟 K 线,再写入两个输出表。此时尚未涉及中间表、节点选择、订阅顺序和故障恢复。
调用 submit() 后,Orca
开始把这条业务链路转换成可以在集群中运行的结构。系统会先检查图结构和表结构,再生成物理图;需要中间承接的位置会插入私有流表,需要跨任务传递数据的位置会插入 Orca
Channel;按 symbol 并行的部分会被拆分成多个可调度的流任务。
随后系统会保存图和表的元数据,让这条链路可以被查询、停止、重启和恢复。调度器再根据表的位置、节点资源、计算组和共享状态等条件,为每个任务选择运行节点。
部署时,Orca 会按顺序创建目标节点上的流表、引擎和订阅关系:先把中间和下游链路准备好,最后再打开 source
订阅。这样数据开始写入 Trade 时,后面的 1 分钟和 5 分钟分支已经可以承接数据并持续计算,结果最终进入
output_1min 和 output_5min。
如果开启检查点,屏障会和数据一起从 source 向下游传播;跨任务的链路由 Orca Channel
处理,保证各个任务保存的是同一批处理进度。后续恢复时,Orca 可以基于最近一次成功的检查点恢复图内状态。
这个例子可以概括 Orca 的分工:开发者编写的是业务逻辑,Orca 负责把它变成实际运行所需的表、引擎、订阅、任务、通道、元数据和恢复逻辑。
