基础概念

本文介绍 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 的函数调用(如 streamTablesubscribeTablecreateReactiveStateEngine 等),并自动封装为可调度的流任务分发至各个节点执行。

DStream API 主要包含以下几类调用:

  • 节点定义类:如 source, buffer, sink, timeSeriesEngine, reactiveStateEngine 等;
  • 节点修饰类:如 setEngineName, parallelize, sync 等;
  • 边操作类:如 mapfork 等。

接口列表见 Orca 声明式 API。用户的思维模式应类似构造链式数据结构:不断向尾部添加或修饰节点,逐步完成整个流图的拼装。

流图

使用 DStream API 构建的计算流程被抽象为一个有向无环图,即流图(Stream Graph),其中每个节点代表一个流表或引擎,每条边表示节点之间的数据传递关系(如级联或订阅)。

流图可分为逻辑流图和物理流图:

  • 逻辑流图是用户调用 DStream API 一步步构建的流图,它代表实际业务逻辑。如图 1-1 所示。
  • 物理流图是用户提交(submit)逻辑流图时,系统经过添加私有流表、应用调度优化、拆分子图与并行任务等过程得到的流图,它定义了实际的流任务。如图 1-2 所示。
1. 图 1-1 逻辑流表
2. 图 1-2 物理流表

为了管理流图的生命周期,Orca 定义了一系列状态和状态转移。分别是:

  • building:任务已调度。当前系统中任务未分发、正在执行级联或构建订阅。
  • running:任务已构建。流计算任务正常运行。
  • error:可恢复错误。例如任务出现资源上的异常,需要重新调度。
  • failed:不可恢复错误。例如任务本身逻辑有误,需要用户介入修改脚本,或是运行时内存溢出。
  • destroying:用户请求销毁流图。当前系统正在销毁。
  • destroyed:流图已销毁。状态机终止。

流表

Orca 支持两类流表,用于数据中转与存储:

  • 私有流表是非持久化的流表,仅供当前流图使用,用于缓存中间结果、并行度匹配、多下游订阅等,不支持直接通过 SQL 语句查询。销毁流图时,会销毁流图中的所有私有流表。
  • 公共流表是持久化流表,可同时被多个流图订阅,用于与外部交互、持久化输出或作为数据源,可直接通过 SQL 语句查询。当无任何流图订阅某公共流表时,该表将自动销毁并释放其全限定名。

子图与流任务

当用户提交逻辑流图时,系统会对逻辑流图进行如下处理:

  1. 检查环:通过拓扑排序检查是否成环,若存在环则报错;
  2. 添加流表:添加私有流表缓存中间结果。比如某引擎存在多个输出或上下游节点并行度不同等情况。因为引擎只能级联输出,不可以被订阅,所以当有多个下游出现时,必须添加流表进行中转。上下游节点并行度不同时,需要中间流表,将上游数据汇总后按照下游并发度重新分发,这个过程称为 shuffle。
  3. 优化:删除多余的新增私有流表。
  4. 分割子图:将整个图拆分成并行度相同的、尽量长的子图(Subgraph)。这一步的目的是尽可能减少需要并行执行的任务数。
  5. 拆分计算任务:在子图内按照并行度拆分出流任务(Stream Task)。流任务是分发执行的单位,每个任务都会被调度到一个线程上运行,即一个流任务内部的引擎和流表构成级联关系,由一个线程按顺序执行。而流任务之间通过订阅传递数据。任务数的总量决定了系统线程的使用上限。Orca 会尽量减少任务数量以提升资源利用率。
  6. 添加 Channel :为了在流任务执行 Checkpoint 时对齐 Barrier(在数据流中插入的标记,详见 Checkpoint 机制),系统需要添加 Channel 。

如图 1-3 所示,红色方框代表子图,蓝色方框代表流任务:

3. 图 1-3 子图与流任务