如何通过Apache NiFi更新Redis中指定Key下JSON的指定字段且保留其他字段
如何在Apache NiFi中实现Redis JSON值的部分更新(保留原有字段)
完全可以实现这种需求,核心思路是先获取Redis中key='x'的原有JSON内容,将其与FlowFile中的新JSON做合并(新字段覆盖旧字段,保留未被更新的原有字段),最后将合并后的完整JSON写回Redis。以下是具体实现方案:
方案一:使用NiFi原生处理器组合
步骤1:获取Redis中原有JSON值
使用FetchRedisKey处理器,配置Redis连接信息,指定要获取的Key为x。该处理器会生成一个包含Redis中原有JSON内容的FlowFile。
步骤2:合并两个JSON内容
- 将
FetchRedisKey输出的FlowFile与原始FlowFile(包含{"key2":"value3"})通过MergeContent处理器合并(可基于自定义属性做binning,确保两个FlowFile被分到同一组),合并格式选择JSON Merge或直接合并为包含两个JSON对象的数组。 - 使用
JoltTransformJSON处理器,配置Jolt规范为merge操作,实现新旧JSON的合并:
该规范会将新JSON中的字段覆盖原有JSON的对应字段,同时保留未被更新的[ { "operation": "merge", "spec": { "*": "&" } } ]key1等字段。
步骤3:将合并后的JSON写回Redis
使用PutRedisKey处理器,指定Key为x,将合并后的FlowFile内容写入Redis,完成部分更新。
方案二:使用ExecuteScript自定义脚本(更灵活)
如果需要更灵活的逻辑,可以用ExecuteScript处理器编写Groovy脚本,直接完成"获取-合并-写入"的全流程:
import groovy.json.JsonSlurper import groovy.json.JsonBuilder import redis.clients.jedis.Jedis def flowFile = session.get() if (!flowFile) return // 读取FlowFile中的新更新数据 def newJson = new JsonSlurper().parseText(flowFile.read().getText("UTF-8")) // 连接Redis并获取原有JSON def jedis = new Jedis("你的Redis地址", 6379) // 如果Redis有密码,添加 jedis.auth("你的密码") def existingJsonStr = jedis.get("x") def existingJson = existingJsonStr ? new JsonSlurper().parseText(existingJsonStr) : [:] // 合并JSON:新字段覆盖旧字段,保留未更新字段 existingJson.putAll(newJson) def mergedJson = new JsonBuilder(existingJson).toString() // 写回Redis jedis.set("x", mergedJson) jedis.close() // 可选:更新FlowFile内容为合并后的结果 flowFile = session.write(flowFile, { out -> out.write(mergedJson.getBytes("UTF-8")) } as OutputStreamCallback) session.transfer(flowFile, REL_SUCCESS)
注意:需在NiFi的
lib目录下放入Redis的Jedis客户端jar包,确保脚本能正常连接Redis。
内容的提问来源于stack exchange,提问作者Fatemeh
相关产品推荐
相关产品推荐

