第一章:DataX 概述与核心概念
了解 DataX 是什么、能做什么、基本架构与设计思想,为深入使用打下基础。
1.1 DataX 简介与应用场景
| 概念名称 | 说明 | 注意事项 |
|---|---|---|
| DataX | 阿里巴巴开源的异构数据源离线同步工具,致力于实现包括关系型数据库、HDFS、Hive、ODPS、HBase、FTP 等各种异构数据源之间稳定高效的数据同步功能。 | DataX 是离线同步工具,不适用于实时数据同步场景。 |
| 开源协议 | DataX 遵循 Apache License 2.0 开源协议,可免费用于商业和非商业用途。 | 使用时需遵守开源协议要求,保留版权和许可声明。 |
| 同步能力 | 支持多种数据源之间的双向同步,如 MySQL ↔ Oracle、MySQL → Hive、SQLServer → PostgreSQL 等。 | 所有同步任务需通过配置 JSON 文件定义,不提供图形化界面(官方版本)。 |
| 应用场景 - 数据迁移 | 在系统升级、数据库替换、云迁移等场景中,用于将旧系统数据迁移到新系统。 | 需确保目标端具备足够的存储空间和写入权限。 |
| 应用场景 - 数仓构建 | 将业务数据库(OLTP)数据抽取到数据仓库(OLAP)中,用于报表分析和决策支持。 | 建议结合调度系统(如 Airflow、Azkaban)实现定时抽取。 |
| 应用场景 - 备份与归档 | 定期将生产数据库中的历史数据同步到归档库或文件系统中。 | 可结合 where 条件实现增量或分片归档。 |
| 核心特点 - 高性能 | 采用多线程模型和插件化架构,支持并发读写,提升同步速度。 | 性能受源端和目标端 I/O 能力、网络带宽限制。 |
| 核心特点 - 稳定性 | 支持错误记录跳过、速率控制、任务重试等机制,保障任务稳定性。 | 错误容忍需在 setting 中显式配置,否则默认失败即中断。 |
1.2 DataX 架构解析(Job、Task、Framework)
| 概念名称 | 说明 | 注意事项 |
|---|---|---|
| Framework | DataX 的核心运行框架,负责任务的调度、资源管理、生命周期控制和容错处理。 | 用户无需直接操作 Framework,其行为由 JSON 配置驱动。 |
| Job(作业) | 一个完整的数据同步任务,由一个 JSON 配置文件定义,包含读、写、通道、限速等配置。 | 一个 Job 对应一次执行,可通过调度系统实现周期性运行。 |
| Task(任务) | Job 被拆分后的最小执行单元,Framework 根据并发度将 Job 拆分为多个 Task 并行执行。 | Task 数量通常等于 channel 数量,过多可能导致源端压力过大。 |
| Job 初始化 | Framework 解析 JSON 配置,进行语义检查、参数校验、全局分片策略制定。 | 若配置语法错误或插件不存在,Job 将在初始化阶段失败。 |
| Task 切分(Split) | 根据 reader 的切分策略(如按主键分片、按文件分片),将 Job 拆分为多个 Task。 | 不同 reader 插件切分策略不同,需确保切分字段能均匀分布数据。 |
| TaskGroup | 一组 Task 的集合,由一个线程执行,用于控制并发粒度和资源隔离。 | 默认每个 TaskGroup 包含多个 Task,可通过配置调整。 |
| 调度控制 | Framework 负责调度 TaskGroup 的执行,监控进度,处理失败重试或终止。 | 支持失败继续(errors 配置),但不支持断点续传。 |
| 运行模式 | 单机多线程模式运行,不依赖 Hadoop 或 Spark 等分布式框架。 | 适合单机可承载的同步任务,超大规模同步需自行分片调度。 |
1.3 插件机制与数据通道模型
| 概念名称 | 说明 | 注意事项 |
|---|---|---|
| 插件机制 | DataX 采用插件化架构,reader 和 writer 均为独立插件,支持动态扩展。 | 插件需遵循统一接口规范,编译后放入 plugin/ 目录方可使用。 |
| reader 插件 | 负责从指定数据源读取数据,如 mysqlreader、oraclereader、hdfsreader 等。 | 每个 job 中只能配置一个 reader。 |
| writer 插件 | 负责将数据写入目标数据源,如 mysqlwriter、hdfswriter、txtfilewriter 等。 | 每个 job 中只能配置一个 writer。 |
| datax-plugin 开发规范 | 插件需实现 Reader 和 Writer 接口,包含 init、prepare、split、startRead/startWrite 等方法。 | 开发者需自行处理连接管理、异常捕获、类型转换等逻辑。 |
| 数据通道(Channel) | 数据在 reader 和 writer 之间的传输通道,每个 channel 对应一个独立的数据流。 | channel 数量直接影响并发度和性能。 |
| Buffer 缓冲区 | 每个 channel 内部使用阻塞队列作为缓冲区,暂存 reader 读取的数据,供 writer 消费。 | 缓冲区大小影响内存占用和吞吐量,默认 32768 条记录。 |
| Record 与 Type | 数据以 Record 形式在通道中传输,包含多种类型字段:String、Long、Double、Date、Boolean、Bytes。 | 类型不匹配可能导致转换异常或数据截断。 |
| Transformer | 可选组件,用于在同步过程中对数据进行转换(如字段拼接、加密、脱敏)。 | 需在配置中显式定义 transformer 节点,支持自定义实现。 |
| 插件加载机制 | 启动时扫描 plugin/reader/ 和 plugin/writer/ 目录,加载所有插件配置和类。 | 插件目录结构和 plugin.json 文件必须符合规范,否则加载失败。 |
第二章:环境准备与快速入门
搭建运行环境,完成第一个数据同步任务,建立直观认知。
2.1 安装与环境配置(Java、Python、DataX 包)
| 概念/步骤 | 说明 | 注意事项 |
|---|---|---|
| Java 环境 | DataX 基于 Java 开发,需安装 JDK 1.8 或以上版本。 | 必须安装 JDK,仅 JRE 无法运行。可通过 java -version 验证。 |
| Python 环境 | DataX 启动脚本使用 Python 编写,需安装 Python 2.6+ 或 Python 3.x。 | 大多数 Linux/Unix 系统默认已安装 Python,Windows 需手动安装并配置环境变量。 |
| 下载 DataX | 从官方 GitHub 仓库(https://github.com/alibaba/DataX)下载打包好的发布版本。 | 建议使用稳定 release 版本,避免使用未测试的开发分支。 |
| 解压安装包 | 将下载的 datax.tar.gz 解压到指定目录,如 /opt/datax。 | 确保解压路径无中文和空格,避免运行时路径解析错误。 |
| 目录结构 | 解压后主要目录:bin/(启动脚本)、plugin/(reader 和 writer 插件)、job/(示例任务配置)、lib/(依赖 jar 包) | 不要随意删除或修改 plugin 目录下的内容,否则插件将无法加载。 |
| 权限配置 | 确保 datax.py 脚本具有可执行权限。 | Linux/Unix 下执行 chmod +x bin/datax.py,Windows 下通常无需处理。 |
| 环境变量(可选) | 可将 DataX 的 bin 目录加入系统 PATH,方便全局调用。 | 非必需,可通过完整路径执行,如 python /opt/datax/bin/datax.py ...。 |
| 验证安装 | 执行 python datax.py job/job.json 运行示例任务,验证安装是否成功。 | 首次运行可能较慢(JVM 启动),成功后应看到同步统计信息。 |
2.2 第一个 DataX 任务:从本地文件同步到控制台
| 概念/配置项 | 说明 | 注意事项 |
|---|---|---|
| 示例配置文件 | 使用 DataX 自带的 job/job.json,实现从本地 CSV 文件读取数据并输出到控制台。 | 该文件用于验证安装和基本功能,不涉及数据库连接。 |
| reader: txtfilereader | 读取本地文本文件(如 CSV、TSV),支持字段分隔符、编码、列类型等配置。 | 文件路径需为绝对路径或相对于 DataX 根目录的相对路径。 |
| writer: streamwriter | 将数据打印到标准输出(控制台),常用于测试和调试。 | 可配置 print 是否打印内容,encoding 指定输出编码。 |
| 配置结构 | JSON 文件包含 job → content → [{reader, writer}] 和 setting。 | content 数组中可定义多个读写对(分片任务),但初学者通常只用一个。 |
| 执行命令 | python bin/datax.py job/job.json | 必须在 DataX 根目录下执行,或使用完整路径。 |
| 输出内容 | 控制台会逐行打印从文件读取的记录,并显示同步统计:读取记录数、写入记录数、速度等。 | 若未看到数据输出,检查 streamwriter 的 print 是否为 true。 |
| 修改测试文件 | 可编辑 job/reader.txt 文件,添加自定义数据,验证同步结果。 | 确保字段数与配置中 column 数量一致,避免解析错误。 |
| 自定义配置 | 可复制 job.json 为 myjob.json,修改 path、column、delimiter 等参数进行测试。 | JSON 语法必须正确,建议使用 JSON 校验工具检查格式。 |
2.3 运行原理与日志解读
| 日志/阶段 | 说明 | 注意事项 |
|---|---|---|
| 启动阶段日志 | 显示 JVM 启动参数、DataX 版本、Python 执行路径等信息。 | 若此阶段失败,检查 Java/Python 是否安装正确。 |
| 任务初始化 | 解析 JSON 配置,加载 reader/writer 插件,进行语义检查。 | 若插件名错误或配置缺失,会在此阶段报错并终止。 |
| Job 切分 | 根据 channel 数量和 reader 切分策略,将任务拆分为多个 Task。 | 日志中会显示切分后的 Task 数量,如 Splits: 3。 |
| Task 执行 | 每个 Task 启动独立线程,reader 读取数据,通过 Channel 传输给 writer。 | 多 Task 并行执行,提升整体吞吐量。 |
| 数据读取日志 | 显示 reader 读取的记录数、速度、平均读取时间等。 | 若读取速度慢,检查源文件 I/O 或数据库查询性能。 |
| 数据写入日志 | 显示 writer 写入的记录数、速度、批次大小等。 | streamwriter 会打印每条记录(若 print=true)。 |
| 同步统计摘要 | 任务结束后输出汇总信息:任务耗时、读取记录数、写入记录数、读取速度、写入速度、错误记录数 | 正常任务应满足:读取数 = 写入数,错误数 = 0。 |
| 错误日志 | 若同步失败,会输出异常堆栈(Exception Stack),定位到具体类和行号。 | 重点关注 Caused by 部分,通常指向根本原因(如连接失败、SQL 错误)。 |
| 日志文件位置 | 运行日志默认输出到控制台,也可重定向到文件:python bin/datax.py job.json > run.log 2>&1 | 生产环境建议记录日志以便审计和排查。 |
| 常见错误类型 | Database connection failed(数据库连接问题)、Column count mismatch(字段数不匹配)、Type conversion error(类型转换失败)、File not found(文件路径错误) | 根据错误类型检查网络、权限、配置、数据格式等。 |
第三章:DataX 配置文件详解
掌握 JSON 配置结构,理解核心模块的配置方式。
3.1 配置文件整体结构(job、content、setting)
| 配置项 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
| job | { "job": { ... } } | 配置文件的根对象,包含任务的全部定义。 | 必须存在,且为最外层 JSON 对象。 |
| content | "content": [ { "reader": {}, "writer": {} } ] | 定义数据同步的读写任务列表,每个元素是一个 reader-writer 对。 | 数组中可包含多个读写任务(用于分片),但通常只定义一个。 |
| reader | "reader": { "name": "xxx", "parameter": { ... } } | 定义数据源读取插件及其参数。 | name 指定插件名(如 mysqlreader),parameter 包含具体读取参数。 |
| writer | "writer": { "name": "xxx", "parameter": { ... } } | 定义数据写入插件及其参数。 | name 指定插件名(如 txtfilewriter),parameter 包含具体写入参数。 |
| setting | "setting": { "speed": {}, "errorLimit": {} } | 定义任务级别的全局控制参数,如速度限制、错误容忍等。 | 可选,但建议配置以控制并发和容错。 |
| speed | "speed": { "channel": N, "byte": B } | 控制任务并发通道数和/或字节速率。 | channel 控制并发 Task 数,byte 控制每秒字节数(可选)。 |
| errorLimit | "errorLimit": { "record": N, "percentage": P } | 设置任务允许的最大错误记录数或百分比。 | 若超过限制,任务失败;record=0 表示不允许任何错误。 |
代码示例:
job 配置示例:
{
"job": {
"content": [...],
"setting": { ... }
}
}
content 配置示例:
"content": [
{
"reader": { ... },
"writer": { ... }
}
]
reader 配置示例:
"reader": {
"name": "mysqlreader",
"parameter": { ... }
}
writer 配置示例:
"writer": {
"name": "txtfilewriter",
"parameter": { ... }
}
setting 配置示例:
"setting": {
"speed": { "channel": 3 },
"errorLimit": { "record": 0 }
}
speed 配置示例:
"speed": { "channel": 5 }
errorLimit 配置示例:
"errorLimit": { "record": 100 }
3.2 reader 模块通用配置规范
| 配置项 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
| name | "name": "plugin_name" | 指定使用的 reader 插件名称。 | 必须与 plugin/reader/ 目录下的插件名一致。 |
| parameter | "parameter": { ... } | 包含该 reader 插件的所有具体参数。 | 不同插件参数不同,需参考具体插件文档。 |
| connection | "connection": [ { "jdbcUrl": [], "table": [] } ] | 通用连接配置结构,包含 JDBC 连接串和表名。 | jdbcUrl 和 table 均为数组,支持多表或多库。 |
| table | "table": ["table1", "table2"] | 指定要读取的表名列表。 | 表名区分大小写,需确保数据库中存在。 |
| column | "column": ["col1", "col2", "*"] | 指定要读取的字段列表,* 表示所有字段。 | 建议明确列出字段,避免使用 * 影响性能和可读性。 |
| where | "where": "id > 1000" | 在 SQL 查询中添加 WHERE 条件,用于过滤数据。 | 条件中避免使用分号 ;,且需确保语法正确。 |
| querySql | "querySql": "SELECT ..." | 自定义查询 SQL,优先级高于 table 和 column。 | 使用 querySql 时,table 和 column 将被忽略。 |
| splitPk | "splitPk": "id" | 指定用于数据分片的主键字段,提升并发读取效率。 | 应选择分布均匀的数字型主键,避免热点。 |
代码示例:
name 配置示例:
"name": "oraclereader"
parameter 配置示例:
"parameter": {
"username": "test",
"password": "123456"
}
connection 配置示例:
"connection": [
{
"jdbcUrl": ["jdbc:mysql://localhost:3306/mydb"],
"table": ["user_info"]
}
]
table 配置示例:
"table": ["orders"]
column 配置示例:
"column": ["id", "name", "age"]
where 配置示例:
"where": "status = 'active'"
querySql 配置示例:
"querySql": "SELECT id, name FROM user WHERE dt='2025-01-01'"
splitPk 配置示例:
"splitPk": "user_id"
3.3 writer 模块通用配置规范
| 配置项 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
| name | "name": "plugin_name" | 指定使用的 writer 插件名称。 | 必须与 plugin/writer/ 目录下的插件名一致。 |
| parameter | "parameter": { ... } | 包含该 writer 插件的所有具体参数。 | 不同插件参数不同,需参考具体插件文档。 |
| connection | "connection": [ { "jdbcUrl": "...", "table": [] } ] | 通用连接配置,包含目标数据库连接信息和表名。 | table 为数组,支持写入多张表(通常只写一张)。 |
| table | "table": ["target_table"] | 指定数据写入的目标表名。 | 表必须预先存在(部分 writer 支持自动建表需显式配置)。 |
| column | "column": ["col1", "col2"] | 指定要写入的目标字段列表。 | 字段顺序需与 reader 输出或 preSql/postSql 匹配。 |
| writeMode | "writeMode": "insert" 或 "writeMode": "delete" 或 "writeMode": "truncate" | 定义写入模式:insert(插入)、delete(先删后插)、truncate(清空表后插入) | delete 和 truncate 会清除原表数据,慎用。 |
| batchSize | "batchSize": 1024 | 每批次提交的记录数,影响写入性能和事务大小。 | 过大可能导致内存溢出或事务超时,过小影响吞吐。 |
| preSql | "preSql": ["DELETE FROM ..."] | 任务开始前执行的 SQL 语句(数组)。 | 常用于清理临时表或准备数据环境。 |
| postSql | "postSql": ["ANALYZE TABLE ..."] | 任务成功后执行的 SQL 语句(数组)。 | 用于更新元数据、统计信息或通知下游。 |
代码示例:
name 配置示例:
"name": "postgresqlwriter"
parameter 配置示例:
"parameter": {
"username": "admin",
"password": "pass"
}
connection 配置示例:
"connection": [
{
"jdbcUrl": "jdbc:postgresql://localhost:5432/mydb",
"table": ["user_export"]
}
]
table 配置示例:
"table": ["log_data"]
column 配置示例:
"column": ["uid", "uname"]
writeMode 配置示例:
"writeMode": "delete"
batchSize 配置示例:
"batchSize": 512
preSql 配置示例:
"preSql": ["TRUNCATE TABLE temp_data"]
postSql 配置示例:
"postSql": ["UPDATE stats SET updated=1"]
3.4 setting 模块:速度控制与出错控制配置
| 配置项 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
| speed.channel | "speed": { "channel": 5 } | 设置并发通道数(即并发 Task 数)。 | 值越大并发越高,但可能增加源端压力,建议根据源端能力调整。 |
| speed.byte | "speed": { "byte": 10485760 } | 设置每秒最大传输字节数(单位:字节)。 | 与 channel 可同时配置,系统会取更严格的限制生效。 |
| speed.record | "speed": { "record": 1000 } | 设置每秒最大传输记录数。 | 用于控制写入速率,防止目标端过载。 |
| errorLimit.record | "errorLimit": { "record": 100 } | 允许的最大错误记录数。 | 0 表示不允许任何错误,任务遇到第一条错误即失败。 |
| errorLimit.percentage | "errorLimit": { "percentage": 0.02 } | 允许的错误记录百分比(相对于总记录数)。 | 与 record 可同时设置,满足任一条件即停止任务。 |
| restore.isRestore | "restore": { "isRestore": true } | 是否开启任务断点续传(恢复模式)。 | DataX 默认不支持断点续传,此功能需特定插件或自定义实现。 |
| restore.restoreColumnName | "restore": { "restoreColumnName": "id" } | 指定用于恢复的列名(如主键)。 | 仅在支持恢复的插件中有效。 |
| log | (非标准配置) | 控制日志输出级别和路径(通常通过启动脚本或 JVM 参数控制)。 | DataX 配置中无直接 log 参数,日志行为由框架控制。 |
代码示例:
speed.channel 配置示例:
"speed": { "channel": 3 }
speed.byte 配置示例:
"speed": { "byte": 5242880 }
speed.record 配置示例:
"speed": { "record": 500 }
errorLimit.record 配置示例:
"errorLimit": { "record": 0 }
errorLimit.percentage 配置示例:
"errorLimit": { "percentage": 0.05 }
restore.isRestore 配置示例:
"restore": { "isRestore": false }
第四章:常用 Reader 插件详解
学习主流数据源作为数据读取端的配置方法。
4.1 streamreader:流式数据读取
| 参数名称 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
| name | "name": "streamreader" | 指定使用 streamreader 插件。 | 常用于测试和调试,不连接真实数据源。 |
| column | "column": [ { "type": "string", "value": "hello" }, ... ] | 定义生成的字段类型和默认值。支持 string, long, double, date, bool, bytes。 | value 为每条记录的固定值,type 必须正确指定。 |
| sliceRecordCount | "sliceRecordCount": 1000 | 每个 Task 生成的记录总数。 | 总记录数 = sliceRecordCount × channel 数量。 |
| orderMode | "orderMode": "none" 或 "orderMode": "hash" | 控制生成数据的顺序模式。none 为无序,hash 为按哈希排序。 | 测试性能时建议使用 none。 |
代码示例:
name 配置示例:
"name": "streamreader"
column 配置示例:
"column": [
{ "type": "long", "value": 1 },
{ "type": "string", "value": "datax" }
]
sliceRecordCount 配置示例:
"sliceRecordCount": 500
orderMode 配置示例:
"orderMode": "none"
4.2 mysqlreader:MySQL 数据读取
| 参数名称 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
| name | "name": "mysqlreader" | 指定使用 mysqlreader 插件。 | 需确保 plugin/reader/mysqlreader/ 目录存在。 |
| connection.jdbcUrl | "jdbcUrl": ["jdbc:mysql://host:port/db"] | MySQL 数据库的 JDBC 连接地址。 | 必须包含数据库名,支持 SSL 参数如 ?useSSL=true。 |
| connection.table | "table": ["table1"] | 要读取的表名。 | 表名区分大小写,需确保有读取权限。 |
| username | "username": "your_user" | 数据库用户名。 | 建议使用只读账号,避免权限过高。 |
| password | "password": "your_pass" | 数据库密码。 | 密码明文存储,注意配置文件权限安全。 |
| column | "column": ["id", "name"] 或 "column": ["*"] | 指定读取的字段列表。 | 推荐明确列出字段,避免 * 影响性能。 |
| where | "where": "id > 1000" | 过滤条件,附加到 SQL 的 WHERE 子句。 | 条件中避免使用分号 ;,防止 SQL 注入风险。 |
| querySql | "querySql": "SELECT ..." | 自定义查询 SQL,优先级高于 table 和 column。 | 使用时 table 和 column 将被忽略。 |
| splitPk | "splitPk": "id" | 用于数据分片的主键字段,提升并发读取效率。 | 必须是表的主键或唯一索引,且为整数类型最佳。 |
| fetchSize | "fetchSize": 1024 | 每次从数据库 fetch 的记录数,影响内存和网络。 | 过大可能导致内存溢出,过小影响性能。 |
代码示例:
name 配置示例:
"name": "mysqlreader"
connection.jdbcUrl 配置示例:
"jdbcUrl": ["jdbc:mysql://192.168.1.100:3306/test_db"]
connection.table 配置示例:
"table": ["user_info"]
username 配置示例:
"username": "datax_user"
password 配置示例:
"password": "secret123"
column 配置示例:
"column": ["id", "user_name", "age"]
where 配置示例:
"where": "status = 1 AND dt = '2025-01-01'"
querySql 配置示例:
"querySql": "SELECT id, name FROM user WHERE dept_id IN (1,2,3)"
splitPk 配置示例:
"splitPk": "user_id"
fetchSize 配置示例:
"fetchSize": 512
4.3 oraclereader:Oracle 数据读取
| 参数名称 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
| name | "name": "oraclereader" | 指定使用 oraclereader 插件。 | 需确保 Oracle JDBC 驱动(ojdbc)在 lib/ 目录。 |
| connection.jdbcUrl | "jdbcUrl": ["jdbc:oracle:thin:@//host:port/service"] | Oracle 的 JDBC 连接串,支持 SID 或 Service Name。 | 推荐使用 @//host:port/service 格式。 |
| connection.table | "table": ["EMPLOYEE"] | 要读取的表名(Oracle 表名通常大写)。 | 注意大小写敏感性,建议大写。 |
| username | "username": "scott" | Oracle 用户名(Schema 名)。 | 需有对应表的 SELECT 权限。 |
| password | "password": "tiger" | 用户密码。 | 明文存储,注意安全。 |
| column | "column": ["ID", "NAME"] | 指定读取的列名。 | 列名建议大写以匹配 Oracle 默认行为。 |
| where | "where": "DEPT_ID = 10" | 过滤条件。 | 支持 Oracle 函数如 TO_DATE。 |
| querySql | "querySql": "SELECT ..." | 自定义 SQL 查询。 | 优先级最高,忽略 table 和 column。 |
| splitPk | "splitPk": "EMP_ID" | 分片字段,用于并发读取。 | 必须是数字类型主键或唯一索引。 |
| fetchSize | "fetchSize": 1000 | 每次 fetch 的记录数。 | 根据网络和内存调整,避免性能瓶颈。 |
代码示例:
name 配置示例:
"name": "oraclereader"
connection.jdbcUrl 配置示例:
"jdbcUrl": ["jdbc:oracle:thin:@//192.168.1.101:1521/ORCL"]
connection.table 配置示例:
"table": ["SALES_DATA"]
username 配置示例:
"username": "admin"
password 配置示例:
"password": "secure_pass"
column 配置示例:
"column": ["EMP_ID", "EMP_NAME", "SALARY"]
where 配置示例:
"where": "CREATE_TIME > TO_DATE('2025-01-01', 'YYYY-MM-DD')"
querySql 配置示例:
"querySql": "SELECT * FROM EMPLOYEE WHERE ROWNUM <= 1000"
splitPk 配置示例:
"splitPk": "ID"
fetchSize 配置示例:
"fetchSize": 500
4.4 sqlserverreader:SQL Server 数据读取
| 参数名称 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
| name | "name": "sqlserverreader" | 指定使用 sqlserverreader 插件。 | 依赖 sqljdbc4.jar 或更高版本。 |
| connection.jdbcUrl | "jdbcUrl": ["jdbc:sqlserver://host:port;DatabaseName=db"] | SQL Server 的 JDBC 连接串。 | 必须包含 DatabaseName 参数。 |
| connection.table | "table": ["Users"] | 要读取的表名。 | 表名区分大小写,需确保存在。 |
| username | "username": "sa" | 登录用户名。 | 建议使用专用账号。 |
| password | "password": "pass123" | 登录密码。 | 明文存储,注意权限。 |
| column | "column": ["Id", "Name"] | 读取的字段列表。 | 字段名区分大小写。 |
| where | "where": "Status = 1" | WHERE 过滤条件。 | 时间字符串需用单引号包围。 |
| querySql | "querySql": "SELECT TOP 1000 ..." | 自定义查询 SQL。 | 优先级高于 table 和 column。 |
| splitPk | "splitPk": "Id" | 分片主键字段。 | 推荐使用自增主键。 |
| fetchSize | "fetchSize": 1024 | 每次 fetch 的记录数。 | 根据网络延迟和内存容量调整。 |
代码示例:
name 配置示例:
"name": "sqlserverreader"
connection.jdbcUrl 配置示例:
"jdbcUrl": ["jdbc:sqlserver://192.168.1.102:1433;DatabaseName=TestDB"]
connection.table 配置示例:
"table": ["OrderInfo"]
username 配置示例:
"username": "datax_user"
password 配置示例:
"password": "P@ssw0rd"
column 配置示例:
"column": ["OrderId", "Customer", "Amount"]
where 配置示例:
"where": "CreateTime > '2025-01-01'"
querySql 配置示例:
"querySql": "SELECT Id, Name FROM Users WHERE Age > 18"
splitPk 配置示例:
"splitPk": "OrderId"
fetchSize 配置示例:
"fetchSize": 512
4.5 postgresqlreader:PostgreSQL 数据读取
| 参数名称 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
| name | "name": "postgresqlreader" | 指定使用 postgresqlreader 插件。 | 需确保 postgresql-*.jar 在 lib/ 目录。 |
| connection.jdbcUrl | "jdbcUrl": ["jdbc:postgresql://host:port/db"] | PostgreSQL 的 JDBC 连接地址。 | 可附加参数如 ?ssl=true&charset=UTF8。 |
| connection.table | "table": ["products"] | 要读取的表名。 | 表名区分大小写,小写无需引号。 |
| username | "username": "pg_user" | 登录用户名。 | 需有 SELECT 权限。 |
| password | "password": "pg_pass" | 登录密码。 | 注意配置文件安全。 |
| column | "column": ["id", "name"] | 指定读取的列。 | 列名区分大小写。 |
| where | "where": "category = 'electronics'" | WHERE 条件过滤。 | 字符串用单引号包围。 |
| querySql | "querySql": "SELECT ..." | 自定义 SQL 查询。 | 忽略 table 和 column。 |
| splitPk | "splitPk": "id" | 分片字段,用于并发读取。 | 必须是整数类型主键。 |
| fetchSize | "fetchSize": 2000 | 每次 fetch 的记录数。 | PostgreSQL 对大结果集支持良好,可适当调大。 |
代码示例:
name 配置示例:
"name": "postgresqlreader"
connection.jdbcUrl 配置示例:
"jdbcUrl": ["jdbc:postgresql://192.168.1.103:5432/mydb"]
connection.table 配置示例:
"table": ["sales"]
username 配置示例:
"username": "datax"
password 配置示例:
"password": "secure_pwd"
column 配置示例:
"column": ["pid", "pname", "price"]
where 配置示例:
"where": "created > '2025-01-01'"
querySql 配置示例:
"querySql": "SELECT * FROM products WHERE price > 100"
splitPk 配置示例:
"splitPk": "product_id"
fetchSize 配置示例:
"fetchSize": 1000
4.6 hdfsreader:HDFS 文件读取
| 参数名称 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
| name | "name": "hdfsreader" | 指定使用 hdfsreader 插件。 | 用于读取 HDFS 上的文本或 SequenceFile。 |
| path | "path": "/user/data/*.csv" | HDFS 文件路径,支持通配符。 | 路径必须存在,否则任务失败。 |
| defaultFS | "defaultFS": "hdfs://namenode:9000" | HDFS 的 NameNode 地址。 | 必须与 Hadoop 集群配置一致。 |
| fileType | "fileType": "text" 或 "orc" 或 "parquet" | 文件类型,支持 text、orc、parquet 等。 | 根据实际文件格式选择。 |
| fieldDelimiter | "fieldDelimiter": "," | 字段分隔符(仅 text 文件)。 | 可使用转义字符如 \t、\u0001。 |
| encoding | "encoding": "UTF-8" | 文件编码格式。 | 确保与文件实际编码一致。 |
| column | "column": ["col1", "col2"] | 定义字段结构(text 文件需指定)。 | 对于 orc/parquet,可从元数据读取 schema。 |
| compress | "compress": "gzip" 或 "bz2" | 压缩格式(如 gzip、bz2、lzo)。 | 需确保 DataX 支持该压缩格式。 |
代码示例:
name 配置示例:
"name": "hdfsreader"
path 配置示例:
"path": "/input/orders/part*"
defaultFS 配置示例:
"defaultFS": "hdfs://hdfs-nn:9000"
fileType 配置示例:
"fileType": "text"
fieldDelimiter 配置示例:
"fieldDelimiter": "\\u0001"
encoding 配置示例:
"encoding": "GBK"
column 配置示例:
"column": ["id", "name", "age"]
compress 配置示例:
"compress": "gzip"
4.7 txtfilereader:文本文件读取
| 参数名称 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
| name | "name": "txtfilereader" | 指定使用 txtfilereader 插件。 | 用于读取本地或挂载的文本文件。 |
| path | "path": ["/data/input/*.txt"] | 本地文件路径数组,支持通配符。 | 路径必须为 DataX 运行机器可访问的绝对路径。 |
| encoding | "encoding": "UTF-8" | 文件编码。 | 编码错误会导致乱码或解析失败。 |
| fieldDelimiter | "fieldDelimiter": "," | 字段分隔符。 | 如 "\t" 表示 Tab 分隔。 |
| lineDelimiter | "lineDelimiter": "\n" | 行分隔符。 | Windows 文件通常为 \r\n。 |
| column | "column": ["f1", "f2"] | 字段定义,支持指定类型。 | 若未指定 name,默认为 col0, col1… |
| skipHeader | "skipHeader": true | 是否跳过首行(标题行)。 | 适用于 CSV 文件包含列名的情况。 |
| csvReaderConfig | "csvReaderConfig": { "safetySwitch": false } | 传递给底层 CSV 解析器的参数。 | 高级配置,一般无需修改。 |
代码示例:
name 配置示例:
"name": "txtfilereader"
path 配置示例:
"path": ["/tmp/data.csv"]
encoding 配置示例:
"encoding": "ISO-8859-1"
fieldDelimiter 配置示例:
"fieldDelimiter": "\t"
lineDelimiter 配置示例:
"lineDelimiter": "\r\n"
column 配置示例:
"column": [
{ "name": "id", "type": "long" },
{ "name": "name", "type": "string" }
]
skipHeader 配置示例:
"skipHeader": true
csvReaderConfig 配置示例:
"csvReaderConfig": { "skipEmptyRecords": false }
第五章:常用 Writer 插件详解
学习主流数据源作为数据写入端的配置方法。
5.1 streamwriter:流式数据输出
| 参数名称 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
| name | "name": "streamwriter" | 指定使用 streamwriter 插件。 | 常用于测试、调试和性能压测。 |
"print": true 或 false | 是否将读取的数据打印到控制台。 | 生产环境建议设为 false 避免日志刷屏。 | |
| encoding | "encoding": "UTF-8" | 输出数据的字符编码。 | 确保终端或日志系统支持该编码。 |
| column | "column": ["col1", "col2"] | 明确指定要输出的字段(可选,通常自动获取)。 | 若未指定,会输出 reader 传递的所有字段。 |
| format | "format": "id: %d, name: %s" | 自定义输出格式(部分实现支持)。 | 官方 streamwriter 通常不支持复杂格式化。 |
代码示例:
name 配置示例:
"name": "streamwriter"
print 配置示例:
"print": true
encoding 配置示例:
"encoding": "GBK"
column 配置示例:
"column": ["id", "name"]
5.2 mysqlwriter:MySQL 数据写入
| 参数名称 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
| name | "name": "mysqlwriter" | 指定使用 mysqlwriter 插件。 | 需确保 plugin/writer/mysqlwriter/ 存在。 |
| connection.jdbcUrl | "jdbcUrl": "jdbc:mysql://host:port/db" | 目标 MySQL 数据库的 JDBC 连接地址。 | 必须包含数据库名,可附加参数如 ?useSSL=false&rewriteBatchedStatements=true。 |
| connection.table | "table": ["target_table"] | 要写入的目标表名(数组形式)。 | 表必须预先存在,或通过 preSql 创建。 |
| username | "username": "your_user" | 数据库用户名。 | 建议使用最小权限账号。 |
| password | "password": "your_pass" | 用户密码。 | 密码明文存储,注意配置文件权限(chmod 600)。 |
| column | "column": ["col1", "col2"] | 指定要写入的目标字段列表。 | 字段顺序必须与 reader 输出顺序一致。 |
| writeMode | "writeMode": "insert" 或 "update" 或 "replace" | 写入模式:insert(插入)、update(更新已存在记录,需配合 updateKey)、replace(先删后插) | update 和 replace 需谨慎使用,避免误删数据。 |
| updateKey | "updateKey": "id" | 指定作为更新条件的主键或唯一键字段。 | 仅在 writeMode 为 update 时有效。 |
| batchSize | "batchSize": 1024 | 每批次提交的记录数。 | 建议设置为 512~2048,过大可能导致事务超时或内存溢出。 |
| preSql | "preSql": ["DELETE FROM ..."] | 任务开始前执行的 SQL 语句(数组)。 | 常用于清理目标表或准备环境。 |
| postSql | "postSql": ["ANALYZE TABLE ..."] | 任务成功后执行的 SQL 语句(数组)。 | 可用于通知下游或更新元数据。 |
| session | "session": ["set session sql_mode='STRICT_TRANS_TABLES'"] | 设置数据库会话参数。 | 可用于调整事务行为或 SQL 模式。 |
代码示例:
name 配置示例:
"name": "mysqlwriter"
connection.jdbcUrl 配置示例:
"jdbcUrl": "jdbc:mysql://192.168.1.100:3306/target_db"
connection.table 配置示例:
"table": ["user_export"]
username 配置示例:
"username": "datax_writer"
password 配置示例:
"password": "write_pass"
column 配置示例:
"column": ["uid", "uname", "reg_time"]
writeMode 配置示例:
"writeMode": "insert"
updateKey 配置示例:
"updateKey": "user_id"
batchSize 配置示例:
"batchSize": 512
preSql 配置示例:
"preSql": ["TRUNCATE TABLE temp_data"]
postSql 配置示例:
"postSql": ["UPDATE stats SET updated=1 WHERE tbl='user_export'"]
session 配置示例:
"session": ["set autocommit=0"]
5.3 oraclewriter:Oracle 数据写入
| 参数名称 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
| name | "name": "oraclewriter" | 指定使用 oraclewriter 插件。 | 需确保 ojdbc 驱动在 lib/ 目录。 |
| connection.jdbcUrl | "jdbcUrl": "jdbc:oracle:thin:@//host:port/service" | Oracle 的 JDBC 连接串。 | 推荐使用服务名格式。 |
| connection.table | "table": ["TARGET_TABLE"] | 目标表名(Oracle 表名通常大写)。 | 注意大小写敏感性。 |
| username | "username": "scott" | Oracle 用户名(Schema)。 | 需有 INSERT 权限。 |
| password | "password": "tiger" | 用户密码。 | 注意安全。 |
| column | "column": ["ID", "NAME"] | 目标字段列表。 | 字段名建议大写。 |
| writeMode | "writeMode": "insert" 或 "update" | 写入模式:insert 或 update。 | update 需配合 updateKey 使用。 |
| updateKey | "updateKey": "ID" | 更新操作的匹配字段(主键)。 | 必须是表的主键或唯一键。 |
| batchSize | "batchSize": 512 | 每批次提交的记录数。 | Oracle 对批量插入性能敏感,建议较小批次。 |
| preSql | "preSql": ["DELETE FROM T WHERE FLAG=1"] | 任务前执行的 SQL。 | 可调用存储过程。 |
| postSql | "postSql": ["COMMIT", "ANALYZE TABLE T"] | 任务后执行的 SQL。 | 用于清理或通知。 |
代码示例:
name 配置示例:
"name": "oraclewriter"
connection.jdbcUrl 配置示例:
"jdbcUrl": "jdbc:oracle:thin:@//192.168.1.101:1521/ORCL"
connection.table 配置示例:
"table": ["SALES_EXPORT"]
username 配置示例:
"username": "admin"
password 配置示例:
"password": "secure_pass"
column 配置示例:
"column": ["EMP_ID", "EMP_NAME", "SALARY"]
writeMode 配置示例:
"writeMode": "insert"
updateKey 配置示例:
"updateKey": "EMP_ID"
batchSize 配置示例:
"batchSize": 256
preSql 配置示例:
"preSql": ["EXECUTE PROC_CLEAN()"]
postSql 配置示例:
"postSql": ["INSERT INTO LOG VALUES ('done', SYSDATE)"]
5.4 sqlserverwriter:SQL Server 数据写入
| 参数名称 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
| name | "name": "sqlserverwriter" | 指定使用 sqlserverwriter 插件。 | 依赖 sqljdbc 驱动。 |
| connection.jdbcUrl | "jdbcUrl": "jdbc:sqlserver://host:port;DatabaseName=db" | SQL Server 连接串。 | 必须包含 DatabaseName。 |
| connection.table | "table": ["TargetTable"] | 目标表名。 | 表名区分大小写。 |
| username | "username": "sa" | 登录用户名。 | 使用专用账号。 |
| password | "password": "pass" | 登录密码。 | 注意安全。 |
| column | "column": ["Id", "Name"] | 写入字段列表。 | 顺序需匹配。 |
| writeMode | "writeMode": "insert" 或 "update" | 写入模式。 | update 需 updateKey。 |
| updateKey | "updateKey": "Id" | 更新匹配字段。 | 应为唯一索引。 |
| batchSize | "batchSize": 1000 | 批量提交大小。 | 根据事务日志空间调整。 |
| preSql | "preSql": ["DELETE FROM T"] | 前置 SQL。 | 可执行存储过程。 |
| postSql | "postSql": ["UPDATE STAT SET CNT=..."] | 后置 SQL。 | 用于日志记录。 |
代码示例:
name 配置示例:
"name": "sqlserverwriter"
connection.jdbcUrl 配置示例:
"jdbcUrl": "jdbc:sqlserver://192.168.1.102:1433;DatabaseName=TargetDB"
connection.table 配置示例:
"table": ["OrderExport"]
username 配置示例:
"username": "dx_writer"
password 配置示例:
"password": "P@ssw0rd"
column 配置示例:
"column": ["OrderId", "Customer", "Amount"]
writeMode 配置示例:
"writeMode": "insert"
updateKey 配置示例:
"updateKey": "OrderId"
batchSize 配置示例:
"batchSize": 500
preSql 配置示例:
"preSql": ["EXEC ClearTemp"]
postSql 配置示例:
"postSql": ["PRINT 'Done'"]
5.5 postgresqlwriter:PostgreSQL 数据写入
| 参数名称 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
| name | "name": "postgresqlwriter" | 指定使用 postgresqlwriter 插件。 | 需 postgresql-*.jar 驱动。 |
| connection.jdbcUrl | "jdbcUrl": "jdbc:postgresql://host:port/db" | PostgreSQL 连接地址。 | 可附加参数。 |
| connection.table | "table": ["target_table"] | 目标表名。 | 小写无需引号。 |
| username | "username": "pg_user" | 用户名。 | 需 INSERT 权限。 |
| password | "password": "pg_pass" | 密码。 | 注意权限。 |
| column | "column": ["id", "name"] | 目标字段。 | 顺序匹配。 |
| writeMode | "writeMode": "insert" 或 "update" | 写入模式。 | update 需 updateKey。 |
| updateKey | "updateKey": "id" | 更新键。 | 唯一键。 |
| batchSize | "batchSize": 2000 | 批量大小。 | PostgreSQL 批量性能较好,可适当调大。 |
| preSql | "preSql": ["TRUNCATE t"] | 前置 SQL。 | 清理数据。 |
| postSql | "postSql": ["ANALYZE t"] | 后置 SQL。 | 更新元数据。 |
代码示例:
name 配置示例:
"name": "postgresqlwriter"
connection.jdbcUrl 配置示例:
"jdbcUrl": "jdbc:postgresql://192.168.1.103:5432/target_db"
connection.table 配置示例:
"table": ["sales_export"]
username 配置示例:
"username": "datax"
password 配置示例:
"password": "secure_pwd"
column 配置示例:
"column": ["pid", "pname", "price"]
writeMode 配置示例:
"writeMode": "insert"
updateKey 配置示例:
"updateKey": "product_id"
batchSize 配置示例:
"batchSize": 1000
preSql 配置示例:
"preSql": ["DELETE FROM temp WHERE flag=1"]
postSql 配置示例:
"postSql": ["INSERT INTO log VALUES ('ok', now())"]
5.6 hdfswriter:HDFS 文件写入
| 参数名称 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
| name | "name": "hdfswriter" | 指定使用 hdfswriter 插件。 | 用于将数据导出到 HDFS。 |
| path | "path": "/user/data/output" | HDFS 目标路径。 | 路径必须存在或可通过 mkdir 创建。 |
| fileName | "fileName": "part" | 生成的文件前缀名。 | 实际文件名为 fileName + taskId。 |
| defaultFS | "defaultFS": "hdfs://namenode:9000" | HDFS 的 NameNode 地址。 | 必须与集群一致。 |
| fileType | "fileType": "text" 或 "orc" 或 "parquet" | 输出文件格式。 | text 为 CSV/TSV,orc/parquet 为列式存储。 |
| fieldDelimiter | "fieldDelimiter": "," | 字段分隔符(仅 text)。 | 避免与数据内容冲突。 |
| encoding | "encoding": "UTF-8" | 文件编码。 | 确保下游可读。 |
| writeMode | "writeMode": "append" 或 "nonConflict" | 写入模式:append(追加)、nonConflict(路径不存在时创建) | append 需目标路径存在。 |
| compress | "compress": "gzip" 或 "bz2" | 压缩格式。 | 提升存储效率,增加 CPU 开销。 |
| column | "column": ["col1", "col2"] | 字段定义(text 文件必需)。 | 对于 orc/parquet,可从 schema 获取。 |
代码示例:
name 配置示例:
"name": "hdfswriter"
path 配置示例:
"path": "/output/orders"
fileName 配置示例:
"fileName": "order_data"
defaultFS 配置示例:
"defaultFS": "hdfs://hdfs-nn:9000"
fileType 配置示例:
"fileType": "text"
fieldDelimiter 配置示例:
"fieldDelimiter": "\\u0001"
encoding 配置示例:
"encoding": "GBK"
writeMode 配置示例:
"writeMode": "nonConflict"
compress 配置示例:
"compress": "gzip"
column 配置示例:
"column": ["id", "name", "age"]
5.7 txtfilewriter:文本文件写入
| 参数名称 | 语法 | 用途 | 注意事项 |
|---|---|---|---|
| name | "name": "txtfilewriter" | 指定使用 txtfilewriter 插件。 | 用于导出数据到本地文件。 |
| path | "path": "/tmp/output/" | 本地目标目录路径。 | 路径必须存在且 DataX 进程有写权限。 |
| fileName | "fileName": "data" | 生成的文件名前缀。 | 实际文件名如 data.0。 |
| encoding | "encoding": "UTF-8" | 文件编码。 | 匹配下游系统编码。 |
| fieldDelimiter | "fieldDelimiter": "," | 字段分隔符。 | 如 "\t" 表示 Tab 分隔。 |
| lineDelimiter | "lineDelimiter": "\n" | 行分隔符。 | Windows 系统建议 \r\n。 |
| writeMode | "writeMode": "append" 或 "truncate" | 写入模式:append(追加)、truncate(清空后写入) | truncate 会删除原文件内容。 |
| format | "format": "yyyy-MM-dd HH:mm:ss" | 日期字段的输出格式。 | 用于 Date 类型字段。 |
| nullFormat | "nullFormat": "\\N" | NULL 值的表示字符串。 | 避免空字符串歧义。 |
代码示例:
name 配置示例:
"name": "txtfilewriter"
path 配置示例:
"path": "/data/export/"
fileName 配置示例:
"fileName": "user_export"
encoding 配置示例:
"encoding": "GBK"
fieldDelimiter 配置示例:
"fieldDelimiter": "\t"
lineDelimiter 配置示例:
"lineDelimiter": "\r\n"
writeMode 配置示例:
"writeMode": "truncate"
format 配置示例:
"format": "dd/MM/yyyy"
nullFormat 配置示例:
"nullFormat": "NULL"
第六章:数据类型与字段映射
理解不同数据源之间的类型转换规则与字段匹配机制。
6.1 DataX 通用数据类型
DataX 定义了一套通用的中间数据类型,用于在不同数据源之间进行转换。所有 Reader 将源数据转换为这些类型,Writer 再将其转换为目标数据源的类型。
| DataX 类型 | 说明 | 对应常见数据源类型示例 | 注意事项 |
|---|---|---|---|
| String | 字符串类型,用于存储文本数据。 | MySQL: VARCHAR, CHAR, TEXT;Oracle: VARCHAR2, CLOB;SQL Server: NVARCHAR, TEXT;PostgreSQL: TEXT, VARCHAR;HDFS: 所有文本字段 | 最常用类型,可容纳数字、日期等,但需注意目标端解析。 |
| Long | 64位有符号整数。 | MySQL: BIGINT, INT;Oracle: NUMBER(10,0);SQL Server: BIGINT, INT;PostgreSQL: BIGINT, INTEGER | 注意源数据是否超出范围(-2^63 ~ 2^63-1)。 |
| Double | 双精度浮点数。 | MySQL: DOUBLE, FLOAT;Oracle: NUMBER;SQL Server: FLOAT, REAL;PostgreSQL: DOUBLE PRECISION | 存在精度损失风险,不适用于高精度金额场景。 |
| Boolean | 布尔值(true/false)。 | MySQL: TINYINT(1), BOOLEAN;Oracle: NUMBER(1);SQL Server: BIT;PostgreSQL: BOOLEAN | 源数据通常以 0/1 或 Y/N 表示,需在 reader 中正确解析。 |
| Date | 日期时间类型,包含年月日时分秒。 | MySQL: DATETIME, TIMESTAMP;Oracle: DATE, TIMESTAMP;SQL Server: DATETIME, DATETIME2;PostgreSQL: TIMESTAMP | 默认格式为 yyyy-MM-dd HH:mm:ss,可通过 format 参数调整。 |
| Bytes | 字节数组,用于二进制数据。 | MySQL: BLOB, VARBINARY;Oracle: BLOB, RAW;SQL Server: VARBINARY, IMAGE;PostgreSQL: BYTEA | 处理图片、文件等二进制内容,性能开销较大。 |
6.2 不同数据源间的数据类型映射规则
| 源数据源 → 目标数据源 | 推荐映射规则 | 示例 | 注意事项 |
|---|---|---|---|
| MySQL → Oracle | VARCHAR → VARCHAR2;INT → NUMBER;DATETIME → DATE;TEXT → CLOB | name (VARCHAR) → NAME (VARCHAR2);create_time (DATETIME) → CREATE_TIME (DATE) | Oracle 表名和字段名默认大写,建议在配置中统一。 |
| Oracle → MySQL | VARCHAR2 → VARCHAR;NUMBER → DECIMAL 或 BIGINT;DATE → DATETIME;CLOB → LONGTEXT | EMP_NAME (VARCHAR2) → emp_name (VARCHAR);SALARY (NUMBER) → salary (DECIMAL(10,2)) | NUMBER 若有小数,目标应为 DECIMAL;若为整数,可用 BIGINT。 |
| SQL Server → PostgreSQL | NVARCHAR → TEXT;INT → INTEGER;DATETIME → TIMESTAMP;BIT → BOOLEAN | CustomerName (NVARCHAR) → customer_name (TEXT);IsActive (BIT) → is_active (BOOLEAN) | PostgreSQL 对大小写敏感,建议使用小写字段名。 |
| HDFS (Text) → MySQL | 文本字段 → String → 目标类型 | ”1001” (String) → id (INT);“2025-01-01 12:00:00” (String) → create_time (DATETIME) | 文本中的数字/日期需确保格式正确,否则转换失败。 |
| MySQL → HDFS (ORC/Parquet) | INT → Long;VARCHAR → String;DATETIME → Date | user_id (INT) → user_id (Long);name (VARCHAR) → name (String) | ORC/Parquet 支持强类型,建议在 writer 中明确定义 schema。 |
| PostgreSQL → SQL Server | TEXT → NVARCHAR(MAX);TIMESTAMP → DATETIME2;BOOLEAN → BIT | description (TEXT) → Description (NVARCHAR(MAX));active (BOOLEAN) → Active (BIT) | 注意 SQL Server 的 NVARCHAR 需指定长度或使用 MAX。 |
6.3 字段映射与别名配置
| 配置方式 | 说明 | 注意事项 |
|---|---|---|
| 直接映射 | 字段名完全一致,DataX 自动匹配。 | 最简单,但要求源和目标字段名相同。 |
| 顺序映射 | 字段名不同,但顺序一致,按位置匹配。 | 风险高,一旦顺序错乱会导致数据错位,不推荐。 |
| 显式映射(推荐) | 在 column 中使用对象形式,明确指定 name 和 type。 | 最安全,可处理字段名不一致、顺序不同、类型转换等情况。 |
| 使用 SQL 别名 | 在 querySql 中使用 AS 为字段设置别名,使其与目标字段名一致。 | 灵活,可在 reader 端完成映射,writer 可直接使用标准名。 |
| 忽略字段 | reader 读取多字段,writer 只写入部分字段。 | 确保 writer 的 column 是 reader 输出的子集。 |
| 常量字段 | writer 中添加固定值字段(需 writer 支持)。 | 并非所有 writer 都支持,需查看具体插件文档。 |
代码示例:
直接映射:
// reader
"column": ["id", "name"]
// writer
"column": ["id", "name"]
顺序映射:
// reader
"column": ["user_id", "user_name"]
// writer
"column": ["id", "name"]
显式映射(推荐):
"column": [
{ "name": "user_id", "type": "long", "index": 0 },
{ "name": "user_name", "type": "string", "index": 1 }
]
// writer
"column": ["id", "name"]
使用 SQL 别名:
"querySql": "SELECT user_id AS id, user_name AS name FROM users"
忽略字段:
// reader
"column": ["id", "name", "age"]
// writer - 只写入部分字段
"column": ["id", "name"]
常量字段:
// writer column
"column": ["id", "name", { "value": "default", "type": "string" }]
6.4 常见类型转换错误与解决方案
| 错误现象 | 可能原因 | 解决方案 | 验证方法 |
|---|---|---|---|
| 数字格式错误 NumberFormatException | 源字段包含非数字字符(如空格、逗号、null)。 | 1. 在 querySql 中使用 TRIM()、REPLACE() 清理数据。2. 使用 CASE WHEN 处理 null 或异常值。3. 在 where 条件中过滤脏数据。 | SELECT TRIM(price) FROM sales WHERE price REGEXP '[^0-9.]' |
| 日期解析失败 ParseException | 日期字符串格式与默认 yyyy-MM-dd HH:mm:ss 不符。 | 1. 在 reader 的 column 中指定 type: “date” 和 format。2. 在 querySql 中使用数据库函数转换日期格式。 | "column": [{ "name": "dt", "type": "date", "format": "dd/MM/yyyy" }] |
| 字段长度超限 | 源字符串长度超过目标字段定义。 | 1. 修改目标表结构,扩大字段长度。2. 在 querySql 中使用 SUBSTR() 或 LEFT() 截断。3. 配置 errorLimit 容忍少量错误。 | SELECT LEFT(description, 255) FROM articles |
| 精度丢失 | Double 类型写入 DECIMAL 时四舍五入或溢出。 | 1. 源端使用 DECIMAL 类型读取。2. 在 querySql 中确保数值精度。3. 目标字段定义足够精度,如 DECIMAL(18,4)。 | 避免用 Double 处理金额,应使用 String 或 Long(以分为单位)。 |
| NULL 值处理异常 | 目标字段不允许 NULL,但源数据为 NULL。 | 1. 在 querySql 中使用 IFNULL()、COALESCE() 提供默认值。2. 修改目标表允许 NULL。3. 在 writer 配置 nullFormat。 | "querySql": "SELECT COALESCE(age, 0) FROM users" |
| 编码乱码 | 源文件/数据库编码与配置的 encoding 不一致。 | 1. 确认源数据真实编码(如 UTF-8, GBK, ISO-8859-1)。2. 在 reader 和 writer 中正确设置 encoding 参数。3. 使用支持多编码的工具预处理文件。 | file -i filename.txt 或数据库 SHOW CREATE TABLE 查看编码。 |
6.5 最佳实践:安全高效的数据映射
| 实践 | 说明 | 示例 |
|---|---|---|
| 显式定义字段 | 始终在 reader 和 writer 的 column 中明确列出字段,避免使用 *。 | "column": ["id", "name", "email"] |
| 使用对象形式配置 | 对于复杂映射,使用 {name, type, index} 对象形式,提高可读性和安全性。 | "column": [{ "name": "src_id", "type": "long", "index": 0 }, { "name": "full_name", "type": "string", "index": 1 }] |
| 统一命名规范 | 在团队内约定字段命名规范(如全小写、下划线分隔),减少映射复杂度。 | user_name 而非 UserName 或 userName |
| 预处理脏数据 | 在 querySql 中清洗数据,确保输入到 DataX 的数据是干净的。 | SELECT TRIM(UPPER(email)), COALESCE(age, 0) FROM users WHERE status = 'active' |
| 合理设置 errorLimit | 根据业务容忍度配置错误记录数,避免因少量脏数据导致整个任务失败。 | "errorLimit": { "record": 10 } |
| 测试验证 | 在正式运行前,使用小数据集测试类型映射和字段转换是否正确。 | 同步前100条记录,检查目标端数据是否符合预期。 |
| 文档化映射关系 | 维护一份源表到目标表的字段映射文档,便于维护和审计。 | 创建 Excel 表格,记录字段名、类型、转换规则、负责人等。 |
第七章:性能调优与高级配置
提升同步效率,掌握高可用与容错配置。
7.1 核心性能参数调优
| 参数名称 | 语法位置 | 默认值 | 推荐值 | 用途 | 调优说明 |
|---|---|---|---|---|---|
| job.setting.speed.channel | job.setting.speed.channel | 1 | 4 ~ 32(根据 CPU 核数) | 并发通道数,控制并发读写任务数量。 | 增加 channel 可提升吞吐量,但会增加数据库和网络压力。建议从 4 开始逐步增加,观察数据库负载和网络带宽。 |
| job.setting.speed.record | job.setting.speed.record | - | 10000 ~ 100000 | 限流:每秒传输的记录数。 | 用于控制整体传输速度,避免对源库造成过大压力。与 channel 配合使用,可实现精确限速。 |
| job.setting.speed.byte | job.setting.speed.byte | - | 10485760 ~ 104857600 (10MB ~ 100MB) | 限流:每秒传输的字节数。 | 适用于大字段(如文本、二进制)同步,防止网络拥塞。 |
| job.setting.errorLimit.record | job.setting.errorLimit.record | 0 | 10 ~ 100 | 允许的最大错误记录数。 | 设置为 0 时,任何错误都会导致任务失败。生产环境建议设置较小容忍值,避免因少量脏数据中断任务。 |
| job.setting.errorLimit.percentage | job.setting.errorLimit.percentage | 0.0 | 0.01 ~ 0.05 | 允许的错误记录百分比。 | 与 record 结合使用,更灵活地控制错误容忍度。 |
7.2 Reader/Writer 插件性能参数
| 插件 | 参数名称 | 默认值 | 推荐值 | 用途 | 调优说明 |
|---|---|---|---|---|---|
| mysqlreader / oraclereader / sqlserverreader / postgresqlreader | fetchSize | 1024 | 512 ~ 4096 | 每次从数据库 fetch 的记录数。 | 增大 fetchSize 可减少网络往返次数,提升读取效率。但过大会增加内存消耗,建议根据单条记录大小调整。 |
| mysqlreader / oraclereader / sqlserverreader / postgresqlreader | splitPk | null | 主键或唯一索引字段 | 用于数据分片的字段,实现并发读取。 | 关键参数。必须选择高基数、分布均匀的整数主键(如自增 ID),才能有效分片。避免使用字符串或低基数字段。 |
| mysqlreader / oraclereader / sqlserverreader / postgresqlreader | querySql | null | 自定义高效 SQL | 自定义查询语句。 | 避免 SELECT *,只查询必要字段。使用 WHERE 过滤数据,减少传输量。确保 SQL 有高效索引支持。 |
| mysqlwriter / oraclewriter / sqlserverwriter / postgresqlwriter | batchSize | 1024 | 512 ~ 2048 | 每批次提交的记录数。 | 增大 batchSize 可减少事务提交次数,提升写入性能。但过大会导致事务过长、内存占用高,甚至超时。需根据数据库性能调整。 |
| mysqlwriter / oraclewriter / sqlserverwriter / postgresqlwriter | writeMode | insert | insert | 写入模式。 | insert 性能最佳。避免使用 replace(先删后插)和 update(逐条更新),除非业务必需。 |
| mysqlwriter / oraclewriter / sqlserverwriter / postgresqlwriter | session | null | set autocommit=0 / set unique_checks=0 / set foreign_key_checks=0 | 设置数据库会话参数。 | 在 preSql 执行前关闭唯一性检查、外键约束等,可大幅提升写入速度。任务完成后在 postSql 中恢复。 |
| hdfswriter / txtfilewriter | compress | null | gzip, snappy | 输出文件压缩格式。 | 启用压缩可显著减少磁盘 I/O 和网络传输,但增加 CPU 开销。snappy 速度更快,gzip 压缩率更高。 |
| hdfswriter / txtfilewriter | fieldDelimiter | , | \u0001 (SOH) | 字段分隔符。 | 使用非常见字符(如 \u0001)可避免与数据内容冲突,减少解析错误。 |
7.3 JVM 与系统级优化
| 优化项 | 配置方法 | 说明 | 注意事项 |
|---|---|---|---|
| JVM 内存 | 修改 datax.py 中的 -Xms 和 -Xmx | 增加 DataX 进程的堆内存。 | 默认 -Xms1g -Xmx1g 可能不足。对于大任务,建议设置为 -Xms4g -Xmx8g 或更高。避免设置过大导致 Full GC 频繁。 |
| 垃圾回收器 | 添加 JVM 参数 -XX:+UseG1GC | 使用 G1 垃圾回收器替代默认的 Parallel GC。 | G1 更适合大内存、低延迟场景,可减少长时间停顿。 |
| 操作系统文件句柄 | ulimit -n 65536 | 增加单进程可打开的文件句柄数。 | 处理大量小文件(如 HDFS 分片)时,可能耗尽句柄。需在系统层面调整。 |
| 网络带宽 | - | 确保网络链路充足。 | 数据同步是 I/O 密集型操作,千兆或万兆网络是基础。避免与其他高带宽应用争抢。 |
| 磁盘 I/O | 使用 SSD 或高性能 RAID | 提升本地临时文件读写速度。 | DataX 在转换过程中可能生成临时文件,SSD 可显著提升性能。 |
7.4 高级配置功能详解
| 功能 | 配置方式 | 用途 | 示例与说明 |
|---|---|---|---|
| 动态参数 | 使用 ${param} 占位符 | 实现配置文件模板化,支持外部传参。 | "path": "/data/input/${date}";执行时:python datax.py job.json -p "-Ddate=2025-01-01";可用于动态指定日期、表名等。 |
| 多数据源并行 | 在 connection 数组中配置多个 JDBC URL | 读取多个数据库实例的数据。 | DataX 会自动合并结果。 |
| 复杂 WHERE 条件 | 在 where 中使用子查询或函数 | 实现精细化数据过滤。 | "where": "create_time >= DATE_SUB(NOW(), INTERVAL 1 DAY)";注意:避免在 where 中使用性能差的函数或全表扫描。 |
| 自定义 SQL (querySql) | 在 reader 中使用 querySql | 执行复杂 JOIN、聚合或数据清洗。 | 使用 querySql 后,table 和 column 参数将被忽略。 |
| 脏数据收集 | 配置 errorLimit 并启用日志 | 记录同步失败的记录。 | 失败记录会输出到日志文件,可用于后续分析和重试。建议结合 record 和 percentage 设置合理阈值。 |
| 任务依赖与调度 | 与外部调度系统(如 Airflow, DolphinScheduler)集成 | 实现自动化、定时执行。 | 将 DataX 任务封装为 Shell 操作符,在调度系统中定义依赖关系和执行计划。 |
多数据源并行配置示例:
"connection": [
{
"jdbcUrl": ["jdbc:mysql://host1:3306/db"],
"table": ["t1"]
},
{
"jdbcUrl": ["jdbc:mysql://host2:3306/db"],
"table": ["t1"]
}
]
自定义 SQL 配置示例:
"querySql": "SELECT a.id, b.name FROM table_a a JOIN table_b b ON a.bid = b.id WHERE a.status = 1"
7.5 性能调优步骤与监控
| 步骤 | 操作 | 工具与方法 |
|---|---|---|
| 1. 基准测试 | 使用小数据集运行任务,记录初始性能。 | time python datax.py job.json;观察日志中的 totalCost 和 bytesPerSecond。 |
| 2. 瓶颈分析 | 判断是 CPU、内存、磁盘 I/O、网络还是数据库成为瓶颈。 | top, htop(查看 CPU 和内存使用率);iostat(查看磁盘 I/O);iftop, nethogs(查看网络流量);数据库监控:慢查询日志、SHOW PROCESSLIST。 |
| 3. 参数迭代 | 逐步调整 channel、batchSize、fetchSize 等参数。 | 每次只调整一个参数,对比性能变化。记录每次测试的吞吐量和资源消耗。 |
| 4. 数据库优化 | 确保源端和目标端有合适的索引。 | 为 splitPk 字段创建索引。避免在 WHERE 条件中使用无索引字段。 |
| 5. 监控与告警 | 将 DataX 集成到监控系统。 | 使用 ELK 收集日志,Prometheus + Grafana 监控关键指标(任务时长、错误数、传输速率)。设置失败告警。 |
| 6. 定期复盘 | 回顾历史任务性能,持续优化。 | 分析月度/季度任务报告,识别性能下降趋势,优化配置。 |
7.6 常见性能问题与解决方案
| 问题现象 | 可能原因 | 解决方案 |
|---|---|---|
| 读取速度极慢 | 1. 未配置 splitPk,单线程读取。2. fetchSize 过小。3. 源 SQL 无索引或性能差。 | 1. 为大表配置合适的 splitPk。2. 增大 fetchSize(如 4096)。3. 优化 querySql,添加索引。 |
| 写入速度瓶颈 | 1. batchSize 过小。2. 目标表有过多索引或约束。3. 使用 replace 或 update 模式。 | 1. 增大 batchSize(如 2048)。2. 临时禁用非必要索引和约束。3. 改用 insert 模式。 |
| 内存溢出 (OOM) | 1. JVM 堆内存不足。2. channel 或 batchSize 过大。3. 处理超大字段(如 CLOB)。 | 1. 增加 -Xmx 值。2. 降低 channel 或 batchSize。3. 分批处理大字段或优化数据结构。 |
| 数据库连接超时 | 1. 任务执行时间过长。2. 数据库连接池配置过期。 | 1. 优化任务性能或拆分大任务。2. 在 JDBC URL 中添加超时参数(如 ?connectTimeout=60000&socketTimeout=60000)。 |
| 网络传输慢 | 1. 网络带宽不足。2. 跨地域传输延迟高。 | 1. 升级网络或错峰执行。2. 尽量在同地域机房部署 DataX 和目标存储。 |
第八章:实战案例与最佳实践
结合真实场景,掌握典型同步任务的配置方案。
8.1 实战案例一:MySQL 到 HDFS 的每日增量同步
业务背景
某电商平台需要将 MySQL 订单库中的每日新增订单数据(约 500 万条)同步到 HDFS,供 Hive 数仓进行 T+1 分析。要求保证数据一致性、高效传输,并以 ORC 格式存储。
解决方案
{
"job": {
"content": [
{
"reader": {
"name": "mysqlreader",
"parameter": {
"username": "dx_reader",
"password": "read_pass",
"connection": [
{
"jdbcUrl": ["jdbc:mysql://mysql-host:3306/orders_db"],
"table": ["orders"],
"where": "create_time >= '${date} 00:00:00' AND create_time < DATE_ADD('${date}', INTERVAL 1 DAY)"
}
],
"column": ["order_id", "user_id", "amount", "status", "create_time"],
"splitPk": "order_id",
"fetchSize": 4096
}
},
"writer": {
"name": "hdfswriter",
"parameter": {
"defaultFS": "hdfs://hdfs-nn:9000",
"fileType": "orc",
"path": "/data/ods/orders/dt=${date}",
"fileName": "part",
"column": [
{"name": "order_id", "type": "LONG"},
{"name": "user_id", "type": "LONG"},
{"name": "amount", "type": "DOUBLE"},
{"name": "status", "type": "STRING"},
{"name": "create_time", "type": "TIMESTAMP"}
],
"writeMode": "nonConflict",
"compress": "snappy"
}
}
}
],
"setting": {
"speed": {
"channel": 8
},
"errorLimit": {
"record": 10
}
}
}
}
关键配置说明
- 动态参数
${date}:通过外部传参指定同步日期(如 2025-01-01),实现每日调度。 - 增量同步:利用 where 条件过滤当日数据,避免全量扫描。
- 分片读取:
splitPk: order_id实现 8 个 channel 并发读取。 - 高效存储:HDFS 使用 ORC + snappy 压缩,节省存储空间并提升查询性能。
- 错误容忍:允许 10 条错误记录,防止脏数据中断任务。
8.2 实战案例二:Oracle 全量迁移至 PostgreSQL
业务背景
企业需将旧系统的 Oracle 数据库(包含用户、订单、商品等 20 张表)整体迁移至 PostgreSQL,要求停机时间最短,数据零丢失。
解决方案
采用 双轨并行 + 最终校验 策略:
- 预迁移阶段:使用 DataX 并行迁移所有静态表(如商品、分类)。
- 增量追平阶段:对核心交易表(用户、订单)开启 CDC 或日志解析,实时捕获变更。
- 切换阶段:业务停机,执行最后一次增量同步,确保数据一致后切换应用。
核心配置片段(以用户表为例):
"reader": {
"name": "oraclereader",
"parameter": {
"username": "SRC_USER",
"password": "src_pass",
"connection": [
{
"jdbcUrl": "jdbc:oracle:thin:@//oracle-host:1521/ORCL",
"table": ["USERS"]
}
],
"column": ["USER_ID", "USERNAME", "EMAIL", "REG_TIME"],
"fetchSize": 2048
}
},
"writer": {
"name": "postgresqlwriter",
"parameter": {
"username": "tgt_user",
"password": "tgt_pass",
"connection": [
{
"jdbcUrl": "jdbc:postgresql://pg-host:5432/target_db",
"table": ["users"]
}
],
"column": ["user_id", "username", "email", "reg_time"],
"preSql": ["DELETE FROM users"],
"postSql": ["ANALYZE users"],
"session": [
"SET session_replication_role = 'replica'"
],
"batchSize": 1000,
"writeMode": "insert"
}
}
关键实践
- 会话优化:通过 session 禁用外键和触发器,大幅提升写入速度。
- 预清理:preSql 删除目标表数据,确保干净环境。
- 并发控制:根据 PostgreSQL 写入能力设置 channel=4 和 batchSize=1000。
- 数据校验:迁移后使用 checksum 工具(如 pg_checksums)或自定义脚本验证数据一致性。
8.3 实战案例三:多源异构数据整合到数据湖
业务背景
金融公司需整合来自 MySQL(客户信息)、SQL Server(交易流水)、MongoDB(行为日志)三个系统的数据,统一写入 HDFS 形成宽表,用于风控模型训练。
解决方案
使用 DataX 多作业组合 + ETL 脚本实现:
- 分别配置三个 DataX 作业,将各源数据导出到 HDFS 临时目录。
- 使用 Hive SQL 或 Spark 作业进行关联、清洗、去重,生成最终宽表。
- 将宽表注册到 Hive Metastore,供下游使用。
MySQL → HDFS 配置示例:
"reader": {
"name": "mysqlreader",
"parameter": {
"connection": [{
"jdbcUrl": ["jdbc:mysql://host:3306/crm"],
"table": ["customers"]
}],
"column": ["cust_id", "name", "age", "city"],
"where": "last_updated >= '${last_run}'"
}
},
"writer": {
"name": "hdfswriter",
"parameter": {
"path": "/tmp/staging/customers",
"fileName": "cust",
"fileType": "text",
"fieldDelimiter": "\u0001",
"writeMode": "append"
}
}
关键实践
- 分步处理:先独立抽取,再集中转换,降低单个任务复杂度。
- 格式统一:所有中间文件使用
\u0001作为分隔符,避免解析冲突。 - 增量更新:每个源系统基于 last_updated 字段进行增量同步。
- 元数据管理:记录每次同步的 start_time 和 end_time,便于追溯。
8.4 最佳实践总结
| 类别 | 最佳实践 | 说明 |
|---|---|---|
| 设计阶段 | 明确同步策略 | 区分全量、增量、CDC,选择合适方案。增量优先,减少资源消耗。 |
| 设计阶段 | 合理分片 (splitPk) | 选择高基数、分布均匀的整数主键作为 splitPk,是提升并发性能的关键。 |
| 设计阶段 | 最小化数据集 | 在 querySql 中只 SELECT 必要字段,使用 WHERE 过滤无关数据。 |
| 配置阶段 | 显式声明字段 | 始终在 column 中列出字段,避免 *,提高可维护性。 |
| 配置阶段 | 使用对象映射 | 对于复杂映射,使用 {name, type, index} 明确对应关系。 |
| 配置阶段 | 启用压缩 | HDFS/TXT 输出建议启用 gzip 或 snappy 压缩,节省 I/O。 |
| 执行阶段 | 动态参数化 | 使用 ${param} 实现配置模板化,便于调度系统集成。 |
| 执行阶段 | 设置 errorLimit | 生产环境配置合理的错误容忍阈值(如 record: 10),避免任务轻易失败。 |
| 执行阶段 | 监控与告警 | 集成日志监控系统,实时跟踪任务状态、吞吐量和错误率。 |
| 运维阶段 | 定期性能测试 | 每月对核心任务进行性能压测,及时发现瓶颈。 |
| 运维阶段 | 版本管理 | 将 DataX 作业配置文件纳入 Git 管理,记录变更历史。 |
| 运维阶段 | 文档化 | 维护数据字典和同步流程图,便于团队协作和故障排查。 |
8.5 常见陷阱与规避方法
| 陷阱 | 表现 | 规避方法 |
|---|---|---|
| 字符集乱码 | 中文显示为 ??? 或乱码。 | 1. 确认源数据库和文件的真实编码。2. 在 reader/writer 中正确设置 encoding 参数(如 "encoding": "UTF-8")。 |
| 主键冲突 | mysqlwriter 报错 Duplicate entry。 | 1. 检查 writeMode 是否为 insert 但目标表已存在数据。2. 使用 preSql 清理目标表或改用 replace 模式(谨慎)。 |
| 数据截断 | 字符串被截断或数字精度丢失。 | 1. 检查目标字段长度是否足够。2. 避免用 Double 存储高精度数值,应使用 String 或 Decimal。 |
| 空指针异常 | 日志中出现 NullPointerException。 | 1. 检查 column 配置是否与实际数据匹配。2. 确保 splitPk 字段不为空且为数值类型。 |
| 连接池耗尽 | 数据库报错 Too many connections。 | 1. 降低 channel 数量。2. 优化数据库连接池配置。3. 使用连接代理(如 ProxySQL)。 |
| 大字段性能差 | 同步包含 CLOB/BLOB 的表时极慢。 | 1. 拆分任务,单独处理大字段。2. 考虑是否真的需要同步二进制数据。 |
第九章:插件开发与扩展
了解如何开发自定义 reader 或 writer 插件。
9.1 插件开发概述
DataX 采用模块化插件架构,所有数据源的读写能力都通过独立的插件实现。核心框架负责任务调度、数据传输、错误处理等通用功能,而具体的读写逻辑由插件完成。
插件类型
| 类型 | 职责 | 示例 |
|---|---|---|
| Reader | 从数据源读取数据,转换为 DataX 通用类型,发送给通道。 | mysqlreader, oraclereader, hdfsfilereader |
| Writer | 从通道接收数据,转换为目标数据源的类型,写入目标存储。 | mysqlwriter, hdfswriter, txtfilewriter |
开发准备
- 环境要求:
- JDK 1.8 或更高
- Maven 3.2+
- Git
- 获取源码:
git clone https://github.com/alibaba/DataX.git
cd DataX
- 项目结构:
DataX/
├── datax-core/ # 核心框架
├── datax-common/ # 公共模块
├── datax-reader/ # 所有 Reader 插件
│ ├── mysqlreader/
│ ├── hdfsfilereader/
│ └── ...
├── datax-writer/ # 所有 Writer 插件
│ ├── mysqlwriter/
│ ├── hdfswriter/
│ └── ...
└── target/ # 打包输出目录
9.2 开发自定义 Reader 插件
步骤 1:创建 Reader 模块
使用 Maven 创建新模块:
cd DataX
mvn archetype:generate \
-DgroupId=com.company.datax \
-DartifactId=customreader \
-DarchetypeArtifactId=maven-archetype-quickstart \
-DinteractiveMode=false
步骤 2:实现核心接口
需要实现三个核心类:
1. CustomReader (Job) — 负责全局初始化、预检查和分片:
public class CustomReader extends Reader {
public static class Job extends Reader.Job {
private Configuration originalConfig;
@Override
public void init() {
this.originalConfig = super.getPluginJobConf();
// 参数预检查
this.validate();
}
private void validate() {
String endpoint = originalConfig.getString("endpoint");
Preconditions.checkArgument(StringUtils.isNotBlank(endpoint), "请配置 endpoint");
// 其他参数校验
}
@Override
public List<Configuration> split(int adviceNumber) {
// 数据分片逻辑
List<Configuration> configurations = new ArrayList<>();
for (int i = 0; i < adviceNumber; i++) {
Configuration cfg = originalConfig.clone();
cfg.set("sliceId", i);
configurations.add(cfg);
}
return configurations;
}
@Override
public void destroy() {
// 资源释放
}
}
public static class Task extends Reader.Task {
// 任务执行逻辑见下一步
}
}
2. CustomReader.Task — 负责实际的数据读取:
public static class Task extends Reader.Task {
private Configuration readerSliceConfig;
private CustomClient client; // 假设的客户端
@Override
public void init() {
this.readerSliceConfig = super.getPluginJobConf();
String endpoint = readerSliceConfig.getString("endpoint");
String accessKey = readerSliceConfig.getString("accessKey");
// 初始化客户端
this.client = new CustomClient(endpoint, accessKey);
}
@Override
public void startRead(RecordSender recordSender) {
String sliceId = readerSliceConfig.getString("sliceId");
List<CustomRecord> records = client.queryData(sliceId); // 从数据源读取
for (CustomRecord srcRecord : records) {
// 创建 DataX 记录
Record record = recordSender.createRecord();
try {
// 字段映射
record.addColumn(new StringColumn(srcRecord.getId()));
record.addColumn(new LongColumn(srcRecord.getTimestamp()));
record.addColumn(new StringColumn(srcRecord.getContent()));
// 发送给通道
recordSender.sendToWriter(record);
} catch (Exception e) {
// 捕获异常,记录脏数据
super.getTaskPluginCollector().collectDirtyRecord(record, e.getMessage());
}
}
}
@Override
public void destroy() {
if (client != null) {
client.close();
}
}
}
3. plugin.json — 在 src/main/resources/ 下创建,声明插件信息:
{
"name": "customreader",
"class": "com.company.datax.CustomReader",
"description": "读取自定义数据源",
"developer": "Your Company"
}
9.3 开发自定义 Writer 插件
步骤 1:创建 Writer 模块
mvn archetype:generate \
-DgroupId=com.company.datax \
-DartifactId=customwriter \
-DarchetypeArtifactId=maven-archetype-quickstart \
-DinteractiveMode=false
步骤 2:实现核心接口
1. CustomWriter (Job):
public class CustomWriter extends Writer {
public static class Job extends Writer.Job {
private Configuration writerSliceConfig;
@Override
public void init() {
this.writerSliceConfig = super.getPluginJobConf();
// 参数校验
String server = writerSliceConfig.getString("server");
Preconditions.checkArgument(StringUtils.isNotBlank(server), "请配置 server");
}
@Override
public void prepare() {
// 任务前准备,如创建表、清理数据
String preSql = writerSliceConfig.getString("preSql");
if (StringUtils.isNotBlank(preSql)) {
executeSql(preSql);
}
}
@Override
public void destroy() {
// 任务后清理
String postSql = writerSliceConfig.getString("postSql");
if (StringUtils.isNotBlank(postSql)) {
executeSql(postSql);
}
}
}
public static class Task extends Writer.Task {
// 任务执行逻辑见下一步
}
}
2. CustomWriter.Task:
public static class Task extends Writer.Task {
private Configuration writerSliceConfig;
private CustomClient client;
private List<Record> buffer = new ArrayList<>();
private int batchSize;
@Override
public void init() {
this.writerSliceConfig = super.getPluginJobConf();
String server = writerSliceConfig.getString("server");
this.client = new CustomClient(server);
this.batchSize = writerSliceConfig.getInt("batchSize", 1000);
}
@Override
public void prepare() {
// 单个 task 的准备
}
@Override
public void startWrite(RecordReceiver recordReceiver) {
Record record;
while ((record = recordReceiver.getFromReader()) != null) {
// 类型转换
List<Column> columns = record.getColumnList();
CustomRecord targetRecord = new CustomRecord();
targetRecord.setId(columns.get(0).asString());
targetRecord.setTimestamp(columns.get(1).asLong());
targetRecord.setContent(columns.get(2).asString());
buffer.add(targetRecord);
if (buffer.size() >= batchSize) {
flush(); // 批量写入
}
}
// 处理剩余数据
if (!buffer.isEmpty()) {
flush();
}
}
private void flush() {
try {
client.batchInsert(buffer);
buffer.clear();
} catch (Exception e) {
// 收集脏数据
for (Record rec : buffer) {
super.getTaskPluginCollector().collectDirtyRecord(rec, e.getMessage());
}
buffer.clear();
throw DataXException.asDataXException(CustomWriterErrorCode.WRITE_DATA_ERROR, e);
}
}
@Override
public void destroy() {
if (client != null) {
client.close();
}
}
}
3. plugin.json:
{
"name": "customwriter",
"class": "com.company.datax.CustomWriter",
"description": "写入自定义数据源",
"developer": "Your Company"
}
9.4 插件打包与部署
步骤 1:编译打包
在插件模块根目录执行:
mvn clean package -Dmaven.test.skip=true
生成 target/customreader-0.0.1-SNAPSHOT.jar 和 target/customwriter-0.0.1-SNAPSHOT.jar。
步骤 2:部署到 DataX
# 将 jar 包复制到 DataX 的 plugin 目录
cp target/customreader-0.0.1-SNAPSHOT.jar /path/to/datax/plugin/reader/customreader/
cp target/customwriter-0.0.1-SNAPSHOT.jar /path/to/datax/plugin/writer/customwriter/
# 创建目录并解压(如果 jar 包包含资源文件)
mkdir -p /path/to/datax/plugin/reader/customreader/
unzip target/customreader-0.0.1-SNAPSHOT.jar -d /path/to/datax/plugin/reader/customreader/
步骤 3:验证插件
查看插件列表:
python datax.py --help
应能看到 customreader 和 customwriter。
9.5 测试与调试
编写测试配置 test_custom.json:
{
"job": {
"content": [
{
"reader": {
"name": "customreader",
"parameter": {
"endpoint": "http://api.example.com",
"accessKey": "xxxx",
"column": ["id", "ts", "data"]
}
},
"writer": {
"name": "customwriter",
"parameter": {
"server": "custom-server:8080",
"batchSize": 500,
"column": ["id", "ts", "data"]
}
}
}
],
"setting": {
"speed": {
"channel": 2
}
}
}
}
执行测试:
python datax.py test_custom.json
调试技巧:
- 日志分析:查看
log/datax.log,重点关注 ERROR 和 WARN。 - 断点调试:在 IDE 中导入项目,设置断点,以
com.alibaba.datax.core.Engine为入口运行。 - 单元测试:为 Job 和 Task 类编写 JUnit 测试,验证参数解析和核心逻辑。
9.6 最佳实践与注意事项
| 实践 | 说明 |
|---|---|
| 参数校验 | 在 init() 中对必填参数进行严格校验,使用 Preconditions.checkArgument()。 |
| 异常处理 | 捕获所有异常,通过 TaskPluginCollector 收集脏数据,避免任务直接崩溃。 |
| 资源管理 | 在 destroy() 中释放数据库连接、文件句柄等资源,防止泄漏。 |
| 性能优化 | 使用批量写入(batchInsert),避免逐条提交。合理设置 batchSize。 |
| 线程安全 | Task 类是多线程执行的,确保成员变量线程安全或使用局部变量。 |
| 编码规范 | 遵循 DataX 原有代码风格,使用 SLF4J 打印日志。 |
| 文档化 | 在 plugin.json 中提供清晰的描述,并编写使用文档。 |
第十章:常见问题与故障排查
快速定位并解决运行中的典型问题。
10.1 故障排查方法论
遵循**“由外到内、由简到繁”**的原则进行排查:
| 步骤 | 操作 | 工具与方法 |
|---|---|---|
| 1. 观察现象 | 仔细阅读错误日志,定位错误类型和发生位置。 | 查看 log/datax.log,搜索 ERROR、Exception、Caused by。 |
| 2. 复现问题 | 使用最小化配置复现问题,排除干扰因素。 | 简化 job 配置,只保留核心 reader/writer 和必要参数。 |
| 3. 检查配置 | 逐项核对 JSON 配置文件的语法、参数名、值类型。 | 使用 JSON 校验工具(如 jsonlint.com),对照官方文档检查参数。 |
| 4. 验证环境 | 确认网络、权限、资源是否正常。 | ping, telnet, nc 测试网络连通性;检查文件权限、磁盘空间、内存。 |
| 5. 分段测试 | 将任务拆解,分别测试 reader 读取和 writer 写入能力。 | 可临时将 writer 改为 txtfilewriter 输出到本地,验证 reader 是否正常。 |
| 6. 查阅文档 | 搜索官方文档、GitHub Issues、社区论坛。 | DataX GitHub: https://github.com/alibaba/DataX/issues |
10.2 常见错误分类与解决方案
A. 配置类错误
| 错误现象 | 可能原因 | 解决方案 |
|---|---|---|
| JSON parse error / Unexpected character | JSON 语法错误,如缺少逗号、括号不匹配、使用单引号。 | 1. 使用在线 JSON 校验工具检查语法。2. 确保所有字符串用双引号 ” 包围。3. 检查末尾是否有多余的逗号。 |
| No reader found / No writer found | 插件名称拼写错误或插件未正确安装。 | 1. 核对 name 字段,如 mysqlreader 而非 mysql_reader。2. 检查 plugin/reader/ 或 plugin/writer/ 目录下是否存在对应插件目录和 plugin.json。 |
| Required value: xxx is null / 参数缺失 | 必填参数未配置。 | 1. 对照官方文档检查 reader/writer 的必填参数(如 jdbcUrl, table, username)。2. 注意大小写敏感。 |
| Unknown column | column 配置的字段名在源表或目标表中不存在。 | 1. 登录数据库执行 DESC table_name 确认字段名。2. 注意大小写(Oracle 默认大写,MySQL 默认小写)。 |
B. 连接类错误
| 错误现象 | 可能原因 | 解决方案 |
|---|---|---|
| Communications link failure / Connection timed out | 网络不通或数据库服务未启动。 | 1. 使用 telnet host port 或 nc -zv host port 测试端口连通性。2. 确认数据库服务正在运行。3. 检查防火墙或安全组规则是否放行端口。 |
| Access denied for user / Login failed | 用户名、密码错误或权限不足。 | 1. 确认用户名密码正确。2. 检查用户是否有 SELECT(reader)或 INSERT(writer)权限。3. 对于 MySQL,检查用户是否允许从 DataX 服务器 IP 登录。 |
| Unknown database / Invalid object name | 数据库或表名错误。 | 1. 登录数据库确认库名、表名拼写正确。2. 注意大小写和特殊字符。 |
| Too many connections | 数据库连接数达到上限。 | 1. 降低 job 配置中的 channel 数量。2. 联系 DBA 增加数据库 max_connections 配置。 |
C. 数据处理类错误
| 错误现象 | 可能原因 | 解决方案 |
|---|---|---|
| Data truncation / Data too long for column | 字符串长度超过目标字段定义。 | 1. 修改目标表结构,增大字段长度(如 VARCHAR(512))。2. 在 querySql 中使用 SUBSTR(col, 1, 255) 截断数据。3. 配置 errorLimit 容忍少量错误。 |
| Incorrect datetime value / Cannot convert string to date | 日期字符串格式不匹配。 | 1. 在 reader 的 column 中指定 type: “date” 和 format。2. 在 querySql 中使用数据库函数转换日期格式(如 MySQL 的 DATE_FORMAT(create_time, '%Y-%m-%d %H:%i:%s'))。 |
| Duplicate entry ‘xxx’ for key ‘PRIMARY’ | 主键冲突,目标表已存在相同主键的数据。 | 1. 检查 writer 的 writeMode,如果是 insert,需确保目标表为空或无冲突。2. 使用 preSql 在写入前清理数据(如 DELETE FROM table WHERE dt = '2025-01-01')。3. 改用 replace 模式(先删后插)。 |
| Column count doesn’t match value count | reader 输出的字段数与 writer 配置的字段数不一致。 | 1. 确保 reader.column 和 writer.column 的字段数量一致。2. 如果使用 querySql,确保其返回的列数与 column 配置匹配。 |
D. 资源与性能类问题
| 错误现象 | 可能原因 | 解决方案 |
|---|---|---|
| Java heap space / OutOfMemoryError | JVM 堆内存不足。 | 1. 修改 datax.py 脚本,增加 -Xmx 参数(如 -Xmx8g)。2. 降低 channel 数量或 batchSize/fetchSize。3. 避免同步包含超大字段(如 CLOB)的表。 |
| GC overhead limit exceeded | 垃圾回收占用过多 CPU,系统几乎无响应。 | 1. 增加堆内存。2. 添加 JVM 参数 -XX:+UseG1GC 启用 G1 垃圾回收器。3. 优化任务,减少单次处理的数据量。 |
| 同步速度极慢 | 瓶颈可能在源库、目标库、网络或 DataX 配置。 | 1. 检查 channel 是否过小,尝试增加并发。2. 确认是否配置了 splitPk 实现分片读取。3. 检查数据库是否有慢查询、锁等待。4. 使用 iostat, iftop 监控 I/O 和网络。 |
| 任务长时间卡住 | 可能是数据库锁、网络阻塞或死循环。 | 1. 查看数据库当前会话,检查是否有长时间运行的 SQL 或锁。2. 在 DataX 日志中查看最后执行的操作。3. 尝试中断任务,分析中断点。 |
10.3 日志分析技巧
DataX 的日志是排查问题的最重要依据。关键日志片段解析:
2025-01-01 10:00:00 [ERROR] JobContainer - Exception when job run
com.alibaba.datax.common.exception.DataXException: Code:[DBUtilErrorCode-07], Description:[读取数据库数据失败. 请检查您的配置的 column/table/where/querySql 或者向 DBA 寻求帮助.]. - 执行的SQL为: SELECT * FROM users WHERE status = 1
...
Caused by: java.sql.SQLSyntaxErrorException: Unknown column 'status' in 'where clause'
- 时间戳:
2025-01-01 10:00:00定位问题发生时间。 - 错误级别:
[ERROR]表示严重错误。 - 错误来源:
JobContainer是核心调度组件。 - 错误码:
DBUtilErrorCode-07可在文档中查询具体含义。 - 错误描述:明确指出”读取数据库数据失败”。
- 关键信息:执行的 SQL 语句。
- 根本原因 (Caused by):
Unknown column 'status'—— 字段名错误。
解决:检查 users 表是否存在 status 字段,或是否拼写错误(如 state)。
10.4 高级排查工具
| 工具 | 用途 | 命令示例 |
|---|---|---|
| telnet / nc | 测试网络端口连通性。 | telnet mysql-host 3306 / nc -zv pg-host 5432 |
| tcpdump | 抓取网络数据包,分析通信问题。 | tcpdump -i any host db-host and port 3306 -w capture.pcap |
| 数据库监控 | 查看数据库会话、慢查询、锁。 | MySQL: SHOW PROCESSLIST;PostgreSQL: SELECT * FROM pg_stat_activity;Oracle: V$SESSION |
| jstack | 生成 JVM 线程堆栈,分析死锁或卡顿。 | jstack <DataX_PID> > thread_dump.txt |
| jmap | 生成 JVM 堆内存快照,分析内存泄漏。 | jmap -dump:format=b,file=heap.hprof <PID> |
10.5 预防性最佳实践
| 实践 | 说明 |
|---|---|
| 配置版本化 | 将 DataX job 配置文件纳入 Git 管理,便于回滚和审计。 |
| 参数化配置 | 使用 ${param} 实现动态配置,避免硬编码。 |
| 预检查脚本 | 在执行 DataX 前,编写脚本检查数据库连通性、表结构、磁盘空间等。 |
| 灰度发布 | 新任务先用小数据量测试,验证无误后再上线。 |
| 监控告警 | 将 DataX 任务集成到监控系统(如 Prometheus + Grafana),设置失败告警。 |
| 定期维护 | 清理旧日志、归档历史任务、更新插件版本。 |