计算任务高可用
为保障流计算任务在分布式环境中的稳定运行,Orca 通过流图级别的 Checkpoint 机制,在任何节点宕机、网络或存储故障场景下都能快速恢复计算任务,保证数据不丢失。
Checkpoint 机制
Orca 的 Checkpoint 实现基于 Chandy–Lamport 分布式快照算法,引入 Barrier标记流图内部的数据流的一致性边界,从而获得全局一致快照。
Checkpoint 的核心流程如下:
- Checkpoint Coordinator 组件定期触发 Checkpoint 任务,并向所有 source 节点注入 Barrier。
- source 节点在收到 Barrier 后将其原子性地写入流表,并记录当前流表的数据偏移量(offset)。
- 流图中的引擎或流表在收到上游的 Barrier 时,将自身状态做一次快照,然后将 Barrier 向下游传递。
- 当 sink 节点收到 Barrier,标志其上游所有节点均已完成快照。当所有 sink 节点都收到 Barrier,标志此次 Checkpoint 完成。
端到端一致性
一致性语义包含以下两种类型:
- AT_LEAST_ONCE:数据至少处理一次,即不会丢失,但可能会被重复处理。
- EXACTLY_ONCE:数据仅处理一次,既不丢失,也不重复。
Orca 的端到端一致性是指从数据源( source 节点)到数据接收端( sink 节点)的一致性,包括 source 端一致性,计算任务(除 source、sink 外的中间节点)的一致性和 sink 端一致性。如前文所述,Checkpoint 机制已确保了计算任务的一致性。若要实现端到端的一致性,还需实现 source 端和 sink 端的一致性:
- source 端一致性:source 端的一致性依赖于系统在重启后能够从指定位置恢复数据。由于使用了持久化的流表,可通过偏移量(offset)记录数据的回放位置,因此 source 端天然实现一致性。
- 计算任务一致性 :Orca 通过实现 Chandy-Lamport 分布式快照算法,保证计算任务在分布式环境中的全局一致性。
- sink 端一致性:sink 端一致性由流表类型决定,只有当除了 source 以外的所有公共流表都是 keyedStreamTable 或 latestKeyedStreamTable 时,依靠这两种流表的去重功能,将上游的重复数据过滤,才能实现端到端的EXACTLY_ONCE 一致性,否则只能实现端到端的 AT_LEAST_ONCE 一致性。
EXACTLY_ONCE 由于依靠键值的去重功能实现,因此只能保证内存中数据的主键唯一性。当某一主键被持久化后,keyedStreamTable 和 latestKeyedStreamTable 仍然能够接受新插入的相同主键。
这种方案可以确保内存中数据键值的唯一性,而不是全局唯一性。尽管如此,仍足以应对多路写入或网络延迟可能导致的重复提交问题。
Barrier 对齐
对于 EXACTLY_ONCE 一致性,在引擎或流表传递 Barrier 的过程中,如果有多个上游,则需要等收到所有上游的 Barrier 后再做快照。这个过程称作 Barrier 对齐。
Orca 引入了轻量级中间组件 Channel,该组件位于除 source 节点外的所有流任务中,用于实现 Barrier 对齐:
- Channel 位于引擎或流表的每一个上游通路上,接收到 Barrier 时会暂停该通路上的数据传输;
- 待所有 Channel 均收到 Barrier,则推动该流任务中所有下游引擎或流表依次完成快照。
- 完成后将 Barrier 转发至该任务的输出流表中,从而将其向下游传递。
