使用NiFi ExecuteScript处理器Python实现GitHub API JSON写入FlowFile
解决ExecuteScript(Jython)写入GitHub接口返回数据到FlowFile的方案
以下两种实现适配你提到的两种使用场景:
场景1:脚本内直接创建FlowFile
完整代码如下,仅需保留原有数据获取逻辑、追加写入逻辑即可:
import urllib.request import json from java.nio.charset import StandardCharsets from org.apache.commons.io import IOUtils from java.io import ByteArrayInputStream # --- 你原有获取数据的代码保持不变 --- headers = { 'Content-Type': 'application/json {}'.format('cat'), 'Authorization': 'token TOKEN_TEXT' } page = 1 response_text = [] while page > 0: req = urllib.request.Request(url="GITHUB_API_URL".format(page),\ data=None, headers=headers) with urllib.request.urlopen(req) as resp: data = json.loads(resp.read().decode("utf-8")) if len(data) == 0: break else: response_text.extend(data) page += 1 # --- 原有代码结束,以下是写入FlowFile逻辑 --- try: # 将结果序列化为JSON字符串 output_json = json.dumps(response_text, ensure_ascii=False) # 创建新的FlowFile flow_file = session.create() # 写入内容到FlowFile flow_file = session.write(flow_file, lambda outputStream: IOUtils.copy(ByteArrayInputStream(output_json.encode('utf-8')), outputStream)) # 可选:设置自定义属性,标记数据源、文件类型 flow_file = session.putAttribute(flow_file, 'data.source', 'github_api') flow_file = session.putAttribute(flow_file, 'mime.type', 'application/json') # 将FlowFile传递到成功关系 session.transfer(flow_file, REL_SUCCESS) session.commit() except Exception as e: log.error('写入FlowFile失败: {}'.format(str(e))) session.rollback()
场景2:使用GenerateFlowFile传入的已有FlowFile
仅需修改写入部分的逻辑,替换场景1中的写入代码即可:
try: # 获取传入的FlowFile,无传入时直接退出 flow_file = session.get() if not flow_file: exit() # 序列化结果 output_json = json.dumps(response_text, ensure_ascii=False) # 覆盖写入FlowFile内容 flow_file = session.write(flow_file, lambda outputStream: IOUtils.copy(ByteArrayInputStream(output_json.encode('utf-8')), outputStream)) # 可选更新属性 flow_file = session.putAttribute(flow_file, 'data.source', 'github_api') flow_file = session.putAttribute(flow_file, 'mime.type', 'application/json') # 传递到下游处理器 session.transfer(flow_file, REL_SUCCESS) session.commit() except Exception as e: log.error('处理失败: {}'.format(str(e))) session.rollback()
补充说明
- ExecuteScript处理器默认绑定了
session、log、REL_SUCCESS、REL_FAILURE等全局对象,不需要额外导入即可直接使用 - 上述代码兼容Jython运行环境,不需要额外安装依赖包
- 完整的ExecuteScript使用规范、绑定对象说明可在对应版本的NiFi官方处理器参考文档中检索获取
内容的提问来源于stack exchange,提问作者ambrrrgris
相关产品推荐
相关产品推荐

