只要做过数据平台开发,大概率都遇到过这种场景:业务方丢过来一个Excel,说“帮我导一下”;或者每天凌晨要跑一批数据同步,把A库的数据清洗完灌到B库。以前的小团队做法是写Python脚本,配crontab,但业务逻辑一复杂,脚本就像滚雪球一样越来越难维护。后来换成了Kettle(现在叫PDI),图拖拽式的ETL确实香。可新问题又来了:调度和监控还是靠人肉,任务状态、失败重试、参数传递全在Kettle外面裸奔。
所以就有了这个项目:把Kettle嵌入Spring Boot服务,用Java代码来加载、执行、监控转换任务。这篇文章把我的完整实践过程写出来,包括为什么选这种集成方式、依赖怎么引入、ktr文件怎么管理、参数怎么动态传、日志怎么接到logback里,还有我在生产环境踩过的一堆坑。适合正在做数据集成平台、想把Kettle能力封装成服务接口的Java工程师参考。
1. 为什么要放弃命令行调度,改用Spring Boot集成Kettle
刚开始用Kettle的时候,大家普遍的做法是装一个客户端,在Spoon里拖好转转换(.ktr)或者作业(.kjb),然后手动跑。稍微正规一点,会用Pan或者Kitchen命令配合脚本做定时调度。这种模式的痛点,做久了你就会懂:
- 任务状态完全不可视。跑没跑完、有没有报错,全靠日志文件里翻。
- 参数传递很痛苦。同一个转换,今天跑昨天、明天跑今天,日期参数得在脚本里拼命令行参数,写错一个符号任务就挂了。
- 没有统一权限和并发控制。几个人同时点一个任务,直接撞库表锁。
- 业务系统要触发一个转换,只能通过Shell脚本或者文件系统约定,接口化更不用想了。
把Kettle集成进Spring Boot之后,这些痛点基本都能解决。所有转换变成Java方法的一次调用,参数天然走方法入参,执行状态走回调监听器,日志统一进logback,更关键的是可以对外暴露REST接口,让别的系统按需触发数据加工任务。
我在这要特别强调一下选型逻辑:如果只是个人临时导数据,没必要上Spring Boot集成,打开Spoon拖一下就好。集成适合的是“别人要调用你的ETL能力”这种工程化场景。也就是说,你做的是一个数据服务,而不是一个数据工具。这一点想清楚了,后面每一步设计都不会跑偏。
2. 依赖引入与版本选择:这里藏了第一个大坑
Kettle本身是开源项目Pentaho Data Integration,Java编写,Maven坐标是pentaho-kettle系列。在实际引入Spring Boot项目时,版本选型是第一个决定成败的点。
2.1 版本组合怎么定
我的生产环境使用的是Spring Boot 2.2.x + Kettle 8.2。这套组合跑了很长时间,比较平稳。Kettle 8.3也试过,但我踩到了和某些第三方库的兼容问题,所以生产上一直锁在8.2。
<dependency> <groupId>org.pentaho</groupId> <artifactId>pdi-engine</artifactId> <version>8.2.0.0-342</version> </dependency> <dependency> <groupId>org.pentaho</groupId> <artifactId>kettle-core</artifactId> <version>8.2.0.0-342</version> </dependency> <dependency> <groupId>org.pentaho</groupId> <artifactId>kettle-engine</artifactId> <version>8.2.0.0-342</version> </dependency>实际用起来,kettle-engine基本涵盖了一个转换从解析到执行的全部核心类,kettle-core提供资源库接口、日志体系等基础能力。如果你的ktr或kjb里用到了高版本才有的组件,再按需加其他模块。
2.2 依赖冲突处理
Kettle 8.2会传递引入一堆老版本的第三方库,其中最容易和Spring Boot冲突的是:
guava:Kettle传递来的版本可能很老,和Spring Boot里用的新版本API冲突。jackson:Kettle内部有自己的Jackson版本,如果覆盖不对,JSON组件解析会出问题。commons-lang3、commons-collections4:这类基础库版本冲突容易造成莫名其妙的NoSuchMethodError。
我的处理方式是在pom.xml里用exclusions排除Kettle传递的旧包,然后由Spring Boot的依赖管理统一控制版本。比如:
<dependency> <groupId>org.pentaho</groupId> <artifactId>kettle-engine</artifactId> <version>8.2.0.0-342</version> <exclusions> <exclusion> <groupId>com.google.guava</groupId> <artifactId>guava</artifactId> </exclusion> <exclusion> <groupId>com.fasterxml.jackson.core</groupId> <artifactId>jackson-databind</artifactId> </exclusion> </exclusions> </dependency>提示:排除之后,Kettle里面用到的部分guava类在较高版本中已被移除,我在实际项目中使用了
guava 30.1-jre并额外引入了failureaccess包,才能正常跑通。这一步不加,会在执行转换时遇到ClassNotFoundException: com.google.common.util.concurrent.internal.InternalFutureFailureAccess。
2.3 JDBC驱动也必须单独处理
Kettle执行转换时,数据库连接由它自己的连接体系管理,不走Spring Boot的数据源配置。所以你要把用到的数据库驱动完整放到classpath里。比如MySQL用mysql-connector-java,达梦用DmJdbcDriver,TDengine用它的JDBC驱动。这些驱动Kettle官方依赖里不会全带,必须手动加。
有朋友在“kettle连接达梦数据库”上反复踩坑,我补充一句:达梦的JDBC驱动jar包下载后,放到项目lib目录或打成系统依赖都行,但注意和Kettle 8.2配合时,驱动类名要写成dm.jdbc.driver.DmDriver,URL格式是jdbc:dm://ip:port/schema。别的驱动类似,Kettle不认Spring Boot的数据源代理,直接用DriverManager加载。
3. 核心代码实现:从一个ktr文件的加载到执行
依赖配好之后,进入正题。先写一个最精简的版本:加载classpath下的.ktr文件,执行它,打印输出日志。
3.1 初始化Kettle环境
Kettle执行前必须先初始化环境,这个动作类似DruidDataSource的初始化。一个JVM只初始化一次,不能每个任务都初始化,否则速度和稳定性都会崩。我放到Spring Boot的启动类里:
@Component public class KettleEnvironmentInitializer implements ApplicationRunner { @Override public void run(ApplicationArguments args) { try { KettleEnvironment.init(); } catch (Exception e) { throw new RuntimeException("Kettle环境初始化失败", e); } System.out.println("Kettle环境初始化完成"); } }这里有一个细节:KettleEnvironment.init()内部会去读kettle.properties配置,还会初始化插件注册表。如果你不想把配置信息放在固定用户目录下,可以通过KettleEnvironment.setKettleHome()指定一个目录,把kettle.properties放到这个目录里。Kettle会优先使用该目录下的配置,避免污染系统的用户主目录。
3.2 加载转换并执行
初始化之后,加载和执行的代码就非常简洁了:
import org.pentaho.di.core.KettleEnvironment; import org.pentaho.di.trans.Trans; import org.pentaho.di.trans.TransMeta; public void runTrans(String ktrPath) { TransMeta transMeta = new TransMeta(ktrPath); Trans trans = new Trans(transMeta); trans.execute(null); trans.waitUntilFinished(); if (trans.getErrors() > 0) { throw new RuntimeException("转换执行失败,错误数: " + trans.getErrors()); } }这个例子里的ktrPath我传的是一个classpath之外的绝对路径,因为生产环境中ktr文件不应该打进jar包里,那样改一个字段就要重新发布应用。更合理的做法是把所有ktr/kjb文件放在一个外部目录,通过配置项指定。比如:
kettle: job-location: /data/kettle/jobs/加载的时候拼接:new TransMeta(jobLocation + fileName)。
3.3 用资源库管理还是直接读文件
Kettle支持资源库(Repository)模式,可以把转换脚本存进数据库。Spring Boot集成时,我更建议直接读文件目录。原因有三个:
- 资源库模式需要内置
kettle-repository相关依赖,会增加大量jar,依赖冲突概率更高。 - 文件模式下,ktr/kjb可以直接用Git做版本管理,变更可追溯、可回滚,这比资源库里存二进制或者XML更便于团队协作。
- 你只需要做一个
ktr文件下载/上传的接口,就能实现在线管理,配合Spring Cloud Config之类的配置中心,比连资源库还要输入账号密码、填库表信息要轻量得多。
当然,如果你手里已经有大量存量任务挂在资源库,也有办法。通过KettleEnvironment.getRepository()拿到资源库实例后,用RepositoryDirectory去遍历目录,再加载TransMeta。但这就意味着你的项目必须引入资源库实现和对应数据库驱动。权衡之后,我果断选择了文件目录方案。
4. 动态参数传递:让一个ktr适配多场景
Kettle里常用的参数传递有两种方式:命名参数(Named Parameters)和变量(Variables)。在Spring Boot集成场景,我主要用命名参数。
在Spoon中,你可以在转换的属性里设置参数名和默认值,ktr中引用的方式是${paramName}。例如一个SQL查询步骤中的SQL可以写成:
SELECT * FROM orders WHERE create_date >= '${startDate}' AND create_date <= '${endDate}'在Java代码中动态修改变成:
TransMeta transMeta = new TransMeta(ktrPath); transMeta.setParameterValue("startDate", "2024-11-01"); transMeta.setParameterValue("endDate", "2024-11-30"); transMeta.activateParameters(); Trans trans = new Trans(transMeta); trans.execute(null); trans.waitUntilFinished();注意,这里activateParameters()必须调用,它负责把参数值设置到转换的运行环境中去。我见过不少人在集成时只set参数忘了activate,结果SQL里还是字面量${startDate},一直查不出数据。
还有一个更容易忽视的地方:如果ktr里用了“获取变量”步骤,通过getVariable("startDate")读取参数,那set方式就不能那样用了,要用transMeta.setVariable("startDate", "2024-11-01")。参数(Parameters)和变量(Variables)在Kettle里是两类不同的体系,定义在哪、怎么用就怎么set,混着来容易出问题。
4.1 基于API分页数据的传参技巧
热词里提到“kettle用post组件获取api分页数据”,这种场景在集成时往往要配合循环。比如你要从第三方接口一页页取数据,Kettle的“循环”可以通过作业(kjb)的Simple Evaluation实现,也可以直接在转换里用Table Input配合SQL调用存储过程来搞。最省事的模式是,把当前页参数传给一个REST Client组件,HTTP请求体里动态拼:
{ "page": "${pageNo}", "size": 500 }在Java里每次循环执行前重新setParameterValue并activateParameters,就能实现“Kettle循环API读取”的效果。这种方式比在Kettle内部做循环更容易控制,出错时还能确定是第几页失败的。
4.2 安全垫:参数默认值 + SQL校验
动态传参有个隐患:如果外部传入的SQL值带有单引号,直接拼进Kettle的SQL里,轻则报错,重则破坏查询条件。我建议在传入前做一层过滤,至少在拼模板前把单引号替换掉:
String safeDate = param.replace("'", ""); transMeta.setParameterValue("startDate", safeDate);另外,Kettle转换内部最好给参数设置一个“默认值”,防止漏传参数时直接执行一个${param}字面量的SQL,这种错误极度误导人,排查起来像看天书。
5. 日志输出与执行状态监控:把黑盒变成白盒
Kettle自带一个内存日志系统,但默认只是往控制台打,这对Spring Boot工程来说不够用。我们的目标是:日志进入logback、关键状态能被业务代码感知、失败时能拿到完整堆栈。
5.1 让Kettle日志接入SLF4J
Kettle的日志体系基于KettleLogStore,每条日志都会经过KettleLogLayout格式化。最简单的接入方式,是给Trans对象添加一个KettleLogLayout类型的监听器,把日志流式转发到SLF4J。
Trans trans = new Trans(transMeta); trans.setLogLevel(LogLevel.BASIC); // 使用 KettleLogLayout 收集日志 KettleLogLayout logLayout = new KettleLogLayout(); trans.addLogChannelInterface(new LogChannelInterface() { // 实际上更稳妥的做法是调用 KettleLogStore 的 listener });更可靠的方案是使用Kettle自带的KettleLogStore的监听器机制。在Kettle 8.2中,可以注册一个KettleLoggingEventListener,在事件里把文本写到logback:
KettleLogStore.getAppender().addLoggingEventListener(event -> { String message = event.getMessage(); if (event.getLevel() == LogLevel.ERROR) { log.error(message); } else { log.info(message); } });这一步做好后,Kettle里的“表输出”步骤、SQL执行日志、每一步的耗时日志,都会出现在你的Spring Boot日志文件里,统一接入ELK完全没问题。
5.2 通过TransListener感知任务状态
业务系统调用Kettle接口时,往往需要同步拿到“成功或者失败”。trans.waitUntilFinished()其实已经包含了阻塞等待,但它的结果只有getErrors()。如果想更细粒度地感知每个步骤的完成状态,可以监听:
trans.addTransListener(new TransListener() { @Override public void transStarted(Trans trans) throws KettleException { log.info("转换开始执行"); } @Override public void transFinished(Trans trans) throws KettleException { if (trans.getErrors() > 0) { log.error("转换执行完成,但存在错误"); } else { log.info("转换执行成功"); } } });这里有个使用注意事项:监听器回调是在Kettle的内部线程中触发的,不要在transFinished里直接操作Spring容器中的无状态单例Bean的事务之类的东西,做简单的状态上报没问题,但要干重活还是建议通过消息队列异步化。
5.3 日志堆积问题
Kettle默认会把所有通道的日志缓存到内存里,给开发者用Spoon界面查看。但在Spring Boot这种长跑进程中,默认配置会让内存涨个不停。必须在初始化后手动关闭日志缓冲区:
KettleLogStore.setMaxAge(HOURS.toSeconds(1)); KettleLogStore.setMaxBufferLines(5000);setMaxBufferLines限制缓冲日志条数,setMaxAge限制保留时间。这个不设置的话,长期跑任务的服务OOM只是时间问题。我在生产环境实测过,连续跑一周后堆内存占用从800MB猛涨到2GB,加了这两行之后稳如老狗。
6. 基于数据库资源库和转换文件的几种扩展玩法
集成的基本盘稳定之后,我陆续加了不少能力,这里挑几个比较实用的说一下。
6.1 数据库资源库方式读取
虽然我推荐文件目录,但业务上有时候存量ktr还在数据库资源库里。用Java读取资源库转换的基本步骤:
Repository repo = ((Repository) KettleEnvironment.getRepository()); RepositoryDirectoryInterface dir = repo.findDirectory("/home/myFolder"); TransMeta transMeta = repo.loadTransformation("myTrans", dir, null, true, null); Trans trans = new Trans(transMeta);这个能跑通的前提,是在kettle.properties里配置好了资源库的类型、连接串和账号。注意,数据库资源库如果配置的是MySQL,需要你自己把MySQL驱动放到Kettle的classpath里,否则repo.loadTransformation直接给你来个ClassNotFound。
6.2 定时任务与并发控制
Spring Boot集成的最大优势就是天然能接上调度框架。我用的是@Scheduled+ExecutorService自管理,而不是直接把Kettle塞进Quartz的Job里。原因很简单:可以统一控制并发数,避免两个任务同时跑同一个转换导致数据重复。
我封装了一个简单的执行锁:
private ConcurrentHashMap<String, AtomicBoolean> runningFlag = new ConcurrentHashMap<>(); public boolean tryLock(String jobName) { return runningFlag.computeIfAbsent(jobName, k -> new AtomicBoolean(false)) .compareAndSet(false, true); } public void unlock(String jobName) { runningFlag.get(jobName).set(false); }任务跑之前先tryLock,跑完或异常时unlock。这个方案比@Scheduled的默认行为可靠得多,因为Spring的定时器默认线程池只有一个线程,如果任务耗时超过间隔,后面的任务就会排队堆积,不是我们想要的。
6.3 数据库类型扩展:达梦、TDengine到MySQL迁移
热词里提到的“kettle支持taos数据库迁移到mysql”“kettle连接达梦”,本质是Kettle如何识别新数据库类型。Kettle 8.2默认的数据库类型列表里没有达梦和TDengine,你需要:
- 下载官方JDBC驱动jar。
- 驱动加载方式有两种:扔到JDK的ext目录(不推荐),或者通过代码注册到Kettle的DriverLocator里。
- 更稳妥的方案是直接在ktr文件里用“Generic Database”类型,把驱动类名和URL模板配进去。
我在做TDengine到MySQL的迁移时,用的是Generic方式,在“表输入”里写SQL,然后“表输出”选Generic Database,URL填jdbc:TAOS-RS://ip:6041/dbname,驱动类名填com.taosdata.jdbc.rs.RestfulDriver。跑批量迁移,10万级别的数据量不是问题。
注意:Generic方式下,Kettle无法自动获取表字段元数据,你必须在“表输出”步骤手动指定字段和映射关系。批量建表建议直接写SQL在“执行SQL脚本”步骤里跑。
7. 踩坑实录:这些问题我不希望你再去试一遍
写到最后,把我在集成过程中遇到的几个典型问题整理成清单。每个问题背后都对应我至少一晚上的排查时间。
7.1 依赖冲突导致的NoSuchMethodError
症状:启动时正常,执行转换时抛NoSuchMethodError: com.google.common.util.concurrent.internal.InternalFutureFailureAccess。原因就是guava版本太低或者被排除了相关类。解决办法:把guava升到30.1-jre并加failureaccess依赖。如果还不行,检查是不是guava-android和guava-jre两个模块混用了。
7.2 JavaScript组件无法运行报org.mozilla.javascript类缺失
Kettle里的“执行SQL脚本”如果没有问题,那“Java代码”或者“Modified Java Script Value”这种脚本组件需要额外的rhino引擎。依赖里缺了js相关jar时,直接报类找不到。补充:
<dependency> <groupId>org.mozilla</groupId> <artifactId>rhino</artifactId> <version>1.7.12</version> </dependency>7.3 中文乱码问题
Kettle读取的文件如果没有指定编码,默认可能是ISO-8859-1,落库后中文就变成问号了。在ktr文件的“CSV输入”“文本文件输入”步骤里,把编码明确改成UTF-8。Java侧加载文件时也建议统一用Files.readAllLines(Paths.get(path), StandardCharsets.UTF_8),避免从参数带进来乱码。
7.4 数据库连接空闲超时
Kettle的数据库连接池默认用的是DBCP,空闲一段时间后可能被数据库主动断掉,再次执行任务时报Communications link failure。传统做法是调大数据库的wait_timeout,但治标不治本。我在集成层做了一层连接池预热+失败重试,每次执行前跑一个SELECT 1探活,连接挂了就自动重建连接,这才彻底解决。
7.5 执行转换时线程阻塞不返回
有段时间一个转换任务偶发卡死,排查后发现ktr里面某个步骤开启了一个事务但没提交,导致后续步骤一直等锁。Kettle的日志一旦停在某一步持续不往下走,数据库侧查一下information_schema.innodb_trx,大概率能发现一个长期挂起的事务。直接杀掉对应事务,Kettle任务会立刻恢复报错。
8. 生产落地建议与后续优化方向
现在整套集成方案已经稳定运行在我负责的数据平台里,每天支撑几十个定时任务和若干实时触发任务。如果你是第一次做这个集成,这几个落地建议可以帮你少走一点弯路。
首先,ktr文件尽量做得“小而专”。一个转换只干一件事,比如“同步订单表”“清洗用户表”,这样Java侧的重用和排错都容易。不要搞一个几百步的巨型转换,出了问题连Kettle自身都得卡半天。
其次,执行结果必须持久化。我在核心表里记录了每次任务执行的开始时间、结束时间、错误数、执行人,这样后续做任务统计和性能分析才有数据依据。Kettle本身不提供这块信息,只能靠Java侧自己埋点。
第三,资源库方式别轻易上。除非你的存量任务全在里面,否则文件目录的轻量方式绝对更适合Spring Boot工程。Git版本管理、评审、回滚,这些现代研发流程要的能力文件方式都能给你,资源库反而成了黑盒。
第四,集成层不要绑死Kettle API。我在项目里抽象了一个EtlExecutor接口,Kettle只是其中一个实现。后面如果某个任务换成Flink或者Spark跑,只需要再写一个实现类,对上层业务完全无感知。这种解耦让你不会被厂商锁定。
关于后续扩展,我正计划做的是:把转换文件上传到OSS,通过配置中心下发执行计划,再配合分布式锁把执行能力横向扩展。Kettle在单机场景下依然很能打,但有了Spring Boot这层壳,它就不再是一个孤立的工具,而是一个可以被编排、被监控、被沉淀成数据资产的服务。嵌入过程中踩过的那些坑,现在看都是值得的——因为整套东西跑起来之后,数据团队终于不用再凌晨爬起来看任务跑没跑完,看一眼仪表盘,比什么都清楚。