NiFi ETL实操:插入JSON至SQL表后获取自增ID插入关联表
NiFi实现自增ID关联插入的解决方案
整体流程调整
原流程:InvokeHTTP(获取API数据) → JoltTransformJson(转换结构) → 新增步骤 → PutSQL(插入category表) → 新增步骤 → PutSQL(插入category_ref表)
具体步骤实现
提取source_id到FlowFile属性
在JoltTransformJson之后添加UpdateAttribute处理器,配置如下:- 添加属性:
source_id,值设为${json-path:evaluate('$.source_id')}
作用:将每条记录的source_id存入FlowFile属性,后续可直接引用,避免数据分离。
- 添加属性:
修改category表插入SQL,返回自增category_id
修改插入category表的PutSQL处理器:- SQL语句改为:
INSERT INTO category (source_name) OUTPUT inserted.category_id VALUES (?) - 设置
Output Result Format为JSON
作用:MSSQL通过OUTPUT inserted.category_id返回刚生成的自增ID,PutSQL会将该结果以JSON格式输出为新的FlowFile,且该FlowFile会继承之前设置的source_id属性。
- SQL语句改为:
合并source_id与category_id为插入结构
添加第二个JoltTransformJson处理器,使用以下Jolt规范:[ { "operation": "modify-overwrite-beta", "spec": { "source_id": "${source_id}", "category_id": "@(0,category_id)" } } ]作用:将FlowFile属性中的source_id与返回的category_id合并为
{"source_id":"xxx","category_id":xxx}的JSON结构,满足category_ref表的插入需求。插入category_ref关联表
添加第二个PutSQL处理器,配置如下:- SQL语句:
INSERT INTO category_ref (source_id, category_id) VALUES (?, ?) - 参数映射:分别将JSON中的
source_id和category_id对应到两个占位符。
- SQL语句:
异常处理建议
- 为两个
PutSQL处理器配置Failure关系,路由到LogAttribute或PutFile处理器记录错误信息,便于排查插入失败的记录。
内容的提问来源于stack exchange,提问作者Duong Dai Tay
相关产品推荐
相关产品推荐

