Orca 实时计算平台

概述

随着多集群部署逐渐成为常态,企业对流数据产品提出了更高的要求。多集群的流数据访问、计算和运维需求日益复杂,而传统流计算架构难以应对复杂任务之间的依赖关系表达、资源调度与高可用需求,逐渐成为业务发展的瓶颈。

当前流计算方案面临以下核心问题:

  • 编码复杂,出错率高

    用户需要手动推导不同引擎间的表结构,自行编写并行逻辑、级联关系、资源清理等代码,导致编码繁琐且出错风险高。

  • 暴露底层概念,心智负担重

    缺乏抽象层,比如流表的发布依赖 share 关键字,用户还需理解共享会话机制等系统细节,导致使用门槛较高。

  • 部署与运维复杂

    用户需手动指定流任务部署在哪台物理节点上。节点重启后,整套流计算框架需要重新搭建,缺乏自动恢复机制。

  • 运维操作割裂

    查看流表、写入数据、引擎预热等操作必须连接到特定节点执行,不利于统一管理和自动化调度。

Orca 在 DolphinDB 现有流数据产品之上提供了一层抽象与增强,核心目标是:

  • 计算抽象

    提供一套声明式 API,支持链式编程风格,简化流计算框架构建流程。用户无需关心表结构推导、并行度拆分、级联订阅与资源清理。

  • 自动调度

    系统根据构建的流图,结合集群拓扑与资源状况,自动完成任务的部署与执行,无需用户干预。

  • 计算任务高可用

    通过内置 Checkpoint 机制,当节点发生故障或机器宕机时,可自动恢复至最近一次快照状态,确保计算任务不遗漏任何待处理的数据。

系统架构

Orca 采用典型的 Master-Worker 架构,结合 DolphinDB 的分布式文件系统(DFS)实现流图的自动部署、任务调度与容错恢复。系统中各组件职责明确,通过集中式状态管理,实现了高可用的流计算平台。

整体架构

Orca 架构如图 1-1 所示:

  • 用户通过调用 Orca 接口,提交流图定义;
  • Stream Master 接收逻辑流图,并根据拓扑和资源状况生成物理流图和调度计划;
  • Stream Worker 负责实际构建流表与引擎、运行任务、进行 Checkpoint;
  • 所有核心状态(包括流图结构、调度记录、流表位置、Checkpoint 元信息等)均持久化至 DFS 表;
  • 系统内部通过心跳检测、状态上报与 Barrier 机制实现计算图的高可用运行。
1. 图 1-1 Orca 架构

Stream Master

Stream Master 部署在控制节点(Controller)上,其职责如下:

  • 接收用户请求(流图提交、删除、状态查询等);
  • 接收 Stream Worker 的状态上报;
  • 执行流图状态机(构建 → 运行 → 错误恢复 → 销毁);
  • 分配任务到各个节点,管理并行度与资源隔离;
  • 定时触发 Checkpoint,协调流图内所有任务执行快照;
  • 维护元信息,持久化到 DFS 表。

Stream Worker

Stream Worker 部署在数据节点(data node)或计算节点(compute node),其职责如下:

  • 接收 Stream Master 分发的任务,构建流表、引擎、级联与订阅结构;
  • 执行并监控流任务(如聚合计算、状态计算、指标生成);
  • 完成本地状态快照,上传Checkpoint。

Stream Worker 的数据处理完全在内存中实现,除公共流表以外,不做本地持久化。引擎状态通过 Checkpoint 记录。

未来规划

为持续提升 Orca 的功能完备性、灵活性与性能表现,未来我们将围绕以下方向进行优化:

功能增强

跨集群能力:支持跨 DolphinDB 集群的数据流动,包括跨集群订阅、跨集群流表访问以及跨集群 Join,满足多集群业务整合需求。

更丰富的状态控制能力:支持计算任务的灵活管理,例如暂停计算、重置流图状态等,提升运维与调试便利性。

更精细的参数调整:允许用户调整 DStream API 中自动生成的参数,如订阅参数(subscribeTable)、私有流表配置(enableTableShareAndCachePurge)等,实现更细粒度的性能调优。

声明式 API 的扩展:丰富 API 能力,支持用户在定义流表时显式指定分区数,支持在同一流图中使用 sourceByName 等高级用法,提升图构建灵活性。

可选的高可用协议:引入可选的流数据高可用协议,支持在不同容错级别与延迟表现之间灵活权衡,满足多样化业务需求。

可选的调度算法:提供灵活的流图调度策略,支持高吞吐与低延迟之间的自定义权衡,进一步允许接入用户自定义调度算法。

性能优化

物理流图优化:优化表连接性能,提升跨节点表连接的执行效率,减少数据 shuffle。

运行时优化:用级联代替订阅,进一步减小时延,提升运行效率。

Checkpoint 性能优化:针对 Barrier 对齐进行进一步优化,引入增量快照与异步快照机制,降低 Checkpoint 对计算性能的影响,缩短容错恢复时间。