基础概念
本文介绍 Orca 的基础概念,包括声明式 API 的编程模型、流图的类别与生命周期、流表的作用划分、子图与流任务的拆解等。
全限定名
DolphinDB 支持数据目录(catalog),用于统一组织各类数据库对象,包括对分布式表的支持。现在,Orca 进一步扩展了 catalog 的应用范围,将流图、流表和流引擎也注册进 catalog 中,使用户能够通过一种统一的方式访问各类元素。
为了唯一标识 catalog 中的对象,当前版本采用了 <catalog>.<schema>.<name>
的三段式命名方式,称其为全限定名(Fully Qualified Name,FQN)。
在为对象命名时,用户只需指定 name,系统将自动补全 FQN:
- catalog 由当前 catalog 决定
- schema 根据对象是流图、流表、引擎分别命名为 orca_graph, orca_table, orca_engine。
声明式 API
Orca 提供了一套声明式 API(Declarative Stream API,DStream API),用于简洁地构建实时流计算图。该 API 屏蔽了底层并行调度、订阅关系、资源清理等复杂逻辑,让用户聚焦于业务本身的表达。
用户无需手动推导表结构,编写并行引擎构建逻辑,显式订阅数据,编写任务清理与销毁逻辑。系统会将用户定义的 DStream API 转换为 DolphinDB 的函数调用(如
streamTable、subscribeTable、createReactiveStateEngine
等),并自动封装为可调度的流任务分发至各个节点执行。
DStream API 主要包含以下几类调用:
- 节点定义类:如
source,buffer,sink,timeSeriesEngine,reactiveStateEngine等; - 节点修饰类:如
setEngineName,parallelize,sync等; - 边操作类:如
map、fork等。
接口列表见 Orca 声明式 API。用户的思维模式应类似构造链式数据结构:不断向尾部添加或修饰节点,逐步完成整个流图的拼装。
流图
使用 DStream API 构建的计算流程被抽象为一个有向无环图,即流图(Stream Graph),其中每个节点代表一个流表或引擎,每条边表示节点之间的数据传递关系(如级联或订阅)。
流图可分为逻辑流图和物理流图:
- 逻辑流图是用户调用 DStream API 一步步构建的流图,它代表实际业务逻辑。如图 1-1 所示。
- 物理流图是用户提交(
submit)逻辑流图时,系统经过添加私有流表、应用调度优化、拆分子图与并行任务等过程得到的流图,它定义了实际的流任务。如图 1-2 所示。
为了管理流图的生命周期,Orca 定义了一系列状态和状态转移。分别是:
- building:任务已调度。当前系统中任务未分发、正在执行级联或构建订阅。
- running:任务已构建。流计算任务正常运行。
- error:可恢复错误。例如任务出现资源上的异常,需要重新调度。
- failed:不可恢复错误。例如任务本身逻辑有误,需要用户介入修改脚本,或是运行时内存溢出。
- destroying:用户请求销毁流图。当前系统正在销毁。
- destroyed:流图已销毁。状态机终止。
流表
Orca 支持两类流表,用于数据中转与存储:
- 私有流表是非持久化的流表,仅供当前流图使用,用于缓存中间结果、并行度匹配、多下游订阅等,不支持直接通过 SQL 语句查询。销毁流图时,会销毁流图中的所有私有流表。
- 公共流表是持久化流表,可同时被多个流图订阅,用于与外部交互、持久化输出或作为数据源,可直接通过 SQL 语句查询。当无任何流图订阅某公共流表时,该表将自动销毁并释放其全限定名。
子图与流任务
当用户提交逻辑流图时,系统会对逻辑流图进行如下处理:
- 检查环:通过拓扑排序检查是否成环,若存在环则报错;
- 添加流表:添加私有流表缓存中间结果。比如某引擎存在多个输出或上下游节点并行度不同等情况。因为引擎只能级联输出,不可以被订阅,所以当有多个下游出现时,必须添加流表进行中转。上下游节点并行度不同时,需要中间流表,将上游数据汇总后按照下游并发度重新分发,这个过程称为 shuffle。
- 优化:删除多余的新增私有流表。
- 分割子图:将整个图拆分成并行度相同的、尽量长的子图(Subgraph)。这一步的目的是尽可能减少需要并行执行的任务数。
- 拆分计算任务:在子图内按照并行度拆分出流任务(Stream Task)。流任务是分发执行的单位,每个任务都会被调度到一个线程上运行,即一个流任务内部的引擎和流表构成级联关系,由一个线程按顺序执行。而流任务之间通过订阅传递数据。任务数的总量决定了系统线程的使用上限。Orca 会尽量减少任务数量以提升资源利用率。
- 添加 Channel :为了在流任务执行 Checkpoint 时对齐 Barrier(在数据流中插入的标记,详见 Checkpoint 机制),系统需要添加 Channel 。
如图 1-3 所示,红色方框代表子图,蓝色方框代表流任务:
