VPC環境で利用できます。
NiFiは Dataflowエンジンとして、お互い違うシステム間のデータを処理してリリースする、使用が簡単で安定的なシステムです。NiFiは、データを収集・処理した後、ロードするための ETLツールの一種であり、分散環境において大量のデータを収集・処理し、Flow-Based Programming(FBP)の概念を実装して開発されたオープンソースです。
NiFiはリアルタイム処理に非常に適合しており、データを損失することなく転送できるというメリットがあります。
NiFiの構成要素
NiFiを構成するコンポーネントは大きく4つで、FlowFile、Processor、Connector、Controllerであり、各要素別の説明は次の通りです。

- FlowFile: NiFiが認識するデータの単位
- Content: データ自体を意味
- Attribute: データ関連情報をキー/値のペアで表現したもの
- Processor: FlowFileを収集、変形、保存する機能
- データ処理を完了した後に新たな FlowFileを作成できる
- Processorは複数個が並列で動作できる
- Connector: Processorと Processorを接続して FlowFileを送信する機能
- FlowFileの Queueを意味
- 優先順位および backpressureを設定して負荷を調整
- Flow Controller: 各プロセスを接続し、その間を移動する FlowFileを管理
Ambari UIでの NiFiサービス確認
Ambari UIの Add serviceから NiFiをインストールできます。
インストール時に入力するオプション
- Advanced nifi-properties-envの contentを、以下の値にすべて置き換える必要があります。
vmで openssl randコマンドを実行して出力された値をコピーし、以下の content設定値のうち、最下位の nifi.sensitive.props.keyに入力します。
openssl rand -base64 32
# Core Properties #
nifi.flow.configuration.file=./conf/flow.xml.gz
nifi.flow.configuration.json.file=./conf/flow.json.gz
nifi.flow.configuration.archive.enabled=true
nifi.flow.configuration.archive.dir=./conf/archive/
nifi.flow.configuration.archive.max.time=30 days
nifi.flow.configuration.archive.max.storage=500 MB
nifi.flow.configuration.archive.max.count=
nifi.flowcontroller.autoResumeState=true
nifi.flowcontroller.graceful.shutdown.period=10 sec
nifi.flowservice.writedelay.interval=500 ms
nifi.administrative.yield.duration=30 sec
nifi.bored.yield.duration=10 millis
nifi.queue.backpressure.count=10000
nifi.queue.backpressure.size=1 GB
nifi.authorizer.configuration.file=./conf/authorizers.xml
nifi.login.identity.provider.configuration.file=./conf/login-identity-providers.xml
nifi.templates.directory=./conf/templates
nifi.ui.banner.text=
nifi.ui.autorefresh.interval=30 sec
nifi.nar.library.directory=./lib
nifi.nar.library.autoload.directory=./extensions
nifi.nar.working.directory=./work/nar/
nifi.documentation.working.directory=./work/docs/components
nifi.nar.unpack.uber.jar=false
####################
# State Management #
####################
nifi.state.management.configuration.file=./conf/state-management.xml
nifi.state.management.provider.local=local-provider
nifi.state.management.provider.cluster=zk-provider
nifi.state.management.embedded.zookeeper.start=false
nifi.state.management.embedded.zookeeper.properties=./conf/zookeeper.properties
nifi.database.directory=./database_repository
nifi.h2.url.append=;LOCK_TIMEOUT=25000;WRITE_DELAY=0;AUTO_SERVER=FALSE
nifi.repository.encryption.protocol.version=
nifi.repository.encryption.key.id=
nifi.repository.encryption.key.provider=
nifi.repository.encryption.key.provider.keystore.location=
nifi.repository.encryption.key.provider.keystore.password=
nifi.flowfile.repository.implementation=org.apache.nifi.controller.repository.WriteAheadFlowFileRepository
nifi.flowfile.repository.wal.implementation=org.apache.nifi.wali.SequentialAccessWriteAheadLog
nifi.flowfile.repository.directory=./flowfile_repository
nifi.flowfile.repository.checkpoint.interval=20 secs
nifi.flowfile.repository.always.sync=false
nifi.flowfile.repository.retain.orphaned.flowfiles=true
nifi.swap.manager.implementation=org.apache.nifi.controller.FileSystemSwapManager
nifi.queue.swap.threshold=20000
nifi.content.repository.implementation=org.apache.nifi.controller.repository.FileSystemRepository
nifi.content.claim.max.appendable.size=50 KB
nifi.content.repository.directory.default=./content_repository
nifi.content.repository.archive.max.retention.period=7 days
nifi.content.repository.archive.max.usage.percentage=50%
nifi.content.repository.archive.enabled=true
nifi.content.repository.always.sync=false
nifi.content.viewer.url=../nifi-content-viewer/
nifi.provenance.repository.implementation=org.apache.nifi.provenance.WriteAheadProvenanceRepository
nifi.provenance.repository.directory.default=./provenance_repository
nifi.provenance.repository.max.storage.time=30 days
nifi.provenance.repository.max.storage.size=10 GB
nifi.provenance.repository.rollover.time=10 mins
nifi.provenance.repository.rollover.size=100 MB
nifi.provenance.repository.query.threads=2
nifi.provenance.repository.index.threads=2
nifi.provenance.repository.compress.on.rollover=true
nifi.provenance.repository.always.sync=false
nifi.provenance.repository.indexed.fields=EventType, FlowFileUUID, Filename, ProcessorID, Relationship
nifi.provenance.repository.indexed.attributes=
nifi.provenance.repository.index.shard.size=500 MB
nifi.provenance.repository.max.attribute.length=65536
nifi.provenance.repository.concurrent.merge.threads=2
nifi.provenance.repository.buffer.size=100000
nifi.components.status.repository.implementation=org.apache.nifi.controller.status.history.VolatileComponentStatusRepository
nifi.components.status.repository.buffer.size=1440
nifi.components.status.snapshot.frequency=1 min
nifi.status.repository.questdb.persist.node.days=14
nifi.status.repository.questdb.persist.component.days=3
nifi.status.repository.questdb.persist.location=./status_repository
nifi.remote.input.host=
nifi.remote.input.secure=false
nifi.remote.input.socket.port=
nifi.remote.input.http.enabled=true
nifi.remote.input.http.transaction.ttl=30 sec
nifi.remote.contents.cache.expiration=30 secs
nifi.web.http.host=0.0.0.0
nifi.web.http.port={{nifi_port}}
nifi.web.http.network.interface.default=
nifi.web.https.host=
nifi.web.https.port=
nifi.web.https.network.interface.default=
nifi.web.https.application.protocols=http/1.1
nifi.web.jetty.working.directory=./work/jetty
nifi.web.jetty.threads=200
nifi.web.max.header.size=16 KB
nifi.web.proxy.context.path=
nifi.web.proxy.host=
nifi.web.max.content.size=
nifi.web.max.requests.per.second=30000
nifi.web.max.access.token.requests.per.second=25
nifi.web.request.timeout=60 secs
nifi.web.request.ip.whitelist=
nifi.web.should.send.server.version=true
nifi.web.request.log.format=%{client}a - %u %t "%r" %s %O "%{Referer}i" "%{User-Agent}i"
nifi.sensitive.props.key= [OPENSSL RAND値を入力]
nifi.sensitive.props.algorithm=NIFI_PBKDF2_AES_GCM_256
- Secure Hadoopの場合、上記の contentの値に以下の認証情報を追加し、Advanced nifi-properties-envの contentに置き換えます。
nifi.security.autoreload.enabled=false
nifi.security.autoreload.interval=10 secs
nifi.security.user.authorizer=single-user-authorizer
nifi.security.allow.anonymous.authentication=false
nifi.security.user.login.identity.provider=single-user-provider
nifi.security.user.jws.key.rotation.period=PT1H
nifi.kerberos.krb5.file=/etc/krb5.conf
nifi.kerberos.service.principal=nifi/_HOST@レルム_情報
nifi.kerberos.service.keytab.location=/etc/security/keytabs/nifi.keytab
nifi.kerberos.spnego.principal=HTTP/_HOST@レルム_情報
nifi.kerberos.spnego.keytab.location=/etc/security/keytabs/spnego.service.keytab
NiFi を使用する
Local Fileを HDFSに移動させる Data Flowを作成する方法を説明します。ローカルディレクトリに作成したファイルが NiFi Data Flowを通じて HDFSに移動する全手順と手順別の説明は、次の通りです。
1.GetFileプロセッサの作成
2.ローカルファイルの作成
3.HDFSプロセッサの作成
4.プロセッサの接続
5.実行と結果確認
6.トラブルシューティング
1.GetFileプロセッサの作成
NiFiで GetFileプロセッサを作成する方法は、次の通りです。
- マスターノードの /tmpディレクトリに以下のように nifi-testディレクトリを作成します。
mkdir /tmp/nifi-test
chown nifi /tmp/nifi-test
-
NiFi Web GUIのコンポーネントツールバーからプロセッサをドラッグしてキャンバスの上に移動します。

-
NiFi Web GUIで GetFileプロセッサを作成した後に設定します。

-
GetFileプロセッサを右クリックし、 [Configure] ボタンをクリックします。
-
[PROPERTIES] タブをクリックします。
-
Input Directory項目に nifi-testディレクトリの位置を入力し、 [APPLY] ボタンをクリックします。

2.ローカルファイルの作成
ローカル環境で任意のファイルを作成する方法は、次の通りです。以下のように viを利用して test1.txtファイルを作成します。
[irteamsu@dev-nch271-ncl nifi-test]$ vi test1.txt
3.HDFSプロセッサの作成
NiFiで HDFSプロセッサを作成する方法は、次の通りです。
- 以下のようにローカル環境にあるデータを保存する HDFSディレクトリを作成します。
sudo -u hdfs hdfs dfs -mkdir -p /user/nifi
sudo -u hdfs hdfs dfs -chown nifi /user/nifi
- NiFi Web GUIで PutHDFSプロセッサを NiFi GUIを通じてキャンバスの上に作成します。
- PutHDFSプロセッサを右クリックし、 [Configure] ボタンをクリックします。
- [PROPERTIES] タブをクリックして以下のようにプロパティを作成します。
- Hadoop Configuration Resources : /etc/hadoop/conf/core-site.xml,/etc/hadoop/conf/hdfs-site.xml
- Directory : /user/nifi
- [RELATIONSHIPS] タブをクリックし、failureと successの terminateにすべてチェックを入れます。

4.プロセッサの接続
NiFiで GetFileプロセッサと PutHDFSプロセッサを連携する方法は、次の通りです。
- NiFi Web GUIで GetFileプロセッサにマウスオーバーし、連携アイコンをドラッグして PutHDFSプロセッサに連携します。
- 完成した Data Flowは次の通りです。

5.実行と結果確認
NiFiで Data Flowを実行してローカル環境で実行結果を確認する方法は、次の通りです。
- NiFi Web GUIにアクセスします。
- GetFileプロセッサにカーソルを当てて右クリックし、 [Start] ボタンをクリックします。
- PutHDFSプロセッサにカーソルを当てて右クリックし、 [Start] ボタンをクリックします。
- /tmp/nifi-testディレクトリを確認するには、ローカル環境で以下のように入力します。
[irteamsu@dev-nch271-ncl nifi-test]$ pwd /tmp/nifi-test [irteamsu@dev-nch271-ncl nifi-test]$ ls- NiFi Data Flowによりローカル環境のファイルが HDFSに移動され、/tmp/nifi-testディレクトリにファイルが存在しません。
- HDFSに移動したファイルを確認するにはローカル環境で以下のように入力します。
[irteamsu@dev-nch271-ncl nifi-test]$ sudo -u hdfs hdfs dfs -ls /user/nifi Found 1 items -rw-r--r-- 2 nifi hdfs 4 2023-08-29 17:33 /user/nifi/test1.txt [irteamsu@dev-nch271-ncl nifi-test]$ sudo -u hdfs hdfs dfs -cat /user/nifi/test1.txt- NiFi Data Flowを通じてローカル環境のファイルが HDFSに移動しました。
ローカル環境の /tmp/nifi-testディレクトリ内にファイルが作成されると、そのファイルは自動的に HDFSに移動され、ローカル環境のファイルは削除されます。
6.トラブルシューティング
Data Flowに問題が生じた時は Data Provenanceを確認して問題を解決できます。Data Provenanceを確認する方法は、次の通りです。
- NiFi Web GUIで確認するプロセッサにカーソルを当てて右クリックします。
- [View data provenance] ボタンをクリックします。
- 以下のように Data Provenance情報を確認します。

