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.0NiFi Toolkit 2.11.0,并放在所有节点的 /opt/package 目录下。

安装 Java 环境

Ubuntu 环境安装 Java 21 版本可以直接使用 apt 进行安装:

1
2
apt update
apt install -y openjdk-21-jdk-headless

检查版本

1
java -version
1
2
3
openjdk version "21.0.12" 2026-07-21
OpenJDK Runtime Environment (build 21.0.12+8-1-22.04-Ubuntu)
OpenJDK 64-Bit Server VM (build 21.0.12+8-1-22.04-Ubuntu, mixed mode, sharing)

配置环境变量

1
2
echo "export JAVA_HOME=/usr/lib/jvm/java-21-openjdk-amd64" >> /etc/profile
source /etc/profile

单节点安装

single-user 模式启动单节点 NiFi

  1. 解压 nifi 压缩包,并配置环境变量:
1
2
3
4
5
cd /opt/package/
unzip nifi-2.11.0-bin.zip
mv nifi-2.11.0 /opt/nifi
echo "export PATH=$PATH:/opt/nifi/bin" >> /etc/profile
source /etc/profile
  1. 进入 NiFi 配置目录,开始初始化配置参数:
1
2
cd /opt/nifi
vi conf/nifi.properties

修改如下参数:

1
2
3
4
5
6
7
8
9
10
## 设置 HTTPS 访问地址,默认是 localhost,设置为 0.0.0.0 可以允许所有地址访问
nifi.web.https.host=0.0.0.0
## HTTPS 默认端口是 8443
nifi.web.https.port=8443

## 允许的 header 值列表,用逗号分隔。首次启动会将这里的地址写入到证书中
nifi.web.proxy.host=u68.zerchin.xyz:8443,192.168.2.68:8443

## 默认是单用户认证,不需要修改
nifi.security.user.authorizer=single-user-authorizer
  1. 设置密码,例如这里设置 admin 用户密码是:
1
nifi.sh set-single-user-credentials admin nifi12345678
  1. 启动 NiFi:
1
nifi.sh start

输出如下结果,并且没有报错,说明启动成功。

1
2
JAVA_HOME=/usr/lib/jvm/java-21-openjdk-amd64
NIFI_HOME=/opt/nifi
  1. 查看日志:
1
tail -f logs/nifi-app.log
  1. 等待一段时间,可以通过 https://192.168.2.68:8443 访问 NiFi UI,用户密码是前面设置的 admin/nifi12345678

基于 docker 运行单节点 NiFi

NiFi 官方也构建了 Docker 镜像,如果有 Docker 等容器环境,可以使用 Docker 快速拉起单节点 NiFi,参考如下:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
docker run \
-itd \
--name nifi-single \
-p 8443:8443 \
-e NIFI_WEB_HTTPS_HOST="0.0.0.0" \
-e NIFI_WEB_HTTPS_PORT="8443" \
-e NIFI_WEB_PROXY_HOST="u68.zerchin.xyz:8443,192.168.2.69:8443"
-e SINGLE_USER_CREDENTIALS_USERNAME="admin" \
-e SINGLE_USER_CREDENTIALS_PASSWORD="Rancher12345" \
-e NIFI_SENSITIVE_PROPS_KEY="Rancher12345" \
-e NIFI_JVM_HEAP_INIT="1g" \
-e NIFI_JVM_HEAP_MAX="2g" \
--hostname nifi \
apache/nifi:2.11.0

常见错误

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
2
3
4
subject=CN = localhost
issuer=CN = localhost
X509v3 Subject Alternative Name:
DNS:localhost, DNS:u68, DNS:192.168.2.68, DNS:u68.zerchin.xyz

这里要保证 DNS: 后面有访问的 IP 或者域名,如果配置错误,需要重新生成证书,方法如下:

  1. 停止 NiFi
1
nifi.sh stop
  1. 删除证书文件
1
rm -f conf/keystore.p12 conf/truststore.p12
  1. 编辑 nifi.properties 配置文件,重置证书相关的参数:
1
2
3
4
5
6
7
## 确认该地址填写正确
nifi.web.proxy.host=u68.zerchin.xyz:8443,192.168.2.68:8443

## 清空这几个地方
nifi.security.keystorePasswd=
nifi.security.keyPasswd=
nifi.security.truststorePasswd=
  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
2
3
4
cat /etc/hosts
192.168.2.68 u68 u68.zerchin.xyz
192.168.2.69 u69 u69.zerchin.xyz
192.168.2.71 u71 u71.zerchin.xyz

网络端口

节点之间至少需要开放:

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
2
3
4
5
cd /opt/package/
unzip nifi-2.11.0-bin.zip
mv nifi-2.11.0 /opt/nifi
echo "export PATH=$PATH:/opt/nifi/bin" >> /etc/profile
source /etc/profile
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 或其他标准证书管理方案。

https://issues.apache.org/jira/browse/NIFI-12200

方法一

在 NiFi1.x 版本有提供一个 tls-toolkit 工具,可以帮助我们一键生成证书文件,并自动配置好 nifi.properties 文件,虽然 NiFi2.x 已经不再提供这个工具了,但是自测下来使用没有什么问题。

下载解压 nifi-toolkit 包,并配置环境变量:

1
2
3
4
5
6
cd /opt/package/
wget https://dlcdn.apache.org/nifi/1.28.1/nifi-toolkit-1.28.1-bin.zip
unzip nifi-toolkit-1.28.1-bin.zip
mv nifi-toolkit-1.28.1 nifi-toolkit
echo "export PATH=$PATH:/opt/nifi-toolkit/bin" >> /etc/profile
source /etc/profile

通过 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
2
3
4
5
# ls nifi-cert/
'CN=nifi-ca_OU=nifi.p12' 'CN=nifi-ca_OU=nifi.password' nifi-cert.pem nifi-key.key u68.zerchin.xyz u69.zerchin.xyz u71.zerchin.xyz

# ls nifi-cert/u68.zerchin.xyz/
keystore.jks nifi.properties truststore.jks

这里可以看到,同时也生成了一个nifi.properties配置文件,该文件会自动配置好证书相关的集群配置参数。

将对应目录下的文件发送到对应节点上:

1
2
3
scp nifi-cert/u68.zerchin.xyz/* u68:/opt/nifi/conf/.
scp nifi-cert/u69.zerchin.xyz/* u69:/opt/nifi/conf/.
scp nifi-cert/u71.zerchin.xyz/* u71:/opt/nifi/conf/.

方法二

由于官方不再维护 tls-toolkit 工具,为了防止兼容性问题,也可以自行生成证书,参考如下:

  1. 创建 CA 证书,在其中一台节点上执行,例如这里在 u68 节点上执行:
1
2
3
4
5
6
7
8
9
10
mkdir /opt/nifi-pki
cd /opt/nifi-pki

export STORE_PASS='nifi-store-12345'

openssl genrsa -aes256 -passout pass:${STORE_PASS} -out ca.key 3072

openssl req -new -x509 -days 3650 -key ca.key -passin pass:${STORE_PASS} -sha256 -out ca.cer -subj "/CN=nifi-ca/OU=nifi/"

openssl x509 -in ca.cer -noout -subject -issuer -dates
  1. 创建节点证书,生成 u68 节点证书:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
NODE=u68
IP=192.168.2.68
DOMAIN=u68.zerchin.xyz
STORE_PASS='nifi-store-12345'
## key
openssl genrsa -aes256 -passout pass:${STORE_PASS} -out ${NODE}.key 3072
## CSR
openssl req -new -key ${NODE}.key -passin pass:${STORE_PASS} -out ${NODE}.csr -subj "/CN=${DOMAIN}/OU=NiFi/" -addext "subjectAltName=DNS:${DOMAIN},DNS:${NODE},IP:${IP}"
## 签发
openssl x509 -req -in ${NODE}.csr -CA ca.cer -CAkey ca.key -passin pass:${STORE_PASS} -CAcreateserial -out ${NODE}.cer -days 825 -sha256 -copy_extensions copyall
## 证书链
cat ${NODE}.cer ca.cer > ${NODE}-chain.cer
## 生成 PKCS12
openssl pkcs12 -export -out ${NODE}.p12 -inkey ${NODE}.key -passin pass:${STORE_PASS} -in ${NODE}-chain.cer -name ${NODE} -passout pass:${STORE_PASS}
  1. 其他两个节点也使用这个方法生成证书:

u69:

1
2
3
4
NODE=u69
IP=192.168.2.69
DOMAIN=u69.zerchin.xyz
...

u71:

1
2
3
4
NODE=u71
IP=192.168.2.71
DOMAIN=u71.zerchin.xyz
...

检查证书:

1
2
3
openssl pkcs12 -in u68.p12 -info -nodes

openssl pkcs12 -in u68.p12 -clcerts -nokeys | openssl x509 -noout -text
  1. 创建 Truststore:
1
2
3
4
5
6
7
8
9
keytool \
-importcert \
-noprompt \
-trustcacerts \
-alias nifi-ca \
-file ca.cer \
-keystore truststore.p12 \
-storetype PKCS12 \
-storepass ${STORE_PASS}
  1. 创建管理员客户端证书:
1
2
3
4
5
6
7
8
9
10
11
12
## key
openssl genrsa -aes256 -passout pass:${STORE_PASS} -out admin.key 3072

## CSR
openssl req -new -key admin.key -passin pass:${STORE_PASS} -out admin.csr -subj "/CN=admin/"

## 签发
openssl x509 -req -in admin.csr -CA ca.cer -CAkey ca.key -passin pass:${STORE_PASS} -CAcreateserial -out admin.cer -days 825 -sha256

## 生成浏览器可导入的
openssl pkcs12 -export -out admin.p12 -inkey admin.key -passin pass:${STORE_PASS} -in admin.cer -certfile ca.cer -name admin -passout pass:${STORE_PASS}

macos 14.4.1 版本加密算法有问题,需要更换成下面的

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
openssl req   -new   -key admin.key   -out admin.csr   -subj "/CN=admin/OU=NiFi"

openssl x509 -req -in admin.csr -CA ca.cer -CAkey ca.key -CAcreateserial -out admin.crt -days 825 -sha256 -extfile admin-ext.cnf

openssl x509 -in admin.crt -noout -subject -issuer -purpose

openssl x509 -in admin.crt -noout -text | grep -A3 -E 'Key Usage|Extended Key Usage'

openssl verify -CAfile ca.cer admin.crt

openssl x509 -in admin.crt -noout -issuer

openssl pkcs12 -export -legacy -out admin-macos.p12 -inkey admin.key -in admin.crt -certfile ca.cer -name "NiFi Admin"

sz admin-macos.p12
  1. 确认证书身份:
1
2
3
4
5
6
7
keytool -list -v -keystore u68.p12 -storepass ${STORE_PASS} | grep '^Owner'

keytool -list -v -keystore u69.p12 -storepass ${STORE_PASS} | grep '^Owner'

keytool -list -v -keystore u71.p12 -storepass ${STORE_PASS} | grep '^Owner'

keytool -list -v -keystore admin.p12 -storepass ${STORE_PASS} | grep '^Owner'

可以看到如下信息:

1
2
Owner: OU=NiFi, CN=u68.zerchin.xyz
Owner: OU=nifi, CN=nifi-ca

后续 authorizers.xml 中的 Node Identity 必须与这里显示的 Owner 完全一致。

官方同样要求 Node Identity 使用证书中的完整 DN。

  1. 同步证书到所有节点上:

u68:

1
2
scp u68.p12 root@192.168.2.68:/opt/nifi/conf/keystore.p12
scp truststore.p12 root@192.168.2.68:/opt/nifi/conf/truststore.p12

u69:

1
2
scp u69.p12 root@192.168.2.69:/opt/nifi/conf/keystore.p12
scp truststore.p12 root@192.168.2.69:/opt/nifi/conf/truststore.p12

u71:

1
2
scp u71.p12 root@192.168.2.71:/opt/nifi/conf/keystore.p12
scp truststore.p12 root@192.168.2.71:/opt/nifi/conf/truststore.p12

配置 Zookeeper

这里我们使用内置的 Zookeeper。

配置 zookeeper.properties

设置 zookeeper 节点

1
vi conf/zookeeper.properties
1
2
3
4
5
6
clientPort=2181
dataDir=./state/zookeeper

server.1=192.168.2.68:2888:3888;2181
server.2=192.168.2.69:2888:3888;2181
server.3=192.168.2.71:2888:3888;2181

配置 state-management.xml

找到 idzk-provider 的标签,将其中的 Connect String 修改为包含所有节点的地址

1
vi conf/state-management.xml
1
2
3
4
5
6
7
8
9
<cluster-provider>
<id>zk-provider</id>
<class>org.apache.nifi.cluster.coordination.state.ZookeeperStateProvider</class>
<property name="Connect String">192.168.2.68:2181,192.168.2.69:2181,192.168.2.71:2181</property>
<property name="Root Node">/nifi</property>
<property name="Session Timeout">10 seconds</property>
<property name="Access Control">Open</property>
</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 下

  1. 配置 GetFile Processor

拖拉 Processor 到画布中,搜索 GetFile

右键 Configure,进入 Properties ,在 Input Directory 中输入 /opt/input,配置从该目录监听文件变化。

点击右上角的 ,可以验证字段是否通过,例如这里验证失败,提示目录不存在,所以需要先创建该目录。

  1. 配置 PutFile Processor

再添加一个 Processor,选择 PutFile

同样还是右键 Configure,进入 Properties ,在 Directory 中输入 /opt/output,配置在改目录写入数据。

  1. 配置 Connection

从 GetFile 连接到 PutFile

这里会发现 PutFile 左上角有个黄色感叹号,说明还有配置没有配完。右键 Configure,进入到 Relationships,将 failure 和 success 勾选上 terminate,意思不管写入成功还是失败,都终止二次写入。例如如果 failure选择 retry,则会尝试再次写入。

配置完之后,就可以点击左侧的三角形图标启动 Processor。

  1. 启动 GetFile

这里可以先启动 GetFile,然后尝试在 /opt/input 目录下写入一个文件:

1
echo aaa > /opt/input/a1.txt

接着刷新页面,可以看到,已经有数据到队列中:

右键 List Queue 查看队列,可以看到具体内容,右侧菜单可以查看文件内容,甚至可以下载该文件。

  1. 启动 PutFile

刷新一下页面,可以看到队列空了,并且 PutFile 也能看到处理了数据。

进入节点,查看数据也已经实时写入到 /opt/output 目录下:

1
2
# cat /opt/output/a1.txt 
aaa

同步本地文件到 HDFS

NiFi 2.x 默认发行包不再包含 Hadoop/HDFS Bundle,但 HDFS Processor 仍然保留并持续维护,NiFi 2 把 Hadoop Bundle 改成了一个可选 Maven Profile。可以在这个网站搜索 nifi-hadoop-libraries-narnifi-hadoop-nar 两个包,下载下来,然后放到所有节点的 extensions 目录下即可。

  1. GetFile
  • Input Directory:设置读取的文件夹路径。
  • Keep Source File:读取完文件是否保存源文件,如果设置为 false,读取完会删除源文件。
  1. 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

  1. 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 只有主节点执行。
  1. ConvertRecord

默认 JDBC 读取的是 Avro 格式的数据,这个数据写入到文件无法被正常读取,所以需要先将 Avro 转换成 JSON 再写入到后端

其中 NiFi1.x 版本使用的 ConvertAvroToJSON,该插件在 NiFi2.x 已移除,官方推荐使用 ConvertRecord,配置如下:

  • Record Reader:新增一个 Service,选择 AvroReader
  • Record Writer:新增一个 Service,选择 JsonRecordSetWriter

记得所有的 Controller Service 都要 Enable

  1. UpdateAttribute

NiFi 支持配置和修改 FlowFile 的属性值,可以通过 UpdateAttribute 来实现这个需求:

  • date:通过右上角的加号新增的属性,后面会用到,这里给当前这个 FlowFile 设置一个日期的属性,使用 Expression Language 语法,例如 ${now():format('yyyy-MM-dd')}
  • filename:默认 FlowFile 的文件名是一个随机字符串,这里我们可以通过新增一个名叫 filename 的 key 来修改文件名,例如 ${now():format('HH-mm-ss')}.json

设置好之后,可以在队列里查看 FlowFile 的属性。

  1. 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:输出的路径,默认输出的路径是固定的,可以通过 Expression Language 语法选择前面设定的 date 参数:/nifi-mysql/${date}

  1. 最终效果:

NiFi UI

HDFS 数据写入

Troubleshooting

  1. 在初始化阶段,如果 authorizers.xml 配置错误需要重新配置该文件,则需要把 authorizations.xml 和 users.xml 删除后,再重新启动 NiFi

  2. Processor 中填写的参数,可以鼠标移到 Property 的 key 后面的问号,可以看到是否支持 Expression Language 语法

  3. 读完文件是否要保存会影响之后的文件读取

  4. 文件是否覆盖会影响后续的写入

  5. bootstrap 文件里可以设置 jvm 内存大小

  6. group 可以设置单个 Group 最大数据大小

参考链接: