ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

增量采集实战:水位线、CDC与幂等写入的完整方案

增量采集实战:水位线、CDC与幂等写入的完整方案 做大数据采集这几年我接过最多的需求就是把一张十几亿行的大表从每天全量拉取改成增量同步。这张表往往是订单、流水或者日志每天新增几百万行全量任务越跑越慢跑批时间从半小时拖到三四个小时数据库压力大得前端业务都受牵连。每次这种需求一出来负责采集的人第一个想到的就是“增量采集”。但增量采集真正落地的时候坑并不在于“增量”这两个字怎么写而在于你怎么定义“上次读到哪”、怎么处理乱序和重复、怎么跟你之前跑的全量衔接上。这篇文章就把我做增量采集时用到的策略、踩过的坑、还有可以直接抄走的实现思路完整聊一遍适合正在做数据集成、数据仓库同步、或者自研采集平台的同学参考。1. 整体设计思路:增量采集到底是什么,以及为什么绕不开1.1 全量采集的痛点我记得特别清楚,那张订单表从2000万行涨到1.2亿行的时候,原有的全量采集任务开始变得不可控。任务每天凌晨跑,先TRUNCATE数仓里的目标表,再从源库把整张表读出来写进去。数据量小的时候这套逻辑一点问题没有,20分钟搞定,下游早上来上班刚好能看到数据。但等数据量过亿,全量任务跑完要接近两个小时,期间源库的IO和CPU都被吃得很厉害,线上业务方开始抱怨接口变慢。更要命的是,跑批时间一长,留给下游任务的时间被压缩,链路一抖动,当天数据就出不来,整套调度全乱。全量采集的本质是成本跟总数据量成正比,而你每天都在搬运那些根本没有变化的存量数据。几千万行里面真正新增的可能只有几万行,但全量任务照样把它们全都读一遍、传输一遍、写入一遍。这件事在数据规模小的时候可以忍,规模一大就忍不了。所以增量采集的出现,核心就是为了把采集成本从“总数据量”降成“变化数据量”。1.2 增量采集的定义与适用边界增量采集说白了就一句话:每次只同步从上一次结束位置之后发生变化的那些数据,并且用一个标记把“我已经读到哪了”记录下来。这个标记,业内一般叫水位线(Watermark)或者位点(Offset)。它可以是时间字段的最大值,可以是自增ID的最大值,也可以是Binlog里的文件位置和偏移量。但要注意,增量采集不是万金油,有些场景确实不该用。比如一张配置表,总共几百行,一年也改不了几次,用全量每天TRUNCATE INSERT一分钟内跑完,完全没有必要上增量;再比如一些临时分析表,本身就不要求数据的连续性,全量拉一次就完事。判断标准其实很朴素:变化频率低、数据总量小、全量成本低于增量成本的时候,全量更合适。反过来,千万行以上体量、每天持续增长、字段有更新时间或者严格的顺序标识,这样的表就是增量采集的主场。适合用增量的典型场景包括:业务订单表入数仓、日志采集入湖、MySQL到Kafka的数据订阅、以及数据仓库ODS层到DWD层的分层同步。2. 增量采集的几种典型方案到底怎么选增量采集没有一个统一的实现标准,不同数据源、不同业务容忍度,对应完全不同的方案。我习惯先把候选方案列出来,再按数据准确度要求、源库改造复杂度、运维成本三个维度去做权衡。2.1 基于时间字段的方案最直观、也是最容易上手的方案,就是利用业务表里的时间字段,比如update_time或者create_time,每次记录下上一次同步到的最大值,下一次查询条件就写成where update_time 上次最大值,把新增和更新的数据捞出来。这个方案的优点是真简单,一段SQL搞定,不需要额外权限,也不影响源库。但问题同样明显:第一,业务字段不一定可靠,update_time完全依赖业务代码有没有正确更新;第二,时间字段如果没有索引,大表查询会全表扫描;第三,同一个时间戳内可能出现多条数据,如果用捞数据,恰好等于边界值的那些行很容易被漏掉。我用的时候一般会把查询改成where update_time last_point order by update_time asc limit 5000,并且在代码里允许边界值重复拉取,靠下游去重来兜底。这个方案的适用场景是:业务更新频率可控、字段质量有保障、能容忍一定延迟且不要求捕获删除操作的数据表。比如一些日志表、埋点表,数据只追加不改动,用时间字段做增量非常舒服。2.2 基于自增ID或序列的方案如果表里有一个严格递增的主键ID,那增量采集就更简单了。每次记录已经拉取的最大ID,下一轮查询用where id 最大ID,按ID排序后拉取。方案的时间复杂度很低,而且不会因为数据更新导致重复,因为每一行数据的ID是唯一的,天然就带上了“增量位置”的属性。但这个方案有一个致命限制:它只能捕获新增数据,捕获不了数据更新。如果一张订单表里面,支付状态从“待支付”变为“已支付”,而它的ID没有变,基于ID的增量就完全感知不到这次变化。再加上很多业务表的主键不是自增的,是雪花ID之类的分布式ID,本身就不存在“上一个ID是什么”的连续性。所以这个方案我只在“数据只追加不改动”的场景里用,比如操作日志、埋点流水表。2.3 基于Binlog或CDC的方案前面两种方案说到底都是在用业务字段来间接推测数据变化,而Binlog方案是直接看数据库底层记录了哪些变更,每一步变化都逃不掉。通过解析MySQL的Binlog,我们能拿到真正的INSERT、UPDATE、DELETE事件,这就把增量采集从“查业务表”升级成了“订阅变更流”。实际落地时,不建议所有团队都从零去写Binlog解析器,因为复刻一套解析流程的成本非常高,还要处理不同MySQL版本的事件格式差异。更稳的做法是直接使用现成的开源组件,比如Canal、Debezium、Flink CDC,它们已经实现了Binlog连接、解析、序列化、位点管理这些脏活累活。我在生产环境用得比较多的是Canal加Kafka的组合,Canal把Binlog变更转成JSON消息放到Kafka,下游采集服务消费消息后写入数仓。Binlog方案的优点非常明显:能捕获删除、能捕获更新前后的值、不依赖业务字段质量、对源库的侵入也小(只需要开启Binlog并给账号授权)。缺点也有:需要额外的运维组件,消费端积压时需要监控,Canal在高并发下偶尔也会出现位点偏差。一句话总结,凡是要求数据准确、删除也能捕获、不允许漏数据的场景,直接选Binlog/CDC方案。2.4 基于消息队列削峰填谷的延伸严格来说,消息队列不是独立的增量采集方案,而是增量采集之后的数据传输管道。数据通过Canal或采集程序读出后,可以先丢进Kafka,下游数仓任务再按自己的节奏消费。这么做的好处是解耦:上游采集只管“把变更读出来并发到Kafka”,下游消费速度慢先挂着,不影响继续采集;Kafka里的消息也能回放,数据出问题可以把历史消息重新消费一遍。但引入消息队列也带来一个新问题:消息投递语义。Kafka默认提供的是at-least-once投递,也就是一条消息可能被重复消费。如果下游没有幂等处理,重复消费一份数据就会导致数仓里出现重复行。这一点我在第三章会展开讲,这里先提个醒:用Kafka做增量传输后,“去重”和“幂等”就从可有可无变成了必须面对的事情。2.5 方案选型对照我把三种常见方案放在一起对比过很多次,大致可以总结成下面这样:方案捕获能力是否需要业务改造延迟运维复杂度适用场景时间字段增量新增加更新需要业务保证字段可靠分钟级极低数据追加为主、可容忍漏删除自增ID增量仅新增需要主键单调递增分钟级极低日志流水、埋点Binlog/CDC增量新增、更新、删除全捕获仅需开启Binlog并授权秒级较高,需运维Canal等严格要求数据准确的业务核心表选型的时候不要只看数据量,还要问业务一句“被删除的数据需不需要同步”,这是很多方案翻车的根源。如果删除不需要,时间字段增量方案能省掉大量运维工作;如果删除也必须一模一样,老老实实上Binlog。3. 核心细节解析与实操要点3.1 水位线到底存在哪里增量采集系统的核心状态就是水位线,它存得稳不稳,直接决定采集任务可靠不可靠。我见过有人把水位线存在内存变量里的,任务一重启就丢了,导致重复采集全量数据;也见过有人把水位线存在本地文件的,集群里起了六个实例,每个实例的水位线互相覆盖,数据越采越乱。水位线的存储选择要按集群规模来定。单机任务可以用本地文件或SQLite,但要保证任务失败重启后能读到上次的值;分布式调度场景下,更稳的是存到Redis、ZooKeeper或者一个专门的数据库表里。我的实践是存在MySQL里的一张水位线表,字段包括table_name、last_point、update_time、batch_id,每次推进水位线都走事务更新。这样任何一个采集节点执行任务,都能读到全局唯一的水位线,不会因为你扩容了几个实例就互相打架。另一个容易踩的坑是水位线推进的时机。正确的顺序永远是:先同步数据并确认数据安全落库,再推进水位线。但有些新手图省事,一边读数据一边推位点,结果下游写失败了,位点已经推进了,这一批数据就再也捞不回来了。正确做法是在采集程序里用一个显式的队列,把“当前已处理成功的最大位点”留在内存,等一批数据全部落库后再批量推进水位线。3.2 数据去重与幂等写入增量采集最大的幻觉就是“我做了增量,数据就不会重复”。这个想法很危险,因为增量管道里处处都是重复:任务重跑会重复,消息队列重投放会重复,时间边界处理不当会重复。所以增量采集的下游,必须设计成幂等写入,也就是无论同一份数据被写入多少次,最终表里的结果都只有一份。具体怎么做呢?最常见的做法是利用目标表的唯一键。比如订单表同步到数仓,我们就可以在目标表上建order_id的唯一索引,写入方式用INSERT ... ON DUPLICATE KEY UPDATE或者MERGE语句。这样就算Canal发了两次同样的变更事件,第二次写入的时候命中唯一键冲突,直接更新而不是插入新行,结果依然正确。如果目标环境不支持MERGE类语句,也可以在采集程序里加一个去重逻辑,用Redis维护最近处理过的业务主键集合,重复消息直接丢弃。但这种方案有内存压力而且不便宜,能走数据库唯一索引的话,优先走数据库唯一索引。3.3 时间乱序与迟到数据处理用时间字段做增量采集的人,迟早会遇到一个诡异问题:某条数据的update_time突然变小了,或者一个三天前创建的订单今天才被更新,而增量任务只盯着“最近一小时变化的字段”,很容易把这批迟到数据漏掉。这类问题在业务上很常见,比如线下补单、审批通过后回填状态、系统修正脏数据。它们的特点都是“事件实际发生的时间很早,但进入采集视野的时间很晚”。我用过的一个有效策略是“主增量加补数窗口”:主流程依然按水位线正常采集,但额外跑一个定时补数任务,每天扫描最近三天的增量变更字段,把那些update_time落在这三天内的数据重新捞一遍。三天的窗口可以覆盖绝大多数业务迟到场景,成本也在可控范围内。补数任务的结果同样走幂等写入,所以重复数据不会污染目标表。3.4 删除数据怎么识别增量采集很容易漏掉DELETE操作。基于时间字段的方案里,一条数据被删了,update_time没有任何变化,你下一轮查询根本看不到它;基于自增ID的方案更是连感知都做不到。如果你的下游数仓同步的是整张业务表,删除漏了之后,数仓里就会攒着大量业务上已经不存在的脏数据,对账的时候怎么都对不上。解决删除识别问题有两种路径:一是推动业务方做软删除,也就是删除时不真正DELETE,而是把is_deleted置为1并更新update_time;这样基于时间字段的增量也能捕获到“删除”操作。二是直接用Binlog/CDC方案,Binlog天然包含DELETE语句的解析,变更流里每条被删数据都能拿到主键信息,到数仓里执行对应的删除操作。从我的经验看,凡是核心业务表,不要跟逻辑绕弯子,直接上Binlog方案才是最省心的。4. 实操过程:从全量切换为增量的完整实施步骤4.1 存量数据先行,增量后续衔接任何表要改增量采集,都绕不开一个问题:历史存量数据怎么处理。你不能直接停掉全量任务然后让增量任务从零开始跑,因为增量任务只认“从某个位置之后的数据”,位置之前的数据它完全不关心。所以正确的步骤是先导存量、再启动增量,两件事之间要接得严丝合缝。我实际操作中的做法是:T0时间点先记录一个快照位点,比如对Binlog方案就是记录当前的Binlog文件和位置;然后启动全量导出任务,把存量数据导入数仓;全量导出结束后,增量任务从刚才记录的位点开始消费。这个方案的关键点在于:记录位点的那一刻和全量导出开启的时刻必须保持一致,不能一边导全量一边让增量任务先跑个几分钟,否则两个任务的时间窗口重叠或断开,数据一定出问题。如果数据量特别大,全量导出要跑好几个小时,这段时间里业务产生的Binlog会不会过期?所以我在大表切换前会先检查Binlog保留时长,和DBA确认保留至少24小时以上,最好能覆盖全量导出的完整时间。4.2 调度与位点推进的一个参考实现下面给一个非常简化的Python实现思路,辅助理解调度推进的核心顺序。import json import time import redis from kafka import KafkaProducer r redis.Redis(host127.0.0.1, port6379) producer KafkaProducer(bootstrap_serverskafka:9092) def fetch_increment(last_point): conn get_mysql_conn() sql ( SELECT id, order_id, update_time, data FROM orders WHERE update_time %s ORDER BY update_time ASC, id ASC LIMIT 2000 ) return conn.query(sql, last_point) def process(): last_point r.get(orders:last_point) or 2024-01-01 00:00:00 while True: rows fetch_increment(last_point) if not rows: break # 注意这里先生产消息,但不推进位点 for row in rows: producer.send(ods_orders, json.dumps(row).encode()).get() # 每条消息都确认发送成功后,再批量推进位点 latest rows[-1][update_time] r.set(orders:last_point, latest) time.sleep(1) if __name__ __main__: process()这段代码表达的核心思想是:先发消息、确认消息可靠发送、再推进位点。get()等待Kafka响应是为了确保消息没有发送失败才推进位置。如果你的团队用Canal这种成熟组件,位点管理就交给Canal自己负责,任务的关注点可以转移到消费端逻辑上。4.3 监控与质量校验增量采集上线的下一步,就是建立质量监控。没有监控的增量任务就像闭着眼睛开车,看着仪表盘转速正常,其实数据早就断流了。我在增量采集的监控体系里至少会盯这几个指标:源库到采集端的采集延迟、采集端到Kafka的消费延迟、Kafka积压量、水位线推进速率、以及目标表的写入错误率。除了延迟,我还会做一种“条数对比”校验:每天定时把数仓目标表当天新增的数据量与源库当天新增数据量做个比对。差异超过阈值就告警,这样能很快发现采集过程中的静默丢失。月度层面再跑一次全量对账,比对两张表的主键集合,确认没有漏数据和重复数据。增量采集不是把全量任务停了就万事大吉,校验机制必须跟着同步设计。4.4 回填与补数流程增量跑着跑着,产品方忽然说“之前有一个业务逻辑算错了,需要重跑某一天的数据”,这不叫事故,这叫日常。增量采集系统必须天然支持按时间区间回填重放。我的处理方式是给采集任务设计一个“区间参数”,平时任务按水位线自动跑,回填时手动指定起始时间和结束时间,任务就会把指定时间内的数据重新捞一遍发到Kafka。下游消费时,因为目标表有唯一键,回填的数据会直接覆盖掉旧值,达到“重算”的效果。回填窗口不要太大,一次回填太大容易把源库和Kafka打爆,我一般建议单次回填控制在小时级,分批回填最稳妥。5. 常见问题与排查技巧实录5.1 增量任务静默空转我在接手别人的增量任务时,见过最隐蔽的问题就是:日志上显示跑批成功,每轮调度都在执行,但目标表的数据就是不动。这个问题排查起来很费劲,因为它没有任何显式报错。我当时的排查路径是:先在源库手动执行增量任务用的同一条SQL,发现返回结果正常;再看代码传参,结果发现任务把update_time当成了字符串,和日期类型比较时隐式转换,查出来的数据和预期完全不同。这类问题最常见的原因基本集中在新表上线时字段类型没确认、日期格式化没统一、源库与采集服务时区不一致。遇到数据量异常,先别急着怀疑网络和组件,拿同一条SQL去源库单独跑一下,就能排除掉一大半问题。5.2 采集条数比源库多或者少增量采集数据条数对不上,是线上最容易被业务方发现的问题。我总结的排查顺序非常固定:先看位点值是否异常;再对比源库和目标库按主键去重后的计数;最后看任务有没有被重复启动导致并发跑批。“采集多了”大概率是重复消费导致的,比如任务重跑一次、Kafka重复投递、或者下游没有去重。这个问题的处理手段在第三章已经聊过,核心就是唯一键加幂等。“采集少了”则大概率是位点推进过快,或者增量查询条件漏掉了边界值。这种时候把位点回拨到上一个批次,重跑一次往往就能把缺失的数据补回来。5.3 Binlog同步延迟导致数据不一致用了Binlog方案之后,有一个新的监控对象叫“同步延迟”。当源库发生大事务,比如DBA半夜跑一个大更新,影响了几百万行数据,对应的Binlog事件会瞬间涌出,消费端Canal或者采集程序处理不过来,积压就开始上涨,这个时候实时数据会明显滞后。延迟带来的问题是:查询目标表的时候,最新状态还没同步过去,业务方可能看到旧数据。我的经验是:延迟指标一定要配置告警,阈值根据业务容忍度来设,一般超过十分钟就报警。报警之后先确认是不是大事务导致的积压,如果是,等待积压自然消化即可,不要盲目重启任务。盲目重启可能导致消费位点回退,重复消费积压消息,反而把压力放大。5.4 全量增量衔接丢数案例最后分享一个我真正踩过的坑。当时做订单表从全量切增量,计划是先导全量再跑增量。结果全量导出任务跑完那一刻,增量任务还没有启动,中间隔了两三分钟。全量导出的数据是导出启动时的快照,而增量任务记录位点的时刻又晚了一点,中间的这两三分钟里新增的数据恰好落在这个空档里,既没有被全量带走,也没有被增量覆盖,结果就丢了一批数据。后来我把流程改成了严格一致的方案:在导出全量之前先记录一个Binlog位点,全量导出完毕后,增量任务从记录的位点开始消费。这样不管导出任务跑多久,中间产生的变更都等着增量来接管,空档期被彻底缝合。这个经验教训我后来写进了团队的切换手册,宁可多等一会儿,也不要在两个任务之间留出缝隙。5.5 常见问题速查表问题现象可能原因解决措施任务跑成功但目标表没数据查询条件类型不匹配/时区错误源库手动执行同SQL排查目标表数据比源库多重复消费、缺少幂等建唯一键,重放走覆盖更新目标表数据比源库少位点推进过快、边界遗漏回拨位点,重跑增量区间同步延迟持续上涨大事务Binlog积压、消费端性能不足停止重启,待积压消化并加大消费能力删除操作没有同步时间字段方案无法捕获DELETE软删除改造或改用Binlog方案增量采集从“能跑”到“跑得稳”,中间差了非常多细节。我自己最大的体会是:任何增量方案都必须配备补数和重放的能力,因为线上环境永远会有意外,你不能指望一条增量管道永远不出错。每次改动采集逻辑之前,备份一份水位线、保留一份操作记录,比事后拍脑袋排查要省心得多。数据同步这件事,慢一点不要紧,漏了才是大麻烦。
返回列表