
做数据ETL的人大约都会经历这样一个阶段数据量小的时候全量同步很省心等数据量上来问题就一个接一个冒出来。我在一个订单系统项目里就吃过这个亏几十亿行的明细表每天全量重抽要跑近5个小时源库CPU被打满下游报表凌晨出不来业务方天天催。后来我们认真落地了增量机制把单次抽取时间从5小时压到10分钟以内源库压力基本可以忽略。这篇内容会把增量机制的设计逻辑讲清楚再手把手说一遍在KettlePDI里从水位表设计、转换构建到作业调度的完整实现包括我在项目里踩过、也帮别人填过的坑。内容适合正在维护批量同步任务的开发也适合准备把同步链路从全量改造成增量的团队参考。1. 增量机制解决的根本问题全量同步的三重代价1.1 时间成本全量抽取的复杂度会随着数据翻倍全量同步的思路很简单每次跑批都把源表的所有数据捞出来清空目标表再重写。小数据量下这套做法完全没问题但数据量一旦上去你会先撞上时间这个硬指标。全量抽取的时间复杂度基本是O(n)数据量翻倍同步时间也跟着翻倍。这不是工程优化能解决的是算术规律。我在那个订单项目里第一次估算时还很乐观觉得5小时能接受结果几个月后表从20亿行涨到60亿行跑批时间直接逼近凌晨下游日汇总根本挤不出窗口。增量机制解决的就是这件事让每次跑批只读取“从上一次跑批以来发生变化的数据”。源表再大一天实际变更的行数通常是整体存量的小零头增量抽取的耗时和资源消耗因此能维持在一个基本恒定的低位不会跟着存量规模线性膨胀。1.2 源库与目标库的连带压力全量同步不是“多跑一会”的事很多人觉得全量同步只是“慢一点”慢就慢吧我多留点时间窗口。等你真正在线上跑起来会发现全量同步的破坏力远不止慢。首先是源端数据库的压力一个大查询把几亿行数据读出来、网络持续打满带宽CPU和磁盘IO都会明显抬升。业务高峰期如果和跑批时间重叠系统响应变慢、连接堆积业务方会直接来找你。其次是目标端的压力全量重写意味着大量DELETE和INSERT目标表的索引维护、事务日志、锁竞争都会放大小库还能撑大库可能出现死锁或复制延迟。增量机制把这个压力也一并拆掉了。每次只处理一个小数据窗口源库查询走索引扫描少量数据目标端写入的批次很小事务短、锁范围小整条链路的噪音都低很多。这也是很多团队宁可花时间设计增量方案也不愿意继续压硬件升级的原因。1.3 增量的本质知道变了什么、从哪里继续、如何落库增量机制设计得好不好归根到底取决于三个问题能不能答清楚第一怎么识别一条记录是否发生了变化第二怎么记住上一次已经处理到哪里第三变化的数据到了目标端之后如何处理新记录、更新旧记录、删除消失的记录。这三个问题决定了你选什么增量策略。识别变化靠时间戳、日志、触发器或者全量比对记住进度需要设计水位Watermark落库逻辑则要看目标表结构和数据约束。Kettle里几乎所有增量场景本质上都是把这三个问题翻译成具体的转换和作业步骤。理解了这个框架你再看各种增量方案就不会觉得它们是一堆零散技巧。2. 五种增量策略的选型逻辑时间戳、CDC与对比法的边界2.1 时间戳字段干净可靠但前提苛刻最常用的增量策略是检查源表里有没有一个“最近修改时间”字段也就是update_time或last_modified_time。只要这个字段存在增量SQL就能写成SELECT * FROM orders WHERE update_time ? AND update_time ?两个问号分别代入上次水位和本次运行时间。实现简单、逻辑清晰Kettle里配合表输入步骤的变量替换半小时就能搭起来。但这个策略有几个苛刻前提。第一源表的所有变更必须都体现在该字段上包括业务修改、后台修复、批量刷数。如果某条数据被直接改数据库而不更新update_time它就永远漏掉。第二物理删除操作不会留下任何时间戳痕迹删除的数据无法被感知。第三时间字段最好有索引否则增量SQL会退化成全表扫描。在我实际经验里源表有没有可靠的update_time字段直接决定了时间戳方案能不能用而不是你想不想用。2.2 自增主键只适合“只增不改”的流水类数据当源表没有时间字段但主键是自增ID时可以用ID水位做增量每次记录当前最大ID下一次抽取ID大于上次水位的数据。这个方案完全不需要update_time性能也好但有一个致命限制——它只能发现新增数据无法感知更新和删除。所以自增ID方案只适合日志型、流水型的数据比如点击日志、操作记录、事实表追加。这类数据一旦写入基本不再修改删除也极少发生。如果你把它用在订单表上用户改个收货地址都不会被同步下游分析就会失真。判断的唯一标准是这张表的业务语义是不是“只增不改”。2.3 CDC日志解析不碰源库但工程链最长当源表连时间戳都没有删除又经常发生业务库又不允许加字段、加触发器时就得考虑CDC了。CDC的全称是变更数据捕获核心思路是解析数据库事务日志或者Binlog把每一次insert、update、delete操作记录捞出来转成统一的变更流。MySQL的Binlog、PostgreSQL的逻辑复制、Oracle的LogMiner都是这个思路。Kettle本身不自带CDC功能但你可以把CDC数据源先落地成一张变更日志表再由Kettle增量读取这张表。这个方案的优点是对业务表零侵入不需要改源表结构删除也能感知。缺点是工程链最长需要额外部署日志捕获组件配置权限和维护成本都不低。Kafka、Canal、Debezium这一套组合常见于实时数据管道如果只是每天凌晨跑一次离线增量我一般不建议一上来就上CDC先把时间戳方案的排查做完至少八成以上的场景是够用的。2.4 触发器与变更日志表可控性与侵入性并存触发器方案是在源库业务表上建触发器当发生INSERT、UPDATE、DELETE时把变更的主键和一些必要信息写入一张变更日志表。Kettle每次增量只读变更日志表再回到业务表捞完整数据这个组合能解决删除和物理删除问题逻辑也容易理解。但它需要DBA在源库执行DDL属于对生产环境的侵入。触发器本身有性能开销一个高频写入的大表挂上触发器对交易的扰动不能忽略。还有一点要注意触发器是业务系统的一部分如果源库是第三方厂商管理人家不一定同意你动。我见过一些项目被这个协作问题卡住最后退回时间戳方案。触发器和日志表更适配那些数据变更频繁、但源库可控性较强的内部系统。2.5 全量对比法什么时候它反而是最优解全量对比法不走增量读取而是把源表当前全量和目标表当前全量都读出来在内存里比对找出新增、更新、删除。听起来很笨但在两种场景下反而是最优解一是源表完全没有时间戳也没有自增ID二是表数据量本身不大全量拉取成本可以接受。Kettle里的“合并记录”Merge Rows步骤就是为这个设计的。它要求两个输入流都按同样的字段排序然后一条条比对输出标志位标识每条记录是identical、new、changed还是deleted。实现上完全不依赖源表的任何增量字段所以特别适合一些历史遗留表和老业务系统。代价是每次比对都必须读全量表一旦超过千万级性能就会很难看。2.6 方案对比小结增量策略发现新增发现更新发现删除源库侵入性实施成本时间戳字段支持支持不支持无低自增ID支持不支持不支持无低CDC日志支持支持支持无高触发器日志表支持支持支持有中全量对比法支持支持支持无中数据量敏感选型没有银弹。我的判断顺序是先确认有没有可靠的update_time字段有就直接用没有再看表是不是只增不改如果有删除需求和强一致要求再评估触发器和CDC。全量对比法永远作为兜底留着对付那些无药可救的孤儿表。3. 水位表与时间窗口Kettle增量设计的核心细节3.1 为什么不能拿“当前时间”当水位很多人第一次做增量时会写这样的SQLSELECT * FROM orders WHERE update_time DATE_SUB(NOW(), INTERVAL 1 DAY)意思是同步最近一天的数据。这个写法有硬伤它是固定滑动窗口每天都要重复扫描时间落在24小时窗口内的数据。涨单高峰期某一小时的数据会被连续两天的跑批各抽一次目标端如果没做幂等就会产生重复。反过来如果某个跑批当天因为故障没执行第二天窗口一移动故障期间的数据就漏掉了。真正可靠的做法不是靠NOW()猜窗口而是把“处理到哪了”记录成明确的进度这个进度就是水位。我调试过不少翻车现场最终都指向同一个问题没有水位或者水位更新逻辑放错了位置。Kettle工程化增量任务第一步永远是设计和落水位表。3.2 水位表设计一张表管住所有同步任务水位表的思路很简单在目标库或者一个专门的管理库里建一张小表记录每个同步任务上一次成功处理的时间点。这样设计CREATE TABLE etl_watermark ( source_table VARCHAR(128) NOT NULL, last_load_time DATETIME NOT NULL, update_time DATETIME DEFAULT CURRENT_TIMESTAMP, PRIMARY KEY (source_table) ); INSERT INTO etl_watermark (source_table, last_load_time) VALUES (orders, 2024-01-01 00:00:00);这张表把每个任务的水位独立存放互不干扰。新增一个同步任务时只要插入一行记录初始化一个起始时间。首轮同步通常选全量初始化开始的时间或者一个业务上明确可追溯的起点。水位表需要建在Kettle可以直接访问的库上目标库或管理库都行我习惯建在目标数据仓库里方便同一套连接管理。水位时间的数据类型要和你源表的时间字段类型对齐如果源表是TIMESTAMP就统一用DATETIME/TIMESTAMP避免类型转换陷阱。3.3 表输入中的参数化查询与变量作用域水位表建好之后Kettle里怎么读取和使用靠的是变量系统。典型的做法是在作业开始处用一个小转换读取水位把last_load_time放到变量里主转换中的表输入再引用这个变量。读取水位的转换里用一个表输入步骤执行SELECT last_load_time FROM etl_watermark WHERE source_table orders出来的字段是last_load_time通过“复制记录到结果”步骤交给父作业再在作业层用“设置变量”步骤把它传给下游转换。主转换的表输入里SQL写法如下SELECT order_id, order_no, customer_id, amount, status, create_time, update_time FROM orders WHERE update_time ${LAST_LOAD_RUN} AND update_time ${THIS_RUN_TIME}注意Kettle的表输入步骤默认是支持变量替换的你在SQL里直接写${LAST_LOAD_RUN}就能取到变量值。作业里定义的变量会向下传递给子转换但子转换内部的变量不会自动传回作业方向是单向的。这个作用域关系经常坑人排查时先确认变量是在作业层定义还是在转换层定义。3.4 时间窗口用半开半闭区间保证不重不漏增量时间窗口的写法有一个标准答案窗口下界用大于上界用小于等于也就是(上次水位, 本次运行时间]。这样第一次窗口是(T0, T1]第二次是(T1, T2]T1这个边界点既不会被重复包含也不会被遗漏区间干净闭合。这里有一个微妙的点本次运行时间THIS_RUN_TIME必须在作业开始时固定下来而不是在SQL里写NOW()。Kettle作业从启动到表输入真正执行SQL之间有延迟如果SQL里每次调用NOW()同一批增量数据的窗口上界是漂移的可能出现跑到一半、源库里又有新数据满足条件导致本次窗口包含了计划外的数据。更规范的做法是在作业最开始用触发器或“获取系统信息”步骤取一次当前时间格式化成yyyy-MM-dd HH:mm:ss存成变量整个作业统一用它做窗口上界。3.5 作业顺序先同步成功再更新水位最后一个关键细节也是增量任务最容易翻车的环节水位更新的时机。水位必须在增量数据全部同步成功之后才推进绝不能放在主转换的数据流中间。如果源表数据读出来了、还没写进目标表水位就已经更新一旦写入失败下次跑批就会跳过这些数据造成永久性丢失。所以在Kettle作业里主转换执行成功之后再接一个SQL脚本步骤更新水位UPDATE etl_watermark SET last_load_time ${THIS_RUN_TIME} WHERE source_table orders作业连线默认是“成功才继续”主转换失败时不会走到更新水位这一步下次重跑会从头把失败批次的数据再抽一遍配合目标表的幂等设计结果依然正确。这套“先数据、后水位”的顺序我建议在架构上就定死任何想为了省事把水位更新塞进主转换的尝试都是给自己埋雷。4. 订单表增量同步实战一次完整Kettle作业的搭建过程4.1 案例场景与表结构定义为了把细节讲透我以一个订单表为例跑一遍完整流程。源库是MySQL目标库是PostgreSQL目标表结构和源表保持一致。同步需求是每天凌晨1点增量同步前一天变更的订单数据包括新增订单和订单状态修改。源表DDL如下CREATE TABLE orders ( order_id BIGINT PRIMARY KEY, order_no VARCHAR(32) NOT NULL, customer_id BIGINT NOT NULL, amount DECIMAL(12,2), status TINYINT, create_time DATETIME, update_time DATETIME, KEY idx_update_time (update_time) );index已经按update_time建了这是时间戳增量能跑通的前提。如果源表没有这个索引后续增量SQL会在亿级数据上做全表扫描性能直接崩掉。4.2 转换部分从水位读取到写入目标的步骤链主转换只做一件事根据传入的LAST_LOAD_RUN和THIS_RUN_TIME抽取增量订单写入目标表。流程如下表输入读取增量订单SQL使用第3.3节那段代码条件携带两个变量。插入/更新Insert/Update连接目标表按主键order_id匹配更新其他字段。插入/更新步骤是增量落库的主力。配置时在“关键字段比较”里选order_id在“更新字段”里选order_no、customer_id、amount、status、update_time同时勾选“插入新纪录”。它做的事情是目标表存在相同主键就更新不存在就插入。对增量数据来说新订单会插入状态修改的订单会更新两类场景都覆盖了。这里有个取舍如果一张表经常有万级以上的增量插入/更新步骤因为要逐条判断主键速度不如直接分成“按主键更新”和“按主键插入”两个流。但对大多数每日增量场景插入/更新的清晰度优势更大性能瓶颈通常在数据库连接和提交频率而不是这一层。4.3 作业部分调度、变量传递与失败保护转换搭好之后需要在作业Job层面把它串起来。我沿用的作业结构是START设置定时调度或手动启动。转换“取系统时间”用“获取系统信息”步骤取当前时间格式化后通过“复制记录到结果”输出再通过“设置变量”存为THIS_RUN_TIME。转换“读取水位线”读取etl_watermark表把last_load_time存为LAST_LOAD_RUN。主转换“订单增量同步”执行4.2节的主转换。SQL脚本“更新水位线”执行UPDATE语句推进水位。这个结构的核心是变量传递链。细节在于“设置变量”步骤里有个选项叫“在运行时替换变量”要确保勾选。否则子转换里拿到的可能还是字符串字面量SQL跑出来全是空结果。我踩过这个坑排查了快两小时最后发现是作业层没开变量替换。失败保护完全依靠作业连线的“成功”语义。第4步失败第5步不会执行水位停留在旧值下一次跑批会自动把这个失败批次的区间重新覆盖。目标表有主键、插入/更新是幂等操作所以重跑不会产生重复数据。这就是“先数据、后水位”的价值。4.4 首次全量初始化与增量切换增量任务上线前必须有一个全量初始化过程。直接建空表就开跑增量会导致目标表缺失历史数据。标准的做法是先手动把源表全量导入目标表导入完成后把水位表的last_load_time设置成全量导入开始的那个时间点然后再切换到增量作业。注意这个时间点必须严格等于全量数据的最新变更时间或者至少不晚于它。我习惯在全量导入完成后执行一条SQL取源表当前最大update_timeSELECT MAX(update_time) FROM orders然后把这个值写进水位表。这样增量作业启动后(水位, 当前时间]区间内只有全量之后的新变更不会把全量已导入的数据再抽一遍。全量初始化期间如果源表还在写入怎么办那就把水位设成全量启动时刻接受一次小范围的重复抽取重复数据靠幂等更新吸收掉损失可以忽略。5. 三个绕不开的坑数据漂移、删除同步与幂等重跑5.1 数据漂移业务时间和变更时间是两个维度时间戳增量看起来简单生产环境里第一个坑就是数据漂移。典型场景一个订单昨天23:59:58下单但支付回调在凌晨00:30才完成update_time被更新到00:30。第二天跑增量这条数据会被抽到但如果目标表的下游分区按create_time或者业务日期划分这条记录就落到今天的分区导致昨天报表缺数、今天报表多单。解决数据漂移的关键是把两个时间维度拆开抽取条件用update_time落库分区用业务时间。也就是说判断“是否变化”用变更时间判断“属于哪一天”用create_time或业务发生时间。ETL开发最容易犯的错是期望一个时间字段同时承担两个职责。另一个和漂移相关的问题是主从延迟。增量查询如果连的是从库源库主库已经提交的update_time从库在复制延迟期间还读不到抽数就会漏。我当时的处理办法是水位窗口向下微调比如水位减2分钟作为实际过滤下界允许极小范围内的重复用幂等吸收换来的是不漏数据。这条经验适合对实时性要求不高的离线跑批核心取舍是宁可重复不可缺失。5.2 删除同步软删除、删除日志、定期对账时间戳方案最大的盲区是物理删除。源库如果直接DELETE一行没有触发器、没有日志下一次增量根本感知不到。应对手段按优先级排有三种。第一推动业务系统做软删除。加一个is_deleted字段删除时把状态置为删除并刷新update_time目标表同步后查询条件统一带上is_deleted 0。这是最干净的做法但要说服业务改代码。第二建立删除日志表。在应用层删除操作发起时额外向del_log写一条记录记录主键值Kettle的增量作业里加一个分支读取删除日志并在目标表执行DELETE。这个方案能捕获物理删除但依赖应用层配合适合内部系统。第三定期全量对账。对实在没有软删除、也没有日志的表每周或每月跑一次合并记录对账。Kettle的“合并记录”步骤会把源表主键集合和目标表主键集合做diff输出deleted标志下游接一个“删除”步骤清掉消失的行。代价是要读全量主键所以频率不能太高。我的做法是把它放在月底大屏数据校正前夜执行平时不跑。5.3 幂等重跑唯一键与窗口边界的配合增量任务迟早会遇到重跑。要么是数据有问题要回刷要么是凌晨跑批失败第二天手动补跑。如果设计时不考虑幂等重跑就会制造重复或者脏数据。幂等有三个支柱。第一目标表必须有主键或唯一索引这是插入/更新能去重的物理前提。第二插入/更新步骤以主键作为关键字段重复执行同一窗口时对已存在的记录只是再做一次更新结果不变。第三时间窗口边界严格使用前文说的半开半闭区间水位单调递增避免两个相邻窗口互相覆盖。有一个容易被忽视的小问题水位更新后如果因为时区切换或者手动回改导致时间回退新的窗口上界小于旧水位就会出现窗口倒挂SQL查询条件变成update_time T2 AND update_time T1结果集为空数据全都漏掉。所以每次改水位时我都要顺手检查一下last_load_time是严格递增的。5.4 突发批量更新被时间戳击穿的特殊场景还有一种场景比数据漂移更隐蔽源系统在凌晨跑日结批任务时会用UPDATE语句把大量历史记录的状态批量刷新比如把所有待支付订单一次性置为关闭。这次批量更新会让几百万行记录的update_time同时变成当前时间。第二天增量任务一抽涌入几百万行变更链路直接被打爆。应对这种场景常规增量通道已经不够了需要分流。我建议和源系统约定日结批任务执行时额外往一张批量变更表写一行汇总记录说明改了哪些范围、改了什么状态。Kettle侧单独建一个“批量变更处理”转换处理方式不是按行更新而是按源系统的批量规则重新执行目标端的等价更新。这样既不会漏数也不会把几百万行变更塞进逐行UPDATE通道。6. 增量链路性能与稳定性索引、分批与跨库适配6.1 增量SQL变慢先查执行计划增量任务刚上线时很快跑了几个月越来越慢这种问题我遇到过不止一次。排查路径非常固定先到源库看这条增量SQL的执行计划重点看过滤字段有没有索引、索引是否失效。时间戳增量最忌讳在WHERE条件里对update_time做函数包裹比如TRUNC(update_time)一旦包裹索引就废了。哪怕你对update_time建了索引只要写WHERE TRUNC(update_time) ?优化器也只能全表扫描。另一个排查点是数据分布。如果update_time字段的统计信息过期优化器可能算错行数选了全表扫描。对这种表定期做ANALYZE/OPTIMIZE能见效。还有增量SQL的范围如果宽度过大哪怕走索引数据库也会因为要扫大量数据而选择放弃索引所以在源库侧保证选择性是优化的前提。6.2 大批量增量游标式分页而不是OFFSET第一次初始化已经用全量跑过了这里的问题主要是两种存量表第一次做增量欠账太多或者某个窗口期突然涌入超大变更量导致单次增量要抽取的数据量过大。Kettle里表输入一次性把所有结果加载到内存几百万行很容易OOM。处理大批量增量的标准做法是分批。但分页不能用OFFSETOFFSET在前面数据发生变化时会重复或漏行。更稳的是游标式分页先按ORDER BY update_time, order_id排序取前N条处理完后记录最后一条的update_time和order_id下一次过滤条件变成WHERE update_time ? OR (update_time ? AND order_id ?)直到取完。这个写法能利用组合索引(update_time, order_id)每次都是小范围扫描而且不受源表并发变更干扰。Kettle里实现游标式分页需要配合作业循环主转换接收当前水位和批次边界跑完之后把最后一条记录传回作业更新批次变量再进入下一轮。结构上比单次抽取复杂但对那些确实存在大增量窗口的表这才是能稳定过夜的做法。6.3 跨数据库方言与第三方驱动适配Kettle的优势在于它可以连接几乎所有主流数据库但这也会带来增量SQL的方言问题。同样是分页MySQL用LIMITOracle用ROWNUM或FETCH FIRSTSQL Server用OFFSET FETCH或TOP。Kettle的表输入步骤把SQL原样发送给源库执行方言差异不会自动屏蔽。我建议规则是复杂逻辑尽量放进Kettle的步骤里比如过滤、判空、时间格式化用字段选择或JavaScript代码步骤完成数据库SQL只做最基础的数据读取。这样换源库时不用重写一整套业务逻辑。另一点是第三方驱动的坑。我在一个Access数据源的小项目里用过ucanaccess驱动驱动类名是net.ucanaccess.jdbc.UcanaccessDriverJDBC连接串是jdbc:ucanaccess://文件路径Access的SQL语法限制很多时间函数和类型转换跟MySQL差异非常大。当时我们把增量SQL里所有复杂函数都抽到Kettle的JavaScript代码步骤里处理驱动这块才稳定下来。处理跨库场景时先用一小段数据把驱动连通性和基础SQL验证过再开始设计增量任务能省很多返工时间。6.4 连接参数与提交批次的小型调优清单增量抽取的性能不只在SQLKettle和目标库连接本身的配置也直接影响效率。我固定使用的一组参数是JDBC连接串里加上rewriteBatchedStatementstrueMySQL或reWriteBatchedInsertstruePostgreSQL让批量写入语句被重写为多值INSERT提交速度可以提升数倍。Kettle的插入/更新步骤里“提交记录大小”一般设在500到1000太小频繁commit拖慢速度太大事务回滚成本高一旦失败要重导的数据也更多。源库连接上连接池最大空闲连接可以按同步任务并发数设置但不要无脑调大数据库连接数本身也是资源。目标库在增量写入期间如果表上有大量索引写入性能会被拖累那些只在报表查询时用到的辅助索引增量跑批前可以临时禁用跑完再重建。这属于常规运维操作但要记得写进操作手册否则切换索引状态的人一旦忘了恢复第二天报表查询会直接变慢。最后说一点个人习惯我每次给项目设计增量同步都会先在纸上画清楚“读取水位、抽取数据、更新水位”三步之间的依赖关系再打开Kettle动手。先想清楚水位放在哪个环节、失败怎么保护、重跑是否幂等再谈配置细节。这个习惯帮我避免了很多次半夜被电话叫醒的尴尬也让我在面对新表接入时不用每次都从头想一遍方案。增量机制看起来是一堆组件和变量真正值钱的其实是那套围绕“变化”和“进度”的确定性设计。