
你有没有遇到过这种需求领导丢过来一句话“把这套环境里几十张表同步到另一套环境”你数了数六十多张表还带了一堆业务前缀。最开始我的做法很老实为每张表手写一个 DataX 的 json 配置写到第 7 个的时候就已经开始怀疑人生了。后来我抽了一个周末重新把这个流程梳理了一遍用“元数据表 模板 JSON 调度脚本”的方式把同步任务从“每个表手搓配置”变成了“往数据库里插一行配置记录”再批量自动执行。我自己的体验是这个方案改完之后后面再加表基本就是一条 SQL 的事。先说一下适用范围DataX 是阿里开源的一款离线数据同步工具我这次做的是 MySQL 到 MySQL 的批量表同步适合测试环境数据重建、跨库数据归档、生产环境抽取分析数据这类离线批量场景。如果你要的是秒级实时同步DataX 不合适该上监听 binlog 的实时同步方案如果你要的是“几十张、上百张表规则还各不相同”的批量迁移那 DataX 这套灵活配置方案应该能让你少熬几个夜。1. 从一次痛苦的“挨个导表”说起1.1 这个方案到底解决什么问题我先说这次方案要解决的几个核心痛点你对照一下是不是也遇到过表太多。一次要同步 50 张表每张表写一个 json 配置哪怕你复制粘贴也要花掉一上午。表结构不一致。源库和目标库的字段名不同、字段顺序不同有的目标表还比源表多了几个字段需要在同步时做映射。同步规则分裂。有的表要全量覆盖有的表按日期增量有的表只同步 where 条件过滤后的部分行。手动执行太累。白天跑量大影响线上业务只能半夜人工登录服务器一条一条执行。没人愿意维护。半年后业务加了一张新表你还得重新打开文件改 json改完还要担心弄坏别的配置。这套方案的本质就是把“每个任务的差异点”抽出来存到数据库表里把“公共的同步逻辑”用模板和脚本固化下来。你只需要关心这张表从哪里来、到哪里去、用什么规则同步剩下的交给脚本去生成配置、执行任务、记录日志。1.2 为什么是 DataX 而不是 mysqldump、Canal、Kettle我最早尝试过 mysqldump命令简单但问题也很明显它更适合整库或整表的逻辑备份对单表自定义过滤条件、字段映射这些需求支持很弱而且它是单线程逻辑导出几十 GB 的大表跑起来时间很长。另一个方案是 Canal监听 binlog 做实时同步确实很强大但部署维护成本高对批量初始化存量数据这件事帮助不大。Kettle 我也试过图形化界面拖拽确实直观但想要做到上百张表无人值守批量调度图形界面反而成了自动化脚本的累赘。DataX 的优势正好在这个位置它走的是“读插件 写插件 通道并发”的架构MySQL 读写插件开箱即用配置是标准 JSON天然适合脚本动态生成没有界面就适合塞到定时任务里跑。做离线批量表同步它是我用过的最稳、最省心的那一个。2. 方案设计把“死的配置”变成“活的数据”2.1 批量同步的三种常见实现方式我对批量同步的实现方式做过一个简单的分类你可以判断自己卡在哪个阶段。方式一手工为每张表写独立 JSON。这是最原始的方案。优点是没有额外代码缺点是人肉维护成本极高遇到几十张表基本就是灾难。方式二单个 JSON 模板 运行时参数替换。DataX 支持在 job 配置里用${param}占位符启动时通过-Dparamvalue传入。好处是省去大量重复 json 文件但批量执行时还得在外面套一层循环字段映射、同步规则一变模板还是得改。方式三元数据表驱动脚本自动生成 JSON 并执行。把所有差异点存进数据库脚本每次运行时动态拼接同步配置。这个方案最灵活也是我这次要展开讲的。为什么最终选了方式三因为到最后你会发现所有同步任务之间“不一样的地方”其实就那么几项源表名、目标表名、字段列表、过滤条件、写入模式。把这些差异点放到一张表里管理剩下的就是一套固定的生成和执行逻辑。2.2 最终选型元数据表 模板 JSON 调度任务整个方案的架构很简洁分为三层元数据表存每一个同步任务的“个性化配置”比如源表、目标表、字段映射、拆分键、写入模式。模板 JSON定义单表同步的公共结构相当于一张标准合同字段变量留给脚本填充。调度脚本读取元数据表逐行生成临时 JSON 并调用 DataX 执行记录成功失败状态。用一句话总结把 JSON 当代码来维护不如把元数据当数据库来管理。任务多了以后你不需要记住每个 json 文件放哪只需要查询一张表所有同步任务的规则一目了然。2.3 整体运行流程一次完整的批量同步跑下来流程是固定的连接元数据库读取enabled1的同步配置。按配置的优先级顺序遍历每一条记录。根据当前记录的表名、字段、条件动态生成一份 job JSON 写入临时目录。调用python bin/datax.py 临时配置.json执行同步。根据返回码判断成功还是失败失败任务重试几次并输出日志。全部跑完后把执行状态回写到元数据表方便下次排查。这套流程里最关键的部分就是第 2 步到第 4 步后面我会放出来可以直接用的脚本。3. 搭建批量同步的底层结构3.1 DataX 安装与基础配置DataX 安装本身不复杂但有几个细节容易被坑到。我简单说一下我的操作过程先确认服务器有 JDK 8 及以上环境执行java -version能看到版本号即可。然后去 DataX 的 GitHub Releases 页面下载打包好的datax.tar.gz上传到服务器解压到固定目录比如/opt/datax。解压后目录结构大概是这样的bin/存放启动脚本 datax.pyjob/存放任务配置 json 的目录plugin/各类读写插件按reader和writer分目录lib/DataX 运行时依赖的 jar 包验证安装是否成功直接执行cd /opt/datax python bin/datax.py -h能正常输出帮助信息就说明装好了。有一点我特别提醒DataX 内置的 MySQL 驱动版本比较老连接 MySQL 8 时经常报时区错误或者 “Public Key Retrieval is not allowed” 的问题。我在 jdbcUrl 里固定加上这几个参数之后问题基本就消失了useSSLfalseserverTimezoneAsia/ShanghaicharacterEncodingutf8mb4allowPublicKeyRetrievaltruerewriteBatchedStatementstrue注意characterEncodingutf8mb4还是utf8取决于你源库和目标库的实际字符集尽量保持和数据库端一致否则同步中文时大概率会有乱码问题。3.2 一张表搞定所有同步规则元数据表设计这是整个方案的核心先看建表 SQLCREATE TABLE sync_table_meta ( id INT NOT NULL AUTO_INCREMENT COMMENT 主键, src_table VARCHAR(128) NOT NULL COMMENT 源表名, dst_table VARCHAR(128) NOT NULL COMMENT 目标表名, src_column VARCHAR(2048) NOT NULL COMMENT 源表导出字段逗号分隔, dst_column VARCHAR(2048) NOT NULL COMMENT 目标表导入字段逗号分隔, split_key VARCHAR(64) DEFAULT NULL COMMENT 并发切分字段建议唯一索引或主键, where_cond VARCHAR(512) DEFAULT NULL COMMENT 同步过滤条件不含where关键字, write_mode VARCHAR(16) NOT NULL DEFAULT insert COMMENT insert/replace/update, pre_sql VARCHAR(512) DEFAULT NULL COMMENT 同步前执行SQL, post_sql VARCHAR(512) DEFAULT NULL COMMENT 同步后执行SQL, batch_size INT DEFAULT 1024 COMMENT 批次大小, channel_num INT DEFAULT 2 COMMENT 并发通道数, priority INT DEFAULT 100 COMMENT 执行优先级越小越先执行, enabled TINYINT NOT NULL DEFAULT 1 COMMENT 是否启用, created_at DATETIME DEFAULT CURRENT_TIMESTAMP COMMENT 创建时间, updated_at DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT 更新时间, PRIMARY KEY (id), KEY idx_enabled_priority (enabled, priority) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4 COMMENTDataX同步任务元数据表;解释几个关键字段的用途src_column和dst_column这两列必须一一对应。DataX 的字段映射是按位置对应的也就是说你源表导出第一列会写到目标表导入的第一列。即使源表和目标表字段名不一致也可以通过调整这两列的排列顺序实现映射。split_key大表并发读取时必填要选数值型或日期型的字段首选主键。pre_sql和post_sql同步前和同步后执行的 SQL。比如全量同步时我习惯在pre_sql里放truncate table 目标表确保目标表是干净数据。write_modeDataX 的 mysqlwriter 支持insert、replace、update三种模式。insert是普通插入遇到主键冲突会报错replace相当于INSERT ... ON DUPLICATE KEY UPDATEupdate只更新已存在的行新行不会写入。插入一条真实任务配置的示例INSERT INTO sync_table_meta (src_table, dst_table, src_column, dst_column, split_key, where_cond, write_mode, pre_sql, batch_size, channel_num) VALUES (user_info, app_user_info, id,user_name,nick_name,create_time, id,user_name,nick_name,create_time, id, status 1, replace, truncate table app_user_info, 1024, 4);这样一条记录就代表一个完整的同步任务。3.3 单个同步任务的 JSON 模板拆解不管用什么方式生成配置最终落到 DataX 里的还是标准 JSON 文件。这里放一个最常用的单表同步 JSON 模板后面脚本就是按照这个结构来生成的{ job: { content: [ { reader: { name: mysqlreader, parameter: { username: sync_user, password: 密码, column: [id, user_name, nick_name, create_time], splitPk: id, where: status 1, connection: [ { table: [user_info], jdbcUrl: [jdbc:mysql://10.10.10.10:3306/src_db?useSSLfalseserverTimezoneAsia/ShanghaicharacterEncodingutf8mb4] } ] } }, writer: { name: mysqlwriter, parameter: { username: sync_user, password: 密码, writeMode: replace, column: [id, user_name, nick_name, create_time], preSql: [truncate table app_user_info], batchSize: 1024, connection: [ { table: [app_user_info], jdbcUrl: jdbc:mysql://10.10.10.11:3306/dst_db?useSSLfalseserverTimezoneAsia/ShanghaicharacterEncodingutf8mb4 } ] } } } ], setting: { speed: { channel: 4, record: -1, byte: -1 } } } }几个我用下来的经验说明column数组里直接用字段名就行不建议写*。因为目标表可能加字段、字段顺序可能调整用*会导致同步完数据对不上。splitPk只对读取阶段有用它的作用是让多个并发 channel 按这个字段的值拆分成多个区间并行查询。如果字段上有大量重复值会造成严重的倾斜某个 channel 累死其他 channel 闲得看戏。preSql是在同步开始前执行的注意后面要讲的坑多个 channel 并发时preSql可能会被并发执行直接写truncate table有一定风险。4. 核心环节实现动态生成配置并批量执行4.1 生成 JSON 配置的 Python 脚本我偏好用 Python 来做这个事因为读取元数据表、拼字符串、调子进程都方便。脚本思路很简单连上元数据库查出所有启用的任务逐条生成 JSON 文件然后立即调用 DataX 执行。import os import json import pymysql import subprocess from datetime import datetime # 元数据库连接信息 META_DB { host: 10.10.10.1, port: 3306, user: meta_user, password: os.environ.get(META_DB_PASSWORD, ), database: sync_manager, charset: utf8mb4, } # 源库和目标库连接配置按实际环境填写 SRC_DB { host: 10.10.10.10, port: 3306, user: sync_user, password: os.environ.get(SRC_DB_PASSWORD, ), database: src_db, } DST_DB { host: 10.10.10.11, port: 3306, user: sync_user, password: os.environ.get(DST_DB_PASSWORD, ), database: dst_db, } JOB_DIR /opt/datax/job LOG_DIR /opt/datax/log DATAX_BIN /opt/datax/bin/datax.py def build_jdbc_url(db_config): return ( fjdbc:mysql://{db_config[host]}:{db_config[port]}/{db_config[database]} ?useSSLfalseserverTimezoneAsia/ShanghaicharacterEncodingutf8mb4allowPublicKeyRetrievaltruerewriteBatchedStatementstrue ) def build_reader(row): return { name: mysqlreader, parameter: { username: SRC_DB[user], password: SRC_DB[password], column: row[src_column].split(,), splitPk: row[split_key], where: row[where_cond], connection: [{ table: [row[src_table]], jdbcUrl: [build_jdbc_url(SRC_DB)], }], }, } def build_writer(row): return { name: mysqlwriter, parameter: { username: DST_DB[user], password: DST_DB[password], writeMode: row[write_mode], column: row[dst_column].split(,), preSql: [row[pre_sql]] if row[pre_sql] else [], batchSize: row[batch_size] or 1024, connection: [{ table: [row[dst_table]], jdbcUrl: build_jdbc_url(DST_DB), }], }, } def run_sync(): conn pymysql.connect(**META_DB) os.makedirs(JOB_DIR, exist_okTrue) os.makedirs(LOG_DIR, exist_okTrue) try: with conn.cursor(pymysql.cursors.DictCursor) as cursor: cursor.execute( SELECT * FROM sync_table_meta WHERE enabled 1 ORDER BY priority, id ) rows cursor.fetchall() for row in rows: job { job: { content: [ { reader: build_reader(row), writer: build_writer(row), } ], setting: { speed: { channel: row[channel_num] or 2, record: -1, byte: -1, } }, } } job_file os.path.join(JOB_DIR, fjob_{row[id]}.json) log_file os.path.join(LOG_DIR, fjob_{row[id]}_{datetime.now().strftime(%Y%m%d%H%M%S)}.log) with open(job_file, w, encodingutf-8) as f: json.dump(job, f, ensure_asciiFalse, indent2) print(f开始同步表 {row[src_table]} - {row[dst_table]}) ret subprocess.run( [python, DATAX_BIN, job_file], stdoutopen(log_file, w), stderrsubprocess.STDOUT, ) if ret.returncode 0: print(f同步成功: {row[src_table]}) else: print(f同步失败日志见: {log_file}) finally: conn.close() if __name__ __main__: run_sync()这段脚本的几个关键点我解释一下。第一所有数据库密码都不要写死在脚本里从环境变量读取。这个习惯能避免配置泄漏到版本控制仓库里。第二preSql字段我加了空判断没有配置时给空数组否则传入null会让 DataX 解析报错。第三每个任务的日志单独落一个文件文件名带上任务 ID 和时间。出错后定位日志会非常快不用在一个巨型日志里翻滚半天。4.2 批量执行与失败重试机制实际生产环境里同步任务不可能永远一次成功。我把脚本调成“逐表执行、独立日志、失败重试”的模式。逐表执行的意思是同一批任务里我不开多进程同时跑多个 DataX 实例。原因很实际DataX 本身就是多通道并发如果我再同时跑好几个任务实例数据库连接数、CPU、磁盘 IO 都可能被打满反过来拖慢每个任务。宁可串行稳定跑完也不要并发一起卡死。失败重试我加了一个简单的循环逻辑MAX_RETRY 3 for attempt in range(1, MAX_RETRY 1): ret subprocess.run(...) if ret.returncode 0: break print(f任务失败第 {attempt} 次重试)当然重试只对网络抖动、数据库临时连接不上这类问题有效。如果是因为数据类型不兼容、字段长度不够导致的失败重试多少次都是白搭这时候更重要的是看日志找出真正原因。我还会在脚本里加一个“迁移前记录源表行数、迁移后统计目标表行数”的校验动作思路是-- 同步前查源表 SELECT COUNT(*) AS src_cnt, IFNULL(SUM(LENGTH(id)), 0) AS src_sum FROM src_db.user_info; -- 同步后查目标表 SELECT COUNT(*) AS dst_cnt, IFNULL(SUM(LENGTH(id)), 0) AS dst_sum FROM dst_db.app_user_info;行数一致不算绝对准确但至少能挡住大部分漏数、少数的低级错误。要求更高的时候可以选两三个有校验意义的字段算 checksum两边比对。4.3 调度落地从手动到定时自动脚本跑通了之后接下就是把整个流程交给定时任务。我用的 Linux 服务器直接上 crontab30 2 * * * /usr/bin/flock -xn /tmp/datax_sync.lock -c cd /opt/datax python bin/sync_meta.py log/cron.log 21注意命令里的flock -xn它的作用是拿一个锁文件如果上一个任务还没执行完这个新的定时任务就直接退出避免同一条同步任务被叠加执行。同步任务最怕的就是上一条还没跑完下一条又启动了两边同时写数据轻则数据重复重则主键冲突任务报错。如果你想更精细一点可以按数据源或业务线拆成多个调度任务比如订单库凌晨 1 点跑、用户库凌晨 2 点跑。但我个人建议初期阶段保持一个入口、一套元数据表就好等任务量真的涨上去了再拆分别一开始就给自己加复杂度。5. 参数调优、常见问题与避坑经验5.1 大表同步性能调优关键参数参数这块我见过很多新手一上来就把 channel 调到几十结果数据库直接被干趴下。我这里给一组适合多数场景的参数参考参数建议值说明channel 数量2 ~ 8从 2 开始观察源库慢查询和 CPU 负载再逐步调大batchSize512 ~ 2048单行字段多、字段内容大时用 512行内容小的表可以用 2048jvm 内存-Xms2g -Xmx2g 起修改 bin/datax.py 里的 DEFAULT_JVM或启动时用 -j 参数覆盖splitPk必须选有唯一索引的数值型字段优先没有唯一索引时选重复率低的字段网络参数rewriteBatchedStatementstrue让 MySQL 驱动走真正的批量插入写入性能提升明显把 JVM 内存调大的方法很简单执行时加上参数python bin/datax.py -j -Xms2g -Xmx2g job.json如果一张表超过几亿行单独用一个大 channel 数去同步会非常吃目标库性能。我的处理方式是拆成分区段同步就是利用where_cond比如一次同步一个时间范围或一个 ID 范围几个任务串行跑完。5.2 常见问题速查表我把实际遇到、以及帮同事排查过的问题整理成了一个速查表现象可能原因解决办法连接报 Communications link failure网络不通、MySQL wait_timeout 太短、防火墙拦截先 telnet 测试端口再检查网络策略和数据库连接超时参数提示 Server time zone value 无法识别连接串没指定 serverTimezonejdbcUrl 统一加 serverTimezoneAsia/Shanghai报 Public Key Retrieval is not allowedMySQL 8 默认认证插件问题jdbcUrl 加 allowPublicKeyRetrievaltrueuseSSLfalse中文乱码连接字符集与库表字符集不一致连接串加 characterEncodingutf8mb4主键冲突任务中断writeMode 用了 insert 但目标表已存在相同主键根据业务改成 replace 或 update数据量对不上splitPk 字段有重复值导致区间切分时漏数或重复读换用重复率低的字段做 splitPk同步很慢目标库 CPU 打满channel 或 batchSize 设置过大降低 channel减小 batchSize观察目标库负载内存溢出 OOMJVM 默认堆太小或批量读入数据量过大调大 JVM 内存同时降低 channel 数量preSql 执行时报表不存在多个 channel 并发执行了 truncate把清表动作挪到调度前单独执行或先 drop 再建表5.3 我踩过的几个坑与最后的建议最后分享几个我自己的经验教训。第一个坑就是preSql并发执行。DataX 的多个 channel 在启动阶段会各自执行preSql如果你在preSql里写死truncate table某一次目标表正好没创建或者刚被清空第二个 channel 再执行就会报“table doesnt exist”。我后来改成两步同步前用单独的 SQL 脚本先清表DataX 任务里只保留普通写入逻辑。这样拆分之后问题再也没出现过。第二个坑是字段顺序。我刚开始做表迁移时源表和目标表字段名基本一致就没太在意顺序。后来有一张表目标端多了两个新字段我随手在dst_column里追加了一下结果同步出来的数据全部串列了。DataX 的字段对应关系就是简单的“位置对应”源第 1 列写目标第 1 列不是按字段名匹配。所以修改字段列表时务必两列对齐检查。第三个坑是关于“灵活”的边界。元数据驱动的方案虽好但不要想着一个模板解决所有问题。比如源库有特殊类型字段如 bit 型、json 型、目标库有生成列这些场景下通用模板大概率会出问题。我的建议是通用方案跑通 90% 的常规表剩下的 10% 异形表单独走专用配置不要为了“灵活”牺牲稳定。这个方案后续还可以继续扩展在元数据表里加sync_type字段全量表用全量模式增量表加一个时间条件从上次记录的时间点开始拉取同步完成后更新标记时间。我后面就是这么干的效果很好真正做到了“加表只插一条记录”。如果你想自己动手做一套我建议从 5 张表的小项目开始验证先手动写 JSON 跑通再逐步迁移到元数据表驱动的方式。这套东西看着简单但真正用起来香不香跑过一个月之后你自然就知道了。