ChunJun数据同步进阶:断点续传、脏数据治理与DDL自动同步,一篇讲透
【免费下载链接】chunjunA data integration framework项目地址: https://gitcode.com/gh_mirrors/ch/chunjun
凌晨两点,你盯着监控面板上跑了两天的同步作业,进度条已经走到 90%,突然一个网络抖动,任务挂了。重跑?意味着从头再来几十个小时。更气人的是,如果中途混进来一条字段超长的"坏数据",整个作业可能直接报错退出,而你还得在几十万条日志里把这条"元凶"翻出来。如果你正在用 ChunJun 做数据同步,这些场景完全可以提前规避——这个开源数据集成框架内置的三套"保险机制":断点续传、脏数据治理和 DDL 自动同步,就是为这类问题设计的。本文不堆概念,直接按"遇到什么问题 → 框架怎么解决 → 你该怎么配置"的顺序,带你逐个上手。
先把三件事的"分工"看清楚
在动手之前,先用一句话给三者定个位:断点续传管"进度不丢",脏数据治理管"坏数据不炸",DDL 同步管"表结构不脱节"。三者相互独立,又能组合进同一个任务里协同工作。
| 能力 | 解决的问题 | 生效位置 | 开启成本 |
|---|---|---|---|
| 断点续传 | 长任务失败后从头重跑,浪费大量时间 | Reader 端(RDB 插件) | 配置 restore 块即可 |
| 脏数据治理 | 单条坏数据导致作业整体失败、排障困难 | Reader/Writer 两端 | 配置 dirty 相关参数 |
| DDL 自动同步 | 源库加字段/改类型后目标库手工维护 | 源库日志捕获 + 语法转换 | 随连接器启用 |
这三套机制分布在 ChunJun 的不同模块里,源码位置分别是chunjun-core(断点续传与脏数据核心逻辑)、chunjun-dirty(脏数据插件)、chunjun-ddl(DDL 解析与转换),下文逐一拆解。
断点续传:让任务从"半路"接着跑
它到底在解决什么
离线同步任务动辄跑几个小时甚至跨天。常规做法是"失败了就全量重来",数据量一大,成本根本兜不住。断点续传想做的就一件事:记录这次同步已经读到了哪一行,下次从那一行之后继续,把"重跑整个任务"变成"只补跑一小段"。
它怎么做到的
ChunJun 的实现思路很朴素,却很可靠:基于 Flink 的 checkpoint 机制。每次触发 checkpoint 时,框架会把 source 端最后一条数据的断点字段值保存下来;任务失败重启后,reader 在拼接查询 SQL 时会自动把保存的字段值塞进 where 条件,形如WHERE id > 记录值,从而跳过已读过的数据。关键在于,这个断点字段必须选递增字段(比如自增主键),因为过滤条件是严格大于号。
你该怎么配置
开启断点续传只需要在任务的setting.restore里配三个参数,对应源码中的RestoreConfig类:
| 参数 | 含义 | 是否必填 |
|---|---|---|
| isRestore | 是否开启断点续传 | 开启时填 true |
| restoreColumnName | 断点字段名 | 开启后必填 |
| restoreColumnIndex | 断点字段在 reader column 中的位置 | 开启后必填 |
一个最小可用的 JSON 片段如下:
"setting": { "restore": { "isRestore": true, "restoreColumnName": "id", "restoreColumnIndex": 0 }, "speed": { "channel": 1 } }两个容易忽略的细节,这里提前说明:一是任务必须开启 checkpoint,否则没有可恢复的"存档点";二是下游 writer 若本身不支持事务,则需要目标表具备幂等写入能力,否则断点恢复后可能出现少量重复数据。此外RestoreConfig里还有一个maxRowNumForCheckpoint参数,默认每处理 1 万条数据触发一次 checkpoint,你可以根据业务对"恢复粒度"的要求适当调整。
脏数据治理:把"坏数据"拦在门外,而不是炸掉整个任务
它到底在解决什么
数据同步中最磨人的不是大数据量,而是"一两条脏数据"。字段超长、类型不符、值超出范围……传统做法要么整条报错拖垮作业,要么静默丢弃后无从排查。脏数据治理的思路是:坏数据照常记下来、收集起来、甚至入库归档,同时作业继续往下跑,是否"忍"这些坏数据由你定阈值。
它怎么做到的
整体架构是经典的生产者-消费者模式,核心逻辑在chunjun-core的DirtyManager类中:
- 收集:reader 和 writer 在读写数据时一旦捕获异常,就调用
DirtyManager.collect(),把数据内容、异常原因、字段名、任务信息封装成一条脏数据记录; - 下发:脏数据被投递进一个内存队列(dirty-queue),同时用一个独立的异步线程池驱动消费者;
- 消费:消费者
DirtyDataCollector循环从队列取数据,真正"怎么处理"由子类实现——日志打印、写 MySQL、自定义插件各有各的逻辑; - 判负:当消费失败条数或总条数超过你设置的上限时,抛出
NoRestartException,任务以失败告终(而不是无限重试)。
你该怎么配置
配置入口有两层。任务级配置对应源码中的DirtyConfig类,核心参数如下:
| 参数 | 含义 |
|---|---|
| maxConsumed | 最大消费条数上限,超过后消费者终止 |
| maxFailedConsumed | 最大失败消费条数上限 |
| type | 脏数据插件类型,如 log、mysql |
| printRate | 每多少条脏数据在日志打印一次 |
| pluginProperties | 插件自定义参数 |
更常用的做法是在 ChunJun 启动参数-confProp中直接声明,无需改任务 JSON:
chunjun.dirty-data.output-type = log # 或 jdbc chunjun.dirty-data.max-rows = 1000 # 脏数据总量上限 chunjun.dirty-data.max-collect-failed-rows = 1000 chunjun.dirty-data.jdbc.url = jdbc:mysql://localhost:3306/db chunjun.dirty-data.jdbc.username = root chunjun.dirty-data.jdbc.password = root chunjun.dirty-data.jdbc.table = chunjun_dirty_data chunjun.dirty-data.log.print-interval = 500选择jdbc类型时,框架会把脏数据持久化到chunjun_dirty_data表,表结构在脏数据插件设计文档里给出了完整的建表语句(含 job_id、job_name、dirty_data、error_message、field_name、create_time 等字段,并为 job_id 建了索引)。脏数据插件本身位于chunjun-dirty模块,目前内置chunjun-dirty-log和chunjun-dirty-mysql两个实现,想接自己的存储(比如 Elasticsearch、Kafka),照着这两个子类的写法扩展即可。
DDL 同步:表结构变更,让框架替你"翻译"和"执行"
它到底在解决什么
数据同步任务一般只在启动时确定表结构。源库业务表今天加了个字段、明天改了字段类型,目标库如果不跟着改,后续同步轻则字段对不上、重则直接写失败。靠 DBA 手工比对两边表结构,在大库、多表场景下既慢又容易漏。DDL 同步的目标是:源库的表结构变更自动捕获,自动"翻译"成目标库的方言,自动执行。
它怎么做到的
这条链路分三步走:
- 捕获:从源库日志中读取 DDL 语句。MySQL 走 binlog,Oracle 走 LogMiner,框架的 CDC 连接器负责这一步;
- 解析与转换:这是
chunjun-ddl模块的核心工作。它基于 Calcite 对 DDL 语句做语法解析,生成统一的语义模型,再经由转换器输出目标数据库方言的 SQL。模块内chunjun-ddl-mysql和chunjun-ddl-oracle各自实现了建表、删表、加列、改列、索引变更等几十种操作的解析与反向生成; - 执行:转换后的 DDL 由目标端的 DDL 处理器(如
chunjun-restore-mysql中的MysqlDDLHandler)落地执行,并对执行结果做缓存与幂等处理,避免重复执行时报错。
你该怎么配置
DDL 同步通常不需要单独开一张配置表,而是随 CDC 连接器(binlog、logminer 等)一起启用。使用前建议先浏览chunjun-ddl模块的源码和测试用例,了解当前支持的操作类型覆盖范围——比如 MySQL 侧的SqlAlterTableAddColumn、SqlCreateTable等已有对应测试,接新场景前先确认你的 DDL 类型在支持列表里,这是最稳妥的上手方式。
实战:把一个"抗造"的任务组装出来
把三件事串起来看,一个生产级任务的骨架大致长这样:
{ "job": { "content": [ { "reader": { "name": "mysqlreader", "parameter": { "column": [{"name": "id", "type": "int"}], "username": "root", "password": "root", "connection": [{"jdbcUrl": ["jdbc:mysql://localhost:3306/test"], "table": ["t_user"]}] } }, "writer": { "name": "hdfswriter", "parameter": {} } } ], "setting": { "restore": { "isRestore": true, "restoreColumnName": "id", "restoreColumnIndex": 0 }, "speed": { "channel": 1 } } } }配合启动时的-confProp脏数据参数,这个任务就同时具备了"失败断点续跑"和"坏数据隔离"能力。如果源端开启 binlog 并启用 DDL 同步,连表结构变更都能自动跟上。
新手最容易踩的四个坑
- 断点字段不是递增的。选了非递增字段做
restoreColumnName,where 条件用>过滤必然漏数据。选字段前务必确认源表该列单调递增,这是断点续传成立的前提。 - 忘了开 checkpoint。断点续传依赖 checkpoint 作为"存档点",Flink 配置里没开 checkpoint,restore 配置写得再对也白搭。
- 脏数据阈值设置成 0 或负数。很多人以为 0 表示"零容忍",实际语义要看插件实现——某些配置下
printRate <= 0表示完全不打印,而条数限制设为负数表示"容忍所有异常不失败"。配置前先读对应文档,别凭直觉。 - 拿到 DDL 就全量放开。DDL 同步虽能自动执行,但删除表、改主键这类高危操作落地到目标库前,建议先在小范围表上验证转换结果,确认
chunjun-ddl对该语法支持良好后再放开,避免线上被"自动"坑一把。
总结:从"能跑"到"跑得稳"
断点续传、脏数据治理、DDL 自动同步,本质上是把数据同步从"能用"推向"用得放心"的三层保障:进度可恢复、异常可容忍、结构可跟随。它们都不复杂——要么是 JSON 里加几行配置,要么是启动参数里多几个键值,但带来的可靠性提升是实打实的。
想深入了解的同学,建议按这个顺序读仓库里的资料:断点续传的原理与参数见 docs/docs_zh/拓展功能/断点续传介绍.md,增量同步(配合 endLocation 指标做跨作业增量)可参考 docs/docs_zh/拓展功能/增量同步介绍.md,脏数据插件设计文档在 docs/docs_zh/拓展功能/脏数据插件设计.md。跑通示例任务的话,chunjun-examples/json和chunjun-examples/sql目录下有大量现成配置可以改改直接用。数据同步这条路,先把"跑得稳"做到位,再谈"跑得快"。
【免费下载链接】chunjunA data integration framework项目地址: https://gitcode.com/gh_mirrors/ch/chunjun
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考