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

Apache NiFi ExecuteStreamCommand生成多个独立FlowFile的方案咨询

解决方案

核心原因是ExecuteStreamCommand处理器默认会将脚本所有stdout输出内容封装为单个FlowFile,不会主动识别多段输出的边界,所以需要人为添加分隔标记,再配合NiFi的拆分处理器实现多FlowFile输出。

步骤1:修改Python脚本,添加输出分隔符

选择一个绝对不会出现在你的JSON内容中的字符串作为分隔符,每个DataFrame转成JSON输出后,写入一次分隔符即可,示例代码如下:

import pandas as pd
import sys

input = sys.stdin.readlines()
# 自定义分隔符,确保不会和JSON内容冲突,也可使用ASCII不可见控制字符如\x1E更稳妥
SEPARATOR = "---X-NIFI-FLOWFILE-SEPARATOR---"

# 保留你原本拆分files的逻辑
# searching for subfiles and saving them to a list with files ...

for dataFrame in files:
    df = pd.DataFrame(dataFrame)
    # 保留你原本的DataFrame处理逻辑
    # several improvements of DataFrame ...
    json_output = df.to_json(orient='records', date_format='iso', date_unit='s')
    # 先写JSON内容,再写分隔符
    sys.stdout.write(json_output)
    sys.stdout.write(SEPARATOR)
sys.stdout.flush()

步骤2:配置NiFi后续处理器拆分FlowFile

在ExecuteStreamCommand处理器的下游添加SplitContent处理器,配置如下参数:

  • Byte Sequence to Split On:填入你刚才定义的分隔符字符串,比如---X-NIFI-FLOWFILE-SEPARATOR---
  • Trim Byte Sequence:设置为true,拆分后自动去掉分隔符,不会保留在最终的JSON FlowFile中
  • 其余参数保持默认即可

拆分完成后,每个拆分得到的FlowFile就是对应单个DataFrame的JSON内容,直接对接后续的InferAvroSchema处理器即可正常识别不同的Schema。

内容的提问来源于stack exchange,提问作者Nicko

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 05:15:00