做了这些年数据工作,Kettle一直是我处理日常数据同步的首选工具之一。最近好几个项目都在聊“增量同步”,不少同事和朋友问我:用Kettle怎么做增量,而不是每天傻乎乎地全量拉一遍。这确实是很多团队都会遇到的现实痛点——数据量越来越大,全量同步越来越慢,业务方又希望报表数据尽量新鲜。增量同步就是在这样的背景下被反复搬上台面的。
这篇文章我会把用Kettle做增量同步的完整思路、具体配置步骤、常见踩坑记录都整理出来。不绕弯子,直接讲方案、讲SQL、讲参数怎么传、讲定时任务怎么配,也会对比一下时间戳增量、全量比对、日志CDC这几条技术路线各自适合什么场景。适合正在用Kettle做数据抽取、想从全量升级到增量的同学,也适合刚接触Kettle、想搞懂增量同步原理的人。
1. 为什么强调“增量”?先想清楚需求再动手
1.1 增量同步能解决什么问题
增量同步说白了就是“每次只同步变化的那部分数据”,和全量同步是相对的。全量同步每天把源表从头到尾拉一遍,数据量小的时候没感觉,一旦表里积累了几千万行,问题就全冒出来了。
首先最直接的是效率问题。几千万行的表每天全量拉一次,再通过网络传到目标库,再写入目标表,整个流程可能要好几个小时。增量同步只抽当天变化的那几千几万行,分钟级甚至秒级就能搞定。其次是源库压力问题,全量扫描大表会对源数据库造成明显的IO和CPU消耗,生产系统很容易被这种同步任务拖慢。另一个容易被忽略的是时效性,增量同步可以把同步周期从T+1压缩到分钟级,业务方看到的报表数据更新鲜,决策也更及时。
当然,增量同步不是没有代价。它要求源表有可靠的增量字段,或者有对应的变更记录机制。如果表结构设计得不好,比如连个update_time都没有,那就得另想办法。所以动手之前,先想清楚以下几个问题:源表有没有时间字段?数据是只增不改,还是会更新历史记录?目标端允许重复数据吗?同步失败的容忍度是多高?这些问题直接决定了你选哪种增量方案。
1.2 Kettle在同步工具里的定位:不是万能的,但足够能打
这些年市面上同步工具也不少,DataX、Canal、Maxwell、Flink CDC,各有各的场景。有朋友问过我,既然这些工具这么火,为什么还要用Kettle?
我的看法是,Kettle在“定时、批处理、表到表”这个场景下依然是最顺手的选择之一。它的核心优势是可视化操作,拖拽几步就能搭好一个数据流,对非专业开发出身的数据人员非常友好。调试也方便,直接在Spoon里点一下“预览”就能看到数据长什么样。而且Kettle是纯Java写的,跨平台,Windows、Linux都能跑,部署成本低。
下面这张表是我自己常用的选型参考:
| 工具 | 实时性 | 实现难度 | 适用场景 |
|---|---|---|---|
| Kettle | 定时批处理,分钟级起步 | 低,可视化拖拽 | 表到表抽取、清洗转换、多数据源整合 |
| DataX | 离线批量同步 | 中,需要写json配置 | 异构数据源之间的高吞吐批量迁移 |
| Canal | 秒级实时 | 中高,需要运维消息链路 | MySQL Binlog监听,实时同步到ES或消息队列 |
| Flink CDC | 秒级实时 | 高,需要开发能力 | 复杂实时数仓、事件驱动架构 |
如果你只是想把A库的一张表稳定地同步到B库,每天跑几次,或者每个小时跑一次,Kettle完全够用。如果业务要求秒级实时同步,那Kettle确实不太合适,应该去走Canal或Flink CDC那条路。做技术选型最忌讳迷信“别人说哪个好”,适合自己的场景和数据规模才是最重要的。
2. 增量同步的三种经典实现思路
2.1 时间戳增量:最简单也最常用的方案
时间戳增量是所有增量方案里最直观的一种,前提是源表里有一个记录最后修改时间的字段,比如update_time、modified_time。每次同步的时候,只取出update_time大于上次同步时间点的记录。
具体做法是单独建一张控制表,比如sync_control,里面存着每个同步任务的上次运行时间。每次跑同步任务时,先读这个时间,同步完成后把本次最大的update_time写回去,作为下次同步的起点。用SQL来表达大概是这样:
-- 读取上次同步时间 SELECT last_run_time FROM sync_control WHERE sync_name = 'orders_sync'; -- 抽取增量数据 SELECT * FROM orders WHERE update_time > ? AND update_time <= NOW();这个问号就是上次同步时间,在Kettle里可以通过变量或者步骤传进去。这种方案最大的优点是实现简单、逻辑清晰,只要源表时间字段维护得好,基本不会出大问题。缺点也很明显,如果源表没有时间字段,或者有时间字段但没建索引,查询效率会很难看。
2.2 全量比对:适合小表的稳妥方案
有些表确实没有像样的时间字段,但数据量不大,几千行或者几万行,这种时候可以用全量比对方案。思路是每次都把源表整个拉过来,和目标表做一次差异比对,找出新增和修改的数据,只把这些变化写进去。
Kettle里有一个专门干这事的步骤,叫“合并记录”(Merge Rows (Diff))。它会根据指定的关键字段,把“源表数据”和“目标表数据”逐行比较,输出“新增”“删除”“修改”和“不变”几类结果。后面再接上“插入/更新”或“删除”步骤,就能完成同步。
我的建议是,这种方案只适合数据量可控的小表,因为它本质上还是每轮全量拉取,只是写入的时候省了点功夫。要是哪天表长到百万行以上,这个方案会越来越吃力,建议尽快换时间戳方案或者引入CDC机制。
2.3 日志与触发器方案:准实时的另一种可能
比时间戳更高级的思路是,基于数据库的变更日志来同步数据。典型的就是MySQL的Binlog、PostgreSQL的WAL,或者是业务表上建触发器,把增删改操作记录到一张单独的日志表里。
触发器的做法离Kettle并不远。你可以在源库建一张change_log表,然后在业务表上建几个触发器,每当有INSERT、UPDATE、DELETE操作,就往change_log里写一条记录,包含操作类型、主键、变更时间。Kettle这边只需要定时去扫change_log,把新增的变更记录同步到目标库,再把这些记录标记为已处理即可。
这种方案的实时性比时间戳高得多,而且不需要源表有update_time字段。但它也有明显的门槛:必须在源库建触发器,这需要DBA的配合,很多生产环境根本不给建。另外,业务系统如果跨库操作,触发器的管理也会很麻烦。所以这个方案可以作为时间戳方案的补充,但并不适合直接全面铺开。
对大多数Kettle使用者来说,时间戳增量已经能解决80%的问题。下面我就以时间戳增量为重点,把完整的实操流程拆开讲。
3. Kettle增量同步实操:从零配置一套“时间戳增量”任务
3.1 环境准备与安装
工欲善其事,必先利其器。Kettle的正式名称叫Pentaho Data Integration,大家习惯叫Kettle。去官网或者开源镜像站下载对应版本的压缩包即可,注意版本要和你的JDK版本匹配。我长期用的是Kettle 8.3和9.0版本,配JDK 8完全没有问题,再新的版本也支持JDK 11。
下载下来解压之后,Windows下直接运行Spoon.bat,Linux下运行spoon.sh,就能打开图形化界面。它依赖Java环境,所以安装Kettle之前先把JDK装好,具体版本要求官方文档写得很清楚,别装错了。打开Spoon之后,你会看到左侧的“转换”和“作业”两个概念:转换负责具体的ETL操作,作业负责编排调度流程。增量同步任务通常需要两者配合使用。
3.2 核心设计:增量抽取转换怎么做
我以一个电商订单表orders为例,表里有订单号order_id、订单金额amount、下单时间create_time、最后更新时间update_time。目标库是另一台服务器上的MySQL库,需要把每天新增和修改的订单同步过去。
第一步,在源库建一张控制表sync_control,记录同步任务的上次运行时间:
CREATE TABLE sync_control ( sync_name VARCHAR(64) PRIMARY KEY, last_run_time DATETIME ); INSERT INTO sync_control (sync_name, last_run_time) VALUES ('orders_sync', '2024-01-01 00:00:00');第二步,在Kettle里新建一个转换,命名为“订单增量抽取”。转换的第一个步骤是“表输入”,用来读取控制表里的上次运行时间。这里的SQL很简单:
SELECT last_run_time FROM sync_control WHERE sync_name = 'orders_sync';注意这个步骤的输出字段是last_run_time。接下来需要把它的值作为变量传给后面的查询。我最常用的做法是在表输入里直接写子查询,一步到位,不需要额外的变量赋值步骤。主数据查询的SQL可以这么写:
SELECT order_id, amount, create_time, update_time FROM orders WHERE update_time > (SELECT last_run_time FROM sync_control WHERE sync_name = 'orders_sync') AND update_time <= NOW();这样写的好处是逻辑全部收拢在SQL里,Kettle这边不需要处理参数传递,后续维护也直观。
第三步,在表输入后面接一个“插入/更新”步骤,目标表设置为orders。这个步骤的核心配置是“用于查询的关键字”,也就是用来判断记录是否已存在的字段,这里选order_id。然后“更新字段”里勾选amount、update_time,意思是如果记录已存在,就更新这几个字段。配置完成后,“插入/更新”会自己判断:order_id存在就更新,不存在就插入,天然避免主键冲突。
这里有一个非常关键的细节:order_id本身不要放进更新字段里,否则每轮同步都会把主键当成普通字段更新一次,既浪费资源又有风险。
3.3 别忘了维护同步时间:Job的编排逻辑
光有上面的转换还不够,因为控制表里的last_run_time一直没更新,下次同步还是会把同一批数据再抽一遍。所以需要新建一个作业,把整个流程串起来:读取上次时间并抽取数据、写入目标表、更新控制表时间。
我在实操中习惯在同一个转换的最前面读last_run_time,在转换的最后面更新last_run_time。更新控制表用“表输出”步骤即可,SQL大致如下:
UPDATE sync_control SET last_run_time = (SELECT MAX(update_time) FROM orders WHERE update_time <= NOW()) WHERE sync_name = 'orders_sync';这里取MAX(update_time)作为新的同步起点,而不是用当前时间,是因为当前时间可能比最后一条数据的update_time大很多,会导致下轮重复抽取一部分数据。用数据自身的最大时间当游标,可以做到天然幂等。
不过有个坑要提醒一下:如果源表里偶尔有数据延迟写入,比如业务系统事务提交慢,或者某个时间戳是历史补录的,下轮同步可能会漏掉它们。我会在更新控制表之前,把MAX(update_time)多往前推60秒作为安全窗口。也就是说,实际写回控制表的时间是MAX(update_time) - INTERVAL 60 SECOND,下一轮会多带60秒的重叠区间,配合“插入/更新”实现幂等覆盖,既不会漏也不会重。
3.4 定时调度:Windows和Linux怎么配
转换和作业都配好后,下一步就是让它按计划自动跑。Kettle本身不内置调度器,通常依赖操作系统的计划任务。Windows下用“任务计划程序”,调用Kettle安装目录下的Kitchen.bat,后面跟上作业文件路径:
D:\kettle\data-integration\Kitchen.bat /file:D:\etl\jobs\orders_sync.kjb /level:BasicLinux和macOS下用crontab,我自己最常用的是这样一行:
0 1 * * * /opt/kettle/data-integration/kitchen.sh -file=/opt/etl/jobs/orders_sync.kjb -level:Basic >> /var/log/etl/orders_sync.log 2>&1这行的意思是每天凌晨1点执行一次同步任务,日志输出到指定文件。如果要改成每30分钟一次,可以写成*/30 * * * *。注意,日志文件长时间跑下来会越来越大,最好配合logrotate做日志轮转,或者定期手动清理,否则磁盘会吃不消。
另外补一句,-level:Basic是日志级别,还有更详细的Debug或Rowlevel,调试的时候可以用,但生产环境别开,日志量太大反而影响性能。
4. 实战中踩过的坑:Kettle增量同步问题排查实录
4.1 时间字段格式不一致,增量条件直接失效
这是我见过最多的问题,也是新手最容易踩的坑。源库的update_time是DATETIME类型,但通过Kettle查询时,参数会以字符串形式拼进SQL,如果格式对不上,比较结果就不对。比如源库存的是2024-05-20 08:30:00,Kettle传进去的参数却变成了2024-05-20或者其他格式,那update_time > ?的判断就会出偏差。
解决办法有两个:一是在SQL里显式做类型转换,比如MySQL下用STR_TO_DATE(?, '%Y-%m-%d %H:%i:%s');二是在表输出的字段映射里,把日期字段明确指定为目标库的DATETIME类型。最好的方式是在SQL查询里就把数据格式统一,别让Kettle再做隐式转换。
还要特别注意时区问题。如果源库是MySQL,连接串里加了serverTimezone=Asia/Shanghai,而你的服务器是UTC时区,查询出来的时间值可能差8个小时,进而导致增量窗口错乱。这个查问题的时候特别隐蔽,我排查过好几次才定位到。
4.2 漏数据与重复数据,怎么保证“既不漏也不重”
漏数据和重复数据是增量同步的两大天敌,而且往往是同时出现的。先说漏数据,典型的场景是:业务系统在源库里更新了一条记录,事务提交的时间比Kettle读取时间晚,导致这条记录的update_time大于Kettle下一轮的起点时间,于是被漏掉了。
解决漏数据的方法就是我上面提到的“安全窗口”,每次把起点时间往前推几分钟,让窗口有一定重叠区域。重叠区域里已经处理过的记录,靠“插入/更新”步骤来覆盖更新,不会重复。
再说重复数据,这个问题多半出在“手动重跑任务”上。某个任务跑失败了,你改完配置重新执行Job,结果刚才已经同步过的一部分数据又被抽了一遍。如果目标表没有主键约束,重复插入是必然的。所以我的经验是:目标表一定要有主键或唯一索引,而且Kettle里要用“插入/更新”而不是单纯的“表输出”。
4.3 数据量大时,Kettle同步慢得像蜗牛怎么办
增量同步如果跑得慢,先别急着骂Kettle,大概率是SQL或者配置的问题。我总结下来主要有四个优化方向。
第一,源表查询要有索引。WHERE update_time > ?这种条件,如果源表update_time字段没有索引,每轮同步都要全表扫描,增量同步就失去了意义。
第二,合理设置批量提交参数。Kettle连接MySQL时,在JDBC连接串上追加参数rewriteBatchedStatements=true&useServerPrepStmts=false,批量写入性能会有非常明显的提升。我自己测过,同样一份数据,开启批量写之后整体耗时能缩短一半以上。
第三,调大提交批次大小。在“插入/更新”步骤的设置里,把“提交记录数量”从默认的1000调到5000或者10000,减少事务提交次数。
第四,如果源表实在太大,可以考虑按时间范围并行抽取。比如把一天的数据按小时切成24个区间,用多个转换并行跑,最后汇总到目标表。并行度不是越高越好,要根据数据库连接数和服务器资源来定。
4.4 增量任务跑挂了,怎么及时发现问题
增量同步最怕的是静默失败。Job没跑成功,日志也没人看,控制表的时间一直停在那里,业务方看到的报表数据越来越旧,等发现问题的时候已经晚了。
我的做法是给同步任务加一个“健康检查”机制。在控制表里除了last_run_time,再加一个last_success_time,每次作业成功时更新。然后单独写一个监控脚本,每隔几小时检查一下last_success_time和当前时间的差距,如果超过阈值就触发告警,发邮件或者钉钉消息都行。
这一步看起来很简单,但真的能救急。我有个项目就是因为加了这个小机制,在生产环境某次数据库异常后第一时间收到告警,最后在业务方发现之前就把问题处理掉了。
5. 从增量同步到进阶:多表同步与扩展玩法
5.1 多表增量同步时,怎么复用一套配置
实际项目里很少只同步一张表,经常是几十张表一起同步。如果每张表都单独建一个转换、一个Job,维护起来会很痛苦。我的做法是用“配置表驱动”的方式:在源库里建一张sync_table_config,每行定义一个同步任务,包括表名、主键字段、时间字段、目标表名。
然后用Kettle的“表输入”步骤读取配置表,再通过“复制记录到结果”和“执行SQL脚本”等步骤,在循环里动态拼接SQL、动态执行同步。这种方式适合结构相似的表,能省下大量重复工作。
不过要提醒一句,动态SQL的调试成本比较高,如果只有三五张表,还是老老实实分别建转换更实在。不要为了追求“优雅”而牺牲可维护性。
5.2 多表合并抽到一个表,或者输出多个Excel
有些场景不是简单的表到表同步,而是要把多张源表合并后写入一张大宽表,或者反向操作,把一张大表分拆成多个Excel文件。
合并多张表时,Kettle里的“合并记录”步骤非常有用。它可以按照关键字段把两个数据流进行全连接、内连接或外连接,配合“字段选择”把需要的字段挑出来,再写入目标表。拆分成多个Excel则可以用“Excel输出”步骤,根据某个字段的值来定义工作表,一个字段值对应一个工作表,导出后就是多个Excel。这些功能本质上是ETL的基础能力,增量同步的思路同样适用,只是把“抽取源”和“写入目标”换了一下。
5.3 增量思维与后续演进方向
“增量”这个词在很多领域都有类似的意思。比如热词里提到的增量式PID控制算法,它和控制系统的输出增量有关,和我们的增量数据同步虽然领域不同,但核心思维是一致的——不要每次都从头开始,只处理变化的部分。
Kettle里的时间戳增量方案,本质上也是“基于变化量驱动”的一种实现。如果哪天你发现数据量已经大到Kettle处理不过来,或者业务要求实时同步,那就该往Flink CDC、Canal这个方向迁移了。我在项目里就是把Kettle当作离线增量同步的主力,同时保留Flink CDC作为准实时链路的补充,两套体系并存,各管各的场景。
最后再分享一点我个人的体会。增量同步的方案本身不难,难的是把“同步状态”管好,把“异常情况”盯住。我见过太多项目初期是用时间戳增量跑得好好的,跑了两周因为控制表时间没更新,悄悄变回全量都浑然不知。所以我在每一个同步任务里都会把控表更新和监控日志做得扎扎实实,宁可多花十分钟配置,也不要在半夜被业务的告警电话叫醒。如果你正准备用Kettle做增量同步,希望这篇能帮你少踩几个坑,少熬几个夜。