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
相关产品推荐
相关产品推荐

