3个开关治好数据同步的三大顽疾:断点续传、增量同步与脏数据处理实战笔记

发布时间:2026/8/21 18:24:56
3个开关治好数据同步的三大顽疾:断点续传、增量同步与脏数据处理实战笔记 3个开关治好数据同步的三大顽疾断点续传、增量同步与脏数据处理实战笔记【免费下载链接】chunjunA data integration framework项目地址: https://gitcode.com/gh_mirrors/ch/chunjun如果你负责过任何一张每天必须准点产出的业务报表大概率经历过这样的深夜凌晨两点手机铃声准时响起——数据同步任务又挂了。你睡眼惺忪爬起来把跑了十几个小时的任务删掉重跑然后在群里回复一句在重跑了明天看结果。可第二天业务方依然投诉数据不对因为从头重跑不仅慢还可能把已经同步好的数据再冲一遍。这篇文章的主角——开源数据集成框架ChunJun正是为解决这类同步既慢又脆的问题而生的它支持把几十种数据源互相打通更内置了断点续传、增量同步和脏数据处理三把手术刀。本文不讲大道理直接带你按它能解决什么问题 → 怎么配置 → 效果如何三步走完这三件事。三分钟跑通第一个最小案例先建立信心在聊高级功能前先用一个最朴素的任务热热身Stream → Stream即从一个虚拟数据流读数据再打印到终端。这个任务不需要任何数据库适合验证环境是否打通。拿到源码后先编译安装包也可以在发布页直接下载现成的插件包git clone https://gitcode.com/gh_mirrors/ch/chunjun cd chunjun mvn clean package -DskipTests打包产物会落在chunjun-assembly模块的 target 目录下进入该目录执行sh bin/chunjun-local.sh -job chunjun-examples/json/stream/stream.jsonlocal模式不依赖 Flink 和 Hadoop 环境一个 JVM 进程就能跑。任务脚本写好之后你可以在 Flink Web UI 上看到Source: streamreader → Sink: streamwriter的算子拓扑数据像流水一样从源头流向终点。跑通之后信心就有了。接下来我们把场景切换到真实的 MySQL 业务库逐一拆解那三个治顽疾的开关。第一招断点续传——像看书夹书签失败后从上次停下的那行接着读解决什么问题全量同步有个致命伤任务一旦失败就得整表重来。假设你从 MySQL 往 HDFS 同步一张千万级的大表跑了 20 个小时后因为网络抖动失败了从头再跑又是 20 小时期间下游数据一直是缺一段的状态。断点续传要解决的就是这个输不起的问题失败了没关系从失败的位置接着跑已同步的数据不重复搬。怎么配置它的实现思路很像你读书时夹书签ChunJun 基于 Flink 的 checkpoint在检查点保存时记录下 source 端最后一条数据的某个字段值相当于记住页码任务重启后读取时把保存的值拼进查询条件只拉取这个值之后的数据。在任务脚本的setting.restore里打开开关setting: { restore: { isRestore: true, restoreColumnName: id, restoreColumnIndex: 0 }, speed: { channel: 1, bytes: 0 } }三个参数的含义一句话就能说清参数一句话解释isRestore是否开启断点续传默认关闭restoreColumnName用哪个字段当书签要求是递增字段restoreColumnIndex这个字段在 reader 的 column 列表里排第几位从 0 开始效果如何开启后任务失败重启source 端拼接 SQL 时会把 checkpoint 里的 state 作为起点从上次读到的位置继续而不是从零开始。效果直接对标一句大白话上次搬了 80% 的货这次只搬剩下 20%前 80% 原封不动。前提是下游支持事务或具备幂等性——这样即使有重复数据写进去也不会产生脏账。一句话总结选对递增字段、打开开关剩下的交给 checkpoint。第二招增量同步——搬家只搬新添的家具不再整屋重搬解决什么问题断点续传解决中途失败增量同步解决每次全跑。很多表只有 Insert 操作数据却在不断增长。每次同步都把整张表从头扫一遍浪费的时间和资源肉眼可见。增量同步的思路是每次只搬上次搬完以后新增的那部分。怎么配置实现原理其实就是配合增量键在 SQL 里拼接过滤条件把已经读过的数据过滤掉。需要两个关键配置increColumn指定增量字段必须在 column 中存在比如自增 idstartLocation起始位置第一次执行可以不配此时是整表同步后续作业取上一次作业记录的 endLocation以 MySQL 为例reader: { name: mysqlreader, parameter: { column: [{name: id, type: int}, {name: name, type: string}], increColumn: id, startLocation: 2, username: root, password: root, connection: [{jdbcUrl: [jdbc:mysql://localhost:3306/test], table: [baserow]}] } }增量同步依赖 Prometheus 收集指标每次作业跑完会记录一个endLocation指标并上传下一次作业就拿它当过滤依据。比如第一次跑完 endLocation 是 10下一次就会拼出SELECT ... WHERE id 10。效果如何从每次都扫全表变成每次只扫增量同步耗时从量级上被压缩。注意有两个使用边界增量字段只能是数值或时间类型Prometheus 不支持字符串且字段值可以重复但必须递增——因为过滤符号是。划重点如果担心任务启动间隙里又插入了增量键值相同的数据比如 id10 的重复插入记得把useMaxFunc设为true。它会取当前表增量键最大值作为本次 endLocation并把换成同时给结果加 当前最大值的上界把夹缝数据也捞回来。第三招脏数据处理——给数据出口装一道安检门解决什么问题数据源里总有坏孩子格式不对的、超出长度范围的、类型对不上的……过去这些脏数据只会悄悄打印在日志里你根本不知道丢了多少、为什么丢。脏数据处理要解决的是脏数据被谁、在哪、因为什么被拦下拦了多少、拦到多少就该报警喊停。怎么配置ChunJun 的脏数据模块采用经典的生产者-消费者架构任务运行时DirtyManager负责收集脏数据及异常原因下发到dirty-queue队列消费者异步轮询队列把脏数据交给具体的插件去处理——比如打印到日志或者写进 MySQL 表。配置项通过-confProp传给任务chunjun.dirty-data.output-type log # 或 jdbc写进数据库 chunjun.dirty-data.max-rows 1000 # 脏数据总条数上限超过则任务失败 chunjun.dirty-data.max-collect-failed-rows 1000 # 处理失败的条数上限 chunjun.dirty-data.log.print-interval 500 # 每500条打印一次 chunjun.dirty-data.jdbc.url jdbc:mysql://localhost:3306/chunjun_dirty chunjun.dirty-data.jdbc.table chunjun_dirty_data核心逻辑在DirtyManager类里收集脏数据 → 下发队列 → 消费者消费 → 计数判断。当脏数据总数达到 max-rows或处理失败数达到 max-collect-failed-rows任务会抛出NoRestartException直接失败且不重试——避免一个坏数据漩涡把集群资源耗光。效果如何脏数据从黑盒丢失变成全程可审计。选log模式能快速定位问题行选jdbc模式脏数据的作业 ID、算子名、异常数据、报错原因、出现时间会逐条落库配合job_id、operator_name等索引排查效率直接拉满。一句话总结给数据出口装一道安检门谁出问题、出了多少一查便知。新手避坑指南这五个坑笔者都替你踩过以下每一条都是真实踩坑血泪强烈建议收藏。断点续传字段没选对重启后数据静默丢失。过滤条件是字段必须严格递增。用了个会回退或重复的字段重启后不是丢数据就是重复数据。选自增主键或更新时间戳最稳。忘了开 checkpointrestore 形同虚设。断点续传依赖 Flink checkpoint 保存状态任务里没开启 checkpoint失败后 state 为空等于白配。restoreColumnIndex 和 column 顺序对不上。配置里写的是字段在 reader column 列表里的下标从 0 开始。写错位置读到的书签根本不是你以为的那一列。增量同步报错先查 Prometheus。增量依赖 Prometheus Pushgateway 收集 endLocation环境没搭好或指标没上传后续作业就拿不到起点。别在 SQL 上瞎调先确认指标flink_taskmanager_job_task_operator_flinkx_endlocation能查到值。脏数据上限默认太严格一遇到脏数据任务就挂。默认值很保守生产环境记得按业务容忍度调大或设成负数表示容忍所有异常。同时别把output-type设成jdbc却忘了建表建表语句在官方文档里直接有现成的。收尾把要点打包成一张可勾选的清单最后把今天的核心要点整理成一张实战检查清单下次配置任务时对照着打勾即可长任务超过一天开启isRestore: true并指定递增的restoreColumnName确认restoreColumnIndex与 reader column 顺序一致且任务已开启 checkpoint只有 Insert 的表用增量同步increColumn选数值或时间类型字段增量键可能重复时设置useMaxFunc: true防漏数据按业务容忍度配置max-rows与max-collect-failed-rows别用默认值裸奔脏数据要可追溯时选用jdbc模式并按文档建好脏数据表想继续深挖这几处是值得动手的地方断点续传原理与参数细节docs/docs_zh/拓展功能/断点续传介绍.md增量同步的环境搭建与useMaxFunc场景docs/docs_zh/拓展功能/增量同步介绍.md脏数据插件架构生产者-消费者、DirtyManager 源码docs/docs_zh/拓展功能/脏数据插件设计.md源码在 chunjun-dirty/各类任务脚本模板chunjun-examples/json/目录下按数据源分门别类照着改参数即可数据同步这事做得久了就会发现比能跑通更重要的是跑不坏、跑不丢、跑不重。把断点续传、增量同步和脏数据处理这三个开关用起来你会惊喜地发现凌晨两点的电话真的可以不接了。去仓库里把源码 clone 下来挑一个真实的 MySQL 库试试吧效果比读十篇文章都来得实在。【免费下载链接】chunjunA data integration framework项目地址: https://gitcode.com/gh_mirrors/ch/chunjun创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

相关新闻