Skip to main content
Version: 5.1.0

基于 Akka 的分布式批量数据同步架构

概述

TIS 5.1 版本在批量数据同步领域实现了重大架构升级,采用 Akka 分布式框架重新设计了整个批量同步任务的执行引擎。这次架构升级不仅解决了 DataX 原生框架不支持分布式任务并发执行的痛点,更通过 Akka Actor 模型带来了更高的稳定性、可扩展性和容错能力。

本文将从架构设计、核心特性、技术选型等多个维度,帮助用户全面了解 TIS 5.1 批量同步方案的原理和优势。

背景:为什么需要架构重构?

传统架构的局限性

在 TIS 5.0 及之前的版本中,批量数据同步任务调度基于 Jenkins 的 task-reactor 库实现:

  • 单机执行模型:所有任务在单个 JVM 进程内执行,无法横向扩展
  • 资源竞争:大量并发任务导致内存和 CPU 资源争抢
  • 容错能力弱:节点故障会导致所有运行中的任务失败
  • 扩展性差:难以应对突发的任务高峰

DataX 的优势与不足

阿里巴巴开源的 DataX 是数据集成领域非常优秀的端到端同步框架:

优势

  • 丰富的数据源插件生态(支持 MySQL、Oracle、PostgreSQL、MongoDB 等数十种数据源)
  • 成熟稳定的数据传输引擎
  • 灵活的数据转换能力

不足

  • 原生 DataX 工具不支持分布式任务并发执行
  • 无法实现任务级别的负载均衡
  • 集群化部署需要额外的调度框架支持

TIS 5.1 的架构重构正是为了在保留 DataX 优秀特性的基础上,通过 Akka 框架弥补其在分布式执行方面的不足。

核心架构设计

整体架构图

TIS 5.1 整体架构图

三层架构设计

TIS 5.1 采用经典的三层架构,职责清晰,层次分明:

1. DAG 算法层

基于 PowerJob 开源项目的核心算法类,负责工作流的拓扑计算和验证:

  • PEWorkflowDAG:可序列化的 DAG 数据模型(点线表示法)
  • WorkflowDAG:运行时引用模型(引用表示法)
  • WorkflowDAGUtils:核心算法工具类
    • 环检测算法(防止死锁)
    • 就绪节点计算(拓扑排序)
    • 依赖关系分析

2. Akka 执行层

基于 Akka Actor 模型实现的分布式任务执行引擎:

WorkflowInstance Actor(有状态实例)

  • 管理单个工作流实例的完整生命周期
  • 缓存工作流状态,避免重复数据库查询
  • 实现并发控制机制(默认最多 5 个并发任务)
  • 计算就绪节点并分发任务

NodeDispatcher Actor(任务分发器)

  • 创建节点执行记录
  • 将任务路由到 Worker 节点
  • 设置超时监控

TaskWorker Actor(任务执行器)

  • 执行具体的 DataX 同步任务
  • 通过 ClusterRouterPool 实现分布式部署
  • 支持轮询(RoundRobin)负载均衡

ClusterManager Actor(集群管理器)

  • 监听集群节点上下线事件
  • 触发故障恢复机制
  • 实现任务自动迁移

3. 持久化层

基于 MySQL 实现的状态持久化:

  • workflow:工作流定义表(存储 DAG 拓扑结构)
  • workflow_build_history:工作流实例表(运行时状态)
  • dag_node_execution:节点执行详情表(任务级别跟踪)

消息传递拓扑

Actor 消息传递拓扑图

图:Actor 消息传递拓扑 - 展示各个 Actor 之间的消息流转关系

关键设计:TaskWorker 通过 getSender() 机制直接回复 WorkflowInstance,避免消息多层转发,降低延迟。

核心特性解析

1. 统一的集群架构:从单节点到多节点的无缝扩展

TIS 5.1 最大的架构亮点是统一的集群架构设计

单节点部署(默认即集群)

启动第一个节点时,无需任何特殊配置:

java -jar tis-web-start.jar

此时系统会:

  • 自动形成单节点集群
  • seed-nodes 默认指向自己
  • 所有任务在本地执行
  • ClusterRouterPool 自动退化为本地路由

扩展到多节点(零配置)

需要扩展时,在新机器上启动节点,只需指向第一个节点:

java -DAKKA_HOST=192.168.1.11 \
-DAKKA_PORT=2551 \
-DAKKA_SEED_NODES="akka://TIS-DAG-Cluster@192.168.1.10:2551" \
-jar tis-web-start.jar

关键优势

  • 第一个节点无需重启
  • 第一个节点无需修改配置
  • 新节点自动加入集群
  • 任务自动分发到新节点
  • 支持运行时动态扩缩容

单节点到多节点扩展示意图

图:集群扩展演进 - 从单节点到 3 节点集群的无缝扩展过程

2. 智能的任务分发与负载均衡

并发控制机制

工作流实例内部实现了智能的并发控制:

等待队列 → 运行任务集合(最多 5 个)→ 完成
↑ ↓
└────── 任务完成后自动填补 ──┘

设计原因

  • 避免 DAG 中的 100 个表同时开始抽取
  • 防止对源数据库造成巨大压力
  • 通过队列机制平滑任务执行

并发上限(maxConcurrentTasks)由「任务触发参数」中的管道任务并发数控制,详见任务执行参数配置

ClusterRouterPool 负载均衡

TaskWorker 通过 ClusterRouterPool 实现集群级别的负载均衡:

  • 本地池大小:每个节点 10 个 Worker 实例
  • 集群总大小:最多 100 个 Worker 实例
  • 路由策略:轮询(RoundRobin)
  • 动态发现:新节点加入后自动纳入路由

以上两个容量值由「任务触发参数」中的单节点最大Worker数集群最大Worker总数控制,可按需调整。

任务分发与负载均衡示意图

图:任务分发与负载均衡 - 100个表抽取任务通过并发控制和负载均衡分发到多节点执行

3. 状态缓存优化:消除重复查询

传统架构的性能瓶颈

  • 每次处理消息都要从数据库加载状态
  • 重复加载 DAG 定义和计算
  • 大量数据库查询成为性能瓶颈

Akka 架构的优化

WorkflowInstance Actor 是有状态实例:

  • 每个工作流实例只需 2-3 次数据库查询
  • DAG 定义和运行时状态缓存在内存中
  • 不同工作流实例完全并行,无锁竞争
  • 同一工作流内串行处理,天然避免并发问题

性能提升:在 100 个节点的工作流中,数据库查询次数从数千次降低到 2-3 次。

4. 强大的容错与故障恢复

故障检测机制

ClusterManager Actor 监听集群事件:

  • MemberUp:节点上线,触发任务重新平衡
  • MemberRemoved:节点下线,触发任务恢复
  • UnreachableMember:节点不可达,启动故障检测

自动恢复流程

故障恢复流程图

图:故障自动恢复流程 - Worker 节点宕机后 8 步自动恢复过程

5. 灵活的流程控制

DAG 流程控制能力

  • 多根节点支持:允许多个独立任务链并行执行
  • 节点禁用:支持动态禁用某些节点
  • 失败跳过:节点失败时可选择跳过继续执行
  • 条件分支:未来可扩展支持条件判断节点

就绪节点算法

基于拓扑排序的智能算法:

  1. 计算所有前置依赖已完成的节点
  2. 自动跳过禁用节点
  3. 支持失败节点的跳过策略
  4. 递归处理依赖链

DAG 执行流程示意图

图:DAG 执行流程 - 菱形结构工作流按依赖关系分批执行,展示等待、运行、完成状态变化

任务执行参数配置

上述并发控制、负载均衡、超时监控等核心特性,均由 TIS 控制台中的「任务触发参数」配置驱动。该配置对应插件 LocalDataXJobSubmitParams(父类 DataXJobSubmitParams),为单实例配置——全系统仅允许保存一个配置实例,Akka 运行期各组件在初始化和任务分发时读取该配置生效。

配置入口:TIS 控制台 → 「设置任务触发参数」对话框:

任务触发参数配置对话框

图:任务触发参数配置 - 控制 Akka 运行期的并发度、Worker 容量、超时与内存规格(图中数值仅为示例)

参数总览

参数字段名默认值校验规则作用范围
名称namelocal_submit_params必填、标识唯一配置实例标识
管道任务并发数pipelineParallelism1不小于 1单个工作流实例
单节点最大Worker数maxInstancesPerNode55 至 10Akka 集群单个节点
任务超时时间taskExpireHours10 小时1 至 24 小时单个 DataX 表任务
内存规格memorySpec1024MB必填每个 fork 的 DataX 进程
ForkJvmforkJvm只读子任务执行方式
集群最大Worker总数maxTotalNrOfInstances100不小于单节点值整个 Akka 集群

其中 ForkJvm集群最大Worker总数 为高级选项,在对话框中默认折叠,点击「高级」开关后显示。

参数详解

1. 名称(name)

配置实例的唯一标识。该配置有单实例约束:系统中已存在实例时,仅允许同名更新,不允许创建第二个不同名的实例。

2. 管道任务并发数(pipelineParallelism)

  • 运行时作用点WorkflowInstanceActor.preStart() 读取该值赋给 maxConcurrentTasks,作为前文并发控制机制中「等待队列 → 运行任务集合」的并发上限。
  • 效果:DAG 中同时处于 RUNNING 状态的 DataX 任务数不超过该值,其余就绪节点在等待队列中排队,任务完成后自动填补空位。
  • 调优建议:根据源数据库的负载能力设置,避免大量表同时抽取对源库造成压力。也可通过 StartWorkflow 消息参数或工作流上下文中的 maxConcurrentTasks 针对单次执行临时覆盖。

3. 单节点最大Worker数(maxInstancesPerNode)

  • 运行时作用点NodeDispatcherActor 创建 TaskWorker 的 ClusterRouterPool 时,同时用作 ClusterRouterPoolSettings.maxInstancesPerNode 和本地 RoundRobinPool 的池大小(见前文ClusterRouterPool 负载均衡)。
  • 效果:每个集群节点最多创建该数量的 TaskWorker Actor 实例。每个实例同一时刻只处理一个任务,因此该值即为单节点的最大并发任务数;新节点加入集群后自动按此上限纳入路由。
  • 约束:最小 5,最大 10。

4. 集群最大Worker总数(maxTotalNrOfInstances)

  • 运行时作用点ClusterRouterPoolSettings.maxTotalNrOfInstances
  • 效果:集群中所有节点的 TaskWorker 实例总数上限,ClusterRouterPool 不会创建超过该总数的 Worker 实例。
  • 约束:不小于 1,且不得小于「单节点最大Worker数」。

5. 任务超时时间(taskExpireHours)

  • 运行时作用点(双重):
    1. NodeDispatcherActor 分发任务时将其写入 TaskExecutionMessage.timeoutMillis,并调度 NodeTimeout 超时监控;
    2. WorkflowInstanceActor 将其设为 ReceiveTimeout,作为工作流实例 Actor 的空闲回收时间。
  • 效果:单个 DataX 表任务执行超过该时长,TIS 会主动关闭甚至 Kill 任务,避免异常任务长期占用 Worker 资源。
  • 约束:1 至 24 小时,默认 10 小时。

6. 内存规格(memorySpec)

  • 运行时作用点:经 DataXJobSubmitParams.getJavaMemorySpec() → 执行链上下文 getJavaMemSpec() 传递,作为 fork 出的 DataX 子任务进程的 JVM 堆内存参数。
  • 效果:决定每个 DataX 子任务进程可使用的内存上限。
  • 调优建议:默认 1024MB;同步大表或宽表时请选择 customize,按实际数据量提高 request/limit(单位均为 MB)。

7. ForkJvm(forkJvm)

  • 运行时作用点:决定 DataXJobSubmit.InstanceType 的选取——为「是」时走 LOCAL 模式,每个子任务通过 Runtime.exec("java -classpath ...") 在独立 fork 的 JVM 进程中执行(拥有独立的内存空间与类加载器);为「否」时走 EMBEDDED 模式,子任务以线程方式在当前 JVM 内执行。
  • 说明:当前版本该选项为只读,固定 fork 独立 JVM。配合「内存规格」实现进程级资源隔离,单个任务内存溢出不会影响 TIS 主进程。

三层限流模型

上述参数在运行时形成三层递进的并发限流:

  1. 工作流级pipelineParallelism 限制单个 DAG 实例内的并发任务数
  2. 节点级maxInstancesPerNode 限制每个集群节点的 Worker 容量
  3. 集群级maxTotalNrOfInstances 限制集群 Worker 总量

系统的实际最大并发度取三层中的最小值;memorySpecforkJvm 则决定每个任务进程的资源配额与隔离方式。集群扩容时,新节点按 maxInstancesPerNode 自动补充 Worker,直至触及集群总量上限。

三层并发限流模型示意图

图:三层并发限流模型 - 任务从工作流的等待队列进入运行集合,再分发到各节点的 TaskWorker,三层限流共同约束系统最大并发度

定时任务触发配置

除手动触发外,TIS 支持为数据管道配置基于 Cron 表达式的定时触发器。该配置对应插件 BatchJobCrontab(全类名 com.qlangtech.tis.dag.BatchJobCrontab),属于管道级配置——每个数据管道最多定义一个定时触发器,配置随管道保存;保存后控制台会实时将调度注册到 Assemble 节点的 DAGSchedulerActor,无需重启服务即可生效。

配置入口:TIS 控制台 → 数据管道详情页 → 「Crontab」配置对话框:

定时触发器配置对话框

图:定时触发器配置 - 为单个数据管道设置 Cron 表达式与启用开关(图中数值仅为示例)

参数总览

参数字段名默认值校验规则作用范围
名称namecrontab只读触发器标识
表达式crontab必填所属数据管道
启用turnOn必选所属数据管道

参数详解

1. 名称(name)

触发器的固定逻辑名称,只读且恒为 crontab。一个数据管道只能定义一个定时任务触发器,重复保存视为对同一触发器的更新。

2. 表达式(crontab)

  • 格式:Spring 风格的 6 段 Cron 表达式(秒 分 时 日 月 周),与 Unix 传统的 5 段格式不同,第一段为秒。常用示例:0 0 12 * * ?(每天中午 12 点)、0 0 2 * * ?(每天凌晨 2 点)、0 15 10 ? * *(每天上午 10:15)。不支持第 7 段「年」字段。
  • 运行时作用点:保存配置时,插件将表达式 POST 到 Assemble 节点的 DAGWorkflowServlet,经 DistributedAKKAJobDataXJobSubmit.handleRegisterSchedule()RegisterSchedule 消息发送至 DAGSchedulerActor;后者使用 cron-utils(CronType.SPRING)计算下一次触发时间,并通过 Akka scheduleOnce 挂载一次性定时任务,每次触发后自动挂载下一次,形成持续的调度链。
  • 效果:到达触发时间时,调度器为该管道创建 TriggerType.CRONTAB 类型的执行上下文,走与手动触发完全相同的 Akka DAG 执行路径,前文的并发控制机制三层限流模型同样生效。
  • 约束
    • 若触发时该管道已有处于 RUNNING / WAITING / QUEUED 状态的工作流实例,本次触发直接跳过,且不会补跑,因此表达式间隔应大于管道的典型执行时长;
    • 表达式非法时调度注册会失败、任务不会触发,请以 Assemble 节点日志中的 registered schedule for pipeline 为准(对话框中的「校验」目前不会真正校验表达式合法性)。

3. 启用(turnOn)

  • 效果:定时触发器的总开关。选「否」时控制台向 DAGSchedulerActor 发送 UnregisterSchedule 消息,立即取消已挂载的定时任务,但保留表达式配置,随时可改回「是」重新启用;删除该插件配置的效果等同于禁用。
  • 说明:配置状态栏会显示触发器当前处于「激活」或「禁用」状态;注册时间、最近触发时间与下次触发时间等权威信息以调度器运行明细(getDAGSchedulerDetail)为准。

调度语义与注意事项

  1. 实现机制:当前版本基于 Akka scheduleOnce + cron-utils 实现,而非 Quartz。调度状态保存在 DAGSchedulerActor 内存中,Assemble 节点重启后会自动从各管道的插件配置重新加载并恢复调度链。
  2. 单触发器约束:每个数据管道只能有一个定时触发器,重复保存即更新表达式与开关状态。
  3. 跳过语义:采用「重叠即跳过、不补跑」策略,适合幂等的全量/增量同步管道;对错过触发必须补跑的场景,需要在管道层自行补偿。

Akka + DataX 融合方案的优势

架构融合点

Akka 与 DataX 融合架构图

图:Akka 与 DataX 融合 - TaskWorker Actor 包装 DataX Job 的完整执行流程

相比原生 DataX 的增强

维度原生 DataXTIS + Akka + DataX
分布式执行❌ 不支持✅ 支持集群部署
任务调度❌ 手动触发定时调度 + 工作流编排
负载均衡❌ 无✅ 自动负载均衡
容错能力❌ 节点故障任务失败✅ 自动故障恢复
可视化监控❌ 命令行日志✅ Web 实时监控
动态扩容❌ 不支持✅ 运行时动态扩容
并发控制⚠️ 需手动控制✅ 智能并发控制

保留 DataX 的优势

TIS 5.1 完全保留了 DataX 的核心能力:

  • 丰富的插件生态:支持 60+ 种数据源
  • 高性能传输引擎:经过阿里巴巴生产验证
  • 灵活的数据转换:支持字段映射、类型转换
  • 完善的错误处理:脏数据处理、限流控制

与开源方案对比

Apache SeaTunnel

Apache SeaTunnel(原 Waterdrop)是另一个流行的数据集成工具。

对比项TIS 5.1Apache SeaTunnel
核心引擎DataX(阿里生产验证)自研引擎
分布式架构Akka ClusterSpark/Flink
部署复杂度⭐⭐ 简单(独立部署)⭐⭐⭐⭐ 复杂(依赖 Spark 集群)
资源占用⭐⭐⭐ 轻量级⭐⭐⭐⭐ 重量级(需 Spark 资源)
批流一体✅ 支持(批量 DataX + 流式 Flink-CDC)✅ 支持
Web 管理界面✅ 完善的 Web Console⚠️ 有限支持
插件生态✅ DataX 60+ 插件✅ 自研插件体系
学习曲线⭐⭐ 平缓⭐⭐⭐ 陡峭(需要 Spark 知识)

TIS 优势

  • 不依赖 Spark/Flink 集群,部署更简单
  • 资源占用更小,适合中小规模场景
  • 完善的 Web 管理界面,开箱即用
  • 统一的批流管理平台

SeaTunnel 优势

  • 更适合大数据生态(已有 Spark 集群)
  • 更强的批流一体能力

Alibaba DataX-Web

DataX-Web 是 DataX 的社区扩展版本,提供 Web 管理能力。

对比项TIS 5.1DataX-Web
分布式执行✅ Akka Cluster⚠️ 简单的多机部署
任务编排✅ DAG 工作流❌ 无工作流概念
负载均衡✅ 自动负载均衡⚠️ 静态分配
容错恢复✅ 自动故障恢复❌ 手动重试
动态扩容✅ 运行时扩容❌ 需重启
批流一体✅ 批 + 流统一平台❌ 仅批量同步
架构成熟度⭐⭐⭐⭐ 生产级⭐⭐⭐ 社区项目

TIS 优势

  • 真正的分布式架构(Akka Cluster)
  • 完善的工作流编排能力
  • 批流一体的统一平台
  • 更强的容错和扩展能力

其他商业方案

方案类型代表产品TIS 5.1 的优势
云厂商 DTS阿里云 DTS、腾讯云 DTS✅ 开源免费
✅ 私有化部署
✅ 无厂商锁定
商业 ETL 工具Informatica、Talend✅ 开源免费
✅ 轻量级架构
✅ 易于定制

实际应用场景

场景 1:大规模表同步

需求:每天凌晨同步 MySQL 中的 200 张业务表到数据仓库

TIS 5.1 方案

  1. 定义包含 200 个 DataX 任务的 DAG 工作流
  2. 设置并发控制为 10(避免对源库压力过大)
  3. 配置 3 个 Worker 节点
  4. 通过定时任务触发配置设置 Cron 调度:0 0 2 * * ?(每天凌晨 2 点)

效果

  • 任务自动分发到 3 个节点的 30 个 Worker
  • 每次并发 10 个表同步,队列平滑执行
  • 单节点故障自动恢复,无需人工干预
  • 全流程可视化监控

大规模表同步执行示意图

图:大规模表同步 - 200 个表通过等待队列、并发运行队列在 3 节点上平滑执行

场景 2:实时扩容应对高峰

需求:双十一期间数据同步任务激增,需要临时扩容

TIS 5.1 方案

  1. 准备 2 台新服务器
  2. 启动 TIS Worker 节点,指向现有集群
  3. 新节点自动加入,无需重启现有服务
  4. 系统自动将任务分发到新节点

操作步骤

# 在新服务器上执行
export AKKA_SEED_NODES="akka://TIS-DAG-Cluster@master-node:2551"
java -jar tis-web-start.jar

效果

  • 0 停机扩容
  • 5 节点集群吞吐量提升到原来的 2.5 倍
  • 高峰过后可以下线节点,自动缩容

场景 3:跨机房容灾

需求:跨机房部署,保证单机房故障不影响服务

TIS 5.1 方案

  1. 机房 A 部署 3 个节点(包含 2 个 Seed Node)
  2. 机房 B 部署 2 个节点(包含 1 个 Seed Node)
  3. 配置 Split Brain Resolver(保留多数派策略)

容灾能力

  • 单节点故障:任务自动迁移到其他节点
  • 单机房故障:保留多数派机房继续服务
  • 网络分区:Split Brain Resolver 自动处理

监控与运维

实时监控面板

TIS 5.1 提供完善的 Web 监控界面:

{/ 图片位置:监控面板截图 /} {/* 图片说明:展示监控面板的主要模块:

  • DAG 拓扑图:实时展示节点执行状态(等待/运行/完成/失败)
  • 队列监控:等待队列长度、运行队列长度统计卡片
  • 集群状态:节点数量、节点列表、节点负载
  • 任务详情:每个节点的执行时间、Worker 地址、执行日志 */}

最佳实践建议

1. 合理设置并发度

  • 小规模(< 50 张表):并发度 3-5
  • 中等规模(50-200 张表):并发度 5-10
  • 大规模(> 200 张表):并发度 10-20

原则:根据源数据库的负载能力设置,避免对源库造成压力。

2. Worker 节点规划

  • 每个 Worker 节点建议 4 核 8G 配置
  • 小规模集群:3 个节点
  • 中等规模集群:5-10 个节点
  • 大规模集群:10+ 个节点

3. Seed Node 选择

  • 选择 3-5 个稳定的核心节点作为 Seed Node
  • Seed Node 应该是长期运行、不频繁重启的节点
  • Seed Node 可以同时作为 Worker 节点

4. 故障恢复策略

  • 配置 Split Brain Resolver:keep-majority(保留多数派)
  • 设置合理的故障检测阈值:acceptable-heartbeat-pause = 3s
  • 启用自动故障恢复机制

总结

TIS 5.1 基于 Akka + DataX 的批量数据同步架构,在以下方面实现了重大突破:

核心价值

  1. 真正的分布式能力

    • 原生 DataX 的分布式短板被完全弥补
    • 线性扩展能力:增加节点直接提升吞吐率
  2. 卓越的稳定性

    • Actor 模型天然隔离故障
    • 自动故障恢复机制
    • 节点下线不影响服务
  3. 极致的易用性

    • 统一集群架构:单节点到多节点无缝扩展
    • 零配置扩容:启动新节点即可加入集群
    • 完善的 Web 管理界面
  4. 生产级的可运维性

    • 实时监控面板
    • 完善的指标和告警
    • 强大的故障诊断能力

适用场景

TIS 5.1 特别适合以下场景:

  • ✅ 中小规模的批量数据同步(10-1000 张表)
  • ✅ 需要工作流编排的复杂同步任务
  • ✅ 对稳定性和容错能力要求高的场景
  • ✅ 希望简化部署和运维的团队
  • ✅ 批流一体的统一数据集成平台

未来展望

TIS 将持续演进批量数据同步能力:

  • 🚀 性能优化:进一步降低资源占用,提升吞吐率
  • 🚀 智能调度:基于任务历史数据的智能并发控制
  • 🚀 条件分支:支持决策节点的条件流程控制
  • 🚀 跨云部署:支持跨云厂商的混合云部署

TIS 5.1 的架构升级,标志着 TIS 在批量数据同步领域迈入了分布式、高可用的新阶段。我们相信,基于 Akka + DataX 的融合方案,将为用户提供更加稳定、高效、易用的数据集成服务。