Apache NiFi 2.0.0缺失PutHDFS处理器,求HDFS数据传输替代方案
Apache NiFi 2.0.0 无PutHDFS处理器的HDFS写入替代方案
一、Python脚本处理器定制实现
NiFi 2.0.0的ExecuteScript或InvokeScriptedProcessor支持Python脚本,可通过Python HDFS客户端库直接实现写入逻辑,步骤如下:
环境准备:在NiFi运行节点的Python环境中安装轻量易用的
hdfs客户端库:pip install hdfs若NiFi使用独立Python环境,需在
ExecuteScript的Python Path属性中指定库所在路径。脚本示例(ExecuteScript处理器):
选择Python作为脚本语言,编写读取FlowFile内容并写入HDFS的逻辑:from org.apache.nifi.processor.io import StreamCallback import io from hdfs import InsecureClient class HdfsWriter(StreamCallback): def process(self, inputStream, outputStream): # 读取FlowFile内容 content = io.TextIOWrapper(inputStream, encoding='utf-8').read() # 初始化HDFS客户端(替换为你的NameNode地址与端口) client = InsecureClient('http://namenode-host:50070', user='hadoop') # 写入指定HDFS路径(可通过NiFi属性动态替换路径/文件名) hdfs_path = '/user/nifi/data/${filename}' with client.write(hdfs_path, overwrite=True, encoding='utf-8') as writer: writer.write(content) # 输出原内容,保留FlowFile向下游传递(按需选择) outputStream.write(content.encode('utf-8')) flowFile = session.get() if flowFile is not None: flowFile = session.write(flowFile, HdfsWriter()) session.transfer(flowFile, REL_SUCCESS)若HDFS启用Kerberos认证,需改用
KerberosClient并配置krb5.conf及对应票据。
二、其他替代处理器/变通方法
1. ExecuteStreamCommand执行HDFS CLI命令
利用Hadoop原生hdfs dfs命令,通过标准输入写入HDFS,配置如下:
- 处理器:
ExecuteStreamCommand - 核心配置项:
Command:/path/to/hadoop/bin/hdfs(指定hadoop客户端的hdfs命令路径)Command Arguments:dfs -put - /user/nifi/data/${filename}(-表示从标准输入读取内容)Redirect Standard Input: 勾选(将FlowFile内容作为命令的标准输入)
- 注意:确保NiFi运行用户拥有HDFS目标路径的写入权限,且Hadoop客户端配置文件(core-site.xml、hdfs-site.xml)已通过
HADOOP_CONF_DIR环境变量指定。
2. 基于WebHDFS的InvokeHTTP处理器
通过HDFS的WebHDFS API实现写入,配置InvokeHTTP处理器:
HTTP Method: PUTRemote URL:http://namenode-host:50070/webhdfs/v1/user/nifi/data/${filename}?op=CREATE&overwrite=true- 替换
namenode-host为你的NameNode地址,/user/nifi/data/${filename}为目标路径(支持NiFi属性动态生成)
- 替换
Send Message Body: 勾选,选择FlowFile Content作为请求体Response Code: 设置成功响应码(如201、200),并配置对应关系传递- 注意:若HDFS启用Kerberos,需在
InvokeHTTP中配置SPNEGO认证。
内容的提问来源于stack exchange,提问作者Filbadeha
相关产品推荐
相关产品推荐

