现成插件覆盖不了你的内部系统时,就要自己写。DataX 的插件接口不复杂,但打包和目录约定容易踩坑。本节给一条可落地的开发路径。
Reader 继承 Reader 实现 Job 和 Task:Job 做切分(split),Task 做实际读取(startRead)。Writer 类似,Task 实现 startWrite。两者都通过框架传入的 Configuration 拿参数,通过 Channel 读写 Record。你只管「这一片怎么读、这一批怎么写」。
// 自定义 Reader 的 Task 骨架 public class MyReaderTask extends Reader.Task { @Override public void startRead(RecordSender sender) { // 1. 从配置拿到本分片的范围 // 2. 拉取数据,逐条 buildRecord // 3. sender.sendToWriter(record) } }
插件是个 maven 模块,依赖 datax-core(provided),打出 jar 放到 plugin/reader/yourreader/ 下,并附 plugin_job_template.json 和 plugin.xml。目录名就是 job.json 里的 "name"。我们通常会写一个示例 job.json 一起提交,方便别人照抄。
# 7.1 自定义插件(Reader/Writer/Transformer)开发指南 plugin/reader/myreader/ myreader-1.0.jar plugin_job_template.json plugin.xml # 引用时 name 必须一致 "reader": { "name": "myreader" }

把冲突版本的第三方库打进插件 jar,导致和核心或其他插件打架。正确做法是在 pom 里把这类依赖标成 provided,只打自己的代码。Plugin Loader 的隔离能缓解,但治本还是别带多余版本。
自定义插件写完后,我们必做三件事:小样本跑通、大表分片并发验证、失败重跑幂等验证。三项全过才进生产。这条清单让我们内部插件至今没在生产出过资损类事故。
| 验收项 | 目的 |
|---|---|
| 小样本跑通 | 接口正确 |
| 大表分片 | 并发正确 |
| 重跑幂等 | 失败安全 |
我们写插件时把第三方依赖标 provided,只打自己代码,避免版本冲突。Plugin Loader 的隔离能缓解,但治本还是别带多余版本。这套规范让插件开发可复用、可审计。
preSql 里带 truncate 的任务,上线前必须二次确认目标表名。
orc 加 snappy 是 Hive 落地的常见稳妥组合,省空间且查询快。
writeMode 必须和数据更新语义对齐,不能凭感觉选。
任务的读写速率差,比绝对速率更能说明瓶颈在哪一段。
Redis 同步重跑要防重复 key,靠固定模板才能幂等。
多租户隔离交给编排层,DataX 保持简单最稳妥。
测试样本先小后大,几分钟校验能省下几小时排错。
把复杂 join 留在计算引擎,DataX 只做贴源搬运。
日志里周期打印的读写速率,是定位瓶颈的第一手材料。
DataX 的设计哲学是把连接差异收敛到插件,让核心只管调度与缓冲。
Web 平台解决协作与可观测,不提升同步能力本身。
源库索引评审应作为同步查询上线的前置环节。
任务可重跑幂等,是生产上线前的硬指标。
splitPk 的列若分布不均,分片会倾斜,部分 Task 拖慢整体。
权限收得越紧,凭据泄露的爆炸半径越小。
Hive 表的分区设计直接影响下游查询性能,写入时就该想清楚。
脏数据阈值设得太高,会掩盖源端的数据质量问题,反而埋雷。
生产环境的稳定性,常常取决于部署习惯而非某个高级特性。
全量基线加增量补充,是批流配合的常见稳妥组合。
关系型 Writer 的批量提交大小,要在往返开销和回滚成本间权衡。
把 DataX 当搬运工而非加工车间,链路才简单可排查。
反压机制保护内存,看到任务变慢应去优化下游而非加并发。
rowkey 的散列前缀设计,能避免 HBase 写入热点。
Channel 是有界缓冲,填满即触发反压保护内存。
当官方插件覆盖不到你的私有存储,就需要自研插件。DataX 的扩展点清晰:实现 Reader 抽象、打包进 plugin/reader 即可被框架加载。
// 自定义 reader 的核心:继承 Reader,实现 Job 与 Task 两个内部类 public class MyStoreReader extends Reader { public static class Job extends Reader.Job { @Override public void init() { /* 解析配置:连接、表、列 */ } @Override public List<Configuration> split(int adviceNumber) { // 按 adviceNumber 把任务切成多个 Task 配置 return configurations; } } public static class Task extends Reader.Task { @Override public void startRead(RecordSender sender) { // 真正抽取数据,逐行 buildRecord 并 sendToWriter while (hasNext()) { sender.sendToWriter(buildRecord()); } } } }
# 编译后把 jar 与依赖放到 plugin/reader/mystorereader/,并写 plugin_job_template.json mkdir -p plugin/reader/mystorereader cp target/mystorereader.jar plugin/reader/mystorereader/ cp lib/*.jar plugin/reader/mystorereader/ # 框架会自动发现并加载 python bin/datax.py job/mystore_to_hdfs.json
自研插件要填三块:Job.init 解析配置、Job.split 决定并发切分(返回多个 Task 配置,对应 channel 数)、Task.startRead 真正读数据并 sendToWriter。Writer 插件对称:Job.split 通常不分片、Task.startWrite 接收 Record 并批量写。Transformer 则实现 transform 做行内转换。三者都靠 Plugin Loader 注入依赖。
只改数据源读取逻辑时,可只写 Reader;需要新目标则只写 Writer;需要行内加工则只写 Transformer。扩展粒度很细,不必全写。
| 扩展点 | 实现类 | 核心方法 |
|---|---|---|
| Reader | Reader.Job/Task | split / startRead |
| Writer | Writer.Job/Task | split / startWrite |
| Transformer | Transformer | transform |
💡 关键直觉:DataX 的扩展模型是「框架管调度,你管读写」。你只要按约定把数据「读出来 / 写进去」,并发、限速、容错都由框架兜住,自研成本被压到最低。
⚠️ 常见坑:依赖未放进插件目录或作用域写错,运行时 NoClassDefFoundError;自研插件应依赖 datax-common 并设为 provided,避免与框架核心 jar 冲突。