恒美微站
首页
关于我们
建站服务
主题模板
案例展示
资讯中心
联系我们
基于Kettle引擎的Web拖拽式数据集成平台设计与实现
首页
资讯中心
/
基于Kettle引擎的Web拖拽式数据集成平台设计与实现
基于Kettle引擎的Web拖拽式数据集成平台设计与实现
发布时间:2026/9/10 11:50:47
简介这是一套基于KettlePentaho Data Integration实现的Web版数据集成平台项目包面向数据工程师、数据分析师及有ETL开发需求的开发者提供浏览器端可拖拽的可视化操作界面帮助用户在不编写代码的前提下完成数据采集、转换与加载任务显著降低数据集成门槛。资源共1645个文件约160.93MB其中Java文件912个、Vue组件74个、XML配置123个、Properties配置163个并含42个JavaScript文件与12个Kettle转换定义ktr覆盖后端逻辑、前端交互、元数据配置及ETL任务定义等核心环节。已有726人学习浏览。平台内置数据源与数据集管理、转换和作业设计、执行监控、版本控制、权限角色等模块支持关系型数据库、文件系统、Web服务等数据源接入同时包含Dockerfile、docker-compose及Shell脚本便于本地部署与二次开发。开发者可借此研究Kettle如何被封装为Web服务结合Java与Vue源码理解前后端调用与部署流程最终用户则可直接部署使用获得无需安装桌面客户端的数据集成工具。1. 从桌面 ETL 到 Web 拖拽平台Kettle 是怎么被搬到浏览器里的Kettle 的正式名字叫 Pentaho Data IntegrationPDI绝大多数人用它是通过 Spoon 桌面客户端画转换Transformation和作业Job画完保存成 .ktr 和 .kjb 文件再用 kitchen.sh 或 pan.sh 跑批。这种模式本身没问题但一旦环境里有多个人维护接口、几十张数据表、几套目标库桌面客户端就成了瓶颈谁的电脑上没装 Kettle 谁就干不了活流程版本散落在各自磁盘里一个需求的改动要等负责人把文件拷出来。所以很多团队开始做「Web 版数据集成平台」核心诉求是把 Spoon 的拖拽能力搬进浏览器把 Kettle 的引擎留在后端当执行器。基于 Kettle 实现的 Web 平台工程上最稳的路线是前端提供拖拽画布生成一套中间描述JSON 或直接生成 XML后端把描述解析成 Kettle 引擎里的 TransMeta / JobMeta再调用其执行 API 跑起来。关键点在于 Kettle 的核心引擎kettle-engine本来就是可以脱离 Spoon 独立嵌入的它不依赖图形界面你完全可以在 Java Web 工程里直接 new 一个转换并执行。这篇博客就把这条路线上的几个关键环节拆开讲引擎嵌入、流程描述、拖拽映射、任务调度与运行期调优。适合需要自建数据集成服务、或者想把自己手头那一堆 Kettle 脚本 Web 化的开发团队。2. 引擎先落地把 Kettle 核心嵌入 Web 工程的前置工作要做 Web 版数据集成平台最先要解决的不是前端拖拽而是后端能不能稳定地把 Kettle 引擎跑起来。Kettle 的引擎部分是一个纯 Java 类库集合Spoon 只是它的一个图形客户端壳kitchen 和 pan 也只是命令行调用入口整个执行能力都集中在 kettle-engine-core、kettle-core 几个包里。把这几块依赖引入 Spring Boot Web 工程就等于给你的平台装了一个能解析和执行 ETL 逻辑的引擎。2.1 依赖引入与版本选择Kettle 9.x 的 Web 工程改造基础引入 Kettle 依赖的难点在于它的仓库和坐标比较特殊Maven 中央仓库里没有需要添加 Pentaho 的公共仓库。常见做法是直接用 9.4 或 8.3 版本9.x 之后 pentaho-kettle 的包名和结构基本稳定社区里基于 9.x 做 Web 工程的案例也最多。maven 配置如下repositories repository idpentaho-releases/id urlhttps://repo.hds.com/artifactory/pentaho/url /repository /repositories但这里有个实际会遇到的问题Pentaho 的 repository 有时候访问不稳定而且从 9.1 开始部分构件没有对外发布源码包。更常见的做法是先用 Kettle 安装包里的 lib 目录生成依赖。比如你下载了 pdi-ce-9.4.0.0-343它的 lib 下有全部运行需要的 jar可以按下面方式批量导入本地仓库mvn install:install-file \ -Dfilekettle-core-9.4.0.0-343.jar \ -DgroupIdorg.pentaho.di \ -DartifactIdkettle-core \ -Dversion9.4.0.0-343 \ -Dpackagingjar然后项目里声明dependency groupIdorg.pentaho.di/groupId artifactIdkettle-core/artifactId version9.4.0.0-343/version /dependency dependency groupIdorg.pentaho.di/groupId artifactIdkettle-engine/artifactId version9.4.0.0-343/version /dependency需要说明几点。第一这种方式适合整体嵌入就是把你下载的 Kettle 全家桶 jar 都丢进 WEB-INF/lib 或通过 maven install-file 逐个引入然后手动做依赖仲裁。第二Kettle 9.x 依赖的 commons-、guava、jackson 等版本比较老如果工程本身是 Spring Boot 2.7很可能出现类冲突最常见的就是 slf4j 版本不匹配导致 Kettle 日志完全不输出。第三JDK 版本用 8 或者 11Kettle 9.x 在 JDK 17 下会有反射访问限制加载插件时容易抛 InaccessibleObjectException。和另外一些参考资料里的结论一致的做法是生产环境跑 Web 集成平台最好单独指定一个 JDK 8 的运行时。2.2 初始化环境与插件注册像 Spoon 一样启动引擎Kettle 在运行前必须做两件事初始化 KettleEnvironment以及设置插件注册表。Spoon 之所以能识别各种输入输出步骤是因为它启动时扫描了 plugins 目录而嵌入式模式下这些插件需要显式加载。最标准的初始化代码是import org.pentaho.di.core.KettleEnvironment; import org.pentaho.di.core.plugins.PluginRegistry; import org.pentaho.di.core.plugins.JobTypePluginType; import org.pentaho.di.core.plugins.StepPluginType; import org.pentaho.di.core.plugins.TransformationPluginType; public class KettleEngineInitializer { public static void init() throws Exception { // 完整初始化会读取 kettle.properties、加载数据库驱动、初始化变量等 KettleEnvironment.init(); PluginRegistry registry PluginRegistry.getInstance(); // 注册步骤插件与作业项插件保证后续加载ktr时能识别“表输入”“CSV文件输入”等节点 registry.registerPluginType(StepPluginType.class); registry.registerPluginType(JobTypePluginType.class); registry.registerPluginType(TransformationPluginType.class); } }这段代码里的KettleEnvironment.init()会读取$KETTLE_HOME/.kettle/kettle.properties没有该文件时使用系统默认值。需要注意的是如果你的平台要部署在 Linux 服务器上Kettle 会把~/.kettle目录建在运行用户的主目录下多任务并发执行时所有日志默认写到同一个pdi.log。所以我在做这类 Web 集成平台时会把.kettle目录改到应用配置的固定路径下并在启动时设置export KETTLE_HOME/opt/dataintegration/config export KETTLE_JNDI_ROOT/opt/dataintegration/config/simple-jndi放在 Java 里则更直接——调用KettleEnvironment.init()之前先做System.setProperty(user.home, /opt/dataintegration/home); System.setProperty(KETTLE_HOME, /opt/dataintegration/config);这里的逻辑在于Kettle 会从这两个变量中寻找 jdbc.properties 属性、共享的数据库连接定义和最近使用列表。为每个租户或每个环境指定独立的 KETTLE_HOME可以避免多个部署实例争夺同一个配置文件这在做企业级 web 开发时是一个很隐蔽但很常见的坑。2.3 转换执行的最小路径加载 ktr 并运行引擎初始化之后最核心的一件事是通过Trans类加载一个 ktr 文件并执行。这个能力决定了 Web 平台的执行器形态前端上传或者后端保存一个文件执行器把文件路径交给引擎即可。代码如下import org.pentaho.di.core.exception.KettleException; import org.pentaho.di.core.logging.LoggingObjectType; import org.pentaho.di.core.logging.SimpleLoggingObject; import org.pentaho.di.repository.Repository; import org.pentaho.di.trans.Trans; import org.pentaho.di.trans.TransMeta; public class KettleRunner { public void runTransformation(String ktrPath, String paramKey, String paramValue) throws KettleException { TransMeta meta new TransMeta(ktrPath); Trans trans new Trans(meta); // 给转换传入运行期变量 trans.setVariable(paramKey, paramValue); trans.setVariable(internal.job.run, Y); // 日志级别可以按需设置BASIC / DETAILED / DEBUG / ROWLEVEL SimpleLoggingObject loggingObject new SimpleLoggingObject(web-runner, LoggingObjectType.TRANS, null); trans.setLog(loggingObject); // 异步执行不阻塞当前线程 trans.startThreads(); trans.waitUntilFinished(); if (trans.getErrors() 0) { throw new KettleException(转换执行报错错误数 trans.getErrors()); } } }这里有几个参数值得展开。setVariable设置的变量只能在转换内部通过${var}引用而TransMeta上还有setParameterValue两者作用域不同Parameter 是定义在 ktr 上的具名参数Variable 是 Kettle 级别的变量。按我平时调测的经验后台接口把外部传入的查询条件、表名、时间窗口统一通过 Variable 传入最省事因为不用预先在 ktr 里定义 Parameter。startThreads()是异步执行适合 Web 接口场景如果直接调用名字看起来像阻塞的execute()实际也是异步的所以要配合waitUntilFinished()。还有执行结果判断不能只看exitStatus必须检查trans.getErrors()因为 Kettle 在部分步骤出错时默认是记录错误行而不是立刻抛出异常。下载安装教程里大家看到的 Kitchen 脚本本质也是这一套逻辑只是外壳封装了命令行参数解析。自己嵌入引擎后你比 Kitchen 多出来的控制权在于可以让你每个转换有独立的日志目录、独立的变量作用域以及把执行状态实时回传给前端轮询接口。3. 拖拽画布的前后端契约如何用 JSON 表达一张 Kettle 流程图引擎能跑 ktr 文件之后剩下的核心问题就是前端拖拽画布怎么产生数据。这里有两种路线。第一种是前端直接生成 ktr 的 XML后端原样交给TransMeta加载第二种是前端生成一份中间 JSON后端把它翻译成TransMeta的 Java 对象。两种在真实项目里都有但我更推荐第二条路线原因很实际XML 是 Kettle 的私有格式字段序列、步骤 id 的生成规则藏在实现里前端直接拼容易拼出引擎识别不了的结构通过 JSON 定义自己平台的数据结构后端在做校验、权限控制、参数注入时都有明确的切入点。3.1 流程的中间数据结构设计我用过的方案是把一条数据集成流程定义为Pipeline包含步骤列表与连线列表。每个步骤至少包含步骤类型、步骤唯一 ID、名称、坐标以及该类型特有的配置项。连线则包含来源步骤、目标步骤和字段映射。下面是一个简化但可运行的 JSON 片段描述了一张「两个表输入分别查订单表和用户表经过排序后合并最后写到 CSV」的流程{ pipelineId: order_user_merge_001, name: 订单用户合并导出, steps: [ { id: step_001, type: TableInput, name: 读取订单表, config: { connectionName: mysql-orders, sql: SELECT order_id, user_id, amount FROM orders WHERE create_time ${beginDate}, rowsPerFetch: 5000 } }, { id: step_002, type: TableInput, name: 读取用户表, config: { connectionName: mysql-orders, sql: SELECT user_id, user_name FROM users } }, { id: step_003, type: SortRows, name: 订单表按用户ID排序, config: { fieldName: user_id, ascending: Y } }, { id: step_004, type: SortRows, name: 用户表按用户ID排序, config: { fieldName: user_id, ascending: Y } }, { id: step_005, type: MergeJoin, name: 按用户ID合并, config: { joinType: INNER, keyFields: [user_id] } }, { id: step_006, type: CsvOutput, name: 写出CSV, config: { fileName: /data/output/order_user_${lastRunTime}.csv, delimiter: ,, encoding: UTF-8 } } ], connections: [ { from: step_001, to: step_003 }, { from: step_002, to: step_004 }, { from: step_003, to: step_005 }, { from: step_004, to: step_005 }, { from: step_005, to: step_006 } ] }这段结构里config中大量使用${beginDate}、${lastRunTime}这样的占位符是为了让同一个流程可以被不同定时任务用不同时间参数触发而后端在把 JSON 翻译成TransMeta时只需要把运行参数注入为 Kettle 变量。连线部分需要考虑一个分支规则Kettle 的TransMeta中每一步的输入输出是通过hops关联的多个上游连接到一个下游步骤时下游步骤要正确设置接收的行集。我的做法是翻译时按连线数组的顺序创建 hop字段是否匹配由引擎在运行时报错时给出因为很多字段映射是运行时通过第 N 个字段名对齐的调试时看trans.getErrors()附带的信息即可。3.2 JSON 到 TransMeta 的翻译器设计既然有了 JSON后端就要做一层翻译器。翻译器的主类大概长这样import org.pentaho.di.trans.TransMeta; import org.pentaho.di.trans.step.StepMeta; import org.pentaho.di.trans.step.StepInterface; import org.pentaho.di.trans.step.StepDataInterface; import org.pentaho.di.trans.step.errorhandling.StreamInterface; import org.pentaho.di.core.plugins.PluginRegistry; import org.pentaho.di.core.plugins.StepPluginType; public class PipelineTranslator { private final PluginRegistry registry PluginRegistry.getInstance(); public TransMeta translate(JsonNode pipelineJson) throws Exception { TransMeta transMeta new TransMeta(); transMeta.setName(pipelineJson.get(name).asText()); JsonNode steps pipelineJson.get(steps); for (JsonNode stepNode : steps) { String type stepNode.get(type).asText(); // 通过插件工厂创建步骤实例 StepMeta stepMeta new StepMeta(type, stepNode.get(name).asText()); // 这里要拿到插件对应的 StepDataInterface // 再把 config 中的字段填充进去 transMeta.addStep(stepMeta); } JsonNode hops pipelineJson.get(connections); for (JsonNode hopNode : hops) { StepMeta from transMeta.findStep(hopNode.get(from).asText()); StepMeta to transMeta.findStep(hopNode.get(to).asText()); transMeta.addTransHop(new TransHopMeta(from, to)); } return transMeta; } }这里翻译器最麻烦的不是 step 的创建而是 config 里的字段如何映射到步骤的StepMetaInterface上。常见做法是前端定义的字段名尽量贴近 Kettle 插件 XML 的属性命名比如 TableInput 步骤的 SQL 字段在 ktr 里叫sql连接名在 ktr 里叫connectionSortRows 步骤里字段是fieldName排序方向是ascending。翻译器里按步骤类型写 if-else 分支把 JSON 字段逐个 set 到StepMetaInterface的对应方法上这是最朴素也最简单的做法。相比之下用反射自动注入属性会因为插件内部 setter 命名不统一而频繁踩空。3.3 前端拖拽节点与 Kettle 步骤类型的映射表前端画布上能拖的节点不能是 Kettle 全量几百个步骤那对用户是灾难。我一般只暴露十来个高频节点表输入、表输出、插入更新、更新、删除、CSV 输入、Excel 输入输出、字段选择、排序、去重、字符串操作、过滤记录、合并记录、Switch/Case、脚本组件Java 脚本或 JavaScript。先做到交互范围内的稳定再逐步扩展。下面这个映射表可以当作平台设计初期的前缀模型前端节点名称Kettle 步骤类型需要预置的关键配置表输入TableInput连接名、SQL、每次读取行数表输出TableOutput连接名、目标表、提交频率插入更新InsertUpdate连接名、目标表、更新字段列表范围字段选择SelectValues选择字段列表、移除字段、元数据调整排序SortRows排序字段与升降序去除重复记录UniqueRows去重字段、是否忽略大小写过滤记录FilterRows判断条件和条件逻辑合并记录MergeRows旧数据源、新数据源、匹配关键字Switch/CaseSwitchCase判断字段、对应值与目标步骤映射类型转换Normalise需转换的字段与类型映射表格里的每一行都要对应前端组件的一份配置表单表单字段名直接绑定到 JSON 的 config这样一个节点从画布生成到后端翻译再到引擎执行链路是完整的。给熟手提个醒SelectValues这个步骤物美价廉很多新手做关联前不知道先裁剪字段导致合并记录时两边字段数不齐运行时出现「输入行字段数不一致」错误。在画布设计时给每个连线提供「查看字段流向」的预览能力非常有用这相当于把你平台变成了一个简化版 Spoon。4. 把平台跑起来Spring Boot Web 工程中的常见实现套路拖拽画布产生 JSON后端翻译成 TransMeta这一层跑通后剩下的是所有 Java Web 工程的常规问题接口怎么暴露、任务怎么调度、日志怎么采集、多人同时建流程时 Kettle 引擎会不会冲突。这章讲讲我实际搭这类平台时用的工程结构和代码。4.1 工程结构分层与 REST 接口设计一个可维护的 Web 版数据集成平台建议按下面的分层组织web-data-integration/ ├── controller/ # REST 接口层 │ ├── PipelineController.java │ ├── TaskController.java │ └── LogController.java ├── service/ # 业务逻辑层 │ ├── PipelineTranslateService.java │ ├── TaskExecuteService.java │ └── ScheduleService.java ├── runner/ # Kettle 执行引擎相关 │ ├── KettleEngine.java │ ├── PipelineRunner.java │ └── TaskLogListener.java ├── repository/ # 元数据存储 └── model/控制器只负责参数接收与返回统一结构真正跑 Kettle 的线程要放到一个独立的执行器池中否则接口调用会占用 Tomcat 线程流程跑 5 分钟前端请求阻塞 5 分钟很容易触发网关超时。常见做法是使用ThreadPoolTaskExecutor核心线程数设成与 CPU 核数相关最大线程数看数据库连接池大小而定。核心的控制器代码如下RestController RequestMapping(/api/pipeline) public class PipelineController { private final PipelineTranslateService translateService; private final TaskExecuteService executeService; PostMapping(/run) public ResultString run(RequestBody String pipelineJson, RequestParam(required false) MapString, String vars) { String taskId UUID.randomUUID().toString().replace(-, ); // 异步提交到执行线程池 executeService.submitTask(taskId, pipelineJson, vars); return Result.success(taskId); } GetMapping(/task/{taskId}/status) public ResultTaskStatus status(PathVariable String taskId) { return Result.success(executeService.getTaskStatus(taskId)); } }接口设计上run接口返回 taskId 后立刻结束前端用定时轮询/task/{taskId}/status刷新状态。这里不要用 WebSocket 推送每个日志行调试初期前端根本跟不上日志速度滚动获取最近 200 行反而更实用。TaskStatus 对象里至少要有状态等待中/运行中/成功/失败/取消、当前步骤名、已处理行数、错误计数、总耗时、最近更新日志列表。4.2 执行器的线程池与变量隔离Kettle 引擎对多线程并发的支持并不是说你起十个线程就能同时跑十个转换而不互相影响它的资源和插件注册表是全局的。为了隔离运行环境我给每个任务都做了变量上下文先收集执行参数再构造TransMeta把运行时需要的连接信息、文件路径、临时目录都放到variables中。下面是线程池定义与任务提交的核心代码Configuration public class ExecutorConfig { Bean(pipelineExecutor) public ThreadPoolTaskExecutor pipelineExecutor() { ThreadPoolTaskExecutor executor new ThreadPoolTaskExecutor(); executor.setCorePoolSize(Runtime.getRuntime().availableProcessors() * 2); executor.setMaxPoolSize(20); executor.setQueueCapacity(200); executor.setThreadNamePrefix(pipeline-run-); executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); executor.initialize(); return executor; } }执行服务里的核心方法Service public class TaskExecuteService { Resource(name pipelineExecutor) private ThreadPoolTaskExecutor executor; private final PipelineTranslateService translateService; private final ConcurrentHashMapString, TaskStatus taskStatusMap new ConcurrentHashMap(); public void submitTask(String taskId, String pipelineJson, MapString, String vars) { TaskStatus status new TaskStatus(WAITING); taskStatusMap.put(taskId, status); executor.execute(() - { status.update(RUNNING); try { TransMeta transMeta translateService.translateToTransMeta(pipelineJson); // 注入变量注意这里要逐项设置 vars.forEach((k, v) - transMeta.setVariable(k, v)); transMeta.setVariable(internal.task.id, taskId); Trans trans new Trans(transMeta); trans.addTransListener(new TaskProgressListener(taskId, taskStatusMap)); trans.startThreads(); trans.waitUntilFinished(); if (trans.getErrors() 0) { status.failed(转换错误数 trans.getErrors()); } else { status.success(); } } catch (Exception e) { status.failed(e.getMessage()); } }); } }这些代码看起来简单但踩坑往往在细节处。ConcurrentHashMap保存任务状态在任务量大时内存会积压我一般加一条上限检查超过 500 个运行完的任务就清理最早的状态记录。日志监听器TaskProgressListener需要实现TransListener接口Trans内部有addStepListener、addTransListener两种维度前者能拿到每个步骤的行数和耗时后者能感知整个转换的开始与结束实际日志采集两种都要因为用户看平台上的展示时既想看见整体又想看见单个步骤的瓶颈。4.3 定时调度与平台内置调度器选型做数据集成平台总不能每个任务都靠手动点按钮触发。业界常见方案是集成 XXL-Job 或者 Quartz但如果你只想在一个独立 Web 工程里做成轻量任务调度直接用 Spring 的Scheduled配 cron 也可以。普通做法是在数据库表task_schedule里存任务的cron表达式、流程 JSON 引用、是否启用等由一张任务轮询表维护状态。以下是一个最小定时触发实现CREATE TABLE task_schedule ( id BIGINT PRIMARY KEY AUTO_INCREMENT, pipeline_json TEXT NOT NULL, cron_expr VARCHAR(64) NOT NULL, vars_json TEXT, enabled TINYINT DEFAULT 1, last_run_time DATETIME, next_run_time DATETIME, create_time DATETIME );Component public class ScheduleTrigger { Scheduled(fixedDelay 30000) public void scanDueTasks() { // 每次扫描找出该执行的任务交给执行线程池 ListTaskSchedule dueTasks taskDao.findDueTasks(new Date(), 10); for (TaskSchedule task : dueTasks) { executeService.submitTask(genTaskId(), task.getPipelineJson(), parseVars(task.getVarsJson())); taskDao.updateLastRunTime(task.getId(), new Date()); } } }用Scheduled做调度器的局限是它跑在单机进程内任务多时无法水平扩展cron 更新后要等下一个固定轮询周期才生效没有失败重试和分片。如果你的平台定位到了多团队共用的程度还是建议替换为 XXL-Job 这类独立调度组件把这里的scanDueTasks换成回调接口即可。但作为第一版明确边界之后用这个方案能快速落地。5. 上生产前必须调的三类参数与日志排错技巧最后一章不讲大框架讲几个把平台推到生产环境时一定会碰到的实际问题数据库连接参数、大表抽取的内存控制以及日志到底从哪里看。5.1 数据库连接的三个必调参数Web 平台里每个流程都可能连不同数据库Kettle 的数据库连接全部走DatabaseMeta但 Web 场景下连接池通常不在 Kettle 里配而是在平台层面统一管理。给表输入、表输出配置连接时务必让用户能设置以下三个参数它们对性能影响非常直接参数作用推荐初始值rowsPerFetch表输入每次从数据库游标取多少行5000 到 10000不要超过 50000fetchSizeJDBC 驱动层面的抓取行数MySQL 需要设1000以上否则默认全量进内存commitSize/batchSize表输出按多少行提交一次事务常见 1000 或 5000太大回滚代价高在 JSON 配置里rowsPerFetch直接对应 TableInput 步骤的配置fetchSize需要在数据库连接 URL 上追加参数比如 MySQL 的jdbc:mysql://host:3306/db?useCursorFetchtruefetchSize1000否则 MySQL 驱动会忽略 fetchSize。大表抽取场景我踩过明显的坑不设置useCursorFetchtrue几百万行的表直接把 JVM 堆打满设置后内存占用降到几百 MB。5.2 多输出与合并场景的常见写法一个表输入输出多个 Excel 文件这个热搜场景在 Kettle 里一般用「表输入 字段选择 分组字段的 Switch/Case 两个 Excel 输出」实现Switch/Case 的判断字段作为输出文件名的变量。一个更省内存的做法是表输入读完后用「克隆行」Clone Row复制流再接不同条件的过滤记录与控制输出步骤。这里不展开节点细节但给出一个调优要点这类场景不要把整表查出来再分流应在 SQL 层就按分区条件拆成多条表输入各查各的分区再连接到对应输出。这与业务库的分区裁剪原理一致源库和 Kettle 两侧的消耗同时降下来。合并两张表输出一个 CSV 的常见写法则更直接两个表输入分别做排序接 Merge Join再接 CSV 输出。这也就是第 3 章 JSON 示例里那张流程的由来。它最容易出的问题是两边排序字段的排序规则不一致Merge Join 要求两边排序的 collation 与大小写行为完全一致否则会出现「数据应该合并却各自散开」的假象表现为输出行数等于两边行数和。复现时看一下 Merge Join 步骤的输出行数如果明显异常先检查两边排序步骤配置的字段是否一字不差。5.3 从 Web 平台快速定位 Kettle 错误日志与状态码Web 平台里看错误要比桌面端困难因为你看不到 Spoon 的步骤运行信息面板。我建议平台在建流程时把每个步骤的日志级别默认设为DETAILEDROWLEVEL只在临时排查时开启因为行级日志会记录每一条流水行数一大磁盘就爆。日志采集要做到两个粒度转换级存Trans的日志到独立文件步骤级通过addStepListener把每个步骤的错误数、行数、时间刷到任务状态里。最后一条经验是看到Step was interrupted这类错误不一定代表数据出问题很多时候是中途手动停止了转换或连接被系统回收导致的误报先看周围步骤的错误数再来判断别把平台时间浪费在翻堆栈上。提示对日志文件切分建议在 logback 里用按天滚动保留 30 天并把 Kettle 的日志独立命名为kettle-task.log和 Web 应用业务日志分开。这样排查问题时只看任务日志文件不会有 Spring Boot 每 10 秒一条的空闲连接打印来干扰判断。本文还有配套的精品资源点击获取