
DataX Transformer 详解内置 UDF 手册、Job 配置与脏数据计量实践【免费下载链接】DataXDataX是阿里云DataWorks数据集成的开源版本。项目地址: https://gitcode.com/gh_mirrors/da/DataX导读DataX 是阿里巴巴 DataWorks 数据集成的开源版本完整支持 E(Extract)、T(Transformer)、L(Load) 三阶段。本文聚焦其 T 阶段——transformer模块系统讲解 Transformer 的运行模型、六大内置 UDFdx_substr、dx_pad、dx_replace、dx_filter、dx_digest、dx_groovy的参数语义、空值与边界处理规则并结合完整 Job 配置示例与计量/脏数据机制帮助你掌握在数据同步管道中灵活裁剪列、转换列、过滤行的实战能力。文中所有实现细节均以当前仓库源码为准可直接对照验证。Transformer 定义与运行模型在数据同步、传输过程中存在用户对数据传输进行特殊定制化的需求场景包括裁剪列、转换列等可以借助 ETL 的 T 过程实现Transformer。DataX 包含了完整的 E(Extract)、T(Transformer)、L(Load) 支持Transformer 就处于读取E与写入L之间的传输管道中对逐条 Record 进行加工。从源码结构看Transformer 的编程接口定义在 transformer 模块 中Transformer.java抽象基类核心方法abstract public Record evaluate(Record record, Object... paras)入参为行记录record与函数参数paras返回处理后的record若返回null则代表过滤该行。transformerName用于标识 transformer 名称其唯一性在 DataX 中检查。ComplexTransformer.java进阶抽象类evaluate方法额外接收MapString, Object tContexttransformer 运行的配置项供需要上下文信息的复杂 transformer 使用。内置 Transformer 的注册与加载统一由 core 模块的 TransformerRegistry.java 负责其在静态初始化块中注册了全部 6 个原生 transformerregistTransformer(new SubstrTransformer()); registTransformer(new PadTransformer()); registTransformer(new ReplaceTransformer()); registTransformer(new FilterTransformer()); registTransformer(new GroovyTransformer()); registTransformer(new DigestTransformer());从该类的checkName方法可以看出命名约束原生内置transformer 名称必须以dx_开头而用户自定义非原生transformer 则不允许使用dx_前缀防止冲突。除内置注册外loadTransformerFromLocalStorage还支持从DATAX_STORAGE_TRANSFORMER_HOME目录延迟加载用户自定义 transformer每个目录需包含transformer.json描述文件与实现类实现插件化的能力。UDF 手册六大内置 Transformer 详解以下 UDF 均以name字段在 Job 的transformer数组中声明参数通过columnIndex与paras传入。各 UDF 的完整实现分别位于 core/src/main/java/com/alibaba/datax/core/transport/transformer 目录下以下语义描述与源码实现保持一致。1. dx_substr按位置截取子串参数3 个第一个参数字段编号对应 record 中第几个字段即columnIndex。第二个参数字段值的开始位置包含该位置。第三个参数目标字段长度。返回从字符串的指定位置包含截取指定长度的字符串。如果开始位置非法抛出异常。如果字段为空值直接返回即不参与本 transformer。举例dx_substr(1,2,5) column 1的value为dataxTesttaxTe dx_substr(1,5,10) column 1的value为dataxTestTest实现细节参见 SubstrTransformer.java实现采用 JavaString.substring(startIndex, startIndex length)注意其起始下标从 0 开始与文档示例中的从 1 计数语义一致若startIndex length 原字符串长度则截取到字符串末尾。空值oriValue null时直接返回原 record不参与处理startIndex超出字符串长度时会抛出TRANSFORMER_RUN_EXCEPTION异常。2. dx_pad头/尾补字符参数4 个第一个参数字段编号对应 record 中第几个字段。第二个参数l、r指示是在头部进行 pad还是尾部进行 pad。第三个参数目标字段长度。第四个参数需要 pad 的字符。返回如果源字符串长度小于目标字段长度按照位置添加 pad 字符后返回。如果长于直接截断都截右边。如果字段为空值转换为空字符串进行 pad即最后的字符串全是需要 pad 的字符。举例dx_pad(1,l,4,A), 如果column 1 的值为 xyz Axyz 值为 xyzzzzz xyzz dx_pad(1,r,4,A), 如果column 1 的值为 xyz xyzA 值为 xyzzzzz xyzz实现细节参见 PadTransformer.java空值按空字符串处理pad 类型仅接受l/r忽略大小写否则抛异常当length 原串长度时直接substring(0, length)截断即都截右边doPad方法按需循环填充 pad 字符不足目标长度时对 pad 串做尾部截取最后按lfinalPad oriValue或roriValue finalPad拼接。3. dx_replace按位置替换子串参数4 个第一个参数字段编号对应 record 中第几个字段。第二个参数字段值的开始位置。第三个参数需要替换的字段长度。第四个参数需要替换的字符串。返回从字符串的指定位置包含替换指定长度的字符串。如果开始位置非法抛出异常。如果字段为空值直接返回即不参与本 transformer。举例dx_replace(1,2,4,****) column 1的value为dataxTestda****est dx_replace(1,5,10,****) column 1的value为dataxTestdatax****实现细节参见 ReplaceTransformer.java空值直接返回原 recordstartIndex超出长度抛异常若startIndex length 原串长度则结果 原串前缀 替换串即替换到末尾否则结果 前缀 替换串 后缀。4. dx_filter行过滤关联 filter 暂不支持即多个字段的联合判断函参太过复杂用户难以使用。若需多字段联合判断可借助dx_groovy实现。参数第一个参数字段编号对应 record 中第几个字段。第二个参数运算符支持以下运算符like,not like,,,,,!,。第三个参数正则表达式Java 正则表达式、值。返回如果匹配正则表达式返回 Null表示过滤该行不匹配表达式时表示保留该行注意是该行。对于都是对字段直接 compare 的结果。like、not like是将字段转换成 String然后和目标正则表达式进行全匹配源码中使用String.matches(value)。、、、、!、对于 DoubleColumn 比较 double 值对于 LongColumn 和 DateColumn 比较 long 值其他 StringColumn、BooleanColumn 以及 ByteColumn 均比较的是 StringColumn 值。如果目标 column 为空null对于 null的过滤条件将满足条件被过滤! null的过滤条件null 不满足条件不被过滤like字段为 null 不满足条件不被过滤not like字段为 null 满足条件被过滤。举例dx_filter(1,like,dataTest) dx_filter(1,,10)实现细节参见 FilterTransformer.java比较运算按列类型分流——DoubleColumn走asDouble()数值比较LongColumn/DateColumn走asLong()数值比较DateColumn 按时间戳 long 值比较其余列StringColumn、BytesColumn、BoolColumn走asString()字典序compareTo比较。空值处理上有细微差别/系列比较时空值不参与比较直接保留空也属于无穷小/无穷大/!时仅当目标值为字符串null忽略大小写时才参与判定时 null 被过滤、!时 null 被保留否则空字段不参与过滤。5. dx_digest哈希摘要参数3 个第一个参数字段编号对应 record 中第几个字段。第二个参数hash 类型md5、sha1。第三个参数hash 值大小写toUpperCase大写、toLowerCase小写。返回返回指定类型的 hashHex。如果字段为空则转为空字符串再返回对应 hashHex。举例dx_digest(1,md5,toUpperCase), column 1 的值为 xyzzzzz 9CDFFC4FA4E45A99DB8BBCD762ACFFA2实现细节参见 DigestTransformer.java基于 Apache Commons Codec 的DigestUtils.md5Hex/sha1Hex实现参数校验严格hash 类型仅接受md5/sha1、大小写参数仅接受toUpperCase/toLowerCase均忽略大小写否则抛出参数非法异常。6. dx_groovyGroovy 动态脚本参数第一个参数groovy code。第二个参数列表或者为空extraPackage。备注dx_groovy只能调用一次不能多次调用源码中采用单例懒加载 synchronized双重检查锁实现见 GroovyTransformer.java。groovy code 中支持java.lang、java.util的包可直接引用的对象有record以及 element 下的各种 columnBoolColumn.class、BytesColumn.class、DateColumn.class、DoubleColumn.class、LongColumn.class、StringColumn.class。不支持其他包如果用户有需要用到其他包可设置 extraPackage注意extraPackage 不支持第三方 jar 包。groovy code 中返回更新过的 Record比如record.setColumn(columnIndex, new StringColumn(newValue));或者返回 null。返回 null 表示过滤此行。用户可以直接调用静态的 Util 方式GroovyTransformerStaticUtil其实现见 GroovyTransformerStaticUtil.java目前提供的方法列表md5(String):Stringsha1(String):String举例groovy 实现的 subStrString code Column column record.getColumn(1);\n String oriValue column.asString();\n String newValue oriValue.substring(0, 3);\n record.setColumn(1, new StringColumn(newValue));\n return record;; dx_groovy(record);groovy 实现的 ReplaceString code2 Column column record.getColumn(1);\n String oriValue column.asString();\n String newValue \****\ oriValue.substring(3, oriValue.length());\n record.setColumn(1, new StringColumn(newValue));\n return record;;groovy 实现的 PadString code3 Column column record.getColumn(1);\n String oriValue column.asString();\n String padString \12345\;\n String finalPad \\;\n int NeedLength 8 - oriValue.length();\n while (NeedLength 0) {\n \n if (NeedLength padString.length()) {\n finalPad padString;\n NeedLength - padString.length();\n } else {\n finalPad padString.substring(0, NeedLength);\n NeedLength 0;\n }\n }\n String newValue finalPad oriValue;\n record.setColumn(1, new StringColumn(newValue));\n return record;;实现细节GroovyTransformer通过GroovyClassLoader在运行时将用户代码包装成一个继承自Transformer的RULE类并解析、实例化参见getGroovyRule方法自动为脚本注入GroovyTransformerStaticUtil静态导入、com.alibaba.datax.common.element.*、DataXException、Transformer与java.util.*等 import再拼接用户传入的 code 作为evaluate方法体extraPackage中的 import 语句会被拼接在代码最前面。编译失败或实例化失败会分别抛出TRANSFORMER_GROOVY_INIT_EXCEPTION异常。Job 定义完整配置示例在 Job 的content[].transformer数组中按顺序声明多个 UDFDataX 会按声明顺序对每条 Record 依次执行。以下示例配置了 4 个 UDFdx_substr、dx_replace、dx_digest、dx_groovy读者源使用streamreader、结果写入streamwriterchannel 数为 1{ job: { setting: { speed: { channel: 1 }, errorLimit: { record: 0 } }, content: [ { reader: { name: streamreader, parameter: { column: [ { value: DataX, type: string }, { value: 1724154616370, type: long }, { value: 2024-01-01 00:00:00, type: date }, { value: true, type: bool }, { value: TestRawData, type: bytes } ], sliceRecordCount: 100 } }, writer: { name: streamwriter, parameter: { print: false, encoding: UTF-8 } }, transformer: [ { name: dx_substr, parameter: { columnIndex: 5, paras: [ 1, 3 ] } }, { name: dx_replace, parameter: { columnIndex: 4, paras: [ 3, 4, **** ] } }, { name: dx_digest, parameter: { columnIndex: 3, paras: [ md5, toLowerCase ] } }, { name: dx_groovy, parameter: { code: //groovy code//, extraPackage: [ import somePackage1;, import somePackage2; ] } } ] } ] } }配置要点说明nameUDF 名称必须是TransformerRegistry中已注册的名称内置 6 个即上文的dx_substr、dx_pad、dx_replace、dx_filter、dx_digest、dx_groovy。columnIndex目标列编号对应 record 中第几个字段示例中 reader 定义了 5 列编号 1~5注意列编号从 1 开始。paras除columnIndex外该 UDF 的其余参数数组均为字符串形式由实现类内部按需解析如Integer.valueOf(...)转整数。code/extraPackage仅dx_groovy使用分别传入 groovy 代码与可选的 import 列表。UDF 的执行管线位于 core 模块的 TransformerExecution.java 与 TransformerExchanger.javaReader 产出的 Record 经 Transformer 链逐条处理后再交给 Writer处理发生在传输交换层因此转换的吞吐与通道channel配置直接相关。计量与脏数据Transform 过程涉及到数据的转换可能造成数据的增加或减少因此更加需要精确度量包括Transform 的入参 Record 条数、字节数。Transform 的出参 Record 条数、字节数。Transform 的脏数据 Record 条数、字节数。如果是多个 Transform某一个发生脏数据将不会再进行后面的 transform直接统计为脏数据。目前只提供了所有 Transform 的计量成功、失败、过滤的 count以及 transform 的消耗时间。涉及运行过程的计量数据展现定义如下Total 1000000 records, 22000000 bytes | Transform 100000 records(in), 10000 records(out) | Speed 2.10MB/s, 100000 records/s | Error 0 records, 0 bytes | Percentage 100.00%注意这里主要记录转换的输入输出需要检测数据输入输出的记录数量变化。涉及最终作业的计量数据展现定义如下任务启动时刻 : 2015-03-10 17:34:21 任务结束时刻 : 2015-03-10 17:34:31 任务总计耗时 : 10s 任务平均流量 : 2.10MB/s 记录写入速度 : 100000rec/s 转换输入总数 : 1000000 转换输出总数 : 1000000 读出记录总数 : 1000000 同步失败总数 : 0注意这里主要记录转换的输入输出需要检测数据输入输出的记录数量变化。从实现看Transformer 的计量统计接入 DataX 的 communication 统计体系见 CommunicationTool.java转换的输入输出记录数、字节数以及 transform 消耗时间均通过该工具类写入任务组/作业的统计通信对象中最终汇总为运行过程与最终作业两级报表。实践提示使用dx_filter时若配置了errorLimit.record 0被过滤行不计入脏数据过滤是正常语义而 transformer 运行抛出的异常如参数非法、起始位置越界会作为脏数据记录并受errorLimit约束。多 transformer 串联时若中间某个发生脏数据后续 transformer 不会再执行该条 Record 直接计入脏数据统计。若想观测转换前后记录数的变化例如确认过滤掉了多少行可重点关注运行过程报表中Transform X records(in), Y records(out)与最终报表中转换输入总数/转换输出总数的差值。扩展你自己的 Transformer除内置的 6 个 UDF 外DataX 的 transformer 模块为自定义扩展预留了清晰的接口参见 Transformer.java 与 ComplexTransformer.java。从 TransformerRegistry.java 的加载逻辑可以梳理出自定义 transformer 的接入方式继承com.alibaba.datax.transformer.Transformer或ComplexTransformer实现evaluate(record, paras)方法并在构造器中调用setTransformerName(...)设置名称自定义名称不能以dx_开头。将实现类按目录结构放置到 transformer 存储目录下对应DATAX_STORAGE_TRANSFORMER_HOME并在该目录中提供transformer.json描述文件声明class与name字段。作业启动时loadTransformerFromLocalStorage会扫描存储目录通过JarLoader加载实现类并注册注册时若名称重复会抛出TRANSFORMER_DUPLICATE_ERROR。在 Job 的transformer数组中用自定义name即可引用。ComplexTransformer相比普通Transformer多接收一个tContextMapString, Object参数可读取 transformer 运行时的配置项适合需要上下文信息的复杂场景。小结本文围绕 DataX 的 Transformer 机制从运行模型、六大内置 UDF 的参数语义与源码实现、完整 Job 配置示例到计量与脏数据统计进行了系统性梳理。核心要点可归纳为内置 UDF 覆盖截取、补齐、替换、过滤、摘要、脚本六类常见转换需求每个 UDF 的空值处理与边界规则已在源码中明确固化使用前务必核对多 transformer 按声明顺序串行执行脏数据会中断后续处理并单独计量。掌握这些规则后你可以在任意 reader/writer 组合之间自由嵌入转换逻辑并在统计报表中准确核对转换对数据量的影响。【免费下载链接】DataXDataX是阿里云DataWorks数据集成的开源版本。项目地址: https://gitcode.com/gh_mirrors/da/DataX创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考