NiFi 数据流处理平台
概述
Apache NiFi 是一个面向 DataFlow(数据流) 的数据集成和数据处理平台。
现实环境中的数据通常不会只存在于一个系统中,而是需要不断在不同系统之间:
- 采集(Ingest)
- 传输(Transport)
- 转换(Transform)
- 过滤(Filter)
- 路由(Route)
- 拆分(Split)
- 合并(Merge)
- 持久化(Store)
- 分发(Deliver)
NiFi 的目标,就是对这些数据在系统之间的流动进行自动化和统一管理。这里的重点并不是单纯“传输数据”,也关心数据中处理过程中的状态,例如处理成功还是失败,历史追踪等。NiFi 更准确的定位是一个用于构建、管理、监控和追踪数据流的可视化 DataFlow 平台。
基本概念
一个易于使用、功能强大、可靠的处理和分发数据框架,主要用于数据的同步传输,支持灵活的数据格式转换,同时可以设置定时调度任务
NiFi 支持数据路由、转换和系统中介逻辑的强大且可伸缩的有向图。
包括以下功能:
- 基于浏览器的用户界面
- 数据来源跟踪
- 丰富的配置
- 可扩展的设计
- 安全通信
核心概念
NiFi 的基本设计概念与基于流程的变成的主要思想密切相关。以下是一些主要的 NiFi 概念以及它们如何映射到 FBP:
| NiFi 术语 | 描述 |
|---|---|
| FlowFile | 数据在 NiFi 中传输时封装的对象,分为属性(attribute)和内容(content),其中属性是键值对的头信息,内容为字符串 |
| FlowFile Processor | 数据处理器组件,通过选择不同的处理器,对数据进行不同的读写或者转换清洗等操作 |
| Connection | 处理器直接的连接,单个处理器可以有多个连接完成数据的分流 |
| Flow Conntroller | 流控制器管理器中的资源分配 |
| Process Group | 处理组,将多个处理器连接的链路封装起来作为一个组管理 |
FlowFile 是数据,Processor 负责处理数据,Connection/Queue 负责传递和缓冲数据,Flow Controller 负责调度整个 DataFlow,Repository 则负责保证数据状态、内容和历史记录能够持久化和追踪。
NiFi 架构
NiFi 在宿主操作系统上的 JVM 中运行。NiFi 在 JVM 上的主要组件如下。
部分组件说明:
FlowFile Repository:FlowFile 存储库是用于存储正在传输时候的数据对象,主要存储数据状态信息,存储在磁盘。
用于跟踪 NiFi 对当前流中处于活动状态的 FlowFile 的了解情况。该存储库的实现是可插拔的。默认方案是在指定磁盘分区上使用持久化预写日志 (Write-Ahead Log)。
Content Repository:内容存储库主要存储数据内容,也是 FlowFile 的主要存储地址,支持多磁盘。
用于存放给定 FlowFile 的实际内容字节。存储库的实现是可插拔的。默认方法是一种相当简单的机制,它将数据块存储在文件系统中。可以指定多个文件系统存储位置,以便使用不同的物理分区,从而减少对单个卷的争用。
Provenance Repository:来源数据库,存储不同的数据来源信息。
溯源存储库用于存储所有溯源事件数据。该存储库采用可插拔结构,默认实现方式是使用一个或多个物理磁盘卷。每个位置的事件数据都已建立索引并可进行搜索。
从 NiFi 1.0 版本开始,NiFi 采用 Zero-Leader Clustering 模式。NiFi 集群中的每个节点都对数据执行相同的任务,但各自处理的数据集不同。Apache ZooKeeper 会选举一个节点作为集群协调器,并自动处理故障转移。所有集群节点都会向集群协调器报告心跳和状态信息。集群协调器负责节点的连接和断开。此外,每个集群还有一个主节点,也由 ZooKeeper 选举产生。作为数据流管理器,您可以通过任何节点的 UI 与 NiFi 集群进行交互。您所做的任何更改都会复制到集群中的所有节点,从而允许多个入口点。
基于上述架构,NiFi 提供了一些非常关键的数据流能力,包括:
- Guaranteed Delivery:通过持久化 Repository 提高数据处理可靠性;
- Buffering / Back Pressure:上下游处理能力不匹配时进行缓冲和流量控制;
- Prioritized Queuing:Queue 中的 FlowFile 可以按照不同策略决定处理顺序;
- Data Provenance:记录数据从进入 NiFi 到离开 NiFi 的处理轨迹;
- Visual Command and Control:通过 Web UI 实时构建和修改 DataFlow;
- Extensible Architecture:Processor、Controller Service、Reporting Task 等均可扩展;
- Scale Out:通过 NiFi Cluster 横向扩展数据处理能力。
这些能力共同决定了 NiFi 是一个面向数据流生命周期管理的数据集成平台。
环境准备
| Item | Version |
|---|---|
| Ubuntu | 22.04.5 |
| NiFi | 2.11.0 |
| Java | 21 |
| 节点 | IP |
|---|---|
| u68 | 192.168.2.68 |
| u69 | 192.168.2.69 |
| u71 | 192.168.2.71 |
提前在 NiFi 官方页面下载最新的版本,这里选择 2.11.0 版本,其中 Source 是源码,需要下载的是 NiFi Standard 2.11.0 和 NiFi Toolkit 2.11.0,并放在所有节点的 /opt/package 目录下。
安装 Java 环境
Ubuntu 环境安装 Java 21 版本可以直接使用 apt 进行安装:
1 | apt update |
检查版本
1 | java -version |
1 | openjdk version "21.0.12" 2026-07-21 |
配置环境变量
1 | echo "export JAVA_HOME=/usr/lib/jvm/java-21-openjdk-amd64" >> /etc/profile |
单节点安装
single-user 模式启动单节点 NiFi
- 解压 nifi 压缩包,并配置环境变量:
1 | cd /opt/package/ |
- 进入 NiFi 配置目录,开始初始化配置参数:
1 | cd /opt/nifi |
修改如下参数:
1 | ## 设置 HTTPS 访问地址,默认是 localhost,设置为 0.0.0.0 可以允许所有地址访问 |
- 设置密码,例如这里设置 admin 用户密码是:
1 | nifi.sh set-single-user-credentials admin nifi12345678 |
- 启动 NiFi:
1 | nifi.sh start |
输出如下结果,并且没有报错,说明启动成功。
1 | JAVA_HOME=/usr/lib/jvm/java-21-openjdk-amd64 |
- 查看日志:
1 | tail -f logs/nifi-app.log |
- 等待一段时间,可以通过
https://192.168.2.68:8443访问 NiFi UI,用户密码是前面设置的admin/nifi12345678
基于 docker 运行单节点 NiFi
NiFi 官方也构建了 Docker 镜像,如果有 Docker 等容器环境,可以使用 Docker 快速拉起单节点 NiFi,参考如下:
1 | docker run \ |
常见错误
HTTP ERROR 400 Invalid SNI
如果访问 UI,发现无法进入登录页面,并且报错 HTTP ERROR 400 Invalid SNI,说明 HTTPS 证书和访问的地址不匹配,可以用下列命令进行检查。
1 | openssl s_client -connect 127.0.0.1:8443 -showcerts 2>/dev/null | openssl x509 -noout -subject -issuer -ext subjectAltName |
1 | subject=CN = localhost |
这里要保证 DNS: 后面有访问的 IP 或者域名,如果配置错误,需要重新生成证书,方法如下:
- 停止 NiFi
1 | nifi.sh stop |
- 删除证书文件
1 | rm -f conf/keystore.p12 conf/truststore.p12 |
- 编辑
nifi.properties配置文件,重置证书相关的参数:
1 | ## 确认该地址填写正确 |
- 接着重新启动 NiFi 即可
1 | nifi.sh start |
集群部署
前面的示例部署的是单节点 NiFi。对于生产环境或者需要更高吞吐能力的场景,可以将 NiFi 部署为 Cluster。
NiFi 使用 Zero-Leader Clustering,每个 Node 运行相同的 Flow,并各自处理自己的数据。Cluster Coordinator 和 Primary Node 通过 Leader Election 自动选举。默认实现使用 ZooKeeper。
同样是需要按照前面的步骤[安装 Java 环境](##安装 Java 环境)。
域名解析
三个节点都需要能够解析另外两个节点,正式环境推荐 DNS,测试环境可以修改 hosts 文件:
1 | cat /etc/hosts |
网络端口
节点之间至少需要开放:
| Port | 用途 |
|---|---|
| 8443 | NiFi HTTPS |
| 11443 | NiFi Cluster Node Protocol |
| 6342 | Cluster Load Balance |
| 2181 | ZooKeeper Client |
| 2888 | ZooKeeper Quorum |
| 3888 | ZooKeeper Leader Election |
官方默认 NiFi HTTPS 为 8443,Cluster Load Balance 为 6342;Node Protocol 默认未启用,本文配置为 11443。Embedded ZooKeeper 通常使用 2181/2888/3888。
解压 NIFI 包
在所有节点,解压 nifi 压缩包,并配置环境变量:
1 | cd /opt/package/ |
1 |
自定义证书
安全 NiFi Cluster 中每个节点都需要自己的 TLS Identity。NiFi 官方要求安全 Cluster 使用 mTLS 进行节点身份验证,每个 Cluster Node 应有自己的 Keystore,并通过 Truststore 建立共同信任关系。也就是每个节点的 Keystore 共同信任一个 Truststore。
NiFi 1.x 曾提供
tls-toolkit.sh用于自动生成 NiFi TLS 证书、Keystore 和 Truststore。该模块在 NiFi 1.24.0 中被标记为 Deprecated,并在 NiFi 2.0 中正式移除。因此 NiFi 2.x 的nifi-toolkit中不再包含tls-toolkit.sh。对于开发 Cluster,Apache NiFi 当前文档提供 OpenSSL + Java keytool 作为手工生成 Keystore/Truststore 的方案;生产环境则建议使用企业 CA、cert-manager 或其他标准证书管理方案。
方法一
在 NiFi1.x 版本有提供一个 tls-toolkit 工具,可以帮助我们一键生成证书文件,并自动配置好 nifi.properties 文件,虽然 NiFi2.x 已经不再提供这个工具了,但是自测下来使用没有什么问题。
下载解压 nifi-toolkit 包,并配置环境变量:
1 | cd /opt/package/ |
通过 tls-toolkit.sh 命令生成证书文件
1 | tls-toolkit.sh standalone -C 'CN=nifi-ca, OU=nifi' -n 'u68.zerchin.xyz,u69.zerchin.xyz,u71.zerchin.xyz' -K nifi12345678 -S nifi12345678 -P nifi12345678 -o nifi-cert |
参数说明:
C:指定 DN- n:主机名列表,以逗号分隔
K&S&P:分别是 key 密码、Keystore 密码、trustStore密码o:证书输出目录
检查:
1 | # ls nifi-cert/ |
这里可以看到,同时也生成了一个nifi.properties配置文件,该文件会自动配置好证书相关的集群配置参数。
将对应目录下的文件发送到对应节点上:
1 | scp nifi-cert/u68.zerchin.xyz/* u68:/opt/nifi/conf/. |
方法二
由于官方不再维护 tls-toolkit 工具,为了防止兼容性问题,也可以自行生成证书,参考如下:
- 创建 CA 证书,在其中一台节点上执行,例如这里在 u68 节点上执行:
1 | mkdir /opt/nifi-pki |
- 创建节点证书,生成 u68 节点证书:
1 | NODE=u68 |
- 其他两个节点也使用这个方法生成证书:
u69:
1 | NODE=u69 |
u71:
1 | NODE=u71 |
检查证书:
1 | openssl pkcs12 -in u68.p12 -info -nodes |
- 创建 Truststore:
1 | keytool \ |
- 创建管理员客户端证书:
1 | ## key |
macos 14.4.1 版本加密算法有问题,需要更换成下面的
1 | openssl req -new -key admin.key -out admin.csr -subj "/CN=admin/OU=NiFi" |
- 确认证书身份:
1 | keytool -list -v -keystore u68.p12 -storepass ${STORE_PASS} | grep '^Owner' |
可以看到如下信息:
1 | Owner: OU=NiFi, CN=u68.zerchin.xyz |
后续 authorizers.xml 中的 Node Identity 必须与这里显示的 Owner 完全一致。
官方同样要求 Node Identity 使用证书中的完整 DN。
- 同步证书到所有节点上:
u68:
1 | scp u68.p12 root@192.168.2.68:/opt/nifi/conf/keystore.p12 |
u69:
1 | scp u69.p12 root@192.168.2.69:/opt/nifi/conf/keystore.p12 |
u71:
1 | scp u71.p12 root@192.168.2.71:/opt/nifi/conf/keystore.p12 |
配置 Zookeeper
这里我们使用内置的 Zookeeper。
配置 zookeeper.properties
设置 zookeeper 节点
1 | vi conf/zookeeper.properties |
1 | clientPort=2181 |
配置 state-management.xml
找到 id 为 zk-provider 的标签,将其中的 Connect String 修改为包含所有节点的地址
1 | vi conf/state-management.xml |
1 | <cluster-provider> |
设置唯一标识符
u68:
1 | mkdir -p state/zookeeper && echo "1" > state/zookeeper/myid |
u69:
1 | mkdir -p state/zookeeper && echo "2" > state/zookeeper/myid |
u71:
1 | mkdir -p state/zookeeper && echo "3" > state/zookeeper/myid |
配置 nifi.properties 集群参数
各个节点配置需要根据具体条件去配置,参考如下:
| 参数项 | u68 | u69 | u71 |
|---|---|---|---|
| nifi.state.management.embedded.zookeeper.start | true | true | true |
| nifi.state.management.provider.cluster | zk-provider |
zk-provider | zk-provider |
| nifi.cluster.is.node | true | true | true |
| nifi.cluster.node.address | 192.168.2.68 |
192.168.2.69 | 192.168.2.71 |
| nifi.cluster.node.protocol.port | 11443 | 11443 | 11443 |
| nifi.cluster.flow.election.max.candidates | 2 (3个节点满足2个即可投票) |
2 | 2 |
| nifi.zookeeper.connect.string | 192.168.2.68:2181,192.168.2.69:2181,192.168.2.71:2181 | 192.168.2.68:2181,192.168.2.69:2181,192.168.2.71:2181 | 192.168.2.68:2181,192.168.2.69:2181,192.168.2.71:2181 |
| nifi.web.https.host | u68.zerchin.xyz | u69.zerchin.xyz | u71.zerchin.xyz |
| nifi.web.https.port | 8443 | 8443 | 8443 |
| nifi.security.keystorePasswd | nifi-store-12345 | nifi-store-12345 | nifi-store-12345 |
| nifi.security.keyPasswd | nifi-store-12345 | nifi-store-12345 | nifi-store-12345 |
| nifi.security.truststorePasswd | nifi-store-12345 | nifi-store-12345 | nifi-store-12345 |
| nifi.sensitive.props.key | nifi12345678 | nifi12345678 |
设置 admin 密码
1 | nifi.sh set-single-user-credentials admin nifi12345678 |
启动 NiFi 集群
三台节点分别执行:
1 | nifi.sh start |
查看日志:
1 | tail -f /opt/nifi/logs/nifi-app.log |
浏览器访问
登录任意一个节点,例如 https://u69.zerchin.xyz:8443,使用前面设置的 admin 的用户和密码进行登录。
右上角设置按钮 -> Cluster,可以看到集群状态:
DEMO 示例
需求:监听 /opt/input 目录下的文件,并写入到 /opt/output 下
- 配置 GetFile Processor
拖拉 Processor 到画布中,搜索 GetFile
右键 Configure,进入 Properties ,在 Input Directory 中输入 /opt/input,配置从该目录监听文件变化。
点击右上角的 √,可以验证字段是否通过,例如这里验证失败,提示目录不存在,所以需要先创建该目录。
- 配置 PutFile Processor
再添加一个 Processor,选择 PutFile
同样还是右键 Configure,进入 Properties ,在 Directory 中输入 /opt/output,配置在改目录写入数据。
- 配置 Connection
从 GetFile 连接到 PutFile
这里会发现 PutFile 左上角有个黄色感叹号,说明还有配置没有配完。右键 Configure,进入到 Relationships,将 failure 和 success 勾选上 terminate,意思不管写入成功还是失败,都终止二次写入。例如如果 failure选择 retry,则会尝试再次写入。
配置完之后,就可以点击左侧的三角形图标启动 Processor。
- 启动 GetFile
这里可以先启动 GetFile,然后尝试在 /opt/input 目录下写入一个文件:
1 | echo aaa > /opt/input/a1.txt |
接着刷新页面,可以看到,已经有数据到队列中:
右键 List Queue 查看队列,可以看到具体内容,右侧菜单可以查看文件内容,甚至可以下载该文件。
- 启动 PutFile
刷新一下页面,可以看到队列空了,并且 PutFile 也能看到处理了数据。
进入节点,查看数据也已经实时写入到 /opt/output 目录下:
1 | # cat /opt/output/a1.txt |
同步本地文件到 HDFS
NiFi 2.x 默认发行包不再包含 Hadoop/HDFS Bundle,但 HDFS Processor 仍然保留并持续维护,NiFi 2 把 Hadoop Bundle 改成了一个可选 Maven Profile。可以在这个网站搜索 nifi-hadoop-libraries-nar 和 nifi-hadoop-nar 两个包,下载下来,然后放到所有节点的 extensions 目录下即可。
- GetFile
- Input Directory:设置读取的文件夹路径。
- Keep Source File:读取完文件是否保存源文件,如果设置为 false,读取完会删除源文件。
- PutHDFS
- Hadoop Configuration Resources:core-site.xml 和 hdfs-site.xml 配置文件路径,例如我这里填写的是
/opt/hadoop/etc/hadoop/core-site.xml,/opt/hadoop/etc/hadoop/hdfs-site.xml - Conflict Resolution Strategy:同一个文件的写入策略,可以设置为 replace 同步源文件内容。
- Directory:输出的路径
离线同步 MySQL 数据到 HDFS
- JDBC 配置
首先拉取一个 ExecuteSQL Processor ,对其中的 Database Connection Pooling Service 进行配置,点击该配置会跳转到 JDBC 配置页面,参数如下:
- Database Connection URL:MySQL 连接配置,写法是
jdbc:mysql://<URL>:<Port>/<Database> - Database Driver Class Name:MySQL JDBC 驱动,直接写
com.mysql.jdbc.Driver - Database Driver Locations:mysql-connector 的 本地 jar 包地址,需要在所有节点都存在这个 jar 包,可以在这个地址,找到 FILES 里的 jar 包进行下载。
- Database User:MySQL 用户名
- Password:MySQL 用户密码
配置完成,要点击 Enable 启动该服务。
接着回到 ExecuteSQL 的配置页面,参数如下:
- Database Connection Pooling Service:上一步配置 MySQL 连接器的地方。
- SQL Pre-Query:前置 SQL 命令,只执行一次,通常可以写
user <database>; - SQL Query:真正监听的 SQL 命令,例如 SELECT 查询。
这里要修改一下 Scheduling:
- Run Schedule :多久运行一次,默认 0 sec,会导致一直在执行 SQL 查询语句,每次查询都会进入队列,导致队列溢出,所以通常这里可以按需配置,例如设置一天(86400 sec)。
- Execution:默认所有节点执行,SQL 查询是一次性的,多次查询会有多份数据,这里可以设置为 Primary node 只有主节点执行。
- ConvertRecord
默认 JDBC 读取的是 Avro 格式的数据,这个数据写入到文件无法被正常读取,所以需要先将 Avro 转换成 JSON 再写入到后端
其中 NiFi1.x 版本使用的 ConvertAvroToJSON,该插件在 NiFi2.x 已移除,官方推荐使用 ConvertRecord,配置如下:
- Record Reader:新增一个 Service,选择
AvroReader - Record Writer:新增一个 Service,选择
JsonRecordSetWriter
记得所有的 Controller Service 都要 Enable
- UpdateAttribute
NiFi 支持配置和修改 FlowFile 的属性值,可以通过 UpdateAttribute 来实现这个需求:
- date:通过右上角的加号新增的属性,后面会用到,这里给当前这个 FlowFile 设置一个日期的属性,使用 Expression Language 语法,例如
${now():format('yyyy-MM-dd')} - filename:默认 FlowFile 的文件名是一个随机字符串,这里我们可以通过新增一个名叫 filename 的 key 来修改文件名,例如
${now():format('HH-mm-ss')}.json
设置好之后,可以在队列里查看 FlowFile 的属性。
- PutHDFS
Hadoop Configuration Resources:core-site.xml 和 hdfs-site.xml 配置文件路径,例如我这里填写的是
/opt/hadoop/etc/hadoop/core-site.xml,/opt/hadoop/etc/hadoop/hdfs-site.xmlConflict Resolution Strategy:同一个文件的写入策略,可以设置为 replace 同步源文件内容。
Directory:输出的路径,默认输出的路径是固定的,可以通过 Expression Language 语法选择前面设定的 date 参数:
/nifi-mysql/${date}
- 最终效果:
NiFi UI
HDFS 数据写入
Troubleshooting
在初始化阶段,如果 authorizers.xml 配置错误需要重新配置该文件,则需要把 authorizations.xml 和 users.xml 删除后,再重新启动 NiFi
Processor 中填写的参数,可以鼠标移到 Property 的 key 后面的问号,可以看到是否支持 Expression Language 语法
读完文件是否要保存会影响之后的文件读取
文件是否覆盖会影响后续的写入
bootstrap 文件里可以设置 jvm 内存大小
group 可以设置单个 Group 最大数据大小
参考链接:
架构:https://nifi.apache.org/nifi-docs/overview.html#nifi-architecture
语法:https://nifi.apache.org/nifi-docs/expression-language-guide.html
集群:https://nifi.apache.org/nifi-docs/administration-guide.html#clustering
中文文档:https://nifichina.github.io/1-%E5%9F%BA%E7%A1%80%E6%96%87%E6%A1%A3/1-Overview.html
docker启动参数:https://github.com/apache/nifi/blob/main/nifi-docker/dockerhub/sh/start.sh