简介:这是一份关于Kettle(Pentaho Data Integration)实现结果集循环获取并传递至下一转换的技术文档,面向有ETL开发需求的工程师,重点解决在Job中通过JavaScript循环处理结果集变量、再交由下一转换继续加工的问题。文档以实际项目为背景,完整展示了从作业到转换的完整链路:用previous_result.getRows()获取上一转换结果集,再用parent_job.setVariable()将id、name等字段写入变量,并在转换var.ktr中读取变量,最终输出到本地txt文件。压缩包内仅含1个PDF文件,大小约130KB,内容涵盖t1.ktr与var.ktr的配置说明、关键JavaScript代码及输出结果演示,便于对照实践。目前已有2365人浏览学习,适合初次接触Kettle循环传参的开发者参考。通过阅读这份文档,读者可快速掌握previous_result、setVariable、循环控制变量等核心写法,并将相关思路复用到自己的ETL作业中,减少变量传递错误与排错时间。
1. 循环结果集是什么:一次一个客户ID,喂给下游转换去处理
接到一个很常见的需求:上游给了一张 3000 行的客户清单,下游要照着每个客户 ID 去调明细接口、更新汇总表或者生成单独文件。你的第一反应多半是拖一个“循环”步骤,但 Kettle 转换里根本没有这玩意儿——这是 Kettle 新手最容易翻车的地方。Kettle 的转换是静态数据流,行在步骤之间流动,不会自己回头反复执行;真正能承担“循环获取结果集中的数据并传入转换里面”这个动作的,是结果集机制加上循环执行步骤的组合,把查询出来的多行结果,按行取出来,再以参数形式喂给另一个转换去处理。这篇笔记就是写给要做跑批、增量抽取、动态 ETL 的人,目标是让你照着一套最小可跑的配置,把循环逻辑搭起来,并知道坑在哪。
2. 结果集产生与传递:表输入、复制行到结果、从结果获取行的最小链路
2.1 行集(Row Set)与结果集(Result):作用域不同,这是到处找不到“循环”的根源
先区分两个容易混的概念。行集(Row Set)是转换内部两根连线之间的东西,表输入把 SQL 查出来的数据一条条流到下一个步骤,它只在当前转换里存在,别的转换无法直接消费。结果集(Result)则是挂在 Kettle 作业上下文里的数据块,转换里用“复制行到结果”把行集快照放进去,作业后续的作业项再用“从结果获取行”把它取出来。很多人在转换里翻遍步骤列表找“循环”找不到,就是因为这一步根本不给转换用,它属于作业和跨转换协作的范畴。
结果集本质上把所有行一次性装进内存,字段名、字段类型、行数这些元数据会跟着快照一起保留。这也是为什么“从结果获取行”取出来之后,字段能直接用于下一步比较和参数映射。代价是结果集不能太大,几千行没什么问题,几十万行全量驻留在内存里,即使后面循环逻辑没问题,Java 堆也扛不住。所以我的习惯是:结果集只用来放“要循环的 ID 清单”或“分页批次号”,而不是放明细大表。
2.2 最小链路:表输入 → 复制行到结果 → 从结果获取行
先用一个最简单的例子把链路跑通。假设有一张客户快照表,需要把客户 ID 清单挂到结果集里。
SELECT client_id, client_name FROM ods.client_snapshot WHERE data_date = '${etl_date}' ORDER BY client_id;这是表输入的 SQL。注意最后一行WHERE data_date = '${etl_date}',表输入步骤里必须勾选“替换 SQL 中的变量”,否则 Kettle 会把${etl_date}当普通字符串拼进 SQL,查询直接报错或返回空。这个开关在表输入的“选项”页里,不同版本叫法略有出入,但“变量替换”这个关键词不变。
表输入后面接“复制行到结果”,这个步骤不需要配什么额外字段,默认把输入行的全部字段连同元数据写入作业结果。然后新建一个作业(Job),把包含这个表输入的转换放进去,作业项高级设置里记得勾选“将结果集传递给下一步骤”,这样下一个作业项才能用“从结果获取行”取数据。很多第一次做结果集传递的人漏掉这个勾选,结果就是下游读到的永远是空,日志里还看不到明显报错,属于典型的 Kettle 玄学问题。
2.3 结果集的字段命名与空结果集表现
复制行到结果时,字段名尽量用统一的小写风格。Kettle 内部对大小写敏感,尤其在做变量映射的时候,CLIENT_ID和client_id会被当成两个参数,报错不会报,传过去永远是空值,排查起来极费时间。我在项目里定的规约是表输入 SQL 里就给字段起别名,全部小写下划线,例如SELECT client_id FROM ...,避免后面所有转换去迁就大小写。
空结果集是这个机制里一个埋了很久的坑:当表输入查不到任何行时,结果集本身是空的,但某些版本里“从结果获取行”仍然会输出一行全 NULL 的数据。这会导致下游转换被白白执行一次,生成空文件或者写一条脏数据。解决办法在链路源头加判断:用表输入查一个COUNT(*)判断行数,然后通过“过滤行”把 0 的情况直接短路;或者干脆在 SQL 里用HAVING COUNT(*) > 0包一层,让无数据时不产生结果集。后面避坑章节会再展开。
3. 用 Execute Row 实现循环取数:父转换 → 子转换的参数映射三步配置
3.1 三种常见循环方案的选型:Execute Row、作业跳转、表驱动批量
Kettle 里做“循环获取结果集并传入转换”这件事,常见做法有三种。第一种是用“循环执行行”步骤,英文叫 Execute Row,它是放在父转换里的步骤,每收到一行输入,就调用一次子转换,并把当前行的字段映射成子转换的参数。这是本标题最直接对应的方案,配置清晰,也方便在日志里看每一轮迭代。
第二种是作业跳转回环:在作业里用“评估结果”和跳转箭头组成循环,跑完一轮判断是否结束,没结束就跳回上一步。这个方案在老项目里常看到,但作业画布上的跳转一多就变成蜘蛛网,维护成本高,而且稍不注意就是死循环,作业一直跑不完也不报错。我一般不建议新项目用。
第三种是表驱动批量:压根不循环,把“每次传一个 ID 下去处理”改成“把整个结果集灌进临时表,然后一条 SQL 关联操作”。这其实是最推荐的方式,性能最好,只是标题既然点名要循环,说明业务往往是真的绕不开逐行处理,比如调外部接口,或者每个 ID 要跑一段异构逻辑。我的选型标准是:能写进一条 SQL 的绝不循环;必须循环时,1000 行以内随便写,超过 5000 行就要考虑分批。
3.2 最小可跑配置四步曲:查询转换 → 父子转换 → 参数映射
整体结构由三个转换加一个作业组成:T_Query 负责把结果集写进作业上下文;T_Loop 是父转换,从结果集取行,用 Execute Row 步骤逐行调用子转换;T_Process 是真正干活的子转换,接收参数并完成业务。
第一步,建 T_Query:表输入查客户 ID 清单,接“复制行到结果”。这一步和第 2 章的最小链路一模一样。
第二步,建 T_Loop:第一步放“从结果获取行”,拉到所有行后进入 Execute Row 步骤。Execute Row 的配置里要指定子转换来源,可以指向 T_Process.ktr 文件,也可以从资源库加载。关键在字段映射区,把输入流里的client_id字段映射到子转换的命名参数client_id。
第三步,建 T_Process:打开转换设置,在“参数”页里显式声明命名参数client_id,默认值填-1或者空字符串。这一步很多人直接跳过,结果单独运行 T_Process 时提示参数未定义,挂到 Execute Row 下面时又发现参数传不过去,原因就是没注册。T_Process 里的表输入 SQL 应该长这样:
SELECT order_id, order_amount FROM ods.t_order WHERE client_id = '${client_id}';注意这里${client_id}必须靠表输入的“替换 SQL 中的变量”开关生效,否则查询条件里就是一个字面量${client_id},查出来永远是空。
第四步,建作业 J_Main:依次放置 T_Query 和 T_Loop。T_Query 执行完,作业上下文里有了结果集,T_Loop 取出来逐行循环。这里不需要额外写循环次数,Execute Row 取到几行就循环几次,输入行流结束了自然停止。
Execute Row 三个关键配置项我一般这么设:
| 配置项 | 推荐值 | 说明 |
|---|---|---|
| 子转换来源 | 文件路径或资源库 | 选文件路径便于挪环境,资源库便于团队协作 |
| 字段映射 | client_id → client_id | 映射名必须和子转换参数名完全一致 |
| 错误处理 | 失败继续或中止 | 调外部接口建议“失败继续”,写库建议“中止” |
3.3 用 kitchen 命令行调起整条作业:参数覆盖与日志定位
开发环境用 Spoon 点运行没问题,跑批场景下最终都要落到命令行。Kettle 提供的 kitchen 是作业的执行器,进入 Kettle 安装目录后执行:
./kitchen.sh -file=/opt/etl/j_loop.kjb \ -param:etl_date=2025-06-01 \ -level:Basic \ -logfile=/opt/etl/log/j_loop.log-param:用来传作业级参数,这里的etl_date会覆盖转换里同名变量;-level:Basic是常用日志级别,循环场景下用Basic比Detailed精简很多,否则 Execute Row 每轮迭代的字段值会刷屏,日志文件一天几个 GB。如果只想在命令行直接跑单个转换,用 pan.sh,参数和 kitchen 基本一致。
日志里要确认循环真的在推进,盯两个地方:一是Execute Row步骤的分组日志出现了多次,二是T_Process的日志里打印出的client_id在变化。我在 T_Process 第一步固定放一个“写日志”步骤,输出当前处理的 client_id,循环卡住时看日志就知道卡在第几个 ID 上。
4. 循环体转换内部设计:三个必调参数与动态 SQL 拼接
4.1 命名参数三件套:显式声明、默认值、作用域
子转换接收循环传进来的值,有三个参数上的细节决定你能不能跑得顺。第一是命名参数必须显式声明,只在 SQL 里写${client_id}不算定义,还得在转换设置的“参数”页里注册同名参数,Kettle 才会把它当参数处理。第二是默认值要给出,我在 T_Process 里给client_id默认-1,这样在 Spoon 里单独调试这个转换时不会因为参数缺失直接报错,-1查不到数据正好可以验证 SQL 逻辑。第三是作用域,循环传入的参数如果想在子转换的后续步骤里继续传递,比如再传给嵌套的转换,要把它提升为全局变量或者至少在作业级可见,只停留在“转换内有效”的话,跨转换就取不到了。
| 设置项 | 值 | 作用 |
|---|---|---|
| 参数名 | client_id | 必须与 Execute Row 映射名一致 |
| 默认值 | -1 | 单转换调试兜底 |
| 作用域 | 作业内/全局 | 确保嵌套子转换可读 |
4.2 循环体里的动态 SQL 与文件路径拼接
循环体里最容易出问题的是参数没有真正下推到查询。比如查出 1000 个 ID,每次循环都把整个订单表拉出来再过滤一个 ID,跑得慢不说,数据库压力也大。正确做法是让 SQL 里的 WHERE 条件直接消费参数,也就是第 3 章那个${client_id}的写法,让数据库只返回当前这一行的数据。
有时候参数不只是用来过滤,还要拼文件名、表名,常见做法是在子转换里用 JavaScript 步骤组装变量。比如按客户 ID 生成独立 CSV:
// 从父转换继承的 client_id,拼成输出路径 var clientId = parent_job.getVariable("client_id", "-1"); var filePath = "/data/output/client_" + clientId + ".csv"; setVariable("file_output", filePath, "r");这段脚本的关键在于parent_job.getVariable能拿到作业上下文里的变量,setVariable第三个参数"r"表示作用域为根(全局)。如果直接写getVariable而不带前缀,取到的可能是子转换本地变量,经常是空。文件路径拼好之后,后面的“文本文件输出”步骤文件名字段里写${file_output}即可。要动态切换表名时,同样先把表名设进变量,再在表输入里引用,但表输入对表名的变量替换支持有限,我遇到这种情况一般会改用“动态 SQL 表输入”或者直接写 JavaScript 里执行 JDBC,这属于另一个话题,不展开。
4.3 输出落地:更新表、写文件、调用接口的取舍
循环体干完活总要有出口。写数据库更新,小数据量没问题,但逐行 update 性能很差,我一般会先更新到临时表,循环结束后再一条 SQL 关联刷进目标表,把行循环和写库解耦。写文件的话,注意文件名里带上批次时间戳和 client_id,避免多轮循环互相覆盖。调 REST 接口是循环用得最多的场景,接口返回 JSON 时,直接用 Kettle 的 JSON 步骤解析,不要塞进 JavaScript 里手拆字符串,拆到一半字段顺序变了就翻车。接口调用还要处理失败行和超时,每轮结果先落到内存里,循环结束统一写失败清单,而不是在循环体里一报错就中断整条作业。
5. 五个高频踩坑点:变量不生效、空结果集执行一次、时区与驱动报错
5.1 变量不生效或串值:下一轮循环沿用了上一轮的 client_id
现象:日志里出现字面量${client_id},或者第一轮循环传参正确,从第二轮开始client_id一直停留在上一个值。原因是多方面的:子转换没有注册同名命名参数,Kettle 找不到参数定义时会把${client_id}当普通文本;Execute Row 的字段映射大小写不一致;循环体里某个“设置变量”步骤在上一轮写入了值,而下一轮 Execute Row 没有覆盖它。
解决:先统一字段和参数全小写,并在子转换设置里显式注册参数。同时在 T_Process 第一步加一个“设置变量”步骤,把所有循环参数强制置为默认值,再让 Execute Row 传入真实值。这个手法看起来多余,但在变量作用域混乱的老作业里特别管用。
5.2 结果集为空仍执行一次:日志里多出一条空行
现象:表输入 SQL 查不到任何数据,但 T_Loop 还是调用了一次 T_Process,文件生成了,表里多了一条不该存在的记录。原因是从结果集获取行在空结果集时仍可能输出一行全 NULL 数据,Execute Row 拿到这一行就照常执行。解决思路在源头拦截:T_Query 里加一个行数判断,或者在外层 SQL 里这么包一层:
SELECT client_id, client_name FROM ods.client_snapshot WHERE data_date = '${etl_date}' AND EXISTS (SELECT 1 FROM ods.client_snapshot WHERE data_date = '${etl_date}') ORDER BY client_id;EXISTS子查询有数据时才返回结果集,没数据时整个查询直接返回空,从根上断掉空结果集触发循环的可能。另外,空字符串参数也会踩 Kettle 的空值转换坑:Kettle 默认不会把传入的空串转成null,下游 SQL 里IS NULL判断永远不成立,需要时在子转换里用“字符串操作”步骤把空串处理掉。
5.3 MySQL 时区报错与 Oracle 驱动不匹配
现象:连 MySQL 8 时启动转换直接报The server time zone value '锟斤拷...',这是 MySQL 驱动无法识别系统时区的经典报错,解决方式是在 JDBC 连接串里显式指定时区:
jdbc:mysql://localhost:3306/ods?useSSL=false&serverTimezone=Asia/Shanghai&allowPublicKeyRetrieval=true其中serverTimezone=Asia/Shanghai是必须项,allowPublicKeyRetrieval=true解决的是 MySQL 8 的公钥检索报错。
另一个高频问题是 Oracle 驱动版本对不上,热词里那个ojdbc6.jar 11.2.0.4是老项目里最常见的组合。现象是连接池报ORA-28040: No matching authentication protocol,或者干脆Class not found oracle.jdbc.OracleDriver。检查 Kettle 的 lib 目录里是不是塞了多个 ojdbc 版本,新旧驱动冲突是常见祸根。保留和数据库版本匹配的那一个,把其他的移出目录,重启 Spoon 再试;驱动可用但版本过旧时,改走升级到 ojdbc8 的路线。
5.4 逐行循环越来越慢:Java 堆上涨与日志刷屏
现象:前 200 行挺快,到 500 行以后每轮明显变慢,任务管理器看到 Java 进程内存持续上涨。原因是结果集全量驻留内存,子转换每次循环都新建数据库连接又没有及时释放,以及过高的日志级别把每轮字段全部打了出来。解决优先级从低到高分别是:日志级别降到Minimal或Error;子转换里数据库连接改成连接池模式,关闭自动提交,减少事务开销;确认结果集里只保留必需字段,别把整行大字段全塞进去。到了几千行以上,第 6 章的分批循环改造才是治本方案。
5.5 循环中途报错:不知道错在第几行
现象:作业跑了一半失败,日志定位到子转换内部,但看不到当时处理的是哪个 ID,排查像大海捞针。原因是没有把当前循环行号或主键打出来。解决:T_Process 第一步骤写日志输出client_id,这一步成本极低,但能把“循环到谁挂了”直接钉死在日志里。更稳妥的做法是加一个失败记录表,循环体里用“数据库查询”或“写日志”把当前行记录下来,作业失败后去表里面看最后一条就是断点位置。
6. 进阶改造:把逐行循环换成页面式批量循环
6.1 分批循环改造:状态表 + 分页游标
循环必然逐行的时候,还有一种折中的优化思路:不按行循环,按批循环。比如一次从结果集里取 1 列 ID 清单,但每次循环处理 5000 行,由子转换内部批量处理,而不是一行调一次。这能显著减少 Execute Row 的调用次数和 JDBC 往返。做法是先建一张状态表记录last_id,循环体每次执行前取一批大于last_id的数据,处理完更新last_id,直到取不到新数据为止。
以 Oracle 为例,分页 SQL 长这样:
SELECT * FROM ( SELECT t.pk, t.data FROM ods.t_order t WHERE t.pk > ${last_id} ORDER BY t.pk ) WHERE ROWNUM <= 5000;MySQL 就把尾段换成LIMIT 5000即可。此时循环的粒度从“行”变成“批”,变量传递也从client_id变成last_id,执行次数能降低几个数量级。对比一下两种模式:
| 指标 | 逐行循环 | 每批 5000 行循环 |
|---|---|---|
| 1 万行调用次数 | 1 万次 | 2 次 |
| JDBC 往返 | 1 万次查询 | 2 次批量查询 |
| 耗时(参考) | 60 分钟以上 | 10 分钟以内 |
| 适用场景 | 调外部接口、强逐行逻辑 | 批量计算、批量写临时表 |
我早年做过一个库存同步作业,1 万行逐行更新跑了快两个小时;改成 5000 行一批之后十分钟出头。后来凡是要写循环,我第一反应都是先问一句:能不能不循环?能就用集合操作,必须循环就优先考虑按批循环。
6.2 验证循环数据完整性的三个检查点
改造完循环,验证和原逻辑一致性比跑通更关键。我习惯检查三个点。第一是总数核对:循环前用COUNT(*)记录结果集行数,循环后统计处理成功的行数,两边对不上说明有行被跳过或重复处理。第二是幂等验证:同一批数据重跑一遍,目标表数据不能翻倍;在目标表主键或唯一索引上做约束,比事后写 SQL 去重更早发现问题。第三是断点续跑:状态表里记录last_id和处理状态,作业中途失败重跑时,从状态表继续而不是从头再来,这个能力在跑批场景里几乎是刚需,早加早舒坦。
最后说个经验之谈:Kettle 里凡是名字里带“循环”的配置,先检查它是不是在错误的作用域里干活。变量传不进去、空结果集多执行一次、越跑越慢,这三大经典问题我基本都见过至少一轮。把这几个点摸透了,循环结果集这个能力才真正算是你的,而不是靠运气跑通。希望帮到你。
本文还有配套的精品资源,点击获取