ETL全量与增量同步方案:从基线建立到断点修复的实战指南

发布时间:2026/9/18 19:18:06
ETL全量与增量同步方案:从基线建立到断点修复的实战指南 简介全量与增量是ETL数据抽取中的核心话题这份PDF学习笔记面向大数据开发工程师、数据仓库工程师以及正在设计数据同步方案的读者通过直观对比帮助读者快速理解两者的适用边界与取舍原则。文档系统梳理了数据采集、数据同步、Cube构建、数据备份四个典型场景下的差异全量抽取简单直接但数据量大增量抽取相对复杂但对业务性能压力小全量同步难捕获物理删除需借助操作日志增量同步逻辑复杂易造成数据不一致全量构建计算量大但查询无需合并Segment增量构建计算量小但需处理Segment合并备份方面则对比了全量、增量、差异三种方式在恢复速度与空间开销上的权衡。资源为1个PDF文件压缩包仅79KB篇幅精炼但覆盖关键知识点。已有3389人学习适合快速建立全量与增量选择的判断框架也可作为内部培训或技术分享的参考材料。1. ETL 全量与增量同步方案决定数据上限凌晨 2 点 17 分数仓调度报警。一张 2 亿行的订单表全量抽取跑了 73 分钟还未进入清洗层ODS 磁盘先红了。这不是第一次——每次大促后的周一全量链路都像定时炸弹加并行度、换大内存机器治标不治本。问题不在资源而在于“全量”这两个字本身它把整张表像历史书一样一遍遍重演。全量与增量不是两个孤立的 ETL 同步模式而是调度层默认的两种世界观。全量负责建立基线、恢复崩溃、修复脏数据增量负责在两次全量之间用最小代价传递变更。真正决定数仓吞吐上限的不是引擎跑多快而是这两个模式怎么接、怎么切、怎么兜底。后面所有讨论都围绕 ETL 概念里这个最常见的二选一场景展开什么时候必须全量什么时候用增量以及两者在真实任务里的边界参数与组合方案。2. ETL 全量策略基线建立与重置语义全量不是简单“把 A 表搬到 B 表”它的本质是重置语义目标表在某个时间点必须与源表完全一致且不依赖任何历史状态。所以全量任务里最关键的一行不是select *而是先drop或truncate。一个常见事故是增量任务连续跑了一个月某天源表历史数据被订正增量没接住于是决定全量重刷。此时目标表里还留着旧数据全量脚本直接覆盖写新旧数据在同一主键下“混搭”。正确做法是先删后插或者先写临时表再原子替换。2.1 用 Spark 重刷全量基线的标准写法大批量全量同步常见做法是走 Spark 这类分布式引擎而不是用 JDBC 一条条塞。原因有两个一是源库连接数打满会拖垮线上二是全量天然适合“读整表、写整表”的批处理语义。from pyspark.sql import SparkSession spark ( SparkSession.builder.appName(full_restore_order) .config(spark.sql.shuffle.partitions, 200) .config(spark.sql.adaptive.enabled, true) .enableHiveSupport() .getOrCreate() ) # 1. 先读源表只拿业务需要字段并统一字段类型 df ( spark.read.format(jdbc) .option(url, jdbc:mysql://prod-db:3306/erp?useSSLfalse) .option(dbtable, t_order) .option(user, reader) .option(password, ***) .option(fetchsize, 5000) .load() .selectExpr(id, user_id, cast(amount as decimal(18,2)), gmt_modified) ) # 2. 写临时目录避免直接覆盖目标表 df.write.mode(overwrite).parquet(/tmp/dw/t_order_stage) # 3. 原子替换先删目标分区文件再 rename target_path /user/hive/warehouse/dw.db/t_order spark.sql(fTRUNCATE TABLE dw.t_order) spark.read.parquet(/tmp/dw/t_order_stage).write.mode(overwrite).parquet(target_path)这段代码有两个值得注意的参数。fetchsize5000控制 JDBC 一次性拉取的行数对 MySQL 这类引擎来说值太大会让内存暴涨太小会增加往返次数建议 100010000 之间按源库规格调。spark.sql.shuffle.partitions200影响后续聚合和 Join 的并行度如果只做纯抽取200 够用如果后面带 groupBy建议按每内核 24 个分区换算。空洞TRUNCATE是安全性的核心。覆盖写不能保证物理清理残留文件旧 schema 的列可能残留在 parquet 的 footer 里导致后续SELECT *行为不确定。先TRUNCATE再写等同于把表重置为“空状态”。这也是全量备份场景下最常见的坑看着行数对了但残留分区导致查询读到双份数据。2.2 分区交换不停服的全量替换如果目标表是数百 GB 的分区表每次全量都TRUNCATE再写代价是查询服务会看到中间状态。常见做法是“写新分区再交换挂载”。-- 假设 dw.ods_order 是按 dt 分区的外部表 ALTER TABLE dw.ods_order DROP IF EXISTS PARTITION (dt 2025-11-17); ALTER TABLE dw.ods_order ADD PARTITION (dt 2025-11-17) LOCATION /data/dw/ods_order/full_build/dt2025-11-17;这里的核心是先把全量结果完整写入一个新路径写完后用ADD PARTITION把路径挂载到表上。整个过程不触碰旧分区业务查询在挂载前后看到的都是完整数据。与INSERT OVERWRITE相比这种方式的好处是如果新数据校验失败可以直接不执行挂载命令旧分区还在原地待命。对比维度TRUNCATE 重写分区交换停机窗口需要不需要失败回滚丢旧数据保留旧分区适合规模小表GB 级大表百 GB 级实现复杂度低中分区交换在数据仓库的数仓建模里更常用但实时链路做“全量重建”时同样适用。原则只有一条宁可多占一份磁盘也要让“替换”这一动作具备原子性。3. ETL 增量策略水印、日志与版本号全量基线建立之后日常同步必须切换到增量模式否则就是每天把整张表重读一遍。增量同步的本质是回答三个问题从哪个点开始拿数据拿哪些数据怎么处理更新和删除三种主流方案水印字段modified_time字段、CDC 日志binlog/WAL、版本号快照。对大多数 MySQL 到数仓的同步需求水印方案最容易落地CDC 方案则适合删除操作频繁或字段更新被隐藏的场景。3.1 水印表设计把同步位点留给数据库不留给内存水印同步的基础不只是一条where条件还需要一张“同步位点表”记录每个同步任务已经读到的时间点。很多人把位点记在配置中心或环境变量里进程一重启就丢后果是重复跑一大段或漏掉一小段。-- 位点表结构 CREATE TABLE sync_checkpoint ( job_name VARCHAR(64) PRIMARY KEY, last_sync DATETIME, updated_at DATETIME ); -- 增量抽取读取位点关闭旧数据 SELECT id, order_no, amount, gmt_modified FROM erp.t_order WHERE gmt_modified 2025-11-17 02:30:00 AND gmt_modified 2025-11-17 03:00:00;增量 SQL 的时间范围要写成(last_sync, current_window_end]左侧开区间右侧闭区间且窗口固定在一个确定长度内不能自然是“从当前时间往前推 5 分钟”这种随时间漂移的写法。窗口漂移会导致上一轮读到02:29:59的数据这一轮又从02:35:00开始中间 5 分钟成为盲区。更隐蔽的问题在事务时间戳。MySQL 的gmt_modified由应用层写入提交时刻晚于写入时刻会出现“晚写入早提交”的乱序。解决办法是窗口内多查一遍WHERE gmt_modified 2025-11-17 02:30:00 AND gmt_modified 2025-11-17 03:00:00 UNION WHERE gmt_modified 2025-11-17 02:30:00 AND gmt_modified 2025-11-17 03:00:00 AND id IN (SELECT id FROM erp.t_order WHERE gmt_modified 2025-11-17 02:30:00)不推荐用UNION解决正确解法是把窗口重叠率设在 10% 以上也就是说上一轮的结束时间03:00:00提前 10% 的窗口时长作为下一轮起点然后再按主键去重。代价是多读少量重复行换来“不漏数据”的确定性。3.2 CDC 日志增量删除操作与精确一致性水印字段的软肋是删除源表记录被物理删除后水印无法感知。此时需要切换 CDC 方案常见的是 MySQL binlog row 模式。binlog 会记录每一行的before image与after image包含删除事件这是水印方案做不到的。表水印与 CDC 的核心差异能力时间戳水印binlog CDC新增支持支持更新支持支持删除不支持持久标记支持历史重放按时间窗口按 binlog 位点对源库影响一次全表扫描索引较低靠 dump 线程运维成本低中高需要管理 binlog 文件使用planetscale或maxwell这类 mysql 增量同步工具连接 binlog 时需要记录两个位点binlog_file和binlog_pos不能只记时间。时间戳在 MySQL 里是会话时区相关的跨时区部署会出偏移位点则是唯一确定事件位置的逻辑指针。同步完成后位点更新最好放在目标端事务里和写入数据一起提交避免“数据写了但位点没更新”造成重复消费。3.3 增量任务里的版本号陷阱有一部分表没有gmt_modified只有流水号seq。用seq last_seq拉增量看起来没问题但源库做全量备份恢复时seq会倒退。此时增量任务必须能被手动重置到全量后的新基线。一个经验是纯自增 ID 只能保证新增不能保证更新。如果表有更新需求即便有主键 ID也得额外建last_updated_version字段由应用层在 UPDATE 时更新。否则会在 ETL 面试题里背过“用 ID 做增量为什么丢更新”后又在生产环境踩一遍。4. ETL 工程落地双写、重跑与一致性校验设计好全量和增量模式后最关键的环节是调度层的容错。增量任务每天跑 48 个小批次任意一批失败都要能断点续跑而续跑的核心是“幂等”。幂等不是靠任务框架保证而是靠任务自身的写入逻辑保证。4.1 批次号回放让重跑不重数据给每个增量批次一个全局唯一批次号batch_id目标表结构里加batch_id和sync_time两列。每次回放时把本批次号对应的旧数据删除后再插入。DELETE FROM dw.ods_order WHERE batch_id 20251117030000; INSERT INTO dw.ods_order SELECT id, order_no, amount, gmt_modified, 20251117030000 AS batch_id, NOW() AS sync_time FROM stg_order_incr WHERE data_date 2025-11-17;这里的规则是先按batch_id删除再插入。重复执行同一批次不会产生双份记录。如果没加batch_id重跑就是简单INSERT同一条订单出现两遍如果加batch_id但重跑时换了一个新批次号旧数据依旧残留。所以批次号必须由调度系统根据逻辑时间生成不能取new Date()。4.2 重试参数先看任务类型再决定次数全量任务与增量任务的重试策略完全不同。全量任务重试成本极高随机重试三次等于把源库读三遍增量任务重试成本低但重试会累积延迟。我常用的做法是给增量任务设置较短的依赖等待和快速重试给全量任务设置“失败后置为等待人工确认”。参数增量场景建议值全量场景建议值单次失败重试次数30 或 1重试间隔30 秒手动超时时间窗口长度的 1.5 倍无固定值失败处理重跑当前批次回滚到旧分区增量任务重试时还要处理“半成功”状态即一批数据写入一部分后失败。切换目标表为merge写入模式或者依赖上面的batch_id清理都能保证半途状态可恢复。但不建议为了省事把增量任务也设计成全量重刷那样等于把全量和增量的最优路径都毁掉了。4.3 双跑期间的一致性校验从全量切换到增量或者反过来都有一个“双跑过渡期”。这个阶段要避免直接信任下游报表建议加一道校验行数与校验和。Checksum 校验的实现不复杂关键是选什么字段。source_sql SELECT id, amount, gmt_modified FROM erp.t_order target_sql SELECT id, amount, gmt_modified FROM dw.ods_order def checksum(conn, sql): cur conn.cursor() # 按主键排序后做拼接避免不同批次顺序影响结果 cur.execute(f SELECT SUM(CRC32(CONCAT_WS(|, id, amount, gmt_modified))) FROM ({sql}) t ) return cur.fetchone()[0] src_hash checksum(mysql_conn, source_sql) dst_hash checksum(hive_conn, target_sql) assert src_hash dst_hash, checksum mismatch, trigger rebuild校验时最容易犯的错误是只对比行数。行数一致但内容不一致的情况在 ETL 里极其常见比如金额字段源库是DECIMAL(10,2)目标表变成FLOAT精度丢失后行数不变数据已经错了。把金额、时间、状态字段拼进CRC32里能显著提高发现问题的概率。注意CRC32本身有碰撞几率正式衡量用XXHASH64或MD5会更好但代价是计算时间成倍增加。5. ETL 兜底技巧用校验快照修复增量断点增量任务跑太久之后最怕出现“不知道断点在哪”的混沌状态。这里给一个我在有赞、饿了么的朋友圈里常用的兜底方案周期性校验快照。即在每 N 个增量窗口结束后额外生成一张轻量级快照表记录每张关键表的源端校验值。具体做法是在全量基线建立时生成第一份快照之后每 10 个增量批次更新一次快照。快照表只存主键、最大时间戳、行数、checksum不存业务明细。增量断点发生后先不急着回放而是读取上一份快照与源库当前状态做对比-- 快照表结构简化版 CREATE TABLE dw.checkpoint_snapshot ( table_name VARCHAR(64), snapshot_time DATETIME, row_count BIGINT, max_modified DATETIME, checksum_agg BIGINT, PRIMARY KEY (table_name, snapshot_time) ); -- 修复前校验找出断点区间的真实增量范围 SELECT COUNT(*), MAX(gmt_modified), SUM(CRC32(CONCAT_WS(|, id, amount))) FROM erp.t_order WHERE gmt_modified 2025-11-17 02:30:00 AND gmt_modified 2025-11-17 03:00:00 HAVING COUNT(*) ! ( SELECT row_count FROM dw.checkpoint_snapshot WHERE table_name t_order AND snapshot_time 2025-11-17 02:30:00 );校验结果显示不一致就用这段区间重新跑一轮增量补齐结果显示一致则直接把断点推进到快照时间避免无脑重放整段 binlog。这个技巧的价值在于把“全量”和“增量”真正焊接到一起。日常的增量不必战战兢兢因为快照提供了周期性的全量参考系而当增量真的断裂时也不需要立刻重建几百 GB 的全量表只要回滚到最近一份校验快照按区间补齐即可。一个同步系统的健壮性不在于全量跑得多快也不在于增量延迟多低而在于它们之间的切换足够平滑、恢复路径足够短。把这套快照机制做成自动调度任务比任何单一引擎参数调优都更能护住数仓的底线。本文还有配套的精品资源点击获取