Apache SeaTunnel 与 DataX
概述
SeaTunnel 与 DataX 是两类常用的数据同步引擎:DataX 是阿里离线批量同步工具,SeaTunnel 是 Apache 社区离线 + 流式统一的同步平台。两者都采用插件化架构,支持多源多目标、限流、断点与脏数据处理。本文逐一讲透。
一、工具定位
| 工具 | 定位 | 特点 |
|---|---|---|
| DataX | 离线批量同步 | 稳定、插件多 |
| SeaTunnel | 离线 + 流式 | 统一、易用 |
| 对比 | DataX | SeaTunnel |
|---|---|---|
| 实时 | 否 | 支持(流式) |
| 配置 | JSON | 配置/UI |
| 部署 | 单机 | 集群 |
| 生态 | 插件 | 插件 + 引擎 |
选型:
纯离线 → DataX
离线+实时 → SeaTunnel二、插件化架构
2.1 插件模型
同步作业 = Reader(读源)+ Transformer(转换) + Writer(写目标)
插件:
Source(数据源读取)
Transform(转换/过滤)
Sink(目标写入)| 插件 | 示例 |
|---|---|
| Source | MySQL/JDBC/Hive/Kafka/S3 |
| Transform | 过滤/映射/补列 |
| Sink | MySQL/ES/Hive/Doris/ClickHouse |
2.2 DataX 插件
reader:mysqlreader、hdfsreader、mongodbreader...
writer:mysqlwriter、hdfs writer、eswriter...
通道(Channel):
并发执行
限速2.3 SeaTunnel 插件
connector:
jdbc、kafka、hive、s3、doris、clickhouse...
引擎:
可基于 Spark/Flink 或自研 Zeta插件化收益:
新源/目标加插件即可
统一抽象三、多源读写
3.1 DataX 多源
json
{
"job": {
"content": [{
"reader": {
"name": "mysqlreader",
"parameter": {
"username": "root",
"password": "***",
"column": ["id", "name"],
"splitPk": "id",
"connection": [{
"jdbcUrl": ["jdbc:mysql://mysql:3306/db"],
"table": ["orders"]
}]
}
},
"writer": {
"name": "hdfswriter",
"parameter": { "defaultFS": "hdfs://nn:9000", "path": "/data/orders", "fileType": "parquet" }
}
}],
"setting": { "speed": { "channel": 4 } }
}
}| 配置 | 说明 |
|---|---|
| reader | 源插件与参数 |
| writer | 目标插件与参数 |
| splitPk | 分片字段 |
| speed.channel | 并发通道 |
3.2 SeaTunnel 多源
conf
env { parallelism = 4 }
source {
Jdbc { url = "jdbc:mysql://mysql:3306/db"
user = "root" password = "***"
table = "orders" }
}
sink {
Doris { fenodes = "doris:8030"
table = "dws.orders" }
}多源支持:
多 source 合并
多 sink 分发
条件路由四、限流
4.1 作用
防止同步压垮源库/目标:
控制速率、保护业务4.2 DataX 限速
| 参数 | 说明 |
|---|---|
| speed.channel | 并发数 |
| speed.byte | 字节速率 |
| speed.record | 行速率 |
json
"setting": { "speed": { "channel": 4, "byte": 10485760 } }4.3 SeaTunnel 限流
conf
env { parallelism = 2 }
transform { ... }
# 通过并行度与配置限速| 方式 | 说明 |
|---|---|
| 并行度 | 控制并发 |
| 限流插件 | 速率控制 |
| 队列 | 缓冲 |
五、断点续传
5.1 问题
大任务中断 → 全部重来5.2 DataX 断点
基于 splitPk 分片:
任务按主键分片
失败 → 重跑未完成分片| 说明 | 要点 |
|---|---|
| 分片 | splitPk 字段分片 |
| 重跑 | 增量区间 |
| 幂等 | 覆盖/去重 |
5.3 SeaTunnel 断点
支持 checkpoint/offset:
流式:位点恢复
批量:分片状态| 机制 | 说明 |
|---|---|
| checkpoint | 状态快照 |
| 位点 | Kafka offset 等 |
| 恢复 | 从断点继续 |
六、脏数据处理
6.1 脏数据
同步失败的记录:
类型不匹配
唯一键冲突
转换失败6.2 处理策略
| 策略 | 说明 |
|---|---|
| 跳过并记录 | 脏数据日志 |
| 重试 | 限次重试 |
| 置脏表 | 单独存储 |
| 告警 | 超阈值通知 |
DataX 脏数据处理:
errorLimit(容忍上限)
超过 → 任务失败| 参数 | 说明 |
|---|---|
| errorLimit.record | 允许脏行数 |
| errorLimit.percentage | 允许比例 |
SeaTunnel 脏数据:
写入指定表/侧输出
统计与监控6.3 质量保障
| 手段 | 说明 |
|---|---|
| 类型校验 | 写入前校验 |
| 主键冲突 | 更新策略 |
| 监控 | 脏数据率 |
七、性能调优
7.1 调优维度
| 维度 | 说明 |
|---|---|
| 并行度 | 通道/并行数 |
| 批大小 | 批量写入 |
| 网络 | 带宽 |
| 资源 | 内存/CPU |
| 索引 | 目标写入优化 |
7.2 DataX 调优
| 参数 | 说明 |
|---|---|
| channel | 合理并发(与源/目标匹配) |
| batchSize | 批量 |
| 内存 | JVM 堆 |
并发建议:
源/目标支持度决定
过高反而慢(锁/连接)7.3 SeaTunnel 调优
| 参数 | 说明 |
|---|---|
| parallelism | 并行度 |
| sink 批量 | 目标批量参数 |
| 引擎 | Zeta 内存配置 |
调优方法:
压测 → 观察吞吐
逐步加并发
找瓶颈(源/转换/目标)八、选型与使用建议
| 场景 | 工具 |
|---|---|
| 纯离线批量 | DataX |
| 离线+实时统一 | SeaTunnel |
| 多源多目标 | SeaTunnel |
| 已有插件需求 | 看生态 |
| 建议 | 说明 |
|---|---|
| 先全量后增量 | 组合方案 |
| 限流保护 | 生产必配 |
| 断点设计 | 大任务 |
| 脏数据监控 | 质量 |
| 幂等 | 可重跑 |
常见问题速查
| 问题 | 原因与处理 |
|---|---|
| 同步慢 | 并发/批量/瓶颈 |
| 中断重来 | 分片续传 |
| 脏数据 | 校验/置脏 |
| 压垮源库 | 限流 |
| 数据重复 | 幂等 |