ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Spring Boot集成Kettle:ETL从手工操作到服务化

Spring Boot集成Kettle:ETL从手工操作到服务化 最近做数据抽取需求时踩了一圈Spring Boot集成Kettle的坑从依赖冲突到驱动缺失、从资源库配置到定时跑批最后总算理出一条能稳定复现的路径。这篇就聊点实际操作不扯虚的把Spring Boot和Kettle怎么真正玩到一起讲清楚。Kettle这名字老用户可能更熟的是它的图形化工具Spoon正经全称是Pentaho Data IntegrationPDI做ETL数据抽取、转换、加载的开源神器。而Spring Boot负责把接口、调度、监控这些上层能力串起来两者结合以后数据抽取就不再是“打开Spoon手动点运行”而是变成应用内可调用、可调度、可监控的一个服务能力。这套玩法尤其适合有大量定时跑批、数据同步、清洗任务同时又不想单独维护一堆调度脚本的团队。文章不挑读者刚接触Kettle的小白能跟着做通一遍有Spring Boot经验但没碰过Kettle的人也能避开我踩过的坑。全文偏实操你需要准备的是一个能跑起来的Spring Boot项目我用的JDK8和Spring Boot 2.x一个Kettle解压包剩下就是耐心。1. 为什么要把Kettle嵌进Spring Boot1.1 先搞明白Kettle在项目里扮演什么角色Kettle本质是一个用Java写的ETL引擎核心能力是把数据从一个地方搬到另一个地方途中还能做清洗、转换、校验、分发。它有两种顶层设计对象转换Transformation和作业Job。转换负责“数据流”从输入步骤到输出步骤一条线跑完作业负责“控制流”可以编排多个转换、脚本、发邮件、判断分支跑批场景基本都靠作业来组织。很多团队最初用Kettle是直接在Spoon里手动拖步骤、点运行生产环境则用kitchen.sh执行作业和pan.sh执行转换配Linux crontab来跑批。这么干的问题很明显跑批状态散落在服务器任务日志不好集中看改个阈值或接口地址要去改ktr文件而且很难和内部系统打通比如“业务方点个按钮立即触发一次全量同步”这种需求crontab根本做不了。把Kettle集成到Spring Boot之后上面这些事就变成常规操作了。转换和作业仍然由Kettle的图形化界面设计但执行、调度、监控、参数下发都由应用层接管业务系统可以像调用一个普通Service一样发起数据同步任务。1.2 三种集成方案选哪个更靠谱网上能搜到的主流做法有三种我先对比一下再给结论。方案实现方式优点缺点命令行调用在Java里用ProcessBuilder执行pan.sh或kitchen.sh不需要处理Kettle类库隔离性好每次执行都要起JVM性能差日志解析靠字符串不靠谱跨平台坑多Java API内嵌引入Kettle的jar在应用里直接new TransMeta、JobMeta执行性能好能精细化获取日志、传参数、拿状态依赖多而杂容易冲突需要花时间调classpath独立ETL服务把Kettle封装成单独的微服务通过REST或消息队列触发职责清晰能独立扩展不影响主业务架构复杂小项目可能过度设计我自己的实践结论是优先选Java API内嵌。除非你的ETL任务量非常大、需要独立扩容否则内嵌方式够用而且最灵活。原因就一条Kettle本身就是Java写的官方提供的KettleEnvironment、TransMeta、JobMeta这些类就是留给开发者做二次集成的与其在外面包一层命令行不如直接调用内核。命令行方式看着简单实际运行时会遇到环境变量、classpath、权限、路径各种问题调试一次就想骂人。1.3 集成后能做什么解决了什么问题集成之后我们能获得几样立竿见影的能力接口触发通过Spring Boot的Controller暴露一个POST /sync/user内部执行Kettle作业业务系统随时调用。统一调度不再依赖Linux crontab直接在应用里用Scheduled或接入分布式调度框架跑批规律、失败重试都好管。参数注入同步日期、目标表名、增量查询的条件值可以在调用时动态传入不用每次改ktr文件。日志聚合Kettle执行日志可以回流到应用日志里配合ELK或Spring Boot Admin做监控。状态反馈知道每个步骤处理了多少行、用了多长时间是成功还是失败方便业务方看到同步结果。这和单纯用Spoon手动跑完全不同ETL从“手工活”变成了“服务能力”这也是项目做数据平台化必走的一步。2. 环境准备与Spring Boot集成前的搭建2.1 Kettle版本怎么选怎么下载安装Kettle的官方名称已经改成PDIPentaho Data Integration去官网或SourceForge能找到历史版本。社区版虽然功能上有阉割但主要ETL功能都在个人使用和内部系统集成足够了。版本选择上我强烈建议用PDI 9.x比如9.3或9.4这两个版本相对稳定对JDK8兼容好网上踩坑案例也多。不建议一上来就追最新版因为Kettle迭代时经常调整jar包里的包名和类名Spring Boot集成的开源资料大多基于9.x用新版本容易遇到类找不着、方法签名变了的问题。我自己用的是pdi-ce-9.4.0.0-343。下载后不用安装系统服务解压即可。关键目录作用如下目录/文件作用>repositories repository idpentaho-releases/id urlhttps://repo.hortonworks.com/content/repositories/releases//url /repository repository idoracle-releases/id urlhttps://repo.oracle.com/maven/public//url /repository /repositories properties kettle.version9.4.0.0-343/kettle.version /properties dependencies !-- Kettle核心 -- dependency groupIdpentaho-kettle/groupId artifactIdkettle-core/artifactId version${kettle.version}/version /dependency dependency groupIdpentaho-kettle/groupId artifactIdkettle-engine/artifactId version${kettle.version}/version /dependency !-- Kettle界面依赖某些API需要 -- dependency groupIdpentaho/groupId artifactIdpentaho-metastore/artifactId version${kettle.version}/version /dependency /dependencies如果某些包用groupId:pentaho-kettle拉不到就需要去PDI安装目录的lib里翻。比如kettle-json、pentaho-vfs之类将它们用systemPath引入dependency groupIdpentaho/groupId artifactIdkettle-json/artifactId version1.0/version scopesystem/scope systemPath${project.basedir}/lib/kettle-json-1.0.jar/systemPath /dependency说实话这部分是最磨人的Kettle的依赖树很深包括commons-vfs2、guava、jfreechart、jackson都和Spring Boot自带版本有重叠。我建议在集成前先把>Configuration public class KettleConfig { PostConstruct public void init() { KettleEnvironment.init(); } }这么写看着没问题但在Spring Boot里埋了个坑KettleEnvironment.init()执行时会加载很多系统配置和插件如果环境里缺少某个数据库方言驱动或者classpath里有冲突的jar它可能不报错但后面执行转换时突然抛异常。我建议改成显式初始化同时加一个判断把Kettle日志也接好Configuration public class KettleConfig { Bean public KettleEnvironment kettleEnvironment() { if (!KettleEnvironment.isInitialized()) { KettleEnvironment.init(); } return KettleEnvironment.getInstance(); } }另外Kettle默认的日志输出是写到控制台集成进Spring Boot后最好把它的日志接到SLF4J不然排查问题时要切两个日志窗口。做法是给org.pentaho.di.*日志输出器换成Slf4jLoggingObject或者直接在转换执行时指定日志级别和回调这部分放到第3节的实操代码里一起说。到这里环境基本就绪。接下来是核心操作加载ktr文件、执行、传参数、拿状态。3. 核心功能实现加载、执行与参数传递3.1 用Java加载转换和作业先分清两种方式Kettle的转换/作业可以存在两种地方文件.ktr/.kjb和资源库数据库资源库/文件资源库。集成时我建议从文件开始因为简单、依赖少。加载转换的代码很直接// 从文件加载转换 TransMeta transMeta new TransMeta(/data/sync/sync_user.ktr); Trans trans new Trans(transMeta); trans.prepareExecution(null); trans.start(); trans.waitUntilFinished(); if (trans.getErrors() 0) { throw new RuntimeException(Kettle转换执行失败); }这里有一个非常关键的点TransMeta是元数据而Trans是运行时实例。同一个TransMeta可以创建多个Trans实例并行跑但如果你在整个应用里共享一个TransMeta多个线程同时调用时很可能会互相污染变量空间。我踩过的坑就是第一次跑正常第二次跑发现数据库连接被复用、上一步的参数串进来了。原因就是我把TransMeta当成单例缓存了。正确做法是每次任务创建独立的TransMeta或者确保每次执行前都调用trans.reset()。加载作业和加载转换类似JobMeta jobMeta new JobMeta(/data/sync/sync_job.kjb, null); Job job new Job(null, jobMeta); job.start(); job.waitUntilFinished();作业里面可以嵌套转换所以跑批调度统一交给作业最合适。3.2 给Kettle传参数实现动态每批次同步实际业务里不可能写死一个SQL里的时间条件。Kettle支持变量variables、**参数parameters和属性properties**三种动态值我用得最多的是参数parameters。在Spoon里设计转换时可以通过“转换设置 参数”定义参数名。Java端用如下方式传参TransMeta transMeta new TransMeta(/data/sync/sync_user.ktr); Trans trans new Trans(transMeta); // 方式一设置变量 trans.setVariable(startDate, 2024-01-01); trans.setVariable(endDate, 2024-12-31); // 方式二设置参数 trans.setParameterValue(targetTable, ods_user_daily); trans.prepareExecution(null); trans.start(); trans.waitUntilFinished();如果作业里包含多个转换参数怎么传这里有个容易忽视的细节作业的参数需要单独设置。JobMeta里可以定义参数然后在Job上执行job.setParameterValue(...)作业会把它传给内部转换的同名参数。但有些情况下内部转换不认你还需要在作业的“参数”页签里做“参数映射”或者在调用转换的前一步使用“设置变量”这一步。我一般用一套约定调用方只传粒度较粗的参数比如日期、批次号Kettle内部通过“JavaScript步骤”或“查询步骤”把参数映射成具体的SQL条件。比如在SQL步骤里这样引用SELECT * FROM user WHERE create_time ${startDate} AND create_time ${endDate}注意这里的${}是变量引用不是参数引用。如果你设置的是参数需要先把它转成变量。最简单的方式是在转换的“参数”页签中勾选“将其作为变量定义”这样参数会自动变成变量SQL才能引到。这个细节让我当时查了一个下午。3.3 配置资源库把ktr/kjb统一管理起来文件方式虽然简单但生产环境有变更麻烦总不能每次把ktr文件拷到服务器上。更好的方式是使用数据库资源库把转换和作业的定义存到数据库表里Java端通过资源库对象加载。数据库资源库本质就是一些元数据表。通过Kettle的Spoon先创建一个数据库资源库并连接然后往里面导入转换和作业。Java端加载时需要设置资源库的连接信息DatabaseMeta databaseMeta new DatabaseMeta(); databaseMeta.setName(kettle_repo); databaseMeta.setDatabaseType(MYSQL); databaseMeta.setAccess(Const.ACCESS_TYPE_NATIVE); databaseMeta.setHostname(localhost); databaseMeta.setPort(3306); databaseMeta.setDBName(kettle_repo); databaseMeta.setUsername(root); databaseMeta.setPassword(password); KettleDatabaseRepository repository new KettleDatabaseRepository(); KettleDatabaseRepositoryMeta repositoryMeta new KettleDatabaseRepositoryMeta(); repositoryMeta.setName(my_repo); repositoryMeta.setDatabaseConnection(databaseMeta); repository.init(repositoryMeta); // 从资源库加载作业/转换 RepositoryDirectoryInterface directory repository.loadRepositoryDirectoryTree(); JobMeta jobMeta repository.loadJob(path/to/job, directory);这个方案比文件方式好维护但它的缺点是要在数据库里维护元数据对团队运维有要求。我的建议是小项目直接用文件方式文件放到配置中心或者统一目录大项目上资源库同时把ktr/kjb纳入版本管理两者不冲突。3.4 用Spring Boot的定时任务实现自动跑批热词里有个“kettle设置自动跑批”在Spring Boot里最常见的做法就是Scheduled。代码不复杂但要注意并发和重复执行的问题。Component public class KettleScheduler { Scheduled(cron 0 0 2 * * ?) public void runDailySync() { executeJob(/data/sync/daily_sync.kjb); } private void executeJob(String kjbPath) { try { JobMeta jobMeta new JobMeta(kjbPath, null); Job job new Job(null, jobMeta); job.setVariable(syncDate, LocalDate.now().minusDays(1).toString()); job.start(); job.waitUntilFinished(); if (job.getErrors() 0) { log.error(跑批失败 kjbPath); } } catch (KettleException e) { log.error(Kettle作业执行异常, e); } } }Scheduled默认是单线程调度的要小心两点如果上一次任务还没跑完下一次触发又开始了可能造成数据重复处理。建议在任务方法上加分布式锁或者用一个AtomicBoolean做进程内互斥。多任务时最好配置一个线程池否则一个耗时任务会把其他定时任务阻塞。更稳妥的做法是给Spring Boot加一个调度框架比如xxl-job或Quartz。Kettle集成和Quartz天然契合因为Quartz本身就是Kettle作业调度用的默认调度器。如果团队已经有xxl-job那就更方便把“执行Kettle作业”封装成Java方法通过xxl-job的XxlJob注解触发日志可以统一到调度中心的日志平台。3.5 捕捉Kettle的日志和状态监控别再靠猜集成后最容易被忽略的是日志。Kettle的日志有自己一套体系如果你直接调用trans.start()日志会打到控制台在Spring Boot里就是系统stdout很难和业务日志对应。我建议把Kettle的监听器挂到应用日志上。Kettle提供了KettleExecutionListener接口可以实现executionStarted、executionFinished等方法拿到状态同时也能拿到每行处理进度。一个实用实现public class KettleLogListener implements KettleExecutionListener { private final Logger log LoggerFactory.getLogger(KettleLogListener.class); Override public void executionStarted(ExecutionControl control) { log.info(执行启动{}, control.getName()); } Override public void executionFinished(ExecutionControl control) { log.info(执行结束{}总行数{}, control.getName(), control.getTotalSteps()); } }注册监听器方式Trans trans new Trans(transMeta); KettleLogListener listener new KettleLogListener(); trans.setExecutionListener(listener);同时Kettle的每个步骤Step都有自己的进度信息包括读取行数、写入行数、错误行数。我通常会在作业结束时把最终的执行结果成功/失败、耗时、处理行数写进业务表里这样“今天凌晨同步了多少条数据”随时可以查。更进一步的监控可以通过日志采集把Kettle的运行指标推到Spring Boot Admin或者Prometheus。比如在监听器里定时记录StepMetrics然后用Micrometer暴露给端点。这部分属于可扩展内容项目有需要时值得做。4. 实战中高频问题与排查记录4.1 常见错误速查表集成过程至少有一半时间在解决各种“找不到类”“连接不上”的问题。我把最常遇到的列成一张表方便你对照排查。现象根本原因解决办法ClassNotFoundException: org.pentaho.di.core.database.DatabaseMetaKettle核心包没引入确认引入kettle-core用systemPath引入本地jar时路径要正确Could not load step from class ...缺少插件或步骤类检查plugins/目录是否包含对应步骤插件或把整个plugins目录放到应用的classpathAccess denied for user rootlocalhost数据库连接配置错误或驱动版本不匹配检查DatabaseMeta设置驱动jar版本要和数据库对应MySQL8需要用mysql-connector-java8.xUnable to load class for JDBC driver驱动jar没放到执行环境把驱动jar放到PDI的lib目录或引入到应用依赖中KettleEnvironment not initialized未调用KettleEnvironment.init()在Spring Boot启动时初始化返回KettleEnvironment.getInstance()The system cannot find the file specified默认从当前工作目录找文件尽量使用绝对路径或使用System.getProperty(user.dir)拼接与Spring Boot的logback冲突导致的NoSuchMethodErrorKettle依赖的老版slf4j/logback和Spring Boot冲突统一slf4j版本在pom里排除旧版本的slf4j-log4j12等UCanAccess driver not found使用Access数据库时未引入UCanAccess驱动单独加入ucanaccess、jackcess、commons-lang3等依赖版本要配合4.2 UCanAccess驱动与Access数据库的坑热词里提到“kettle ucanaccess 驱动”这个确实坑很多。如果你的数据源是Access.mdb/.accdbKettle默认带的是sun.jdbc.odbc.JdbcOdbcDriver但JDK8以后移除了JDBC-ODBC桥所以必须改用UCanAccess驱动。集成时的注意点UCanAccess需要几个配套jar包括ucanaccess、jackcess、commons-lang3、commons-logging、hsqldb等缺一不可。在Kettle的数据库连接里连接类型选“Access”但Java代码中通过DatabaseMeta创建连接时可能不支持UCanAccess的驱动类更靠谱的办法是先设置DatabaseMeta的自定义驱动和URL。databaseMeta.setDriverClass(net.ucanaccess.jdbc.UcanloadDriver); databaseMeta.setURL(jdbc:ucanaccess:///data/db/test.accdb);你可能会遇到“java.lang.NoClassDefFoundError: com/healthmarketscience/jackcess/...”那基本就是jackcess缺失。在Spring Boot里建议直接用Maven依赖dependency groupIdcom.h2database/groupId artifactIdh2/artifactId version2.1.210/version /dependency dependency groupIdnet.sf.ucanaccess/groupId artifactIducanaccess/artifactId version5.0.1/version /dependencyUCanAccess自带了一个H2的内嵌依赖所以要保证H2版本匹配否则会报错。这属于小众场景但一旦碰到就会卡很久先记下。4.3 并发执行、共享状态与性能调优Kettle转换能不能多线程跑可以但要注意几点。同一个TransMeta不能同时被多个Trans并发执行。原因是步骤对象内部会有状态比如“已打开的连接数”“游标位置”都是实例级别的。我后来统一在每次执行时调用transMeta.cleanup()且新建TransMeta彻底解决问题。Kettle的数据库连接池默认是每个步骤独立建的如果一次跑100个步骤可能瞬间开100个连接。这个在kettle.properties里配置连接池上限也可以在每个数据库连接的“连接池”选项里设置最大连接数。Spring Boot侧也要控制好线程池避免大量任务同时涌进来。内存溢出是常客。Kettle执行时会把数据放入缓冲区当“表输入”和“表输出”之间的行集RowSet过大而目标库写入慢时内存就爆。解决方式是在关键步骤中间加上“延时”或“阻塞数据直到完成”或者调大JVM的Xmx更彻底的是给批量写加上合理提交大小比如每次5000行。性能调优没有银弹我的经验是先看trans.stepPerformanceSnapShot找出耗时最长的步骤再针对它优化。Kettle里最耗时的往往不在计算而在数据库IO。检查SQL是否走了索引、是否一次性取全表、目标表是否有索引锁这些都比换配置项更有效。4.4 对外提供接口时集成服务应该放哪热词里有一个“spring boot对外提供的接口(给第三方)应该放在哪里?是单独服务还是放在对应服务”放到Kettle场景里我的建议是在已有系统里增加一个独立的ETL Controller模块不要让业务Controller直接去触发重型Kettle任务。原因很实际Kettle作业执行是耗时的如果同步几百万数据接口可能要跑几分钟第三方的HTTP调用根本等不了。正确做法是暴露一个POST /api/etl/start接口接收任务编号和参数。接口把请求扔进线程池或者消息队列立刻返回“任务已提交”。任务服务异步执行Kettle执行完把结果更新到任务状态表。第三方通过GET /api/etl/status?taskIdxxx轮询查询进展。这样既保证了第三方调用体验又不阻塞业务主流程。如果团队服务已经拆得很细把Kettle能力独立成一个“数据同步服务”也完全可以但不要为了“干净”而强行拆最终依据还是调用方和部署环境。5. 工程化落地的几个扩展建议5.1 把ktr路径和参数玩成配置化用Spring Boot集成Kettle后一个很容易犯的错是在代码里硬编码ktr/kjb路径、数据库连接信息然后每次环境变了都要改代码重新部署。我建议把这些全部挪到application.yml里。etl: kettle: root-path: /data/etl repo-enabled: false tasks: daily-user-sync: file: job/daily_user_sync.kjb variables: targetTable: ods_user_daily cron: 0 0 2 * * ?再用ConfigurationProperties绑定一个EtlProperties类任务执行时按名字找到文件路径和参数。这样新增一个Kettle任务不需要改Java代码只在配置里加一段运维同事也能看懂。5.2 失败重试、重跑和告警Kettle跑批失败很常见比如源库在凌晨做备份把连接踢了、目标表临时锁死。这时候要判断该不该重试。我的经验是不要盲目重试整个作业要根据失败点决定。如果是网络抖动或临时锁整体重试1-2次是合理的。如果是数据质量问题源数据格式错、约束冲突重试再多次也没用应该告警人工介入。在Spring Boot里实现重试很方便用Spring Retry即可Retryable(maxAttempts 3, backoff Backoff(delay 2000, multiplier 2)) public void executeJobWithRetry(JobMeta jobMeta) { executeJob(jobMeta); }但注意Retryable不适合直接作用于长期运行任务最好包一层让它只检查启动瞬间的异常。作业执行中如果是因为SQL出错Kettle会返回errors 0而不是抛异常所以重试逻辑要同时处理Kettle错误数和Java异常两种情况。告警我建议直接用Spring Boot Admin或结合已有的钉钉/企微机器人。把“作业执行失败”“作业耗时超过阈值”作为告警源简单用ApplicationEventPublisher发一个事件监听器发消息。5.3 结合WebSocket把执行日志推到前端如果你有Web应用需要实时展示ETL任务进度Kettle日志配合WebSocket是个很舒服的方案。Spring Boot集成的配置只要在application.yml里加一段spring: websocket: mapping-path: /ws/etlKettle那边在步骤监听器里通过事件发送日志消息用SimpMessagingTemplate推给指定会话。用户发起同步后页面上能实时看到“正在读取第20000行”这种动态体验比轮询接口好太多。这里不展开代码原理就是WebSocket Kettle步骤级监听值得做。最后说一个我实际用下来很有用的细节Kettle转换里如果有数据库连接在Spring Boot环境里最好把每个数据库连接都设置成“每次执行后关闭”。怎么操作在Spoon里双击数据库连接打开“选项”Tab勾选“连接池”配置时设置“每次运行后释放连接”。不然应用持续运行几天后经常会出现连接耗尽而日志只是报一个“Connection is not available, request timed out”。我当时排查了很久最后发现根本不是Kettle的问题而是连接池泄漏。养成这个习惯能少踩很多莫名其妙的坑。Spring Boot集成Kettle这套组合前期依赖处理确实烦但一旦跑通你会发现自己手里多了一把数据管线的钥匙。它能解决的问题远超“跑一下转换”本身更像是在业务系统和数据世界之间架起了一座可编程的桥。
返回列表