
DataX 插件开发完全指南从框架原理、接口实现到打包测试的实战全流程【免费下载链接】DataXDataX是阿里云DataWorks数据集成的开源版本。项目地址: https://gitcode.com/gh_mirrors/da/DataX本指南面向 DataX 插件开发人员以阿里云 DataWorks 数据集成的开源版本 DataX 的官方插件开发文档dataxPluginDev.md为骨架结合仓库内common、core、streamreader、mysqlwriter、datax-example等模块的真实源码系统阐述一个 DataX 插件的完整开发历程先理解框架 插件的逻辑与物理执行模型再掌握Reader/Writer的编程接口、plugin.json定义、assembly 打包规范、JSON 配置与ConfigurationDSL、数据传输与类型转换、脏数据处理、类加载原理最后通过datax-example编写可本地调试的测试用例并遵循插件文档规范发布。读完本文你将具备从零开发并交付一个 DataX 读写插件的完整能力。一、开发之前理解 DataX 为什么需要插件机制从设计之初DataX 就把异构数据源同步作为自身的使命。为了应对不同数据源的差异同时提供一致的同步原语和扩展能力DataX 自然而然地采用了框架 插件的模式插件只需关心数据的读取或者写入本身——读取插件负责从源端拉取数据并转换为 DataX 内部统一的Record写入插件负责把Record落到目标存储。同步的共性问题由框架统一处理——包括类型转换、性能并发、限速、流控、统计速率、进度、脏数据限制、断点与日志等插件无需重复实现。作为插件开发人员你只需要关注两个问题数据源本身的读写数据正确性这是插件的核心职责如何与框架沟通、合理正确地使用框架这是本文的重点。开工前请先想明白你的插件要解决什么数据源的问题读、写还是两者都做数据源有什么样的物理分片能力可以用来拆分任务在动手 coding 之前回答清楚这些问题能避免路走远了发现没走对的尴尬。二、插件视角看框架执行模型与核心概念2.1 逻辑执行模型插件开发者不必关心框架内部的所有细节但需要理解自己的代码在逻辑上是怎么被调度的以及每一个方法在什么时机被调用。先明确以下概念JobDataX 用以描述从一个源头到一个目的端的同步作业是数据同步的最小业务单元。例如从一张 MySQL 表同步到 ODPS 的一个表的特定分区就是一个Job。Task为最大化并发而把Job拆分得到的最小执行单元。例如读一张有 1024 个分表的 MySQL 分库分表Job可以拆分成 1024 个读Task由若干个并发执行。TaskGroup一组Task的集合特指在同一个TaskGroupContainer中执行的那组Task。JobContainerJob执行器负责Job全局拆分、调度、前置语句和后置语句等工作角色类似 YARN 中的 JobTracker。TaskGroupContainerTaskGroup执行器负责执行一组Task角色类似 YARN 中的 TaskTracker。简而言之Job拆分成Task分别在框架提供的容器中执行插件只需要实现Job和Task两部分逻辑。2.2 物理执行模型框架为插件提供物理上的执行能力线程共有三种运行模式模式说明Standalone单进程运行没有外部依赖Local单进程运行统计信息、错误信息汇报到集中存储Distributed分布式多进程运行依赖DataX Service服务对插件的编写而言这三种模式没有任何区别——只要你避开下面会讲到的一些小错误例如Job与Task之间共享变量插件就能在单机/分布式之间无缝切换。判别规则很简单当JobContainer和TaskGroupContainer运行在同一个进程内时就是单机模式Standalone/Local当它们分布在不同进程中执行时就是分布式Distributed模式。2.3 编程接口Reader/Writer抽象类与内部类Job/TaskJob和Task的逻辑是如何对应到具体代码的首先插件的入口类必须继承Reader或Writer抽象类并分别实现Job和Task两个内部抽象类。这两个实现必须是内部类的形式原因见下文加载原理一节。以 Reader 为例的完整骨架如下public class SomeReader extends Reader { public static class Job extends Reader.Job { Override public void init() { } Override public void prepare() { } Override public ListConfiguration split(int adviceNumber) { return null; } Override public void post() { } Override public void destroy() { } } public static class Task extends Reader.Task { Override public void init() { } Override public void prepare() { } Override public void startRead(RecordSender recordSender) { } Override public void post() { } Override public void destroy() { } } }这两个抽象类的定义位于 common/src/main/java/com/alibaba/datax/common/spi/Reader.java 与 common/src/main/java/com/alibaba/datax/common/spi/Writer.java。从源码看Reader.Job extends AbstractJobPlugin、Reader.Task extends AbstractTaskPlugin其中AbstractJobPlugin/AbstractTaskPlugin又继承自 AbstractPlugin.java后者持有pluginJobConf本插件相关的任务配置、pluginConf插件自身的plugin.json配置等字段并通过getPluginJobConf()、setPluginJobConf()等方法暴露给插件使用。Job接口各方法职责initJob 对象初始化工作此时可以通过super.getPluginJobConf()获取与本插件相关的配置。读插件获得任务配置中reader部分写插件获得writer部分。prepare全局准备工作例如 odpswriter 清空目标表。split拆分Task。参数adviceNumber是框架建议的拆分数一般是运行时所配置的并发度返回值是Task的配置列表。post全局的后置工作例如 mysqlwriter 同步完影子表后的 rename 操作。destroyJob 对象自身的销毁工作。Task接口各方法职责initTask 对象的初始化此时通过super.getPluginJobConf()获取的配置是Job.split方法返回的配置列表中的其中一个即每个 Task 拿到自己那份切片配置。prepare局部的准备工作。startRead从数据源读数据写入RecordSenderRecordSender会把数据写入连接 Reader 和 Writer 的缓存队列Reader 特有。startWrite从RecordReceiver中读取数据写入目标数据源RecordReceiver中的数据来自 Reader 和 Writer 之间的缓存队列Writer 特有。post局部的后置工作。destroyTask 对象自身的销毁工作。两个必须牢记的注意点Job和Task之间一定不能有共享变量——因为分布式运行时无法保证共享变量被正确初始化两者之间只能通过配置文件split返回的配置切片传递依赖。prepare和post在Job和Task中都存在插件需要根据实际操作的作用域全局还是局部决定放在哪一级实现。框架按照下图顺序执行Job和Task的接口上图中黄色表示Job部分的执行阶段蓝色表示Task部分的执行阶段绿色表示框架执行阶段。整体脉络是框架先串行执行Job.init → Job.prepare → Job.split将拆分出的多个 Task 切片交给TaskGroupContainer并行执行每个Task.init → Task.prepare → Task.startRead/startWrite → Task.post → Task.destroy待所有 Task 结束后再回到 Job 侧执行Job.post → Job.destroy。相关类关系如下三、插件定义plugin.json告诉框架如何找到你代码写好了框架是怎么找到插件的入口类的答案是每个插件的资源目录下都有一个plugin.json文件它定义了插件的基本信息与入口类。以 mysqlwriter/src/main/resources/plugin.json 为例{ name: mysqlwriter, class: com.alibaba.datax.plugin.writer.mysqlwriter.MysqlWriter, description: useScene: prod. mechanism: Jdbc connection using the database, execute insert sql. warn: The more you know about the database, the less problems you encounter., developer: alibaba }字段含义name插件名称大小写敏感。框架根据用户在任务配置文件中指定的名称来搜寻插件。十分重要。class入口类的全限定名称框架通过反射实例化插件入口类。十分重要。description描述信息说明使用场景、实现机制等供人阅读。developer开发人员/团队。四、打包发布assembly 与统一目录结构DataX 使用 Mavenassembly插件打包。打包命令如下mvn clean package -DskipTests assembly:assembly每个插件的pom.xml中都声明了 assembly 打包配置例如 mysqlwriter/src/main/assembly/package.xml它把src/main/resources下的plugin.json、plugin_job_template.json输出到plugin/writer/mysqlwriter/把构建产物mysqlwriter-0.0.1-SNAPSHOT.jar也输出到同目录并把runtime作用域的依赖统一收集到plugin/writer/mysqlwriter/libs/——这与下文的插件目录规范完全对应。DataX 插件需要遵循统一的目录结构${DATAX_HOME} |-- bin | -- datax.py |-- conf | |-- core.json | -- logback.xml |-- lib | -- datax-core-dependencies.jar -- plugin |-- reader | -- mysqlreader | |-- libs | | -- mysql-reader-plugin-dependencies.jar | |-- mysqlreader-0.0.1-SNAPSHOT.jar | -- plugin.json -- writer |-- mysqlwriter | |-- libs | | -- mysql-writer-plugin-dependencies.jar | |-- mysqlwriter-0.0.1-SNAPSHOT.jar | -- plugin.json |-- oceanbasewriter -- odpswriter各目录职责${DATAX_HOME}/bin可执行程序目录datax.py启动脚本。${DATAX_HOME}/conf框架配置目录core.json、logback.xml。${DATAX_HOME}/lib框架依赖库目录。${DATAX_HOME}/plugin插件目录下分reader与writer子目录读写插件分别存放。插件目录规范${PLUGIN_HOME}/libs插件的依赖库第三方 jar。${PLUGIN_HOME}/plugin-name-version.jar插件本身的 jar。${PLUGIN_HOME}/plugin.json插件描述文件。尽管框架加载插件时会把${PLUGIN_HOME}下所有的 jar都放到 classpath但仍然推荐把依赖库的 jar 与插件本身的 jar 分开存放便于排查依赖冲突与定位问题。⚠️注意插件的目录名字必须与plugin.json中定义的插件名称一致否则框架无法按名定位插件。五、配置文件任务 JSON、core.json与ConfigurationDSLDataX 使用 JSON 作为配置文件的格式。一个典型的 DataX 任务配置如下{ job: { content: [ { reader: { name: odpsreader, parameter: { accessKey: , accessId: , column: [], isCompress: , odpsServer: , partition: [ ], project: , table: , tunnelServer: } }, writer: { name: oraclewriter, parameter: { username: , password: , column: [*], connection: [ { jdbcUrl: , table: [ ] } ] } } } ] } }5.1core.json框架默认行为框架有 core.json 配置文件指定框架的默认行为主要包括common.column.*列类型转换的默认格式datetimeFormat/dateFormat/timeFormat/extraFormats/timeZone/encoding供类型转换使用core.transport.channel.*传输通道实现类、通道容量capacity: 512、字节容量byteCapacity: 67108864、流控间隔flowControlInterval: 20等core.transport.exchanger.*数据交换器实现类与缓冲区大小bufferSize: 32core.container.*Job/TaskGroup 容器的行为reportInterval、每个 TaskGroup 的默认channel: 5core.statistics.collector.*统计收集器实现类与默认最大脏数据条数maxDirtyNumber: 10。任务配置中可以指定框架中已经存在的配置项任务配置具有更高优先级会覆盖core.json中的默认值。配置传递规则job.content.reader.parameter的 value 部分会传给Reader.Jobjob.content.writer.parameter的 value 部分会传给Writer.Job。Reader.Job和Writer.Job通过super.getPluginJobConf()获取这份配置Task侧拿到的是split之后的一份切片配置。DataX 框架还支持对特定配置项进行RSA 加密加密后的值以*开头。配置项的加密/解密过程对插件完全透明插件仍然以不带*的 key 来查询和操作配置项。5.2 如何设计配置参数任务配置中reader和writer下的parameter部分就是插件的配置参数配置文件的设计是插件开发的第一步应当遵循以下原则驼峰命名所有配置项采用驼峰命名法首字母小写后续单词首字母大写如jdbcUrl、sliceRecordCount。正交原则配置项必须正交功能上没有重复、没有潜规则避免两个参数控制同一件事造成歧义。富类型合理使用 JSON 的类型减少无谓的处理逻辑和出错可能使用正确的数据类型比如布尔值用true/false而不是yes/true/0这类字符串合理使用集合类型比如用数组替代带分隔符的字符串。类似通用遵守同一类型插件的习惯。例如关系型数据库的connection参数都是如下结构多连接、多表共用{ connection: [ { table: [ table_1, table_2 ], jdbcUrl: [ jdbc:mysql://127.0.0.1:3306/database_1, jdbc:mysql://127.0.0.2:3306/database_1_slave ] }, { table: [ table_3, table_4 ], jdbcUrl: [ jdbc:mysql://127.0.0.3:3306/database_2, jdbc:mysql://127.0.0.4:3306/database_2_slave ] } ] }5.3 使用Configuration类与路径 DSL为了简化对 JSON 的操作DataX 在 common/src/main/java/com/alibaba/datax/common/util/Configuration.java 中提供了Configuration类配合一套简单 DSL 使用。该类内部采用结构化树形保存JSON 信息注释中对比了打平 Key方案的劣势并额外维护secretKeyPathSet记录被加密的 keyPath以便在分布式场景下把加密值抛到 DataXServer。Configuration提供了常见的get、带类型get、带默认值get、set等读写操作以及clone、toJSON、merge等方法。配置项读写操作都需要传入一个path作为参数这个path就是 DataX 定义的 DSL语法只有两条子 map 用.key表示path的第一个点省略数组元素用[index]表示。比如操作如下 JSON{ a: { b: { c: 2 }, f: [ 1, 2, { g: true, h: false }, 4 ] }, x: 4 }调用configuration.get(path)的结果如下path结果x4a.b.c2a.b.c.dnull不存在返回 nulla.b.f[0]1a.b.f[2].gtrue注意插件看到的配置只是整个配置的一部分使用Configuration对象时需要时刻清楚当前的根路径是什么getPluginJobConf()返回的根路径即job.content.reader.parameter或job.content.writer.parameter。更多Configuration的操作可以参考ConfigurationTest.java。六、插件数据传输RecordSender/RecordReceiver/Record/Column与一般的生产者-消费者模式一样Reader 插件和 Writer 插件之间通过channel实现数据传输。channel 可以是内存的也可能是持久化的插件不必关心其具体实现——框架默认使用内存通道 MemoryChannel.java其底层是ArrayBlockingQueue配合ReentrantLock与条件变量实现背压流控关闭时向队列放入TerminateRecord通知消费者结束。插件的使用方式Reader 侧通过RecordSender见 RecordSender.java往 channel 写入数据接口提供createRecord()、sendToWriter(record)、flush()、terminate()、shutdown()等方法。Writer 侧通过RecordReceiver见 RecordReceiver.java从 channel 读取数据。channel 中的一条数据是一个Record对象Record中可以放多个Column对象可以简单理解为数据库中的记录和列。Record接口见 Record.java提供如下方法public interface Record { // 加入一个列放在最后的位置 void addColumn(Column column); // 在指定下标处放置一个列 void setColumn(int i, final Column column); // 获取一个列 Column getColumn(int i); // 转换为json String String toString(); // 获取总列数 int getColumnNumber(); // 计算整条记录在内存中占用的字节数 int getByteSize(); }因为Record是一个接口Reader 插件首先要调用recordSender.createRecord()创建一个Record实例然后把Column一个个添加到Record中最后调用recordSender.sendToWriter(record)送出。Writer 插件调用recordReceiver.getFromReader()方法获取Record然后把Column遍历出来写入目标存储。注意该方法的行为约定当 Reader 尚未退出、传输还在进行时如果暂时没有数据getFromReader()会阻塞直到有数据当传输已经结束会返回nullWriter 插件据此判断是否结束startWrite方法。七、类型转换六种内部类型与转换矩阵为了规范源端和目的端类型转换操作、保证数据不失真DataX 支持六种内部数据类型Long定点数Int、Short、Long、BigInteger 等Double浮点数Float、Double、BigDecimal 等String字符串类型底层不限长使用通用字符集UnicodeDate日期类型Bool布尔值Bytes二进制可以存放诸如 MP3 等非结构化数据。对应地common/src/main/java/com/alibaba/datax/common/element 下实现了DateColumn、LongColumn、DoubleColumn、BytesColumn、StringColumn和BoolColumn六种Column实现。Column抽象类见 Column.java除了提供getRawData()、getType()、getByteSize()等数据方法外还提供一系列以as开头的数据类型转换方法asLong()、asDouble()、asString()、asDate()、asDate(String dateFormat)、asBytes()、asBoolean()、asBigDecimal()、asBigInteger()。DataX 的内部类型在实现上会选用不同的 Java 类型以保证不失真内部类型实现类型备注Datejava.util.DateLongjava.math.BigInteger使用无限精度的大整数保证不失真Doublejava.lang.String用 String 表示保证不失真Bytesbyte[]Stringjava.lang.StringBooljava.lang.Boolean类型之间相互转换的关系如下from \ tofrom\toDateLongDoubleBytesStringBoolDate-使用毫秒时间戳不支持不支持使用系统配置的 date/time/datetime 格式转换不支持Long作为毫秒时间戳构造 Date-BigInteger 转为 BigDecimal然后BigDecimal.doubleValue()不支持BigInteger.toString()0 为 false否则 trueDouble不支持内部 String 构造 BigDecimal然后BigDecimal.longValue()-不支持直接返回内部 StringBytes不支持不支持不支持-按照common.column.encoding配置的编码转换为 String默认utf-8不支持String按照配置的 date/time/datetime/extra 格式解析用 String 构造 BigDecimal然后取longValue()用 String 构造 BigDecimal然后取doubleValue()会正确处理NaN/Infinity/-Infinity按照common.column.encoding配置的编码转换为byte[]默认utf-8-true为truefalse为false大小写不敏感其他字符串不支持Bool不支持true为1L否则0Ltrue为1.0否则0.0不支持不支持-这些转换规则与格式配置在 ColumnCast.java 中落地ColumnCast.bind(configuration)会把common.column.datetimeFormat、common.column.dateFormat、common.column.timeFormat、common.column.extraFormats、common.column.timeZone、common.column.encoding等配置注入StringCast/DateCast/BytesCast默认值见 core.json日期yyyy-MM-dd、时间HH:mm:ss、datetimeyyyy-MM-dd HH:mm:ss、extra 格式[yyyyMMdd]、时区GMT8、编码utf-8。框架在启动时完成bind插件在数据转换过程中即可透明地使用这些默认格式。八、脏数据处理收集、上报与限制8.1 什么是脏数据目前主要有三类脏数据Reader 读到不支持的类型、不合法的值不支持的类型转换比如Bytes转换为Date写入目标端失败比如写 MySQL 整型长度超长。8.2 如何处理脏数据在Reader.Task和Writer.Task中通过AbstractTaskPlugin.getTaskPluginCollector()可以拿到一个TaskPluginCollector见 TaskPluginCollector.java它提供一系列collectDirtyRecord重载方法可附带Throwable异常与errorMessage错误信息以及collectMessage(key, value)用于收集自定义信息Job 插件可在 post 阶段通过getMessage()获取。当脏数据出现时只需调用合适的collectDirtyRecord方法把被认为是脏数据的Record传入即可// 方式一仅记录脏数据 collector.collectDirtyRecord(dirtyRecord); // 方式二附带异常 collector.collectDirtyRecord(dirtyRecord, throwable); // 方式三附带错误信息 collector.collectDirtyRecord(dirtyRecord, errorMessage);用户可以在任务配置中指定脏数据限制条数或百分比限制当脏数据超出限制时框架会结束同步任务并退出core.json中默认maxDirtyNumber: 10。插件需要保证脏数据都被收集到其他工作计数、汇报、限流、超限退出交给框架即可。九、加载原理框架如何定位并实例化插件了解加载原理有助于理解为什么Job/Task必须是内部类以及为什么目录名必须等于插件名。框架的加载流程如下框架扫描plugin/reader和plugin/writer目录加载每个插件的plugin.json文件以plugin.json文件中name为 key索引所有插件配置如果发现重名插件框架会异常退出用户任务在reader/writer配置的name字段指定插件名字框架根据插件类型reader/writer和插件名称去对应插件路径下扫描所有 jar加入 classpath根据插件配置中定义的入口类class框架通过反射实例化对应的Job和Task对象。正因如此Job/Task必须声明为入口类的静态内部类框架通过Class.forName(入口类)得到插件入口类后再通过内部类名如入口类$Job、入口类$Task反射出对应的内部类并实例化。同时由于 classpath 的组装依赖目录名所以插件的目录名必须与plugin.json中的name完全一致。十、编写测试用例用datax-example实现本地快速调试框架从datax.py启动时依赖datax.home系统变量传统调试流程需要反复 maven 打包、设置环境变量、再让 JarLoader 加载最新代码非常耗时。仓库为此提供了 datax-example 模块专用于本地调试和复现 BUG它不依赖datax.home而是把 IDE 编译后的target目录作为插件的类加载目录从而把调试流程简化为改代码 → 直接调试两步。使用步骤修改插件 pom在插件的build中补充 resources 配置把src/main/resources也输出到 target确保plugin.json等资源可被加载build resources !--将resource目录也输出到target-- resource directorysrc/main/resources/directory includes include**/*.*/include /includes filteringtrue/filtering /resource /resources plugins !-- compiler plugin -- plugin artifactIdmaven-compiler-plugin/artifactId configuration source${jdk-version}/source target${jdk-version}/target encoding${project-sourceEncoding}/encoding /configuration /plugin /plugins /build在测试模块中调用ExampleContainer.start(jobPath)。ExampleContainer见 ExampleContainer.java内部通过ExampleConfigParser.parse(jobPath)解析任务 JSON然后直接构造Engine并start(configuration)启动同步。参考 datax-example-streamreader 的StreamReader2StreamWriterTestpublic class StreamReader2StreamWriterTest { Test public void testStreamReader2StreamWriter() { String path /stream2stream.json; String jobPath PathUtil.getAbsolutePathFromClassPath(path); ExampleContainer.start(jobPath); } }PathUtil.getAbsolutePathFromClassPath(path)负责把 classpath 下的 JSON 路径解析为绝对路径。更复杂的场景可参考 datax-example-neo4j 的StreamReader2Neo4jWriterTest它启动任务后还会根据 channel 和 reader 的 mock 数据校验结果集是否符合预期示范了启动 校验的完整测试范式。更多细节见 datax-example/doc/README.md。10.1 一个可参考的最小插件实现streamreader仓库中的 streamreader 是理解上述全部概念的理想范本。它的 StreamReader.java 完整演示了Job.init中通过super.getPluginJobConf()获取配置、校验必填参数如sliceRecordCount不能小于 1、解析并回写column列表支持random混淆函数Job.split(adviceNumber)中按adviceNumber克隆出相同数量的任务切片配置for (int i 0; i adviceNumber; i) { configurations.add(this.originalConfig.clone()); }Task.startRead中循环构造RecordrecordSender.createRecord()→record.addColumn(column)→recordSender.sendToWriter(record)buildOneColumn中按STRING/LONG/DOUBLE/DATE/BOOL/BYTES六种内部类型分别构造对应的Column实现类。把这个插件与 streamwriter 配对就能跑通一个最小的端到端同步链路非常适合作为新插件开发的起点模板。十一、Last but Not Least插件文档规范文档是工程师的良知。每个插件都必须在 DataX 官方 wiki 中有一篇文档文档需要包括但不限于以下内容快速介绍介绍插件的使用场景、特点等。实现原理介绍插件实现的底层原理比如mysqlwriter通过insert into和replace into实现插入tair插件通过 tair 客户端实现写入。配置说明给出典型场景下的同步任务 JSON 配置文件介绍每个参数的含义、是否必选、默认值、取值范围和其他约束。类型转换插件是如何在实际的存储类型和 DataX 的内部类型之间进行转换的以及是否存在特殊处理。性能报告软硬件环境系统版本、Java 版本、CPU、内存等数据特征记录大小等测试参数集多组、系统参数如并发数、插件参数如 batchSize不同参数下的同步速度Rec/s、MB/s、机器负载load、cpu 等、对数据源压力load、cpu、mem 等。约束限制是否存在其他使用限制条件。FAQ用户经常会遇到的问题。仓库中所有插件的doc/目录如 mysqlwriter/doc/mysqlwriter.md、streamreader/doc 等都是这一规范的落地示例新插件可以参照这些既有文档的写法组织自己的说明文档。结语从理解框架到交付插件回顾整个插件开发生命周期设计配置 → 实现Job/Task内部类 → 编写plugin.json→ assembly 打包 → 编写任务 JSON 与测试用例验证 → 撰写插件文档。贯穿始终的关键约束是插件只负责数据源读写的正确性而把类型转换、并发调度、限速统计、脏数据限制等共性问题交给框架Job与Task之间只通过配置切片通信、绝不共享变量目录名与plugin.json的name必须一致。把握住这些原则无论是内存流式的 streamreader、JDBC 家族的 mysqlwriter还是面向对象存储、搜索引擎等任意异构数据源的新插件都能快速、正确地接入 DataX 生态。【免费下载链接】DataXDataX是阿里云DataWorks数据集成的开源版本。项目地址: https://gitcode.com/gh_mirrors/da/DataX创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考