简介:基于Kettle实现的Web版数据集成平台源码包,面向需要快速搭建拖拽式ETL工具的数据工程师、后端开发者及企业数据团队。项目将Kettle的抽取转换加载能力封装为浏览器端可视化服务,用户无需编码即可完成数据源接入、转换流程设计、任务调度与执行监控,并支持作业版本控制和用户权限管理。压缩包共1645个文件,以912个Java后端源码、74个Vue前端组件、123个XML配置、12个KTR转换脚本为主,另含Dockerfile、docker-compose、properties及yaml等部署配置文件,整体约160.93MB,目录结构清晰,便于按前后端模块阅读和二次开发。资源还提供中文说明文档、Kettle转换示例及前端css/js资源,可帮助读者从零搭建环境,深入理解Kettle Server与Web端通信机制及拖拽流程编排的实现思路。已有727人学习下载,适合希望将Kettle产品化或学习数据集成平台架构的开发者参考。
1. 从脚本到Web:为什么要把Kettle搬到浏览器里拖拽
很多团队的日常是这样的:数据工程师用Spoon画好转换和作业,导出成ktr/kjb文件,找台服务器配上定时任务跑。业务部门提需求,改一个字段映射就要改ktr、重跑、验证,来回一天。真正让人头疼的不是Kettle本身,而是它作为一个桌面工具,能力全在少数几个人的本地上,流程、脚本、日志都不透明。这个标题要解决的,就是把这套东西包一层Web外壳,用浏览器拖拽节点和连线,后台自动生成Kettle可识别的XML,再交给Kettle引擎去执行。适合正在建设数据中台、想做自助数据集成的工程师和平台负责人,也适合想把分散在各处的Kettle作业收归统一门户的运维团队。
我按自己做这类平台的思路,把整个方案拆成选型、解析、调度、排错四个环节往下讲,全部围绕一个目标:让不碰Spoon的人也能画出能跑的Kettle任务。
2. 选型与架构:Kettle引擎嵌入Web服务的三种姿势与取舍
2.1 先想清楚:你是在改Kettle,还是在封一层壳
这是做Web版数据集成平台最容易跑偏的地方。如果你打算改Kettle源码,给Spoon的拖拽画布换一套Web前端,那工作量会大到失控。Kettle的Spoon是一个基于SWT的桌面应用,它的画布交互、节点编辑、日志视图全部耦合在桌面框架里,把它Web化等于重写客户端。而业界绝大多数“Kettle Web化”方案,其实是封装:前端自己做一套拖拽画布,把画布结构翻译成Kettle标准的ktr/kjb文件,服务端调用Kettle引擎去执行。Kettle本身熟悉并且久经考验的转换和作业模型保留不变,替换掉的只是人与引擎之间的那层桌面交互。搞清楚这一点,架构就不会跑偏。
我一般也建议团队把“Web拖拽生成ktr”与“Kettle执行ktr”当作两个独立模块来开发。前端输出的是标准ktr,那么哪怕将来执行引擎换成Spark或Flink,只需要替换解析层,画布资产不会作废。执行层也一样,Kettle版本升级不影响前端结构。这种边界清晰的拆法,让整个平台不会变成一堆谁都不敢动的旧作业。
2.2 三种嵌入方式:Kitchen命令、Carte服务、内嵌API
Kettle本身提供了几条让外部程序调用它的路径,选哪条直接决定平台的运维模型。
第一种是命令行方式,通过pan跑转换、kitchen跑作业。Web服务端把前端生成的ktr写到临时目录,然后ProcessBuilder拉起pan进程执行。好处是隔离性最强,转换崩了不会拖垮Web进程,内存可以通过JVM参数单独控制;坏处是每次执行都要新起一个JVM,启动耗时少则几秒多则十几秒,高并发下进程开销大。适合任务频率低、单个转换体量大、Web服务内存本身不宽裕的场景。
第二种是Carte服务。Kettle自带一个轻量Web服务Carte,可以在远程机器上启动,接收执行请求并返回运行状态。Spoon里那个“远程执行”功能就是通过它实现的。Web后端把ktr发送给Carte,由远端引擎执行,天然支持分布式。但Carte本身只暴露很简单的HTTP接口,作业信息、执行历史都需要自行对接,而且它默认不带鉴权,暴露到网络上有安全风险,需要自己做一层网关。适合需要集群执行、按任务量横向扩展的场景,也是我建议平台早期就预留的扩展位。
第三种是内嵌API,也是我在这类项目里最常用的方式。在Spring Boot进程里引入kettle-core、kettle-engine等依赖,KettleEnvironment.init()初始化引擎,随后用TransMeta加载ktr,再用Trans执行。没有跨进程传输,执行速度快,日志可以编程方式接管,直接打进数据库或者消息队列,非常契合Web平台对状态透明的要求。缺点是和Web应用共享JVM,大的转换会挤占Web服务的堆内存,必须对并发数做限流。
三种方式不是互斥的。常见做法是默认走内嵌API,遇到大转换或机器资源紧张时,通过Carte把任务扔到独立执行机上。平台在设计任务提交层时,把“执行方式”抽象成一个接口,内嵌和Carte各实现一个,后面加新引擎也容易。别一开始就把所有任务都压在Web服务这一侧。
2.3 平台模块划分与前后端数据流转
按我自己的习惯,这个平台的核心模块拆成五块:画布设计器、解析生成器、任务调度器、执行引擎适配层、日志与告警。前端画布负责提供拖拽节点、连线、属性面板;解析生成器把画布的JSON结构翻译成ktr/kjb;调度器负责定时触发和被外部系统调用;执行引擎适配层封装三种跑法;日志告警负责收集执行状态和步骤指标。
数据的流转链路是这样的:前端把用户画的DAG保存成一个JSON草图,提交给后端,后端解析成标准ktr并存入文件库或数据库字段;定时任务触发时,调度器从库里捞出ktr,通过适配层下发执行;执行过程中引擎回调日志和每一行处理数,后端落库,前端通过WebSocket订阅状态刷新。这套链路成熟、好调试,出问题时能明确知道是哪一层出了问题。
一个容易忽略的模块是参数管理。Kettle作业里常常要传日期、表名这些变量,前端设计器里要有一层全局变量面板,编译ktr时把变量注入到<parameters>节点里,而不是让用户在每一个步骤属性里硬填。我在后期重构时补了这个模块,平台的整体可复用性立刻上了一个台阶。
3. 用Web拖拽生成ktr:从DAG到XML的转换实现
3.1 前端画布的数据模型怎么设计
前端拖拽面板,核心是维护一张有向无环图。每个节点是一个Kettle步骤,比如表输入、表输出、字段选择、排序、去重、Join;每条边是一条数据流HOP。我给前端定义的数据结构长这样,节点和边分开维护:
{ "transName": "order_to_dw", "params": { "bizDate": "2024-01-01" }, "nodes": [ { "id": "node_1", "type": "TableInput", "config": { "connection": "mysql_dw", "sql": "SELECT order_id, user_id, amount FROM orders WHERE dt = '${bizDate}'" }, "x": 120, "y": 80 }, { "id": "node_2", "type": "TableOutput", "config": { "connection": "mysql_dw", "table": "dwd_order", "commitSize": "1000" } } ], "edges": [ { "id": "edge_1", "from": "node_1", "to": "node_2" } ] }这个结构里一个值得注意的设计点是config对象。不同类型的Kettle步骤,配置项千差万别,前端属性面板不可能为每个步骤单独写一套表单。我做法是给每类节点维护一个JSON Schema,渲染属性面板时按Schema生成表单,提交后Schema校验再合并进config。这样新增一种步骤类型,只需要后端加一个转换器,前端加一份Schema,不用改动画布主流程。
节点坐标x和y有两个用途:一是给前端画布做布局还原,二是生成ktr时给Kettle步骤设置<gui>标签里的坐标。Kettle在Spoon里打开ktr时会读取这些坐标来摆放节点,坐标合理的话,运营同学用Spoon排障时看到的面板就是清晰可读的,这点直接影响工具的信任度。
3.2 把DAG翻译成ktr:Java代码怎么组织
ktr文件本质是一个XML描述,包含转换信息、步骤列表、HOP连线列表。前端提交JSON后,后端构建XML并落盘或存库。用Java的标准DOM解析器就可以完成,不需要额外引入模板引擎。核心逻辑是把JSON里的每个节点映射成一个<step>,每条边映射成一个<hop>:
public void buildKtr(TransGraph graph, OutputStream out) throws Exception { DocumentBuilderFactory dbf = DocumentBuilderFactory.newInstance(); DocumentBuilder builder = dbf.newDocumentBuilder(); Document doc = builder.newDocument(); doc.setXmlStandalone(true); Element root = doc.createElement("transformation"); doc.appendChild(root); // info段:转换名称和全局参数 Element info = doc.createElement("info"); Element name = doc.createElement("name"); name.setTextContent(graph.getTransName()); info.appendChild(name); Element parameters = doc.createElement("parameters"); for (Map.Entry<String, String> entry : graph.getParams().entrySet()) { Element param = doc.createElement("parameter"); Element key = doc.createElement("name"); key.setTextContent(entry.getKey()); Element val = doc.createElement("default_value"); val.setTextContent(entry.getValue()); param.appendChild(key); param.appendChild(val); parameters.appendChild(param); } info.appendChild(parameters); root.appendChild(info); // step段:每个节点生成一个Kettle步骤 for (TransNode node : graph.getNodes()) { Element step = doc.createElement("step"); Element stepName = doc.createElement("name"); stepName.setTextContent(node.getId()); step.appendChild(stepName); Element type = doc.createElement("type"); type.setTextContent(StepTypeMapping.getKettleType(node.getType())); step.appendChild(type); // 各种驱动类和连接的配置,每个步骤类型差异很大 Element connection = doc.createElement("connection"); connection.setTextContent(node.getConfig().get("connection")); step.appendChild(connection); Element sql = doc.createElement("sql"); sql.setTextContent(node.getConfig().getOrDefault("sql", "")); step.appendChild(sql); Element lookup = doc.createElement("lookup"); lookup.setTextContent(node.getConfig().getOrDefault("lookup", "")); step.appendChild(lookup); Element show = doc.createElement("show"); show.setTextContent("false"); step.appendChild(show); Element gui = doc.createElement("gui"); Element xloc = doc.createElement("xloc"); xloc.setTextContent(String.valueOf(node.getX())); Element yloc = doc.createElement("yloc"); yloc.setTextContent(String.valueOf(node.getY())); gui.appendChild(xloc); gui.appendChild(yloc); step.appendChild(gui); root.appendChild(step); } // hop段:边的集合,from和to对应步骤的name for (TransEdge edge : graph.getEdges()) { Element hop = doc.createElement("hop"); Element from = doc.createElement("from"); from.setTextContent(edge.getFrom()); hop.appendChild(from); Element to = doc.createElement("to"); to.setTextContent(edge.getTo()); hop.appendChild(to); Element enabled = doc.createElement("enabled"); enabled.setTextContent("true"); hop.appendChild(enabled); root.appendChild(hop); } TransformerFactory tf = TransformerFactory.newInstance(); Transformer transformer = tf.newTransformer(); transformer.setOutputProperty(OutputKeys.INDENT, "yes"); transformer.setOutputProperty(OutputKeys.ENCODING, "UTF-8"); transformer.transform(new DOMSource(doc), new StreamResult(out)); }这段代码里有两个需要留意的参数。第一个是type字段,不能直接使用前端节点的类型名。比如前端叫DataBaseInput,Kettle内部类型是TableInput;前端叫DataBaseOutput,Kettle内部是TableOutput;像StreamLookup、MergeJoin这些名字也不是一眼能对上的。所以必须维护一张映射表,我在工程里把它放在StepTypeMapping里,Kettle各版本对类型名要求严格,写错一个type标签,加载ktr时会直接报StepPluginClass找不到。第二个是<hop>的from和to,必须严格对应<step>的name值,不能填前端生成的node_1这种ID,否则ktr加载时连线找不到端点。
3.3 服务端执行ktr并回传状态
ktr一旦生成,执行逻辑就与前端解耦了。内嵌API执行最简做法如下:
public void runTrans(File ktrFile, Map<String, String> variables) { KettleEnvironment.init(); TransMeta meta = new TransMeta(ktrFile.getAbsolutePath()); Trans trans = new Trans(meta); // 参数注入:把平台侧的运行时间、批次号塞进转换 if (variables != null) { variables.forEach((k, v) -> trans.setVariable(k, v)); } // 设置日志级别:生产环境用BASIC,排查时用ROWLEVEL trans.setLogLevel(LogLevel.BASIC); trans.addTransListener(new TransAdapter() { @Override public void transFinished(Trans trans) { int errors = trans.getErrors(); if (errors > 0) { saveExecResult("failed", trans.getLogChannel().getLogText()); } else { saveExecResult("success", "rows:" + trans.getRowsWritten()); } } }); trans.execute(new String[0]); trans.waitUntilFinished(); }trans.execute是非阻塞的,启动后会立即返回,真正的执行在后台线程里进行,必须调用waitUntilFinished等待完成,否则任务刚提交就退出,后面拿不到执行结果。LogLevel这里有门道,生产环境如果设置成ROWLEVEL,Kettle会打印流经每一行的详细日志,数据量大时日志量会直接撑爆存储,所以默认我用BASIC,只有在具体排查某一条链路时才临时切成ROWLEVEL。
执行结果的回传,我会封装一个ExecResult对象,包含状态、总行数、错误行数、开始结束时间、日志摘要,既有内存对象又有DB落库。前端通过WebSocket拿到推送后刷新任务状态,这就算跑完了一个最小闭环。把这个链路跑通后,后面的调度和监控只是在这个环上加壳。
4. 把任务跑起来:资源库、自动跑批与运行监控
4.1 资源库选型:数据库资源库还是文件仓库
ktr生成后要有一个地方存放它,Kettle官方提供两类存储:文件仓库和数据库资源库。文件仓库就是一坨目录,ktr/kjb按目录层级摆放,简单直接;数据库资源库把转换、作业、步骤参数等元数据写进一系列表里,支持版本管理、多用户共享。做Web平台,一定要用数据库资源库,因为在多用户环境下文件仓库没法解决并发保存、权限控制和版本回滚的问题。
Kettle资源库的表结构由Kettle引擎初始化时自动创建。配置方式是用KettleEnvironment.init()之前先设置kettle.properties,里面指定资源库类型、数据库连接和用户名。Spring Boot工程里我习惯在启动时调用一个初始化Bean,把Kettle环境变量配好。需要建哪些表由Kettle管理,平台侧不要手工改它的表结构,版本升级时Kettle会自行做迁移。资源库选型上我踩过一个坑:Kettle官方资源库默认支持H2、MySQL、PostgreSQL等,但生产上并发的调度任务如果都往里写元数据,连接池配置不对容易出现锁等待,建议给资源库单独用一个数据库实例,并且连接池要比业务库更宽松。
4.2 定时调度的两种实现:Quartz集群与自研调度表
平台的“自动跑批”功能,背后需要一个调度器定时触发ktr执行。这里有两种常见路线:一是直接集成Quartz,二是自己写一张调度表配一个扫描线程。Quartz好处是成熟、支持集群、支持Cron表达式,坏处是每次触发都发一个任务消息,调度历史和失败重试都要自己额外记。自研调度表则更轻量,适合任务数量几百条以内、调度策略不复杂的团队。
我在这类平台里推荐的做法是“自研调度表 + Quartz执行器”混合模式。调度表记任务的Cron、参数、依赖关系;Quartz只负责到点触发一个统一入口,入口从库里查当前批次该跑哪些任务并按依赖排序。这样可以随时在前端修改调度配置,不用重启任何服务。依赖关系的处理是Kettle平台的一个重头戏,因为很多跑批任务要求父任务成功后子任务才能启动,最简单的方式是给调度表加两个字段:parent_task_id和status,子任务在扫描时只取父任务已成功的记录。
// 简化版调度扫描逻辑,实际工程中会配合Quartz的@Scheduled触发 @Scheduled(cron = "0 * * * * ?") public void scanTask() { List<Task> readyTasks = taskMapper.findReadyTasks(); for (Task task : readyTasks) { if (!isParentSuccess(task.getParentTaskId())) { continue; } task.setStatus("running"); taskMapper.updateStatus(task); ExecutorService executor = Executors.newFixedThreadPool(4); executor.submit(() -> { try { runTrans(new File(task.getKtrPath()), task.getParams()); task.setStatus("success"); } catch (Exception e) { task.setStatus("failed"); task.setErrorMessage(e.getMessage()); } }); } }这段代码里要注意的是线程池和状态更新。跑批任务不能串行执行,一个是慢,一个是某个任务卡住会把整条链堵死,所以要用独立线程池去并发放飞。而状态的更新必须放在真正的执行前后,不能在放线程池前就置为success。这个代码里Executors.newFixedThreadPool(4)是一个简化示范,线上一般用带队列和拒绝策略的ThreadPoolExecutor,不然任务量一上来内存直接被打满。
4.3 运行状态监控与失败告警
平台做到这个阶段,如果只是能看到“跑完/跑挂”两个状态,离好用还差很远。我一般会在执行链路里埋三层状态信息:任务级状态、转换步骤级状态、每批次的数据行数。任务级状态就是刚才代码里的success/failed;步骤级状态需要监听Kettle每个步骤的StepPerformanceSnapShot,这个对象记录了每个步骤的输入行数、输出行数、吞吐量;批次行数则是业务器记录的,用于对账。
这三层状态落库后,前端可以展示一张任务实例详情页,列出每个步骤的耗时、输入输出行数、错误数。还有一个容易被忽视的功能是失败重试。Kettle作业本身跑挂,很多情况是数据库连接闪断、目标表锁冲突这类瞬时故障。我习惯给任务实例加一个retry_count,失败后先自动重试两次,第二次仍然失败才置为failed并发告警。重试间隔用指数退避,第一次30秒、第二次60秒,这样比立即重试的成功率高很多。
告警这块,不要把告警直接绑死在钉钉或微信上,平台内部定义一个告警接口,邮件、企业微信、钉钉各自实现,让用户在配置任务时自己选。告警内容至少包含任务名、批次号、失败原因摘要、日志文件链接。没有这个链接,运维每天处理告警一半时间在翻日志目录找对应文件和批次。
5. Kettle Web化避坑与排查:从驱动失效到内存溢出的5条实战记录
5.1 Kettle环境重复初始化导致性能退化
现象:平台运行一段时间后,每个任务执行越来越慢,日志里能看到大量类加载警告,个别任务偶发KettleEnvironment状态错乱。原因:内嵌API调用KettleEnvironment.init()时,如果每次提交任务都初始化一次,Kettle会反复加载插件和类定义,不仅慢,还容易在并发时出现未定义状态。解决:把KettleEnvironment.init()放到Spring Boot的启动Bean里,只初始化一次,所有任务共用引擎实例。
5.2 ojdbc与MySQL驱动缺失导致数据库步骤全部连接失败
现象:本地Spoon能跑的ktr,Web平台一执行就报Driver class not found。原因:Spoon安装目录里自带驱动,而Web工程要自己引入JDBC驱动,这个最容易忘;另一个坑是Kettle 9.x默认用com.mysql.cj.jdbc.Driver,平台里还是老驱动包就只能报ClassNotFound。解决:在pom中显式添加对应的驱动依赖,同时确认Kettle版本对驱动类名的要求;数据库连接串里的useSSL=false等参数也要按目标库实际环境配好,不然部分驱动在所有连接上都报SSL警告并被Kettle当作错误处理。
5.3 ktr里硬编码绝对路径,换环境必翻车
现象:ktr用Spoon保存后,很多步骤会把文件路径、驱动路径写成绝对路径,Web平台迁移服务器或部署到测试环境,任务大面积失败。原因:Kettle的设计者没有强制使用变量,Spoon在保存时把用户原样输入的路径写进XML。解决:平台保存ktr前做一次路径规范化,把常见的/data/...、D:\...替换成${BASE_DIR}/...,在执行时集中注入变量。转义坑在于路径中可能含中文和空格,ktr是XML格式,写路径时要先做XML转义,否则加载ktr直接抛解析异常。
5.4 大转换并发执行把Web服务JVM挤爆
现象:Web平台上线初期正常,某天多个10分钟级大转换同时触发,Web服务响应变慢,最终整机OOM。原因:内嵌模式转换任务与Web服务共享堆内存,大批量并发没有限流。解决:在调度层加并发信号量,最大同时运行任务数按堆内存设置,建议初始值取内存GB数的一半,比如8G内存最多同时跑4个任务。超过阈值的新任务进入等待队列,而不是直接拒绝。另一个缓解手段是给大转换单独配置Carte执行机,让大任务默认走Carte通道,平台执行器适配层做路由。
5.5 参数和文件里的中文乱码
现象:链路中有中文表名、中文数据源或输出文件的Scenario,在Web页面正常显示,执行时全部变问号,表都找不到。原因:ktr文件保存时用了本地默认字符集,或Kettle在不同环节切换了文件编码。解决:生成ktr时强制使用UTF-8写文件,Transformer设置OutputKeys.ENCODING为UTF-8;执行时给JVM指定-Dfile.encoding=UTF-8;数据库连接串里也显式定义字符编码参数。这个坑在Windows开发机上报得最多,Linux服务器上反而少,做一个编码统一规范能省不少血泪时间。
6. 进阶:三个让平台更耐用的技巧与验证方法
平台跑通后,我一般会立刻做三件事:变量化改造、执行能力扩展、数据一致性验证。
变量化改造是把所有环境相关的配置从ktr中剥离开。ktr里只留${CONN_MYSQL}、${BASE_DIR}这类变量,平台将每种环境的连接配置存在一张配置表里,任务下发时自动注入。好处是测试、预发、生产三套环境的ktr完全一致,平台发布时不用逐个改任务。这个习惯是我在一个客户环境里连续改了三轮连接配置后才养成的,现在所有新任务默认变量化。
执行能力扩展是引入Carte集群。内嵌引擎适合几十条任务的场景,后面任务量和数据量上来,单JVM终究会撞到天花板。做法是在平台的任务配置里增加执行组概念,把任务分组绑定到不同Carte节点,平台通过Carte HTTP接口提交ktr并轮询状态。这样既保留内嵌的轻量,又能在高峰期把压力甩到独立机器上。
数据一致性验证,是每个平台上线前必做的功课。我会构造一组脏数据,包含超长字段、空指针、重复主键等边界情况,用Web平台跑一条全链路转换,同时写一个比对脚本统计源表和目标表的行数以及校验字段的哈希,确认转换既不丢行也不重行。大多数平台“翻车”都发生在没有这层验证直接接业务需求的时候。
我现在做这类平台,默认会先花两天时间把执行状态日志和变量化打好地基,再往后加新功能会顺手很多。毕竟平台的核心是把Kettle这件事做得让人愿意天天用,而不是让它成为一个黑匣子,希望帮到你。
本文还有配套的精品资源,点击获取