ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Kettle数据传递实战:变量、结果集与JSON解析全攻略

Kettle数据传递实战:变量、结果集与JSON解析全攻略 你有没有遇到过这种情况一个Kettle转换跑完日志里明明查到了数据可下一个步骤就是用不上或者变量设置好之后SQL里就是替换不进去再或者从接口拿回来的JSON里面字段看得见摸不着解析出来全是空。这三个问题本质上都指向同一件事——Kettle里数据的传递方式没有搞透。变量、结果集、JSON这三样东西贯穿了ETL开发的全过程也是排查效率的分水岭。我这几年用Kettle做数据同步、接口对接、多表抽取绝大多数报错和“数据丢失”都出在这些环节。这篇就把我的实操经验和踩坑过程完整写出来。1. Kettle 数据传递的三张网变量、结果集、数据流字段1.1 先理清一个容易被忽略的前提Kettle里的“数据”可以走三条完全不同的通道。第一条是普通的数据流字段也就是Excel输入、表输入这些步骤往后传的行数据这是最直观的。第二条是变量用${变量名}这种形式在SQL、文件路径、作业参数里被引用本质上是内存中的键值对。第三条是结果集它介于两者之间可以承载多行多列的数据块专门用于作业内不同转换之间的接力。很多新手容易把这三者混在一起尤其是“查询到的数据”到底该走哪条通道完全凭感觉。我见过不少同事把多行查询结果硬塞进变量结果只留下最后一行也见过把明明该用结果集传多行的场景用了一堆字符串拼接最后代码乱成一团。1.2 变量本身的“出身”不同获取方式也不同Kettle中的变量来源大概有三类。第一类是在kettle.properties里定义的环境变量以及通过“设置变量”步骤Set Variables动态写入的变量这类变量有作用域的概念可以指定为JVM级、作业级或转换级。第二类是内置变量比如${Internal.Job.Filename.Directory}表示当前作业文件所在目录${Internal.Transformation.Name}表示当前转换名称这类变量直接可用不需要你手动设置。第三类是从结果集或参数中传入的变量通常在作业执行转换时通过“参数”标签页映射进去。获取方式也对应不同。直接在使用处写${变量名}是最常见的如果想在转换里把变量作为一个字段输出可以用“获取变量”步骤Get Variables如果只是想看某个变量当前的值可以临时建一个“生成记录”步骤加一个“获取变量”步骤输出到日志里看。这三种方式覆盖了绝大多数场景。1.3 变量值到底存的是什么“类型”Kettle的变量本质上全是字符串。你在“设置变量”步骤里把数字字段 123 写进变量取出来也是字符串 “123”。这在拼SQL时通常没问题数据库会做隐式转换但在做字符串比较、日期比较时就会出幺蛾子。比如变量里存的是2024-01-05你拿去和数据库的DATETIME字段比较不同数据库的行为可能不一致。我自己习惯在“设置变量”之前用“字段选择”或“字符串操作”先把格式和类型理干净而不是把问题留给SQL去扛。2. 查询结果转变量三条通道的适用边界2.1 “设置变量”步骤是首选但有一个关键前提把SQL查询结果变成变量标准做法是这样先用“表输入”执行SELECT MAX(update_time) AS max_ts FROM orders然后接一个“设置变量”步骤在字段映射里把max_ts这个字段写入变量max_ts_var作用域选“当前作业”。后续转换再写WHERE update_time ${max_ts_var}就能引用。但这里藏着一个巨大的坑如果“表输入”“设置变量”和后续使用变量的步骤在同一个转换里很可能会出现“变量设置了但后面读不到”的情况。原因是Kettle转换的步骤是并行调度执行的不是像我们想象中那样按连线顺序等前一步跑完再跑下一步。“设置变量”步骤还没执行完后面的“表输入”可能已经把SQL跑完了当然取不到新值。我最早遇到这个问题时排查了很久日志打了一遍又一遍后来才意识到是执行模型的问题。解决方案有两种。第一种是把“设置变量”放到一个单独的转换里再用作业把两个转换串联起来作业项是顺序执行的前一个转换跑完变量已经写入JVM级或作业级作用域后一个转换必然能读到。第二种是干脆不用变量用“数据库查询”步骤把查询结果作为字段直接带入数据流。2.2 “获取变量”步骤把变量从“暗处”拉到“明处”“获取变量”步骤的作用是把一个或多个变量转换为输出流中的字段。什么时候必须用它我遇到比较多的是这种场景希望在转换的每一行数据上都带上一个全局变量比如批次号、业务日期。你可以在“获取变量”步骤里指定变量名、输出的字段名和类型然后和主数据流做“记录关联”或“合并记录”让每一行都带上这个变量值。用这个步骤调试也很方便。我曾经在排查问题时直接在“设置变量”后面拖一个“获取变量”步骤再连一个“写日志”步骤把变量值打印出来。注意“获取变量”步骤本身不产生数据如果前面没有输入流要先加一个“生成记录”步骤产生一行数据。2.3 “数据库查询”步骤流式传递查询结果绕开变量如果只是想在当前转换中根据主数据流的某些字段去查数据库并把查询结果追加到行上用“数据库查询”步骤更合理。比如主数据流是订单ID列表你需要查出每个订单对应的客户名称就可以配置查询SQL为SELECT customer_name FROM customers WHERE id ?Kettle会逐行用输入流字段替换?并执行查询把结果字段追加到输出流。这种方式绕开了变量作用域和执行时序的问题在单转换内的数据处理中比“设置变量”更可靠。但要注意它是逐行查询数据量大时性能会受影响。我一般会在查询SQL里尽量命中索引或者改用“流查询”步骤先把整个维表加载到内存再关联。2.4 实操案例增量抽取中的“最终更新时间”变量我做一个订单增量同步时目标表每天要拉取源表当天新增和修改的数据。完整的作业结构是这样的作业开始→ 转换A查询源表SELECT MAX(update_time) AS max_ts FROM orders接到“设置变量”步骤变量名max_ts作用域选“当前作业”→ 转换B“表输入”写SELECT * FROM orders WHERE update_time ${max_ts}后面接“表输出”写入目标库→ 作业结束转换A里的“表输入”只会产生一行一列数据运行很快转换B的SQL里用${max_ts}引用变量。注意Kettle在转换B启动时会先解析SQL中的变量此时作业已经执行完转换A变量肯定存在所以不会出问题。3. 结果集多行多列数据在转换间的接力棒3.1 结果集和变量的本质区别变量适合传单值结果集适合传“一批数据”。比如一个抽取任务需要先查出所有待同步的表名再逐张表执行同步逻辑。表名可能有很多个变量塞不下结果集就是为这种情况设计的。结果集在Kettle里的实现逻辑很简单上游转换用“复制行到结果”步骤把当前数据流的所有行写入内存中的结果区下游转换用“从结果获取行”步骤把这些行取出来作为输入流继续处理。中间不需要什么连接配置作业会把上一个转换产生的结果集自动传给下一个转换。3.2 正确姿势跨转换传递而不是在同一个转换内闭环新手常犯的错误是在同一个转换里把“复制行到结果”和“从结果获取行”连在一起想着“复制进去再取出来”。这虽然能跑通但毫无意义。真正的用法是在两个不同的转换之间接力。比如作业结构这样→ 转换A表输入查出所有需要同步的表名 → “复制行到结果”→ 转换B“从结果获取行”得到表名字段 → 逐表执行同步需要留意的是从“从结果获取行”读出来的也是普通数据流字段不是变量。如果下游SQL要用表名不能直接在“表输入”的SQL里写SELECT * FROM ${table_name}因为Kettle是在SQL解析时替换变量而table_name在这里是字段不是变量。你需要先在“从结果获取行”后面接一个“Java脚本”或“字符串操作”步骤把表名字段拼接成完整SQL字符串再用“执行SQL脚本”步骤去执行。3.3 关于动态表名同步的一个可靠方案我做过一个归档场景源库里有按月份命名的日志表比如logs_202401、logs_202402一直到当前月需要全部同步到目标库。步骤是转换A查SELECT table_name FROM information_schema.tables WHERE table_name LIKE logs_%复制行到结果。转换B从结果获取行得到字段table_name。后接一个“JavaScript 代码”步骤var sql INSERT INTO target.logs SELECT * FROM source. table_name; row.sql_exec sql;再把sql_exec字段传给“执行SQL脚本”步骤选择“从字段读取SQL”。这样就实现了动态表名的逐表同步。很多人卡在表名不能作为绑定参数的问题上用这个思路可以绕开。3.4 结果集传递的边界条件结果集存放在内存中不适合传递超大结果。我曾经一次传过几十万行的结果集作业直接内存溢出。后来改成按批次处理上游转换每次只查一批数据、复制行到结果作业里配合“作业循环”反复执行每批处理完再拉下一批。另外结果集的字段名和下游步骤期望的字段名必须一致。由于结果是存储在内存里的行集字段类型信息有时会在传递过程中丢失或改变我在多个版本里遇到过数字字段变成字符串的情况。稳妥的做法是在“从结果获取行”后面接一个“字段选择”步骤显式指定类型和名称。4. JSON 里的数据怎么变成变量和行4.1 Kettle 处理 JSON 的步骤选型Kettle处理JSON主要用三个步骤“HTTP client”或“REST client”负责发起请求拿到JSON字符串“JSON input”负责按JSONPath提取字段“JSON Path”步骤则适合快速取单个值。如果你是把JSON字符串作为数据流中的一个字段解析用“JSON input”并把源类型选为“字段”即可如果是解析整个文件源类型选“文件”。最新版Kettle中还有一个“JSON Extractor”步骤功能类似但路径配置方式略有差异。我自己通常用“JSON input”因为它支持一次配置多个输出字段而且对数组展开的支持比较直观。4.2 JSONPath 的选取逻辑和常见写法JSONPath 是定位JSON字段的查询语言类似SQL里的WHERE条件。$ 表示根节点比如接口返回{ status: success, data: { total: 2, items: [ {id: 1, name: 苹果}, {id: 2, name: 香蕉} ] } }要取total字段路径写$.data.total要展开items数组路径写$.data.items[*]Kettle会为数组里的每个元素生成一行输出。如果想同时取数组里的id和数组外的status在“JSON input”的字段配置里加两行路径即可一行$.data.items[*].id一行$.status输出行数会以数组元素数量为准。4.3 从接口JSON中取数并写入变量的完整链路我的一个项目里需要调用内部接口获取某个业务的最新状态码然后根据状态码决定后续同步逻辑。作业结构是转换A用“HTTP client”请求接口响应存在字段body里。接“JSON input”配置一个源字段为body路径为$.status输出字段status。再接“设置变量”步骤把status字段写入变量latest_status作用域选“当前作业”。转换B根据${latest_status}的值走不同分支。比如在“表输入”里写SELECT * FROM sync_config WHERE status_code ${latest_status}。这个链路的关键点在于HTTP请求和JSON解析在同一个转换内通过数据流字段传递不存在时序问题而变量的写入和后续使用拆分到两个转换通过作业串联确保变量已生效。4.4 一个容易踩的坑HTTP响应乱码Kettle的“HTTP client”步骤在部分版本中默认按ISO-8859-1解码响应体如果接口返回的是UTF-8编码的中文JSON你会在后面的“JSON input”里看到中文乱码JSONPath虽然能定位到字段但值已经是乱码字符串。解决办法是在“HTTP client”步骤的配置里找到“编码”或“头信息”相关设置显式指定UTF-8。如果版本里没有这个选项可以在HTTP响应后用“字符串操作”步骤做转码或者在“JavaScript 代码”步骤里手动处理var byteArray org.apache.commons.codec.binary.StringUtils.getBytesUtf8(body); var fixedBody new java.lang.String(byteArray, UTF-8);这种方式比较麻烦所以我一般建议第一步先确认HTTP client的响应编码而不是直接怀疑JSONPath写错了。4.5 JSON数组展开后和上级字段的关联实际业务中经常遇到“数组元素要带上外层公共字段”的需求。比如上面的JSON每个items元素都要带上total的值。按照4.2的写法把$.data.total和$.data.items[*].id配置在同一个“JSON input”里Kettle输出时每一行都会重复带上 total 的值符合预期。需要注意的一点是如果外层字段和数组字段的路径包含不同的层级嵌套某些版本的Kettle在展开时会把外层字段重复填充这没有问题但如果你配置了一个路径本身就指向数组整体而不是数组元素比如$.data.items那么输出的可能是一个JSON数组字符串而不是展开的多行。这是新手容易困惑的地方看到“输出一个字段但内容是一长串JSON”别急着改代码先检查路径末尾是否忘了加[*]。5. 变量、结果集、JSON的坑位排查手册5.1 变量在SQL中替换失败先查这三个位置第一变量名大小写。${max_ts}和${MAX_TS}在Kettle中是两个完全不同的变量设置时写了小写引用时用大写就取不到。第二作用域。在转换里设置的变量默认只在当前转换内有效如果要在作业的其他转换里用必须把“有效范围”改成“当前作业”或更高层级。第三值本身是否为空。如果“表输入”查询结果为空“设置变量”步骤不会执行变量根本不存在SQL引用了不存在的变量会报错或替换为空字符串。我现在的习惯是在“设置变量”之前加一个“空操作”判断如果查询结果为空用“生成记录”步骤造一行默认值保证变量始终有值。5.2 结果集“只有一行”还是“一行都没有”“从结果获取行”读不到数据最常见的原因是上游转换没有真正输出行。比如“复制行到结果”前面连的是“表输入”但表输入因为变量替换失败查出来是空集结果集自然也是空的。排查时先在上游“复制行到结果”之前加一个“写日志”步骤确认有没有行经过。另一个场景是同一个结果集被多个转换读取第二个转换发现读不到。这跟Kettle结果集的消费机制有关——结果集在作业项间传递时每个转换开始时读取一次。如果一个作业项已经消费了结果集并结束后面的作业项再读时结果集可能已经清空。解决办法是尽量让结果集只被一个下游转换消费确实需要多个下游转换的考虑把结果集先写到一个临时表再让多个转换各自读取临时表。5.3 JSONPath 明明没错字段还是空如果JSONPath已经在在线工具里验证过没问题但Kettle解析出来是空我一般按顺序排查四件事源字段是不是字符串类型。如果“HTTP client”的输出字段在元数据里被推断成了字节数组或其他类型“JSON input”可能不认识。可在中间加“字段选择”步骤把它转成String。是不是把文件路径和字段内容搞混了。源类型选“文件”时Kettle会把文件内容当作JSON解析选“字段”时是从某个字段的值里解析。这个选错解析结果往往是空或报错。是不是有BOM头。UTF-8带BOM的JSON字符串在Kettle某些版本里解析第一个字段时会失败因为$前面的BOM字符干扰了JSONPath匹配。可以在前面用“字符串操作”步骤替换掉BOM。是不是大小写和空格问题。很多接口返回的字段名带大小写$.Data.items和$.data.items不同JSON里字段名前后的空格也可能导致匹配失败。5.4 增量同步中取不到最大时间戳的经典问题有人用SELECT MAX(update_time) FROM orders查最大时间戳设置变量然后增量抽取结果发现每次空跑或者重复跑。排查之后发现多数情况是源表的时间字段根本是字符串类型MAX()按字典序返回的不是最新的时间。比如2024-01-05和2023-12-30按字典序比较结果可能不符合预期。解决方法是先确认字段类型。如果是字符串先用STR_TO_DATE、CAST这类函数转成日期再取最大值或者在SQL里统一用ORDER BY update_time DESC LIMIT 1的方式取。另外还要考虑时区问题如果源库和目标库时区不同同步时间戳也需要做偏移转换否则增量窗口边界会漏数据。6. 写少一点感悟写多一点“别踩”真要说有什么总结我觉得是Kettle本身不难难的是理解它“数据到底在哪”这件事。变量像便利贴写下什么后面照着念结果集像托盘一次能端很多碗JSON字段像箱子里的标签得用正确的钥匙打开。三者各管一摊用对了地方ETL就顺了一大半。我早期做增量同步时被“设置变量”步骤的并行特性坑了一整晚日志刷了几百行也没定位出来。后来把作业结构改成“设置变量的转换”和“使用变量的转换”分开问题当场消失。这种“拆分转换用作业串联”的思路后来几乎用在了所有需要跨步骤传值的场景里。如果你现在正遇到变量取不到、结果集读不到、JSON解析全空的问题别急着翻文档先停下来理一下数据现在在哪条通道里它要去哪条通道中间有没有人把它截住了理清了问题就解决了一半。
返回列表