简介:这是一套基于Kettle(Pentaho 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 Integration(PDI),绝大多数人用它是通过 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> <id>pentaho-releases</id> <url>https://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 \ -Dfile=kettle-core-9.4.0.0-343.jar \ -DgroupId=org.pentaho.di \ -DartifactId=kettle-core \ -Dversion=9.4.0.0-343 \ -Dpackaging=jar然后项目里声明:
<dependency> <groupId>org.pentaho.di</groupId> <artifactId>kettle-core</artifactId> <version>9.4.0.0-343</version> </dependency> <dependency> <groupId>org.pentaho.di</groupId> <artifactId>kettle-engine</artifactId> <version>9.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 或者 11,Kettle 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 里叫connection;SortRows 步骤里字段是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/Case | SwitchCase | 判断字段、对应值与目标步骤映射 |
| 类型转换 | 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 Result<String> run(@RequestBody String pipelineJson, @RequestParam(required = false) Map<String, String> vars) { String taskId = UUID.randomUUID().toString().replace("-", ""); // 异步提交到执行线程池 executeService.submitTask(taskId, pipelineJson, vars); return Result.success(taskId); } @GetMapping("/task/{taskId}/status") public Result<TaskStatus> 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 ConcurrentHashMap<String, TaskStatus> taskStatusMap = new ConcurrentHashMap<>(); public void submitTask(String taskId, String pipelineJson, Map<String, 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() { // 每次扫描找出该执行的任务,交给执行线程池 List<TaskSchedule> 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,不要超过 50000 |
fetchSize | JDBC 驱动层面的抓取行数 | MySQL 需要设1000以上,否则默认全量进内存 |
commitSize/batchSize | 表输出按多少行提交一次事务 | 常见 1000 或 5000,太大回滚代价高 |
在 JSON 配置里,rowsPerFetch直接对应 TableInput 步骤的配置;fetchSize需要在数据库连接 URL 上追加参数,比如 MySQL 的jdbc:mysql://host:3306/db?useCursorFetch=true&fetchSize=1000,否则 MySQL 驱动会忽略 fetchSize。大表抽取场景我踩过明显的坑:不设置useCursorFetch=true,几百万行的表直接把 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 的步骤运行信息面板。我建议平台在建流程时把每个步骤的日志级别默认设为DETAILED,ROWLEVEL只在临时排查时开启,因为行级日志会记录每一条流水,行数一大磁盘就爆。日志采集要做到两个粒度:转换级存Trans的日志到独立文件,步骤级通过addStepListener把每个步骤的错误数、行数、时间刷到任务状态里。最后一条经验是:看到Step was interrupted这类错误不一定代表数据出问题,很多时候是中途手动停止了转换或连接被系统回收导致的误报,先看周围步骤的错误数再来判断,别把平台时间浪费在翻堆栈上。
提示:对日志文件切分,建议在 logback 里用按天滚动保留 30 天,并把 Kettle 的日志独立命名为kettle-task.log,和 Web 应用业务日志分开。这样排查问题时,只看任务日志文件,不会有 Spring Boot 每 10 秒一条的空闲连接打印来干扰判断。
本文还有配套的精品资源,点击获取