简介面向需要处理复杂 ETL 流程的数据工程师或 Kettle 学习者这份资源聚焦“循环获取结果集并传入下一个转换”的典型场景先由第一个转换得到结果集再由作业中的 JavaScript 步骤调用获取上一结果集的方法得到数据借助作业变量把表名、行数、编号、名称等写入作业环境并在循环控制变量的驱动下逐条传给第二个转换最终输出到本地文本文件。资料不仅展示了关键脚本代码和“检验字段的值”步骤还说明了两个转换的配置方法、变量传递关系和输出效果能帮助读者理解 Kettle 中结果集、作业变量与循环三者的配合方式。资源包为 1 个 PDF 文档大小约 130KB内容紧凑、便于速查。目前已有 3113 人学习浏览示例由 KevinYang-凯 提供读者可直接复用其中的循环代码和转换配置思路也可作为排查类似场景问题的参考。无论是完整复现示例流程还是在现有作业中加入循环逻辑这份资料都能帮助节省排错时间。1. Kettle循环获取结果集中的数据并传入转换这个需求比你想象的常见做ETL的人迟早会撞上一个需求先把一批主数据查出来再一条一条地传给后面的任务处理。比如批量重算1000张订单的金额、逐个门店刷新月度报表、按客户ID调一遍外部接口并回写结果。KettlePDI天然是流式处理的一百行数据会一起流过每个步骤可业务要的是一行一个独立任务。这个标题讲的就是Kettle里最常见的解法用结果集Result Set把查询结果暂存在内存里作业层用从结果获取行一条条取出转成变量传给转换再通过循环连接一直跑到结果集清空。适合正在用Kettle做批量数据抽取、报表批处理又发愁怎么让每行数据各自跑一遍的人。2. 先把结果集产出表输入 复制行到结果的配合2.1 结果集到底是什么内存行集快照与两道边界Kettle里的结果集不是数据库游标它是一份存在JVM堆内存里的行集快照。转换里执行完查询后通过复制行到结果Copy rows to result / Rows to result把当前数据流里的每一行按字段原样追加到这份内存集合中供同一个作业后续步骤读取。它有两道边界第一生命周期只在当前作业内部作业跑完就释放不能跨作业调用第二大小受JVM最大堆内存限制几十万行宽表数据就有可能把堆打爆。理解这两道边界很重要因为很多新手把结果集当成表来用——想对它做关联、筛选、排序甚至期待它像Oracle的游标一样有懒加载。实际上Kettle的结果集只做一件事在作业步骤之间搬运整批数据。取出一行后这一行就会从结果集里移除靠这个边取边删的机制来支撑循环的终止判断。2.2 最小可跑通的例子一个作业加两个转换先搭一个作业里面放开始步骤、一个查数据的转换、一个处理数据的转换。查数据的转换里放表输入和复制行到结果两个步骤。表输入的SQL按你的业务写比如把待处理的订单号、门店编号、状态码查出来SELECT order_id, store_code, order_status, total_amount FROM t_order WHERE order_status PENDING AND order_date DATE(now()) - INTERVAL 1 day ORDER BY order_id LIMIT 500;这段SQL把今天待处理的500张订单主键查出来字段包括订单号、门店编码、状态、金额。注意LIMIT的作用先控制结果集规模避免后面循环时一次性占满内存。字段顺序本身没有特殊要求但字段名会被后续从结果获取行步骤当成变量名所以尽量避免用中文列名和特殊字符。复制行到结果不需要额外配置只要上游表输入有行输出它就把这些行全部写入内存结果集。跑完这个转换结果集里就有了数据接下来交给作业层的取行步骤。作业跑起来一般走命令行这样能看到完整日志/opt/pdi/data-integration/kitchen.sh -file/opt/etl/jobs/process_order.kjb -levelBasic -logfile/opt/etl/logs/process_order.logkitchen.sh是作业执行入口-file指定作业文件.kjb-level控制日志详细程度调试时我会用Detailed或Debug看得更细-logfile把日志写到文件不然只打在控制台。如果你只想跑单个转换用pan.sh替换kitchen.sh即可。调参时也可以在命令后追加 -param:ORDER_STATUSPENDING 这类键值对给作业注入外部命名参数。2.3 作业里把行变成变量从结果获取行的配置与连接从结果获取行Get rows from result是作业专用步骤。它每次执行从结果集FIFO弹出第一行把这一行的每个字段自动转成作业变量变量名就是字段名值就是字段值。这样后续步骤里直接用 ${order_id} 这种占位符就能拿到当前行的数据。接线是关键从从结果获取行引出三条连接。成功分支接到要循环执行的目标转换上表示还有行可取开始干活循环分支从目标转换接回来回到同一个从结果获取行继续取下一行失败分支接到空操作或退出作业上表示结果集没有行了整个批跑完。我见过很多人把失败分支留空结果作业取完所有行后以红色失败状态结束调度系统就误报错。正确做法是把失败分支接到一个空操作步骤或者直接连成功出口让作业正常绿色收尾。另外成功分支和循环分支不要同时接到同一个步骤上否则会出现取一行跑两遍的重复执行。3. 循环驱动作业级逐行循环和转换内行集循环两条路3.1 作业级循环最贴合这个标题的驱动方式作业级循环的整体流程是从结果获取行弹出一条数据 → 字段变成变量 → 目标转换用变量执行一次 → 返回后通过循环连接回到取行步骤 → 再弹下一条。结果集每取一次就少一行直到取空循环自然收尾。这就是Kettle里实现for-each的常规做法。为什么优先选作业级循环因为每个子转换都是独立执行空间出错时可以单独重跑日志清晰适合循环体很重的场景——比如每行记录要生成一个文件、调一次外部接口、跑一段复杂计算。缺点是每次启动子转换有固定开销转换的加载、初始化、连接池建立都要时间几万行级别的循环会明显变慢这个后面第四章再展开。我一般会把循环体拆成单独的转换文件命名带 clear 语义比如 process_order_one.ktr这样在作业里一眼就能看出这是被循环调用的单元。转换内部不依赖上一轮循环的残留状态每次进去都靠变量重新定位数据这是保证循环可重复执行的前提。3.2 转换内行集循环轻量、同进程、适合流水线处理如果循环体很轻或者你不想为一个简单操作单独建一个转换文件可以用转换内部的从结果获取行Rows from result步骤。做法是在同一个转换里先用复制行到结果把主数据写入结果集然后用从结果获取行把行读出来再进入后续处理步骤。要注意的是一个转换内的步骤是并行流式执行的不会自己等自己。所以转换内做循环通常要配合阻塞直到步骤完成Block until steps finish来强制串行先让主查询全部落入结果集再用从结果获取行交给处理链。这样做的优势是不切换转换上下文处理速度快劣势是变量作用域有限且出错时不容易定位到具体是第几行数据出的问题。还有一类更轻的替代用映射子转换步骤。映射可以把当前行的数据推给一个内嵌子转换处理子转换跑完把结果返回主流程。它本质上是行流处理不产生作业级循环适合每行做一次字段加工、查一次字典表这种轻量逻辑。如果你的循环体只是简单变换优先考虑映射而不是作业级循环。3.3 两种循环的选型对比维度作业级循环转换内行集循环实现层级作业Job转换Transformation每次循环开销高需加载子转换低同进程内流转变量传参通过字段变量传递直观依赖行集输入不转变量出错定位日志按子转换轮次分开难定位第几行出问题适合规模几十到几千行几千到几万行典型场景每行调接口、生成文件字典表关联、字段重算选型时先问一个问题每行数据的处理是不是独立事务如果是作业级循环更安全一个子转换失败不会污染其他行如果只是纯粹的数据变换转换内行集循环效率更高。这个决定不要等到写完才改因为两种结构的作业布局和日志排查方式完全不同。4. 传给转换之后怎么接、怎么查变量作用域与调试三连4.1 三种传参方式选哪种变量、命名参数、直接行集从结果获取行自动生成的字段变量作用域覆盖当前作业和它直接调用的转换。子转换里用 ${order_id} 就能取到值这是最省事的方式。但变量有一个硬约束它只能存字符串。数字、日期、大金额经过变量传递后类型信息会丢失可能造成精度问题——比如把13位时间戳当字符串传给SQL数据库隐式转换后对不上索引。命名参数是更可控的替代。在作业里的转换步骤上配置参数映射把变量值显式传给子转换的正式参数名子转换里用获取变量步骤接收。这样做的好处是参数名录清晰子转换可以独立测试不依赖作业里恰好存在某个变量。缺点是配置步骤多一步小项目里容易嫌麻烦。第三类是直接行集传参子转换入口放一个从结果获取行步骤读取作业传入的行集数据让数据以行流方式进入子转换根本不经过变量。这样能保留字段原始类型也支持一次传入多行做批处理。但当子转换里需要反复引用当前行的值时行流不如变量直观。4.2 表输入里用变量SQL占位与变量替换开关子转换的表输入步骤里SQL可以这样写SELECT detail_id, product_code, quantity, price FROM t_order_detail WHERE order_id ${order_id} AND order_status ${order_status};${order_id} 是数值不用加单引号${order_status} 是字符串必须用单引号包住否则SQL语法错误。这里有个常见坑表输入步骤默认不会自动替换变量必须打开替换SQL中的变量选项不然Kettle会把 ${order_id} 当字面值发给数据库然后报未找到列名之类的错误。如果子转换是通过命名参数接值的表输入SQL里写的是 ${参数名} 前提是子转换的属性里定义了同名参数并且勾选了使用命名参数。参数名、变量名、字段名三者不一致是最容易出问题的位置所以我建议从SQL的列名开始就统一用小写蛇形命名order_id、store_code、total_amount一路保持一致到变量和参数。4.3 调试三连写日志、看轮次、数行数循环场景下最常见的困惑是到底跑了几行、每次传的值对不对。靠肉眼在Spoon里看是不可能的因为作业跑起来很快而且Spoon在后台执行循环时界面不一定实时刷新。我会在目标转换入口放一个写日志步骤把关键字段打印出来2025/01/12 10:23:45 - Write to Log.0 - order_id 10086, order_status PENDING, store_code SH001写日志步骤的字段列表里选好要打印的字段日志级别选Basic就能在控制台看到。跑完以后用日志文件里打印的最后一个订单号比对结果集总行数就能判断循环是不是完整跑完了。脚本里可以顺手加一个行数统计步骤在查数据的转换里用计数或写日志输出结果集总行数两边一对清不清楚一目了然。我还习惯在每个轮次的日志里加一个轮次标记做法是在结果集的查询SQL里直接生成一个序号列ROW_NUMBER() OVER (ORDER BY order_id) AS loop_seq然后写日志时把它一起打出来。这样万一某条数据出错日志里直接看到第107轮跑了单号10086能省掉大量翻日志的时间。5. Kettle结果集循环的五个翻车现场与排查清单5.1 作业只跑了一次就停了结果集是空的现象作业启动后目标转换只执行了一次就直接结束日志里没有报错但也没有继续跑第二轮。原因通常有三个表输入SQL没查出数据查数据的转换里忘了放复制行到结果作业里从结果获取行的连接没接对取空行时走了失败分支。排查顺序建议先看数据库里SQL单独执行有没有结果再看转换有没有正确输出结果集最后检查作业的接线。解决方法是给查数据的转换加一个写日志步骤先把查到多少行打出来如果0行问题就出在SQL上不要先去翻作业。5.2 循环停不下来结果集越循环越大或者重复执行现象日志显示目标转换的执行次数远超结果集应有行数或者作业内存持续上涨跑了很久不结束。原因最常见的是接线错误——成功分支和循环分支都连回了从结果获取行导致取一行数据被处理两次另一个隐蔽原因是目标转换内部也放了复制行到结果每跑一次循环就向结果集追加一批新行结果集永远清不完形成指数膨胀。解决方法是先数作业里往结果集写数据的步骤有几个确认只有前导查询转换里有然后把从结果获取行的连接整理成一条取数循环、一条成功执行、一条失败退出不要绕圈。5.3 变量没传进子转换SQL里一堆空值或原样字符串现象目标转换的表输入执行时日志里显示SQL语句的占位符没有被替换直接以 ${order_id} 的字面形式发给了数据库或者查出来的数据是0行。原因表输入步骤没勾替换SQL中的变量变量名大小写不一致比如SQL输出列叫 ORDER_ID作业里却写 ${order_id}又或者子转换是通过命名参数接值但参数没映射。解决办法是先打开表输入步骤确认替换SQL中的变量勾上了再核对字段名和变量名Kettle的变量替换是大小写敏感的如果是命名参数回到作业的转换步骤属性里把参数映射补上。5.4 大结果集把堆内存打爆OOM和越来越慢的执行现象作业跑到中途报 java.lang.OutOfMemoryError: Java heap space或者明显感觉越到后面循环越慢。原因结果集全量载入JVM内存如果每行数据还带大字段五万行就可能撑爆默认堆大小。解决方向分两层。第一层是调大堆内存在启动脚本之前设置环境变量export PENTAHO_JAVA_OPTIONS-Xms512m -Xmx4g -XX:MaxMetaspaceSize512m-Xms是启动时初始堆大小-Xmx是最大堆大小循环场景直接给到4G以上比较稳妥。Windows环境在spoon.bat或kitchen.bat里用 set PENTAHO_JAVA_OPTIONS... 设置。第二层是改查询策略不要让一次查询把全量主键装进结果集而是按日期、按ID区间切片每次只查1000条处理完再查下一段这样结果集永远保持在小规模。5.5 循环体里的数据库连接反复重建连接池挤兑现象循环里每次执行目标转换日志都出现数据库连接建立、释放的反复记录有时还报连接超时或获取连接失败。原因每个子转换实例独立管理自己的连接循环次数多时连接池被频繁打开关闭数据库端压力很大。解决建议是把数据库连接配置做成共享连接在PDI里多个转换引用同一个连接名称让连接信息统一管理同时尽量在子转换里只保留必要的数据操作不要在循环体里反复做连接初始化类的步骤。如果循环量级已经大到要考虑连接性能不如回到第四章的想法改成批处理。6. 让循环跑得更快的三个进阶技巧批处理、并行和别用循环6.1 技巧一把逐行循环改成行批循环逐行循环最大的浪费在于每次循环都要重新加载子转换。如果结果集里先按批号分好组每轮循环处理一批效率会高很多。做法是在查数据的转换里用 GROUP BY 生成批号比如按订单日期、按门店ID、按100条一个分组循环变量变成批号子转换SQL里用 WHERE batch_id ${batch_id} 拉出整批数据再流式处理SELECT order_id, store_code, total_amount FROM t_order WHERE order_status PENDING AND batch_id ${batch_id} ORDER BY order_id;这样循环轮次从每行一次变成每批一次子转换内部恢复成Kettle擅长的流式多行处理吞吐量能上一个量级。6.2 技巧二循环不一定要串行同一结果集拆给多任务并行如果循环体是调用外部接口接口响应时间长串行循环会很吃亏。可以考虑把结果集在作业里复制出多份分别驱动几条并行处理链每条链各自从结果集取行。这样做的副作用是结果集的并发读取需要Kettle内部协调行不会重复但调试复杂。我更常用的做法是把结果集先导出到一张临时表再按模数分片作业里用两个并行分支分别处理偶数和奇数ID互不干扰。分片条件写进SQL的WHERE里谁都不会碰到对方的数据比直接在结果集上做并发更可控。6.3 技巧三百万行级别的循环是反模式直接用批量SQL替代这个技巧放在最后是因为它最反直觉Kettle结果集循环适合几千行勉强上万行超过十万行就该停下来想想是否真的需要逐行处理。如果是逐行重算金额、逐行更新状态这类逻辑大多能用一条UPDATE联查完成。我曾经用作业级循环逐行重算十万张订单的状态和金额跑了一个通宵后来改成先按订单日期把大表切成十个分片每个分片内用一条SQL做汇总更新再跑一次行批循环做校验二十分钟就收了。循环的价值是保留每行独立处理的灵活性而不是制造一套慢系统。所以我的习惯是接到循环处理的需求先问两个问题——每个处理结果能不能用一条批量SQL表达如果能就绝不上循环如果不能再看是适合作业级循环还是转换内行集循环最后才考虑并行和批处理优化。这个顺序帮我避开过不少性能翻车也减少了半夜被作业失败短信吵醒的次数。希望帮到你。本文还有配套的精品资源点击获取
阅读完成 · 觉得有帮助?