尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Spring Boot集成Kettle:从ktr加载到执行监控的完整实践

Spring Boot集成Kettle:从ktr加载到执行监控的完整实践 只要做过数据平台开发大概率都遇到过这种场景业务方丢过来一个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 IntegrationJava编写Maven坐标是pentaho-kettle系列。在实际引入Spring Boot项目时版本选型是第一个决定成败的点。2.1 版本组合怎么定我的生产环境使用的是Spring Boot 2.2.x Kettle 8.2。这套组合跑了很长时间比较平稳。Kettle 8.3也试过但我踩到了和某些第三方库的兼容问题所以生产上一直锁在8.2。dependency groupIdorg.pentaho/groupId artifactIdpdi-engine/artifactId version8.2.0.0-342/version /dependency dependency groupIdorg.pentaho/groupId artifactIdkettle-core/artifactId version8.2.0.0-342/version /dependency dependency groupIdorg.pentaho/groupId artifactIdkettle-engine/artifactId version8.2.0.0-342/version /dependency实际用起来kettle-engine基本涵盖了一个转换从解析到执行的全部核心类kettle-core提供资源库接口、日志体系等基础能力。如果你的ktr或kjb里用到了高版本才有的组件再按需加其他模块。2.2 依赖冲突处理Kettle 8.2会传递引入一堆老版本的第三方库其中最容易和Spring Boot冲突的是guavaKettle传递来的版本可能很老和Spring Boot里用的新版本API冲突。jacksonKettle内部有自己的Jackson版本如果覆盖不对JSON组件解析会出问题。commons-lang3、commons-collections4这类基础库版本冲突容易造成莫名其妙的NoSuchMethodError。我的处理方式是在pom.xml里用exclusions排除Kettle传递的旧包然后由Spring Boot的依赖管理统一控制版本。比如dependency groupIdorg.pentaho/groupId artifactIdkettle-engine/artifactId version8.2.0.0-342/version exclusions exclusion groupIdcom.google.guava/groupId artifactIdguava/artifactId /exclusion exclusion groupIdcom.fasterxml.jackson.core/groupId artifactIdjackson-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达梦用DmJdbcDriverTDengine用它的JDBC驱动。这些驱动Kettle官方依赖里不会全带必须手动加。有朋友在“kettle连接达梦数据库”上反复踩坑我补充一句达梦的JDBC驱动jar包下载后放到项目lib目录或打成系统依赖都行但注意和Kettle 8.2配合时驱动类名要写成dm.jdbc.driver.DmDriverURL格式是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日志接入SLF4JKettle的日志体系基于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在事件里把文本写到logbackKettleLogStore.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集成的最大优势就是天然能接上调度框架。我用的是ScheduledExecutorService自管理而不是直接把Kettle塞进Quartz的Job里。原因很简单可以统一控制并发数避免两个任务同时跑同一个转换导致数据重复。我封装了一个简单的执行锁private ConcurrentHashMapString, 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 DatabaseURL填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 groupIdorg.mozilla/groupId artifactIdrhino/artifactId version1.7.12/version /dependency7.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这层壳它就不再是一个孤立的工具而是一个可以被编排、被监控、被沉淀成数据资产的服务。嵌入过程中踩过的那些坑现在看都是值得的——因为整套东西跑起来之后数据团队终于不用再凌晨爬起来看任务跑没跑完看一眼仪表盘比什么都清楚。
返回列表