简介:这份PDF资料围绕ETL中的全量与增量策略展开,面向数据仓库、大数据开发及数据同步方向的初中级工程师,帮助厘清两种抽取方式在采集、同步、构建与备份等场景下的差异与取舍。资源包共1个PDF文件,约79KB,内容以文字讲解为主,便于快速通读与查阅。资料从数据采集切入,对比全量抽取简单但数据量大、增量抽取复杂却对业务系统压力更小的特点;进而延伸到数据同步中全量覆盖、异步写与物理删除带来的隐患,以及增量同步长期易出现的数据一致性问题。文中还讨论了全量构建与增量构建在Cube更新、Segment合并及查询性能上的不同,并对比全量备份、增量备份与差异备份在速度、恢复和磁盘占用上的折中关系。目前已有3390人学习,适合需要系统梳理全量与增量选型思路、为实际项目方案做技术储备的读者参考。
1. ETL 全量与增量:为什么你写的同步任务总在凌晨三点崩
凌晨三点,调度平台弹出一串红色告警,一张三千万行的订单表同步任务跑了四个小时还没结束,源库连接数被占满,下游报表全线延迟。这种场景做数据集成的人多少都遇到过,而根子往往不在代码写得多烂,在于一开始就没想清楚:这份数据到底该全量抽,还是增量抽。
ETL 里的全量与增量,说的其实是两种数据搬运策略。全量是每次把源端符合条件的记录整批读出来,覆盖或追加到目标端;增量是只读取上次同步之后发生变化的那部分数据,靠时间戳、自增 ID 或日志位点来界定边界。选错了策略,轻则任务越跑越慢,重则数据重复、丢失、对不上账。这篇文章面向正在用 DataX、Kettle、Flink CDC 或自研脚本做数据同步的工程师,把两种模式的判断依据、落地写法、参数设置和踩坑点讲透,让你下次面对一张新表时能直接拍板用哪种,而不是先跑起来再说。
2. 全量与增量的选型判断:先看这四件事再动手
2.1 数据特征决定策略上限
判断一份数据适合全量还是增量,第一件事是看它有没有可靠的变更标识。所谓变更标识,就是能告诉你「这条记录是新的还是改过的」的字段。常见的有三类:自增主键、更新时间戳、数据库日志位点。
自增主键只适合新增场景。如果源表只插入不更新,比如日志表、流水表,用id > 上次最大 id就能拿到增量,简单可靠。但一旦有更新操作,自增主键就失效了,因为更新不会改变主键值,你根本不知道哪条被改过。
更新时间戳适用范围更广,前提是业务代码在每次写入时都维护这个字段。我见过太多表号称有update_time,结果一查发现大量记录的更新时间还是建表时的默认值,业务用的是ON DUPLICATE KEY UPDATE但没更新这个字段。这种表你按时间戳抽增量,改过的数据永远同步不过去。
日志位点是最可靠的方式,MySQL 的 binlog、PostgreSQL 的逻辑复制槽都属于这类。它不依赖业务字段,数据库层面记录所有变更。代价是需要开启相应配置,对源库有一定压力,而且解析链路更复杂。
第二件事是看数据量级和变更比例。一张一亿行的表,每天只变几百行,全量抽就是纯浪费;反过来一张十万行的配置表,每天改一半,增量抽的复杂度还不如直接全量覆盖来得省心。
第三件事是看目标端的写入能力。全量同步通常用批量覆盖或 truncate + insert,对目标库的写入压力集中在短时间内;增量同步是持续小批量写入,压力平摊。如果目标端是分析型数据库,批量写入效率远高于逐条更新,全量反而更快。
第四件事是看对数据一致性的要求。全量同步天然是一致性快照,只要抽取时加锁或走从库一致性读,拿到的就是某个时间点的完整数据。增量同步如果处理不好边界,容易出现漏抽或重复,需要额外的去重和校验机制。
2.2 三种典型组合的适用场景
把上面的判断落到具体组合上,常见的有三种。
全量覆盖适合小表、配置表、维度表。这类表数据量不大,变更频繁且无规律,每次直接全量拉取覆盖目标端,逻辑最简单,不会出错。典型的是商品类目表、地区编码表、汇率表。
增量追加适合只增不改的流水表。用自增 ID 或时间戳做水位线,每次只抽新数据追加到目标端。关键是要保证水位线的推进是单调的,不能因为任务重跑导致水位线回退。
增量合并(upsert)适合有增有改的业务表。抽取增量数据后,在目标端按主键做 insert or update。这种方式对源端压力小,但要求变更标识可靠,且目标端支持高效的 upsert 操作。
提示:不要在一张表上混用多种策略。我见过有人对同一张订单表,白天跑增量、晚上跑全量,结果目标端出现大量重复记录,排查了一整天才发现是两个任务的边界没对齐。
2.3 用 SQL 快速判断该用哪种策略
动手之前,先跑几条 SQL 摸清数据底细。下面这几条针对 MySQL,其他数据库改下语法即可。
-- 1. 看表有没有可靠的时间戳字段,以及更新是否真的在维护它 SELECT COUNT(*) AS total_rows, COUNT(update_time) AS non_null_update_time, MIN(update_time) AS earliest_update, MAX(update_time) AS latest_update FROM orders; -- 2. 看主键是否连续自增,有没有空洞 SELECT MIN(id) AS min_id, MAX(id) AS max_id, COUNT(*) AS actual_count, MAX(id) - MIN(id) + 1 AS expected_count FROM orders; -- 3. 看最近一天的实际变更量,判断增量比例 SELECT COUNT(*) AS changed_rows FROM orders WHERE update_time >= DATE_SUB(NOW(), INTERVAL 1 DAY);第一条 SQL 告诉你update_time的填充率。如果non_null_update_time远小于total_rows,说明这个字段不可靠,不能用来做增量水位线。第二条看主键有没有大量空洞,如果actual_count远小于expected_count,说明有大量删除操作,纯靠自增 ID 做增量会漏掉删除。第三条算出日变更量,用changed_rows / total_rows得到变更比例,低于 1% 考虑增量,高于 20% 考虑全量。
这三条 SQL 跑完,基本能拍板用哪种策略。我一般会把结果记在同步任务的配置注释里,后面换人维护时不用重新摸一遍。
3. 全量同步的落地写法:从批量分片到断点续传
3.1 全量抽取的分片策略与参数
全量同步最大的坑是「一把梭」——一条SELECT * FROM big_table直接拉,结果源库内存爆了,网络带宽打满,任务还卡死。正确做法是分片抽取。
分片的核心是找一个均匀分布的字段做切分键,通常是自增主键。按主键范围把大表切成 N 片,每片独立抽取,可以并行也可以串行。下面是一个 Python 分片抽取的骨架。
import pymysql from concurrent.futures import ThreadPoolExecutor def get_slice_ranges(conn, table, slice_size=500000): """按主键切分全量抽取范围""" with conn.cursor() as cur: cur.execute(f"SELECT MIN(id), MAX(id) FROM {table}") min_id, max_id = cur.fetchone() ranges = [] start = min_id while start <= max_id: end = start + slice_size - 1 ranges.append((start, end)) start = end + 1 return ranges def extract_slice(conn_params, table, start_id, end_id): """抽取单个分片,带重试""" conn = pymysql.connect(**conn_params) try: with conn.cursor(pymysql.cursors.SSCursor) as cur: cur.execute( f"SELECT * FROM {table} WHERE id BETWEEN %s AND %s", (start_id, end_id) ) for row in cur: yield row finally: conn.close() def full_sync(conn_params, table, slice_size=500000, workers=4): conn = pymysql.connect(**conn_params) ranges = get_slice_ranges(conn, table, slice_size) conn.close() with ThreadPoolExecutor(max_workers=workers) as pool: for start, end in ranges: pool.submit(extract_slice, conn_params, table, start, end)slice_size控制每片的大小,默认 50 万行。这个值太小会导致分片过多、连接频繁创建销毁;太大则单次查询耗时过长、失败重试成本高。我一般根据单行平均大小来调:单行 1KB 左右用 50 万,单行 10KB 以上降到 10 万。workers是并行度,不要超过源库能承受的连接数,生产环境一般 4 到 8 个。
SSCursor是流式游标,避免一次性把结果集加载到内存。这个细节很多人忽略,用默认游标抽大表,Python 进程内存直接飙到几个 G。
3.2 断点续传与幂等写入
全量同步跑到一半失败是常态,网络抖动、源库主从切换、目标端磁盘满,都会中断任务。没有断点续传,每次失败都从头再来,一张大表可能永远跑不完。
断点续传的实现思路是记录每个分片的完成状态。分片粒度要足够细,细到单次失败重试的成本可以接受。下面是一个基于状态表的管理方式。
-- 分片状态表 CREATE TABLE sync_slice_status ( id BIGINT AUTO_INCREMENT PRIMARY KEY, table_name VARCHAR(128) NOT NULL, slice_start BIGINT NOT NULL, slice_end BIGINT NOT NULL, status TINYINT DEFAULT 0 COMMENT '0待处理 1处理中 2已完成 3失败', retry_count INT DEFAULT 0, updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, UNIQUE KEY uk_table_slice (table_name, slice_start, slice_end) );任务启动时先查状态表,跳过已完成的片,只处理待处理和失败的片。每片开始前把状态置为 1,完成后置为 2,失败置为 3 并累加retry_count。retry_count超过阈值(比如 3 次)就告警人工介入,不要无限重试。
写入端要保证幂等。全量覆盖场景下,如果目标端是整表替换,用临时表 + rename 的方式最安全:先把数据写到table_tmp,全部完成后RENAME TABLE table TO table_old, table_tmp TO table。这样切换是原子的,下游不会读到半截数据。如果是分片追加,目标端要按主键做 upsert,避免重跑时产生重复。
注意:
RENAME TABLE在 MySQL 里是原子操作,但如果目标端有外键引用,rename 会失败。这种情况改用INSERT OVERWRITE或先删后插,但要接受短暂的空窗期。
3.3 全量同步的资源控制
全量同步对源库的冲击是集中式的,必须做资源控制。几个关键参数:查询超时时间、单次读取行数、写入批次大小。
查询超时设太短,大分片查不完就超时;设太长,慢查询拖垮源库。我一般设 300 秒,配合分片大小控制单次查询时长在 60 秒以内。单次读取行数用游标的fetchmany控制,每次取 1000 到 5000 行,避免网络往返太频繁。写入批次大小和读取批次对齐,批量 insert 比逐条 insert 快一个数量级。
源库层面,全量抽取尽量走从库,并且加/*+ MAX_EXECUTION_TIME(60000) */这类 hint 限制执行时间。如果源库是 MySQL,还可以在会话级别设置SET SESSION transaction_isolation='READ-COMMITTED',减少锁竞争。
4. 增量同步的落地写法:水位线、去重与迟到数据
4.1 水位线的选择与推进
增量同步的核心是水位线,也就是「上次同步到哪儿了」的标记。水位线选得好,增量就稳;选不好,要么漏数据,要么重复抽。
时间戳水位线最常见,但有两个坑。一是时间戳精度,如果源表update_time只精确到秒,同一秒内多条记录更新,水位线推进到这一秒,下一批可能漏掉同秒的其他记录。解决办法是水位线回退一秒,配合目标端去重。二是时钟回拨,服务器时间被 NTP 校正后可能倒退,导致水位线回退,重复抽取。所以水位线要单调递增,发现新值小于旧值时不推进。
自增 ID 水位线适合只增不改的表,简单可靠。但要注意删除操作,自增 ID 无法感知删除,需要额外的删除标记或定期全量比对。
日志位点水位线最可靠,binlog 的position或 PostgreSQL 的LSN都是单调递增的,且包含所有变更类型。代价是链路复杂,需要解析日志。
下面是一个时间戳水位线的推进逻辑。
def get_watermark(conn, table): """读取当前水位线""" with conn.cursor() as cur: cur.execute( "SELECT watermark_value FROM sync_watermark WHERE table_name = %s", (table,) ) row = cur.fetchone() return row[0] if row else '1970-01-01 00:00:00' def update_watermark(conn, table, new_value): """推进水位线,保证单调递增""" with conn.cursor() as cur: cur.execute( """UPDATE sync_watermark SET watermark_value = %s WHERE table_name = %s AND watermark_value < %s""", (new_value, table, new_value) ) if cur.rowcount == 0: # 新值不大于旧值,不推进,记录告警 print(f"watermark not advanced for {table}: {new_value}")WHERE watermark_value < %s这个条件保证了水位线只前进不后退。如果rowcount为 0,说明新值没有超过旧值,可能是时钟回拨或重复抽取,这时候不推进水位线,但要记录日志排查原因。
4.2 目标端去重与 upsert
增量抽取的数据在目标端写入时,必须处理重复。重复的来源有两个:水位线回退导致的重复抽取,以及任务重跑。处理方式取决于目标端类型。
如果目标端是 MySQL,用INSERT ... ON DUPLICATE KEY UPDATE。如果目标端是 PostgreSQL,用INSERT ... ON CONFLICT ... DO UPDATE。如果目标端是 Hive 或 ClickHouse 这类分析型数据库,通常用ReplacingMergeTree或定期 merge 去重。
-- MySQL upsert 示例 INSERT INTO orders_target (id, user_id, amount, status, update_time) VALUES (%s, %s, %s, %s, %s) ON DUPLICATE KEY UPDATE user_id = VALUES(user_id), amount = VALUES(amount), status = VALUES(status), update_time = VALUES(update_time);ON DUPLICATE KEY UPDATE依赖目标表的主键或唯一索引。如果目标表没有主键,这个语句会退化成普通 insert,重复数据照样进去。所以目标表建表时一定要有主键,哪怕源表没有,也要用业务字段组合出一个唯一键。
批量 upsert 时,把多条记录拼成一条INSERT ... VALUES (...), (...), ... ON DUPLICATE KEY UPDATE,比逐条执行快得多。但要注意单条 SQL 的长度限制,MySQL 默认max_allowed_packet是 4MB,拼太多会报错。我一般每批 500 到 1000 条。
4.3 迟到数据与乱序处理
增量同步还有一个隐蔽的坑:迟到数据。源端的事务提交时间和update_time字段的写入时间可能不一致。比如一个事务在 10:00:00 开始,10:00:05 提交,但update_time写的是 10:00:00。如果水位线在 10:00:03 推进了,这条记录就被漏掉。
处理迟到数据的常见做法是水位线延迟推进,比如每次推进到max(update_time) - 5 分钟,给迟到数据留出窗口。代价是同步延迟增加 5 分钟。另一种做法是定期做全量比对,找出增量漏掉的记录补上。
乱序处理是指同一主键的多条变更记录到达目标端的顺序和源端不一致。如果目标端是 upsert,最终值取决于最后到达的那条,可能不是最新的。解决办法是在目标端按update_time做条件更新,只允许更新的记录覆盖旧的。
INSERT INTO orders_target (id, amount, update_time) VALUES (%s, %s, %s) ON DUPLICATE KEY UPDATE amount = IF(VALUES(update_time) > update_time, VALUES(update_time), amount), update_time = IF(VALUES(update_time) > update_time, VALUES(update_time), update_time);这个写法保证只有update_time更大的记录才能覆盖目标端,避免乱序导致的数据回退。
5. 避坑与排查:全量增量同步的五个血泪教训
5.1 全量任务把源库连接打满
现象:全量同步任务启动后,源库监控显示连接数飙升到上限,业务查询开始超时。
原因:分片并行度设得太高,每个分片一个连接,加上连接池没有复用,瞬间创建大量连接。
解决:把workers降到源库能承受的范围,一般不超过 CPU 核数的两倍。连接池设置max_connections上限,超过就排队等待。全量任务尽量安排在业务低峰期跑。
5.2 增量任务漏抽更新记录
现象:源表某条记录被更新了,但目标端还是旧值,增量任务显示已同步完成。
原因:update_time字段没有被业务代码维护,更新时没有刷新这个字段,增量抽取按时间戳过滤时漏掉了这条。
解决:先跑 2.3 节的 SQL 检查update_time填充率。如果不可靠,改用 binlog 或触发器方案。如果必须用时间戳,推动业务方修复写入逻辑,或者在 ETL 层加全量比对兜底。
5.3 水位线回退导致数据重复
现象:目标端出现大量重复记录,主键冲突告警频繁。
原因:服务器时钟被 NTP 校正后回拨,水位线新值小于旧值,任务重新抽取了已经同步过的数据。
解决:水位线推进加单调递增条件(见 4.1 节代码)。目标端写入用 upsert 保证幂等。监控水位线变化,发现回退立即告警。
5.4 全量切换时下游读到空数据
现象:全量同步完成后,下游报表查询报错,提示表不存在或数据为空。
原因:用DROP TABLE+CREATE TABLE+INSERT的方式切换,中间有空窗期,下游查询正好落在空窗期。
解决:用RENAME TABLE原子切换,或者用临时表写完再 rename。如果目标端不支持 rename,用事务包裹删除和插入,但大表事务会锁很久,不推荐。
5.5 增量任务延迟越来越大
现象:增量同步任务刚开始延迟几秒,跑着跑着延迟变成几小时,且持续增长。
原因:单批处理的数据量超过了处理能力,每批处理时间大于数据产生速度,积压越来越多。
解决:先看单批处理耗时,如果超过批次间隔,要么提高并行度,要么减小批次大小。如果是目标端写入慢,检查目标端索引是否过多、是否有锁竞争。如果是源端抽取慢,检查增量字段有没有索引,没有索引的WHERE update_time > ?会全表扫描。
6. 用校验和监控把同步质量兜住
同步任务跑起来只是开始,能不能持续跑对才是关键。我一般会在同步链路里加两层保障:数据校验和运行监控。
数据校验分两种粒度。行数校验最简单,源端和目标端各跑一次COUNT(*),对比是否一致。但行数一致不代表数据一致,更新丢失的场景下行数完全一样。所以还需要内容校验,常用的是对主键和关键字段做MD5或CRC32聚合。
-- 源端和目标端各跑一次,对比结果 SELECT COUNT(*) AS row_count, MD5(GROUP_CONCAT(id ORDER BY id)) AS id_checksum, MD5(GROUP_CONCAT(CONCAT_WS(',', id, amount, status) ORDER BY id)) AS content_checksum FROM orders WHERE update_time >= '2024-01-01' AND update_time < '2024-01-02';GROUP_CONCAT有长度限制,默认 1024 字节,大表要调大group_concat_max_len或者分批校验。更稳妥的方式是按主键分片,每片单独算校验和,对比不一致的片再细查。
监控方面,我关注四个指标:同步延迟(源端最新数据时间减去目标端最新数据时间)、批次耗时、失败重试次数、水位线推进速度。延迟超过阈值告警,批次耗时突增告警,重试次数超过 3 次告警,水位线长时间不推进告警。
最后一个习惯:每次上线新的同步任务,先跑一周的「影子模式」——同步任务照跑,但下游不切流量,每天对比源端和目标端的数据。一周后确认无误再切正式流量。这个习惯帮我拦下过好几次update_time字段不可靠、目标端主键缺失的问题。同步这件事,宁可上线慢一点,也别让脏数据流到下游。希望帮到你。
本文还有配套的精品资源,点击获取