基于 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 采用经典的三层架构,职责清晰,层次分明:
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 之间的消息流转关系
关键设计: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 流程控制能力
- 多根节点支持:允许多个独立任务链并行执行
- 节点禁用:支持动态禁用某些节点
- 失败跳过:节点失败时可选择跳过继续执行
- 条件分支:未来可扩展支持条件判断节点
就绪节点算法
基于拓扑排序的智能算法:
- 计算所有前置依赖已完成的节点
- 自动跳过禁用节点
- 支持失败节点的跳过策略
- 递归处理依赖链

图:DAG 执行流程 - 菱形结构工作流按依赖关系分批执行,展示等待、运行、完成状态变化
任务执行参数配置
上述并发控制、负载均衡、超时监控等核心特性,均由 TIS 控制台中的「任务触发参数」配置驱动。该配置对应插件 LocalDataXJobSubmitParams(父类 DataXJobSubmitParams),为单实例配置——全系统仅允许保存一个配置实例,Akka 运行期各组件在初始化和任务分发时读取该配置生效。
配置入口:TIS 控制台 → 「设置任务触发参数」对话框:

图:任务触发参数配置 - 控制 Akka 运行期的并发度、Worker 容量、超时与内存规格(图中数值仅为示例)
参数总览
| 参数 | 字段名 | 默认值 | 校验规则 | 作用范围 |
|---|---|---|---|---|
| 名称 | name | local_submit_params | 必填、标识唯一 | 配置实例标识 |
| 管道任务并发数 | pipelineParallelism | 1 | 不小于 1 | 单个工作流实例 |
| 单节点最大Worker数 | maxInstancesPerNode | 5 | 5 至 10 | Akka 集群单个节点 |
| 任务超时时间 | taskExpireHours | 10 小时 | 1 至 24 小时 | 单个 DataX 表任务 |
| 内存规格 | memorySpec | 1024MB | 必填 | 每个 fork 的 DataX 进程 |
| ForkJvm | forkJvm | 是 | 只读 | 子任务执行方式 |
| 集群最大Worker总数 | maxTotalNrOfInstances | 100 | 不小于单节点值 | 整个 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)
- 运行时作用点(双重):
NodeDispatcherActor分发任务时将其写入TaskExecutionMessage.timeoutMillis,并调度NodeTimeout超时监控;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 主进程。
三层限流模型
上述参数在运行时形成三层递进的并发限流:
- 工作流级:
pipelineParallelism限制单个 DAG 实例内的并发任务数 - 节点级:
maxInstancesPerNode限制每个集群节点的 Worker 容量 - 集群级:
maxTotalNrOfInstances限制集群 Worker 总量
系统的实际最大并发度取三层中的最小值;memorySpec 与 forkJvm 则决定每个任务进程的资源配额与隔离方式。集群扩容时,新节点按 maxInstancesPerNode 自动补充 Worker,直至触及集群总量上限。

图:三层并发限流模型 - 任务从工作流的等待队列进入运行集合,再分发到各节点的 TaskWorker,三层限流共同约束系统最大并发度
定时任务触发配置
除手动触发外,TIS 支持为数据管道配置基于 Cron 表达式的定时触发器。该配置对应插件 BatchJobCrontab(全类名 com.qlangtech.tis.dag.BatchJobCrontab),属于管道级配置——每个数据管道最多定义一个定时触发器,配置随管道保存;保存后控制台会实时将调度注册到 Assemble 节点的 DAGSchedulerActor,无需重启服务即可生效。
配置入口:TIS 控制台 → 数据管道详情页 → 「Crontab」配置对话框:

图:定时触发器配置 - 为单个数据管道设置 Cron 表达式与启用开关(图中数值仅为示例)
参数总览
| 参数 | 字段名 | 默认值 | 校验规则 | 作用范围 |
|---|---|---|---|---|
| 名称 | name | crontab | 只读 | 触发器标识 |
| 表达式 | 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)计算下一次触发时间,并通过 AkkascheduleOnce挂载一次性定时任务,每次触发后自动挂载下一次,形成持续的调度链。 - 效果:到达触发时间时,调度器为该管道创建
TriggerType.CRONTAB类型的执行上下文,走与手动触发完全相同的 Akka DAG 执行路径,前文的并发控制机制与三层限流模型同样生效。 - 约束:
- 若触发时该管道已有处于 RUNNING / WAITING / QUEUED 状态的工作流实例,本次触发直接跳过,且不会补跑,因此表达式间隔应大于管道的典型执行时长;
- 表达式非法时调度注册会失败、任务不会触发,请以 Assemble 节点日志中的
registered schedule for pipeline为准(对话框中的「校验」目前不会真正校验表达式合法性)。
3. 启用(turnOn)
- 效果:定时触发器的总开关。选「否」时控制台向
DAGSchedulerActor发送UnregisterSchedule消息,立即取消已挂载的定时任务,但保留表达式配置,随时可改回「是」重新启用;删除该插件配置的效果等同于禁用。 - 说明:配置状态栏会显示触发器当前处于「激活」或「禁用」状态;注册时间、最近触发时间与下次触发时间等权威信息以调度器运行明细(
getDAGSchedulerDetail)为准。
调度语义与注意事项
- 实现机制:当前版本基于 Akka
scheduleOnce+ cron-utils 实现,而非 Quartz。调度状态保存在DAGSchedulerActor内存中,Assemble 节点重启后会自动从各管道的插件配置重新加载并恢复调度链。 - 单触发器约束:每个数据管道只能有一个定时触发器,重复保存即更新表达式与开关状态。
- 跳过语义:采用「重叠即跳过、不补跑」策略,适合幂等的全量/增量同步管道;对错过触发必须补跑的场景,需要在管道层自行补偿。
Akka + DataX 融合方案的优势
架构融合点

图:Akka 与 DataX 融合 - TaskWorker Actor 包装 DataX Job 的完整执行流程
相比原生 DataX 的增强
| 维度 | 原生 DataX | TIS + Akka + DataX |
|---|---|---|
| 分布式执行 | ❌ 不支持 | ✅ 支持集群部署 |
| 任务调度 | ❌ 手动触发 | ✅ 定时调度 + 工作流编排 |
| 负载均衡 | ❌ 无 | ✅ 自动负载均衡 |
| 容错能力 | ❌ 节点故障任务失败 | ✅ 自动故障恢复 |
| 可视化监控 | ❌ 命令行日志 | ✅ Web 实时监控 |
| 动态扩容 | ❌ 不支持 | ✅ 运行时动态扩容 |
| 并发控制 | ⚠️ 需手动控制 | ✅ 智能并发控制 |
保留 DataX 的优势
TIS 5.1 完全保留了 DataX 的核心能力:
- ✅ 丰富的插件生态:支持 60+ 种数据源
- ✅ 高性能传输引擎:经过阿里巴巴生产验证
- ✅ 灵活的数据转换:支持字段映射、类型转换
- ✅ 完善的错误处理:脏数据处理、限流控制
与开源方案对比
Apache SeaTunnel
Apache SeaTunnel(原 Waterdrop)是另一个流行的数据集成工具。
| 对比项 | TIS 5.1 | Apache SeaTunnel |
|---|---|---|
| 核心引擎 | DataX(阿里生产验证) | 自研引擎 |
| 分布式架构 | Akka Cluster | Spark/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.1 | DataX-Web |
|---|---|---|
| 分布式执行 | ✅ Akka Cluster | ⚠️ 简单的多机部署 |
| 任务编排 | ✅ DAG 工作流 | ❌ 无工作流概念 |
| 负载均衡 | ✅ 自动负载均衡 | ⚠️ 静态分配 |
| 容错恢复 | ✅ 自动故障恢复 | ❌ 手动重试 |
| 动态扩容 | ✅ 运行时扩容 | ❌ 需重启 |
| 批流一体 | ✅ 批 + 流统一平台 | ❌ 仅批量同步 |
| 架构成熟度 | ⭐⭐⭐⭐ 生产级 | ⭐⭐⭐ 社区项目 |
TIS 优势:
- 真正的分布式架构(Akka Cluster)
- 完善的工作流编排能力
- 批流一体的统一平台
- 更强的容错和扩展能力
其他商业方案
| 方案类型 | 代表产品 | TIS 5.1 的优势 |
|---|---|---|
| 云厂商 DTS | 阿里云 DTS、腾讯云 DTS | ✅ 开源免费 ✅ 私有化部署 ✅ 无厂商锁定 |
| 商业 ETL 工具 | Informatica、Talend | ✅ 开源免费 ✅ 轻量级架构 ✅ 易于定制 |
实际应用场景
场景 1:大规模表同步
需求:每天凌晨同步 MySQL 中的 200 张业务表到数据仓库
TIS 5.1 方案:
- 定义包含 200 个 DataX 任务的 DAG 工作流
- 设置并发控制为 10(避免对源库压力过大)
- 配置 3 个 Worker 节点
- 通过定时任务触发配置设置 Cron 调度:
0 0 2 * * ?(每天凌晨 2 点)
效果:
- 任务自动分发到 3 个节点的 30 个 Worker
- 每次并发 10 个表同步,队列平滑执行
- 单节点故障自动恢复,无需人工干预
- 全流程可视化监控

图:大规模表同步 - 200 个表通过等待队列、并发运行队列在 3 节点上平滑执行
场景 2:实时扩容应对高峰
需求:双十一期间数据同步任务激增,需要临时扩容
TIS 5.1 方案:
- 准备 2 台新服务器
- 启动 TIS Worker 节点,指向现有集群
- 新节点自动加入,无需重启现有服务
- 系统自动将任务分发到新节点
操作步骤:
# 在新服务器上执行
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 方案:
- 机房 A 部署 3 个节点(包含 2 个 Seed Node)
- 机房 B 部署 2 个节点(包含 1 个 Seed Node)
- 配置 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 的批量数据同步架构,在以下方面实现了重大突破:
核心价值
真正的分布式能力
- 原生 DataX 的分布式短板被完全弥补
- 线性扩展能力:增加节点直接提升吞吐率
卓越的稳定性
- Actor 模型天然隔离故障
- 自动故障恢复机制
- 节点下线不影响服务
极致的易用性
- 统一集群架构:单节点到多节点无缝扩展
- 零配置扩容:启动新节点即可加入集群
- 完善的 Web 管理界面
生产级的可运维性
- 实时监控面板
- 完善的指标和告警
- 强大的故障诊断能力
适用场景
TIS 5.1 特别适合以下场景:
- ✅ 中小规模的批量数据同步(10-1000 张表)
- ✅ 需要工作流编排的复杂同步任务
- ✅ 对稳定性和容错能力要求高的场景
- ✅ 希望简化部署和运维的团队
- ✅ 批流一体的统一数据集成平台
未来展望
TIS 将持续演进批量数据同步能力:
- 🚀 性能优化:进一步降低资源占用,提升吞吐率
- 🚀 智能调度:基于任务历史数据的智能并发控制
- 🚀 条件分支:支持决策节点的条件流程控制
- 🚀 跨云部署:支持跨云厂商的混合云部署
TIS 5.1 的架构升级,标志着 TIS 在批量数据同步领域迈入了分布式、高可用的新阶段。我们相信,基于 Akka + DataX 的融合方案,将为用户提供更加稳定、高效、易用的数据集成服务。