增量采集是个很基础但又特别容易翻车的话题。很多时候面试也好、做方案也好,上来就说“用时间戳字段拉数据”,真正落地才发现要么漏数、要么重复、要么把业务库拖垮。这篇文章把增量采集的核心思路、技术选型和实践细节拆开揉碎讲清楚,希望能给你的数据管道提供一些参考。
大数据采集这个领域,“增量”两个字几乎决定了整个数据管道的效率和成本。全量同步像是拍一张快照,简单直接,可业务数据一多,每天几千万条的变化量,全量跑一遍不仅耗时,还特别浪费计算资源。而增量采集更像是在记流水账,只捕捉“变了什么”,让下游数仓、分析引擎始终拿到最新的数据。
1. 内容整体设计与思路拆解
1.1 为什么增量采集是数据管道的核心话题
增量采集之所以重要,是因为数据采集的上游和下游对它的诉求完全不同。上游是业务库(MySQL、PostgreSQL、Oracle这类OLTP系统),它的第一诉求是业务稳定,不希望采集任务给数据库带来太大压力。下游是数仓或分析引擎,它的诉求是数据要及时、准确、不丢不重。增量采集恰好是这两个诉求之间的平衡点:通过尽可能小的代价,把源端的变化数据同步到目标端。
拿我做过的一个电商订单系统来说,订单表每天新增几十万条记录,同时有不少历史订单会更新状态(比如“待支付”变成“已支付”、“已发货”变成“已签收”)。刚开始用全量同步,每天晚上跑一次,把整张表拉下来覆盖写入。最开始数据量几百GB还能忍,后来订单积累到几个TB,全量同步要跑三个多小时,而且业务侧反映数据时效太差——白天产生的订单变更,要到第二天早上才能看到。这就是全量采集在数据量增长之后的典型瓶颈。
增量采集的核心价值在于:它关注的是“变化”而非“存量”。变化的数据量通常远小于存量数据,采集效率能提升一个数量级,时效性也大大改善。但增量采集的难点在于如何识别变化。不同的识别方式,决定了采集方案的可靠性、实时性和对源端的侵入程度。
1.2 增量采集常见的三种技术路线
增量采集的方案选型,本质上是绕不开“如何识别变化”这个问题的。业界常用的路线主要有三种:基于时间戳、基于Binlog(或Redo Log)解析,以及基于消息队列的变更数据捕获(CDC)。
基于时间戳是最直白、最容易落地的方式。前提是业务表里有一个能表示“最近修改时间”的字段,比如updated_at。每次采集时,把上一次记录的max(updated_at)作为本次查询的起始点,把源表中所有updated_at > 上次水位的数据拉出来。这种方式实现简单,对数据库也只多了一个范围查询的负载,很快就能见效。但它的缺点也很明显:
- 如果业务表没有
updated_at字段,或者更新时没有维护这个字段,就没法用这种方案。 - 时间戳字段存在精度问题,如果秒级精度恰好有同一秒内的大批量更新,容易出现漏数。
- 依赖“先更新数据,再提交事务”的顺序,如果业务上先提交事务、后异步更新时间戳,会丢失部分更新记录。
- 物理删除的数据不会体现在时间戳变化里,下游会残留已删除的脏数据。
基于Binlog解析是更可靠的方案。MySQL的Binlog(二进制日志)记录了所有数据变更操作,包括Insert、Update、Delete以及DDL。采集程序伪装成从库,向主库请求Binlog,解析出每一条变更记录,再写入目标端。这种方式不依赖业务表的任何字段,对业务零侵入,而且天然支持删除操作和事务级一致性。目前主流的Canal、Debezium、Flink CDC都是这个思路。
Binlog方案的劣势在于技术门槛相对高。要处理Binlog的格式(Row还是Statement)、位点管理、网络断线续传、DDL变更的兼容等,需要投入不少精力。但一旦跑通了,效果是另外两种方案没法比的——它可以做到准实时(秒级延迟),并且数据完整性有保障。
基于消息队列的方式比较特殊。它通常不是源端主动推数据,而是业务应用在写入数据库的同时,向Kafka这类消息队列发送一条变更事件,采集任务再消费Kafka里的消息写入目标端。这种方式业务侵入性强,要求研发在代码里做埋点,好处是对数据库完全无压力,还能做到业务系统与数据平台解耦。缺点是如果业务侧忘记埋点或埋点逻辑有bug,数据会悄悄丢失,而且排查起来比较困难。
1.3 技术路线对比怎么选
这三条路线没有绝对的优劣,更多是看场景。
| 维度 | 时间戳增量 | Binlog解析 | 消息队列辅助 |
|---|---|---|---|
| 实现难度 | 低 | 中高 | 中 |
| 对源库影响 | 低(有查询压力) | 极低(类从库拉取) | 极低 |
| 删除操作支持 | 不支持 | 支持 | 取决于埋点 |
| 实时性 | 分钟级 | 秒级 | 秒级 |
| 可靠性 | 中(有漏数和重复风险) | 高(有事务位点) | 中(依赖业务侧配合) |
个人建议:如果业务表有可靠的updated_at字段,数据量在千万级以下,对实时性要求不高(只做T+1同步),优先用时间戳方案,成本最低,运维最简单。如果数据量大、需要准实时同步、业务表频繁更新删除,直接上Binlog方案,一次到位,后续省心。消息队列方案一般用在已有成熟MQ基建、以及业务系统本身就需要发消息的场景,不需要为了采集单独引入一套消息链路。
2. 核心细节解析与实操要点
2.1 时间戳增量的水位线管理
时间戳增量的核心在于“水位线”(Watermark)的管理。水位线记录的是“我已经成功同步到哪个时间点”,下次增量从这个时间点往后取数。听起来简单,实际坑在边界条件。
先说时间精度的问题。很多业务表的时间字段是datetime(0),秒级精度。如果你的采集任务是每5分钟跑一次,某个订单在3分59秒更新,下一次采集在5分00秒拉取时,updated_at刚好是3分59秒,这没问题。可如果同一个秒内有两笔订单,第一笔在3分59秒01毫秒更新,第二笔在3分59秒990毫秒更新,而水位线记录的是3分59秒,第二次采集时用updated_at > '2024-06-01 03:59:59'查询,第二笔订单就会漏掉——因为它的更新时间恰好等于上一次的水位值。
解决方式很简单:水位线存储时精确到毫秒,哪怕源表的updated_at是秒级精度,查询条件也要用>=并把水位线往前拨1秒,再在应用层做去重。更稳妥的做法是:每次取数时把水位线条件设为updated_at > 上次水位 AND updated_at <= 当前时间 - 1秒,避免把正在执行中的事务数据(可能刚更新一半还没提交)拉出来。
另一个关键点是:水位线必须在数据成功写入目标端之后再更新。我之前遇到过一个线上事故:采集任务先更新了ZooKeeper里的水位线,然后才去写目标表,结果某些分区的写入失败,任务重启后从新水位线继续拉取,失败的那一批数据就永远丢了。正确做法是:先把数据写入目标端(最好在同一批事务里),确认写入成功后再推进水位线。如果目标端写入失败,保留原水位线,让任务重试消费。
2.2 Binlog解析的关键机制
Binlog方案的细节比时间戳方案多不少,这里重点讲三个机制:Row格式解析、位点记录、DDL处理。
首先,Binlog必须设置为binlog_format=ROW。这个格式下,Binlog里保存的是每一行数据变更前后的完整值,解析起来最直观,也能正确处理Update的“前镜像”和“后镜像”。如果使用Statement格式,Binlog里存的是SQL语句本身,虽然日志体积小,但解析要额外模拟SQL执行,复杂得多,还容易出错。
其次,位点管理极其重要。Binlog位点通常用binlog文件名 + position偏移量来表示。采集程序每消费一条Binlog事件,都要记录下当前的位点。位点丢失意味着从头重放或从最后重放,前者重复消费,后者丢数据。生产环境建议把位点存储到目标端数据库或ZooKeeper/Etcd里,防止本地磁盘故障导致位点丢失。
DDL处理是Binlog方案最容易被忽略的环节。源库执行ALTER TABLE加了一个新字段,Binlog里会有对应的Query事件。如果采集程序没有解析这个事件,后续Insert语句的字段数和目标端的表结构就对不上了,轻则写入失败,重则字段对应错乱。成熟的同步组件(如Debezium、Canal)会自动把DDL事件同步到目标端,但自研方案必须自己处理这个逻辑。我的经验是:DDL事件一律同步结构变更到目标端,并记录变更日志,方便追溯表结构变化历史。
2.3 业务系统配合的技术规范
不管用哪种方案,源端业务系统的配合程度都直接决定了增量采集的可靠性。这里列几个我踩过坑之后的硬性要求:
- 所有需要做增量采集的表,必须有
updated_at字段,且应用层更新数据时必须更新这个字段,不能依赖数据库的ON UPDATE CURRENT_TIMESTAMP(因为批量更新时容易被跳过)。 - 业务上的“软删除”优于“物理删除”。物理删除在Binlog方案下虽然能识别,但删除前该行在目标端的关联数据(比如维度表的引用)可能没人帮你清理。软删除用
deleted标记,增量采集天然能看到数据变化。 - 大事务要拆小。如果一个事务里更新了几十万行,Binlog会产生几十万条事件,采集端消费压力骤增,还可能造成源库Binlog磁盘空间暴涨和主从延迟。DBA一般会限制大事务,但作为采集方案的负责人,你有义务在业务侧推动这个规范。
3. 实操过程与核心环节实现
3.1 方案选型与整体架构
我在最近一个项目里做了这样一个选型决策:源库是MySQL 8.0,业务表总量大约2000张,核心业务表30多张,数据量较大的表有上亿的行数。业务方要求订单类数据的延迟不超过5分钟,历史数据要支持回补,并且要保证数据不丢。
这个需求下,时间戳方案满足不了5分钟延迟的硬指标,所以直接锁定Binlog方案。组件选了Debezium Embedded Engine,嵌入到我们的采集服务里,通过Java程序直接消费Binlog,把数据写入Kafka,再由Flink作业消费Kafka并写入Hive分区表和ClickHouse。整个架构的链路线是这样的:
源库MySQL → Debezium Embedded → Kafka → Flink → Hive / ClickHouse
为什么不用Canal?Canal本身也很成熟,但它部署上需要一个独立的Server进程,多了一套运维成本。Debezium Embedded把抓取逻辑嵌入应用,直接以库的形式调用API,控制粒度更细,也方便我们统一管理位点和监控指标。Flink的加入则是为了做数据的清洗、转换和分区分发。
3.2 Debezium的配置细节
Debezium的配置是整个链路的第一道关口。这里给出一个实际运行的配置片段,按字段逐个说明:
Properties props = new Properties(); props.setProperty("connector.class", "io.debezium.connector.mysql.MySqlConnector"); props.setProperty("offset.storage", "org.apache.kafka.connect.storage.FileOffsetBackingStore"); props.setProperty("offset.storage.file.filename", "/data/offset/offset.dat"); props.setProperty("offset.flush.interval.ms", "5000"); props.setProperty("database.hostname", "192.168.1.20"); props.setProperty("database.port", "3306"); props.setProperty("database.user", "debezium"); props.setProperty("database.password", "***"); props.setProperty("database.server.id", "10001"); props.setProperty("database.server.name", "order_center"); props.setProperty("database.include.list", "order_db"); props.setProperty("table.include.list", "order_db.t_order,order_db.t_order_item"); props.setProperty("database.history", "io.debezium.relational.history.FileDatabaseHistory"); props.setProperty("database.history.file.filename", "/data/history/dbhistory.dat"); props.setProperty("snapshot.mode", "schema_only_recovery"); props.setProperty("decimal.handling.mode", "string"); props.setProperty("tombstone.on.delete", "false");几个关键点分别说一下:
offset.storage是位点存储。开发环境可以用FileOffsetBackingStore,生产环境强烈建议改成KafkaOffsetBackingStore或自己实现一个基于数据库的存储,避免采集服务重启后位点回退导致重复消费。
database.server.id是给采集程序在MySQL主库上模拟从库的ID,每个采集进程必须唯一,不能和其他从库或采集进程冲突。如果多个采集程序消费同一个源库,server.id必须各不相同;如果同一个库被多个Binlog采集任务盯上了,还会对主库的Binlog产生重复拉取压力,尽量一个库只保留一个采集任务。
snapshot.mode配置的是首次启动时的快照行为。schema_only_recovery的意思是:只在有历史位点时恢复Schema,不重新拉全量数据。配合此配置,首次启动需要先做一次初始化快照(用initial模式),后面重启就不会全量重扫了。
tombstone.on.delete设为false,这样删除操作在Kafka里直接保留一条Tombstone消息,方便Flink侧处理删除语义,而不是把消息标记为过期后自动清理。
3.3 Kafka与Flink侧的处理
Debezium输出到Kafka的消息格式是带有Schema的JSON,每条消息都包含before和after字段,以及op操作类型(c表示Create、u表示Update、d表示Delete)。Flink消费时,重点关注after字段里的数据内容,以及op字段判断操作类型。
Flink作业的核心逻辑可以抽象成一个简单的状态机:
- 收到
c(Create):直接按主键写入目标表。 - 收到
u(Update):先查目标表是否存在该主键,存在则更新,不存在则做一次“拉齐”——这条数据可能是快照期间产生的,目标表还没有对应记录。 - 收到
d(Delete):按主键删除目标表记录,或者在Hive表里写入一条删除标记。
这里有个实际场景容易出错:如果源库执行的是UPDATE ... SET name='x' WHERE id=1,Binlog里产生的是一条Update事件。但如果这个Update根本没有改变任何值(比如设置的值和原来相同),MySQL仍然会产生一条更新日志,这属于正常现象,不处理也没问题,目标端多执行一次无变化的更新而已。
实时流同步的另一个关键点是主键。Binlog里的Update事件里,before和after都包含主键值,Flink侧更新目标表时,必须用主键精确匹配,不要用其他业务字段匹配。如果源表没有主键,建议在采集方案里加一个“虚拟主键”字段(比如将表中所有字段拼接后做哈希),否则重复数据的处理会非常痛苦。
3.4 目标端写入策略
写入Hive时,我推荐按事件时间做分区,而不是按处理时间。这样即使数据延迟到达,也能落回它本应属于的那个分区,保证下游查询的数据分布符合预期。
写入ClickHouse时,要利用好它的ReplacingMergeTree引擎。这个引擎允许重复写入相同主键的数据,后台会按版本号合并去重。实际操作中,我会让Flink作业写数据时带一个业务时间戳或Binlog位点作为版本字段,这样ClickHouse内部合并时保留最新版本。千万别用SummingMergeTree来处理更新类数据,那个引擎主要用于聚合类场景,更新会被错误地累加。
批量写入的批次大小也值得调优。我测试下来,单批次5000~10000条是性能和吞吐的平衡区间。批次太小时网络往返过多;批次太大时,一旦目标端网络抖动,整个批次回退重试的成本很高。另外,写入ClickHouse时建议用async_insert配合wait_for_async_insert=0,能明显降低写入延迟——前提是你接受“数据可能延迟可见”这个代价。
4. 常见问题与排查技巧实录
4.1 增量数据延迟越来越大
这是踩得最多的坑。增量采集任务跑着跑着,Kafka里的Lag越来越大,目标端的数据永远追不上源库。排查思路按下面顺序来:
- 先看源库的Binlog产生速率。如果源库本身有大事务或者大批量更新,Binlog瞬间会产生大量事件,采集端消费能力跟不上,这是最直接的原因。应对办法是压缩消息体积,或者给采集进程扩容(比如Debezium的
max.batch.size调大)。 - 再看Flink作业的并行度。并行度上不去,可能卡在单分区消费上。如果Debezium输出到Kafka时没有按主键做分区(默认按主键哈希),某些表的数据会集中到少数几个Kafka分区,Flink的并行读取就受限了。可以在Debezium的配置里针对大表单独设置分桶键,让数据分散。
- 查有没有某个表的消费有异常。Flink侧可以打印每条消息的处理耗时,如果某条数据太大(比如一个字段存储了超大JSON),序列化和写入目标端的耗时就会很长。
4.2 数据重复消费
重复消费几乎是分布式采集一定会碰到的问题。Flink的Checkpoint机制保证了“至少一次”(At Least Once)的语义,也就是遇到故障恢复时,某些数据会被再次消费。这本身是设计选择,允许重复消费,但目标端必须能幂等处理。
解决重复的核心是“幂等写入”。主键相同的数据多次写入,目标端最终保留的必须是最新一条,而且不能产生脏数据。ClickHouse的ReplacingMergeTree、Kafka的keyed-table、Hive的动态分区加主键去重,都能实现幂等。不要在应用层硬扛重复,那是一定扛不住的。
另一个排查方向是位点回退。如果你的offset.flush.interval.ms设得太大,比如30秒,采集进程在两次flush之间崩溃退出,重启后位点会回退到上一次flush的位置,这段时间产生的数据就会重复消费。把flush间隔调小,或者改用更可靠的位点存储,能显著减少重复。
4.3 DDL变更导致任务失败
这是Binlog方案绕不过去的坎。某天业务方执行了一条ALTER TABLE t_order ADD COLUMN buyer_remark VARCHAR(255),采集任务直接报错或者产出的目标表结构对不上。
最佳实践是让采集端自动消费DDL事件并同步到目标端。Debezium的database.history会记录表结构的变更历史,Flink侧的事件流里也会包含DDL的Schema变更消息。你需要在Flink作业里专门处理这类消息:提取表名和新的Schema定义,动态变更目标端的表结构。
如果目标端是Hive,DDL同步相对简单,因为Hive对列的类型变化容忍度较高;如果目标端是ClickHouse,字段顺序敏感,表结构调整就麻烦得多。我的建议是:在源库和生产环境之间加一层表结构变更审批流程,所有DDL必须先通知数据团队评估,再执行。否则你天天忙着救火,业务方还觉得是数据团队的问题。
4.4 常见故障快速速查表
| 故障现象 | 可能原因 | 排查思路 |
|---|---|---|
| 消费中断且位点丢失 | 位点存储介质损坏 | 检查offset存储文件所在磁盘,恢复最近一次成功位点 |
| 目标端数据比源库少 | 大事务产生的Binlog被跳过 | 检查源库binlog_row_image是否FULL,确保包含完整前后镜像 |
| 字段错位 | DDL变更后Schema信息未同步 | 对比源库和目标端的表结构,重新加载最新Schema |
| 写入目标端超时 | 目标端合并操作过慢 | 检查目标端是否有大查询锁表,分批写入 |
| 延迟堆积但无报错 | 单个分区数据倾斜 | 按主键分桶策略调整,增加目标表分片数 |
| 源库Binlog空间暴涨 | 采集任务停止时间过长 | 优先恢复采集,再清理Binlog,防止文件被覆盖 |
4.5 一个印象深刻的踩坑案例
有一次我们同步一个“商品收藏”表,这个表写频繁但几乎不更新,只有Insert操作。上线几天后突然发现目标端数据量和源库差了10%左右。排查了很久,最后定位到问题:源库的binlog_row_image被设置成了MINIMAL,Binlog里只记录被修改的列,而增量采集默认按全行解析,某些列拿不到值就写了NULL。而我们目标表的字段恰好设置了非空约束,写入时丢了几万条。
这个坑提醒我两点:一是源库Binlog相关的参数必须在上线前检查确认,不能拿默认值想当然;二是目标端对字段要有默认值兜底,空值可以设默认值,不能直接拒绝导致数据丢失。
5. 增量采集的延伸应用与实际体会
5.1 从业务库到数仓之外的增量场景
增量采集的应用场景不只是业务库同步到数仓。我在实际工作中还把它用到了几个延伸的地方:
- 缓存刷新:用户画像数据存在Redis里,每次全量刷新太贵,后来做成基于Binlog的增量更新。用户资料变更时,采集程序感知到后直接更新Redis里的对应key。
- 搜索引擎索引同步:ES索引里的文档需要跟着业务数据变化而更新。用增量采集把变化的数据抛到Kafka,再写一个消费程序调用ES的Bulk API做文档更新,比定时全量重建索引效率高太多了。
- 业务审计与溯源:Binlog里自带操作时间和事务ID,天然适合做审计。我们用它把高敏感表的所有变更记录保存到单独的审计库,保留全量历史,供安全团队查询。
这些场景对数据的时效性和准确性要求各不相同,但增量采集的底层逻辑完全一致——识别变化,传递变化,应用变化。
5.2 我个人在实际操作中的体会
做了这么多次增量采集,最大的体会是:可靠性靠的不是多复杂的代码,而是对细节的敬畏。水位线的正确推进、位点的持久化、DDL变更的应对、目标端的幂等设计,每一个环节都不能想当然。以我的经验,增量采集跑得稳定,80%的功夫花在上线前的参数核查和流程设计上,只有20%是在写代码。
另外,随着数据需求的复杂化,采集方案的设计越来越像“数据契约”的制定。它不只是技术问题,还是协作问题——你需要和DBA确认Binlog保留时长,和业务研发确认更新字段的维护规范,和运维确认监控报警的阈值。多花点时间沟通,远比事后救火省心。