
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-ddlDDL 解析与转换下文逐一拆解。断点续传让任务从半路接着跑它到底在解决什么离线同步任务动辄跑几个小时甚至跨天。常规做法是失败了就全量重来数据量一大成本根本兜不住。断点续传想做的就一件事记录这次同步已经读到了哪一行下次从那一行之后继续把重跑整个任务变成只补跑一小段。它怎么做到的ChunJun 的实现思路很朴素却很可靠基于 Flink 的 checkpoint 机制。每次触发 checkpoint 时框架会把 source 端最后一条数据的断点字段值保存下来任务失败重启后reader 在拼接查询 SQL 时会自动把保存的字段值塞进 where 条件形如WHERE id 记录值从而跳过已读过的数据。关键在于这个断点字段必须选递增字段比如自增主键因为过滤条件是严格大于号。你该怎么配置开启断点续传只需要在任务的setting.restore里配三个参数对应源码中的RestoreConfig类参数含义是否必填isRestore是否开启断点续传开启时填 truerestoreColumnName断点字段名开启后必填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、mysqlprintRate每多少条脏数据在日志打印一次pluginProperties插件自定义参数更常用的做法是在 ChunJun 启动参数-confProp中直接声明无需改任务 JSONchunjun.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 走 binlogOracle 走 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 同步连表结构变更都能自动跟上。新手最容易踩的四个坑断点字段不是递增的。选了非递增字段做restoreColumnNamewhere 条件用过滤必然漏数据。选字段前务必确认源表该列单调递增这是断点续传成立的前提。忘了开 checkpoint。断点续传依赖 checkpoint 作为存档点Flink 配置里没开 checkpointrestore 配置写得再对也白搭。脏数据阈值设置成 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),仅供参考