Hallo zusammen!

Die Aufgabe besteht darin, den oben dargestellten Flow auf N Servern zu implementieren, wobei der Flow testweise ist – es erfolgt die Generierung einer Datei und deren Versand an einen anderen NiFi-Instanz. Die Datenübertragung erfolgt über das NiFi Site to Site-Protokoll.
NiFi Site to Site (S2S) ist eine sichere und leicht konfigurierbare Methode zur Datenübertragung zwischen NiFi-Instanzen. Wie S2S funktioniert, erfahren Sie unter und es ist wichtig, nicht zu vergessen, die NiFi-Instanz so zu konfigurieren, dass S2S erlaubt ist, siehe .
In Fällen, in denen es um die Datenübertragung mittels S2S geht, wird eine Instanz als Client-Instanz und die andere als Server-Instanz bezeichnet. Die Client-Instanz sendet Daten, die Server-Instanz empfängt sie. Es gibt zwei Möglichkeiten, die Datenübertragung zwischen ihnen zu konfigurieren:
- Push. Von der Client-Instanz aus werden die Daten über eine Remote Process Group (RPG) gesendet. An der Server-Instanz werden die Daten über einen Input Port empfangen.
- Herunterladen. Der Server empfängt die Daten über RPG, der Client sendet über den Output Port.
Den Flow zur Implementierung speichern wir im Apache Registry.
Apache NiFi Registry ist ein Unterprojekt von Apache NiFi, das ein Werkzeug zur Speicherung von Flows und zur Versionsverwaltung darstellt. Sozusagen ein GIT. Informationen zur Installation, Konfiguration und Nutzung des Registrys finden Sie unter . Der Flow für die Speicherung wird in einer Prozessgruppe zusammengefasst und in diesem Format im Registry gespeichert. Darauf werden wir später in dem Artikel zurückkommen.
Zu Beginn, wenn N eine kleine Zahl ist, wird der Flow von Hand in akzeptabler Zeit bereitgestellt und aktualisiert.
Aber mit zunehmendem N gibt es immer mehr Probleme:
- Die Aktualisierung des Flows benötigt mehr Zeit. Es müssen alle Server aufgerufen werden.
- Es gibt Fehler bei der Aktualisierung der Templates. Hier wurde aktualisiert, aber dort wurde es vergessen.
- Menschliche Fehler bei der Ausführung einer großen Anzahl ähnlicher Operationen.
All dies bringt uns zu dem Punkt, dass wir den Prozess automatisieren müssen. Ich habe folgende Ansätze ausprobiert, um dieses Problem zu lösen:
- MiNiFi anstelle von NiFi verwenden.
- NiFi CLI.
- NiPyAPI.
Einsatz von MiNiFi.
— Teilprojekt Apache NiFi. MiNiFy ist ein kompakter Agent, der die gleichen Prozessoren wie NiFi verwendet und es ermöglicht, dieselben Flows zu erstellen wie in NiFi. Die Leichtgewichtigkeit des Agents wird unter anderem dadurch erreicht, dass MiNiFy keine grafische Benutzeroberfläche für die Konfiguration von Flows hat. Das Fehlen einer grafischen Benutzeroberfläche bedeutet, dass die Übertragung von Flows nach minifi eine Herausforderung darstellt. Da MiNiFy aktiv im IoT eingesetzt wird, gibt es viele Komponenten, und der Prozess der Übertragung von Flows zu den Endinstanzen von minifi muss automatisiert werden. Eine bekannte Aufgabe, oder?
Zur Lösung dieser Aufgabe kann ein weiteres Teilprojekt helfen — der MiNiFi C2 Server. Dieses Produkt dient als zentrale Anlaufstelle in der Architektur zur Bereitstellung von Konfigurationen. Wie man die Umgebung konfiguriert, ist beschrieben in auf Habré, und die Informationen sind ausreichend, um die gestellte Aufgabe zu lösen. MiNiFi aktualisiert in Verbindung mit dem C2-Server automatisch seine Konfiguration. Der einzige Nachteil dieses Ansatzes ist, dass Vorlagen auf dem C2-Server erstellt werden müssen; ein einfacher Commit im Registry reicht nicht aus.
Die im Artikel oben beschriebene Variante ist funktional und einfach umzusetzen, aber man sollte folgendes nicht vergessen:
- Minifi unterstützt nicht alle Prozessoren von NiFi.
- Die Prozessorversionen in Minifi hinken hinter den Prozessorversionen in NiFi hinterher.
Zum Zeitpunkt der Veröffentlichung war die neueste Version von NiFi – 1.9.2. Die Prozessorversion der neuesten MiNiFi-Version – 1.7.0. Prozessoren können in MiNiFi hinzugefügt werden, aber aufgrund der Versionsabweichungen zwischen NiFi- und MiNiFi-Prozessoren kann dies möglicherweise nicht funktionieren.
NiFi CLI.
Laut Das Tool auf der offiziellen Website ist ein Werkzeug zur Automatisierung der Interaktion von NiFi und NiFi Registry im Bereich der Bereitstellung von Flows oder der Prozessverwaltung. Um zu beginnen, müssen Sie dieses Tool herunterladen. .
Wir starten das Dienstprogramm
./bin/cli.sh
_ ___ _
Apache (_) .' ..](_) ,
_ .--. __ _| |_ __ )
[ `.-. | [ |'-| |-'[ | /
| | | | | | | | | | ' '
[___||__][___][___] [___]', ,'
`'
CLI v1.9.2
Geben Sie 'help' ein, um eine Liste der verfügbaren Befehle zu sehen, verwenden Sie die Tabulatortaste zur automatischen Vervollständigung.
Um den benötigten Flow aus dem Registry zu laden, müssen wir die Bucket-IDs und die Flow-ID kennen. Diese Daten können entweder über die CLI oder über die Web-Oberfläche der NiFi Registry abgerufen werden. In der Web-Oberfläche sieht das so aus:

Mit der CLI funktioniert es so:
#> registry list-buckets -u http://nifi-registry:18080
# Name Id Description
- -------------- ------------------------------------ -----------
1 test_bucket 709d387a-9ce9-4535-8546-3621efe38e96 (empty)
#> registry list-flows -b 709d387a-9ce9-4535-8546-3621efe38e96 -u http://nifi-registry:18080
# Name Id Description
- ------------ ------------------------------------ -----------
1 test_flow d27af00a-5b47-4910-89cd-9c664cd91e85
Wir starten den Import der Prozessgruppe aus dem Registry:
#> nifi pg-import -b 709d387a-9ce9-4535-8546-3621efe38e96 -f d27af00a-5b47-4910-89cd-9c664cd91e85 -fv 1 -u http://nifi:8080
7f522a13-016e-1000-e504-d5b15587f2f3
Ein entscheidender Punkt ist, dass als Host, auf den wir die Prozessgruppe anwenden, jede nifi-Instanz angegeben werden kann.
Die Prozessgruppe wurde mit angehaltenen Prozessoren hinzugefügt, diese müssen aktiviert werden.
#> nifi pg-start -pgid 7f522a13-016e-1000-e504-d5b15587f2f3 -u http://nifi:8080
Gut, die Prozessoren sind gestartet. Allerdings müssen, gemäß den Anforderungen, die NiFi-Instanzen Daten an andere Instanzen senden. Angenommen, wir haben die Push-Methode zur Datenübertragung gewählt. Um die Datenübertragung zu organisieren, muss die Übertragung im hinzugefügten Remote Process Group (RPG), die bereits in unserem Flow aktiviert ist, eingeschaltet werden.

In der Dokumentation der CLI und anderen Quellen habe ich kein Verfahren gefunden, um die Datenübertragung zu aktivieren. Wenn Sie wissen, wie das funktioniert, schreiben Sie bitte einen Kommentar.
Da wir bash verwenden und bereit sind, bis zum Ende zu gehen – lassen Sie uns eine Lösung finden! Wir können die NiFi API nutzen, um dieses Problem zu lösen. Wir verwenden die folgende Methode, wobei wir die ID aus den obigen Beispielen nehmen (in unserem Fall ist das 7f522a13-016e-1000-e504-d5b15587f2f3). Beschreibung der NiFi API-Methoden. .

Im Body muss ein JSON der folgenden Art übergeben werden:
{
"revision": {
"clientId": "value",
"version": 0,
"lastModifier": "value"
},
"state": "value",
"disconnectedNodeAcknowledged": true
}
Die Parameter, die ausgefüllt werden müssen, damit es "funktioniert":
state — Datenübertragungsstatus. Verfügbar sind TRANSMITTING zum Aktivieren der Datenübertragung und STOPPED zum Deaktivieren.
version — Prozessorversion
Die Standardversion wird beim Erstellen 0 sein, diese Parameter können jedoch mithilfe von Methoden abgerufen werden.

Für Bash-Skriptliebhaber mag diese Methode nützlich erscheinen, ich finde jedoch, das ist nicht meine Stärke — Bash-Skripte sind nicht gerade mein Favorit. Der nächste Ansatz ist meiner Meinung nach interessanter und praktischer.
NiPyAPI.
NiPyAPI — eine Bibliothek für die Programmiersprache Python zur Interaktion mit NiFi-Instanzen. enthält die erforderlichen Informationen zur Verwendung der Bibliothek. Der schnelle Einstieg wird beschrieben in auf GitHub.
Unser Skript zur Bereitstellung der Konfiguration ist ein Programm in Python. Lassen Sie uns mit dem Programmieren beginnen.
Wir konfigurieren die Einstellungen für die weitere Nutzung. Folgende Parameter werden benötigt:
nipyapi.config.nifi_config.host = 'http://nifi:8080/nifi-api' #Pfad zur NiFi-API-Instanz, auf der wir die Prozessgruppe erstellen
nipyapi.config.registry_config.host = 'http://nifi-registry:18080/nifi-registry-api' #Pfad zur NiFi-Registry-API
nipyapi.config.registry_name = 'MeinSchönerRegistry' #Name des Registrys, wie er in der NiFi-Instanz angezeigt wird
nipyapi.config.bucket_name = 'BucketName' #Name des Buckets, aus dem wir den Flow abrufen
nipyapi.config.flow_name = 'FlowName' #Name des Flows, den wir abrufen
Ich werde jetzt die Methodennamen dieser Bibliothek einfügen, die beschrieben sind. .
Wir verbinden das Registry mit der nifi-Instanz über
nipyapi.versioning.create_registry_clientAn diesem Schritt kann man auch überprüfen, ob das Registry bereits zur Instanz hinzugefügt wurde, dafür kann man die Methode nutzen
nipyapi.versioning.list_registry_clientsWir finden den Bucket für die weitere Suche nach Flow im Korb.
nipyapi.versioning.get_registry_bucketBasierend auf dem gefundenen Bucket suchen wir den Flow.
nipyapi.versioning.get_flow_in_bucketEs ist außerdem wichtig zu verstehen, ob diese Prozessgruppe bereits hinzugefügt wurde. Prozessgruppen werden nach Koordinaten platziert, und es kann vorkommen, dass sich zwei Komponenten überlappen. Ich habe das überprüft; das kann passieren 🙂 Um alle hinzugefügten Prozessgruppen zu erhalten, verwenden wir die Methode
nipyapi.canvas.list_all_process_groupsund anschließend können wir zum Beispiel nach dem Namen suchen.
Ich werde den Prozess der Template-Aktualisierung nicht beschreiben, sondern nur sagen, dass es keine Probleme mit Warteschlangen gibt, wenn in der neuen Version des Templates Prozessoren hinzugefügt werden. Wenn jedoch Prozessoren entfernt werden, können Probleme auftreten (nifi erlaubt es nicht, einen Prozessor zu löschen, wenn davor eine Warteschlange von Nachrichten besteht). Wenn Sie interessiert sind, wie ich dieses Problem gelöst habe – schreiben Sie mir bitte, und wir können darüber sprechen. Die Kontaktdaten finden Sie am Ende des Artikels. Lassen Sie uns zum Schritt der Hinzufügung der Prozessgruppe übergehen.
Bei der Fehlersuche des Skripts bin ich auf ein besonderes Merkmal gestoßen, dass nicht immer die neueste Version des Flows geladen wird. Daher empfehle ich zunächst, diese Version zu überprüfen:
nipyapi.versioning.get_latest_flow_verWir deployen die Prozessgruppe:
nipyapi.versioning.deploy_flow_versionWir starten die Prozessoren:
nipyapi.canvas.schedule_process_groupIm Abschnitt über die CLI wurde erwähnt, dass die Datenübertragung im Remote Process Group nicht automatisch aktiviert wird? Bei der Implementierung des Skripts bin ich auf dasselbe Problem gestoßen. Zu diesem Zeitpunkt gelang es mir nicht, die Datenübertragung über die API zu starten, also entschied ich mich, den Entwickler der Bibliothek NiPyAPI um Rat/Hilfe zu bitten. Der Entwickler antwortete mir, wir besprachen das Problem, und er meinte, dass er Zeit benötige, um "einige Dinge zu überprüfen". Und dann, nach ein paar Tagen, erhielt ich eine E-Mail mit einer Funktion in Python, die mein Startproblem löste!!! Zu diesem Zeitpunkt war die Version von NiPyAPI 0.13.3, und darin gab es natürlich nichts dergleichen. Aber in die vor kurzem veröffentlichte Version 0.14.0 wurde diese Funktion bereits aufgenommen. Herzlich willkommen,
nipyapi.canvas.set_remote_process_group_transmissionAlso haben wir mit der NiPyAPI-Bibliothek das Registry verbunden, den Flow implementiert und sogar die Prozessoren sowie die Datenübertragung gestartet. Jetzt können wir den Code aufbereiten, verschiedene Prüfungen und Logging hinzufügen, und das war's auch schon. Aber das ist eine ganz andere Geschichte.
Von den Optionen zur Automatisierung, die ich betrachtet habe, erscheint mir die letzte als die funktionalste. Erstens handelt es sich um Python-Code, der es ermöglicht, unterstützenden Programmcode zu integrieren und alle Vorteile der Programmiersprache zu nutzen. Zweitens entwickelt sich das NiPyAPI-Projekt aktiv weiter, und im Falle von Problemen kann man den Entwickler kontaktieren. Drittens ist NiPyAPI ein flexibleres Werkzeug für die Interaktion mit NiFi bei der Lösung komplexer Aufgaben. Zum Beispiel kann es dabei helfen festzustellen, ob die Nachrichtenwarteschlangen im Flow derzeit leer sind und ob die Prozessgruppe aktualisiert werden kann.
Das ist alles. Ich habe drei Ansätze zur Automatisierung der Bereitstellung von Flows in NiFi beschrieben, die Stolpersteine aufgezeigt, mit denen Entwickler konfrontiert werden können, und funktionierenden Code zur Automatisierung der Bereitstellung bereitgestellt. Wenn Sie sich für dieses Thema interessieren, genau wie ich —
Quelle: habr.com
