简介这份PDF面向大数据ETL开发与数据仓库建设人员系统梳理全量与增量两种数据抽取策略的适用边界。内容从数据采集、数据同步、Cube构建到备份恢复四个维度展开对比点明全量抽取简单但数据量大、增量抽取精准但逻辑复杂的核心差异并揭示全量同步中数据物理删除难以捕获的隐患及借助日志文件追踪的思路。读者可据此结合业务场景选择合适方案避免因策略误用造成资源浪费或数据不一致。资源为单个PDF文档压缩包仅79KB便于快速阅读已有3389人浏览学习适合对ETL有一定基础、需要系统理解增量和全量区别的开发者参考。1. ETL 全量与增量选错初始策略后面全在还债一个三千万行的订单表每天凌晨用全量重跑从 MySQL 拖到 Hive 要 40 分钟下游报表只能跟着深夜加班。很多团队把这个问题归结为集群不够快其实根子在 ETL 策略全量加工把所有变更历史都重新搬了一遍。全量抽取逻辑简单、无状态适合小表增量抽取只搬运变化行但需要维护水位线、处理乱序和删除。ETL 里的“全量与增量”不是两个按钮而是你先得搞清楚源表能否支撑增量计算、下游能容忍多大延迟再决定初始化用哪种、长期同步用哪种。这篇文章把概念、实现机制和混合调度串起来讲。2. ETL 全量与增量的概念边界与选型判据2.1 ETL概念全量抽了全部行增量抽了变化行全量抽取不需要知道上次状态每次SELECT * FROM整表导入即可。增量抽取的核心问题变成如何定位“变化”。按实现难度我一般把增量分成三档基于时间戳字段、基于版本号/状态位、基于数据库日志解析。时间戳最简单但要求表里有可靠的update_time并且数据库能索引到该字段状态位适合业务可控的同步但对外系统往往不愿意改表日志解析能捕获 DELETE 和物理删除是实时链路的主流。边界划定后要注意全量和增量不是互斥的同一张表在不同阶段可以用不同策略。第一次同步必须先做全量然后才有资格谈增量位点。很多 ETL 工具会在内部完成这件事比如 Flink CDC 里配置snapshot.modeinitial就是先做全量快照再自动切到 binlog如果一开始就只配增量模式要么丢数据要么直接初始化失败。2.2 选型判据数据量、频率、源表约束判断主表用全量还是增量不能只按行数拍脑袋。我早期吃过亏一张 2000 万行的配置表每天就改 20 行全量跑半小时这明显应该增量反过来一张 10 万行的订单表每秒都在写如果下游要准实时增量比全量难控制得多。所以我会看四个指标源表总行数、每日变更行数、源表字段里有没有可用于增量的标记、下游 SLA。变更行数占比少且下游要求高优先增量行数小且更新频繁反而全量更稳。2.2.1 一个简单的策略判断脚本我常在调研阶段写这种半自动脚本用来和业务方确认边界# etl_strategy.py def choose_strategy(rows: int, change_rows_per_day: int, latency_min: int) - str: 根据预估量输出 ETL 策略方向。 rows: 源表行数 change_rows_per_day: 每天变更行数 latency_min: 下游可接受延迟分钟 change_ratio change_rows_per_day / max(rows, 1) # 每天更新比例低于 5%且对延迟敏感优先增量 if change_ratio 0.05 and latency_min 60: return incremental # 数据量小全量重跑成本低 if rows 500000 and latency_min 60: return full # 变更占比大且行数大增量也没多大用 return full_or_hybrid脚本本身不负责执行只做策略初筛。latency_min是最关键的参数——实时大屏给不到 5 分钟就不适合用离线全量再打补丁T1 报表能接受全量重跑就没必要给每张表维护水位线。阈值 5% 和 50 万不是通用值我通常是按目标仓库的吞吐反推跑一次全量要 30 分钟以上的表才有必要进增量池。2.3 全量与增量对比表维度全量抽取增量抽取数据范围源表全部行上次同步后的变化行实现成本低SQL 简单高需维护水位线或解析日志历史更新覆盖天然覆盖容易漏需额外设计物理删除自然消失时间戳方案无法感知目标表写入可重跑覆盖需要幂等键典型延迟批级分钟级或秒级这张表解决的不只是选型还决定目标表建模。全量多用INSERT OVERWRITE或者目标临时表替换增量必须给目标表设计主键或唯一键用来做 upsert否则增量任务跑第二次目标表就会出现重复数据。这些在 ETL 面试题里经常出现但落地时容易被忽略。另外全量抽取在目标端的表现也要纳入考虑。如果全量只是简单覆盖分区表可以按天分区date20250301写如果不是分区表全量覆盖会锁表。所以选全量时我也会确认目标表是否按日分区没有分区就先做分区改造。增量抽取则不同目标表必须保留当前最新状态通常把源端主键作为目标表唯一键使用ON DUPLICATE KEY UPDATE或 SQL Merge。很多团队把增量表建成了流水表只插入不更新结果下游要按主键取最新值时逻辑复杂等于把问题转嫁给数据应用。所以策略选择实际上会定义目标表模型全量表多是快照表增量表多是镜像表或拉链表。3. ETL 增量抽取的三种机制与 Spark ETL 脚本落点3.1 时间戳增量水位线怎么取才不会漏时间戳增量是做到增量最便宜的方式。只要源表有一个更新时间字段比如update_time每次同步便执行SELECT id, col1, col2, update_time FROM biz_order WHERE update_time :last_ts AND update_time :now_ts ORDER BY update_time;配合脚本把last_ts持久化到 Meta 表或本地文件。我用本地文件维护水位线时常用的逻辑是本次消费完一批取这批数据中最大的update_time作为下一次起点而不是以任务启动时间也不能直接覆盖旧状态。一个常见坑是源库的update_time精度是秒级而业务在同一秒内大量更新水位线一旦等于max(update_time)下一秒从该时间点开始查询会把这一秒里的部分行漏掉。稳妥做法是将水位线向前退一步比如last_ts max_seen - INTERVAL 1 SECOND消费结果写入目标时做 upsert靠幂等兜住重复但不漏数。3.1.1 Python 版本的水位线状态机# sync_order.py import mysql.connector import json from datetime import timedelta def sync_orders(conn, state_file): try: f open(state_file) last_ts json.load(f)[last_ts] except FileNotFoundError: last_ts 2000-01-01 00:00:00 sql SELECT id, update_time FROM orders WHERE update_time %s AND update_time NOW() ORDER BY update_time cursor conn.cursor(dictionaryTrue) cursor.execute(sql, (last_ts,)) rows cursor.fetchall() # 实际生产这里会写入目标仓库并做 upsert for row in rows: upsert_order(row) # 水位线回退 1 秒防止同一个秒级时间戳被边界截断 max_ts max(r[update_time] for r in rows) new_ts max_ts - timedelta(seconds1) if rows else last_ts json.dump({last_ts: str(new_ts)}, open(state_file, w))state_file是每个同步表的位点文件upsert_order为示意实际需要写成事务否则半途失败会导致水位线提前提交。我发现时间戳实现的问题大多不在 SQL 而在提交顺序先提交数据再更新水位线数据重复但可控先更新水位线再提交数据直接丢一批。因此代码里优先保证数据落库位点稍旧不要紧。3.2 基于日志解析的 CDC 增量Binlog 与消息队列时间戳方案对DELETE无解。如果业务表会被物理删除即使删了十行水位线没有感知目标表仍留着旧数据。这时要用数据库日志解析也就是 CDCChange Data Capture。常见的链路是 MySQL binlog - Canal/Debezium - Kafka - Spark Streaming。这种方案不再依赖业务字段binlog 里天然包含 insert、update、delete 三类事件并且带事务坐标或行位点可以做到精确消费。我建议不要在 ETL 里自己解析 binlog优先复用 Maxwell、Canal、Flink CDC 这类增量同步工具。它们把复杂位点和类型映射都处理好了ETL 脚本只需要消费 JSON 事件。下面是一段 Spark ETL 脚本消费 CDC 事件的简化写法from pyspark.sql import functions as F df (spark.readStream.format(kafka) .option(kafka.bootstrap.servers, kafka:9092) .option(subscribe, mysql_biz_order) .load()) events df.selectExpr(CAST(value AS STRING) as json, partition, offset) parsed (events .select(F.from_json(json, schema).alias(data)) .select(data.op, data.after.*, data.ts_ms)) # 根据事件类型直接 upsert这里的关键参数是subscribe的 topic 名以及schema中定义的after结构。Spark ETL 脚本通常只做转换和写出不负责水位线管理binlog 的位点由消费组提交到 Kafka或者由 Flink CDC 的 checkpoint 管。这样做的收益是 ETL 脚本天然支持乱序和重放比时间戳方案健壮得多。CDC 链路还要注意重复投递binlog 解析重跑时如果消费组没有提交 offset下游会看到同样的事件必须让目标写入操作幂等。Flink CDC 通过 checkpoint 可以接近 exactly once但 Canal 到 Kafka 再消费还是要自己处理去重通常用业务主键加事件时间做去重。3.3 MySQL 增量同步工具选型四个方向看参数做增量同步时大家最先搜的是 mysql 增量同步工具。我常备四个方向的备选Canal、Maxwell、Flink CDC、以及常规离线工具 Sqoop 的增量模式。Sqoop 的 lastmodified 只能算弱增量无法捕获删除Canal 与 Maxwell 都需要在 MySQL 开启binlog_formatROW并预留 binlog 保留时长Flink CDC 内置全量和增量自动衔接比较适合新项目。工具增量方式位点推进适合场景Sqooplastmodified 时间戳--last-value手动管理HDFS 离线批增量Canalbinlog 解析meta.dat 自动Java 生态、订阅到 MQMaxwellbinlog 解析自动存到表中直接投 Kafka JSONFlink CDCsnapshotbinlogcheckpoint 推进流批一体、实时数仓选型前确认源库 binlog 的保留时间如果保留 24 小时而增量任务挂了三天恢复时从上一个 checkpoint 继续会找不到 binlog 文件位置。Canal 有默认位点失效策略Maxwell 会配置replica_server_id防止和从库冲突这些参数必须在压测阶段调好而不是上线后再说。4. ETL 全量 增量混合调度的落地脚本与状态机4.1 全量初始化再切入增量的执行顺序一个稳定链路通常分三步第一步全量抽取把目标表建成一份当前快照第二步在初始化结束那一刻记一个位点或消费组 offset第三步切到增量。很多工具已经封装了这件事离线场景里常见的做法是先生成 CSV 或 Parquet再执行增量脚本。这里的关键是全量和增量不能在同一时间并发写同一张目标表否则全量覆盖会把增量刚写入的数据抹掉。我一般用一张sync_state表控制切换。阶段动作状态值初始全量抽取到目标临时表full_running全量完成记录当前源库位点full_done切换临时表替换目标表switch_done增量从位点开始消费incremental这种状态机可以放在调度系统里也可以是一个独立脚本。重点在于全量完成后先记位点再切换表很多人把顺序反过来导致全量和增量之间那一小段时间里发生的变更全部丢失。4.2 用状态文件与调度器维护同步链路写一个手工可控的调度器。状态文件记录每张表当前处于全量还是增量增量水位线单独存放。下面是 Python 控制脚本的核心片段import os, subprocess TABLES [biz_order, biz_user] def run(table): state_file f/data/etl/state/{table}.state # 表没有状态文件先跑全量 if not os.path.exists(state_file): subprocess.run([python, sync_full.py, table], checkTrue) with open(state_file, w) as f: f.write(incremental) return # 已有状态直接跑增量 subprocess.run([python, sync_inc.py, table], checkTrue) for table in TABLES: run(table)sync_full.py里必须做完临时表替换sync_inc.py要带幂等参数比如目标表唯一键。这个脚本没有考虑分布式多 worker 跑同一张表的互斥问题一旦同时触发任务两个进程可能都判断为全量。生产环境我会在目标数据库里执行GET_LOCK或者用调度器的 single run 锁否则状态文件读到旧值就出乱子。4.2.1 失败重跑时看什么增量任务报错后最先查两处水位线时间是否停在失败批次以及目标表是否有 upsert 重复数据。如果水位线已经跳到最后一行但业务数据只写入一半恢复时直接用当前水位线再次跑会静默丢数。我习惯在同步脚本里把每个批次写成带batch_id标记的中间表重跑时先删除该batch_id记录再重放这样幂等边界清晰。4.3 处理更新历史数据导致的漏数ETL 增量最大的隐蔽风险是业务方凌晨修改了历史订单的状态update_time落在过去三天而增量水位线已经跑过那个时间段。如果只按update_time last_ts抽取这次更新永远进不了增量链路。常见解法有三种第一种在源端把这种修改做成新增一条变更明细目标表按主键取最新值第二种缩短水位线补偿窗口每次增量除了抽取新水位线还要额外抽取最近 N 小时的数据并 upsert即“增量里带一个小全量”第三种精度要求高的链路用 CDCbinlog 会记录所有更新不受业务时间影响。我一般做主表的时候选第二种因为它成本低、可控。把增量脚本的过滤条件改成update_time last_ts - INTERVAL 24 HOUR然后目标表按主键 upsert。这样每天会重复处理前一天的数据但换来的是历史更新不再漏。代价是重复数据量会扩大所以参数不能设太大24 小时对多数业务表足够遇到日志类毫秒级高频表就要另想。状态机脚本在实际生产里最容易被忽略的是与表结构的耦合。每张表的全量脚本要传入 schema、分区字段、目标表名增量脚本要传入唯一键和更新时间列最好用一个小的配置表统一管理。另外切换阶段临时表替换如果用RENAME TABLE tmp TO target这是元数据操作耗时很短但要确保没有查询正停在目标表上否则可能出现锁等待。5. ETL 增量延迟验证与一个压测技巧5.1 用乱序事件验证增量链路增量任务上线最怕什么顺顺利利跑两天突然有一天缺数。我常用一个压测技巧在测试环境生成一批update_time乱序的变更事件直接灌进增量链路看下游目标表在水位线推进后是否还能聚到正确结果。很多团队压测只用顺序时间戳永远暴露不了乱序问题。import random, datetime records [] base datetime.datetime(2025, 1, 1, 8, 0, 0) for i in range(5000): t base datetime.timedelta(secondsrandom.randint(0, 3600)) records.append({id: i % 200, value: i, update_time: t}) random.shuffle(records) # 将 records 按乱序写入源表触发增量消费观察点如果下游是一个按主键取最新值的模型最终每个 id 的值应当等于该 id 在update_time最大那条记录如果出现中间值说明水位线提交和批处理顺序有问题。然后用上一章提到的 24 小时补偿窗口重跑看是否能收敛。这个方法的参数是乱序窗口大小我一般设置为同步频率的 5 倍以上比如同步间隔 5 分钟窗口就拉 30 分钟。窗口太短测不出延迟窗口太长会放大重复计算压力。还可以配合一个校验 SQL在目标端做全量对账SELECT sum(value) FROM target_table和源端同一主键的最新值求和做对比。全量对账不需要每天做一天一次即可对不齐再按补数任务重放最近批次。下次压测把乱序、删除、历史更新这三件事同时丢进去水位线还能站住增量链路才算真正能上线。本文还有配套的精品资源点击获取