You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.03 15:06:04