Automatisierung der Datenübertragung in Apache NiFi

Hallo zusammen!

Automatisierung der Datenübertragung in Apache NiFi

Die Aufgabe besteht darin, den oben dargestellten Flow auf N Servern zu implementieren, wobei Apache NiFider 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 Dokumentation. und es ist wichtig, nicht zu vergessen, die NiFi-Instanz so zu konfigurieren, dass S2S erlaubt ist, siehe hier.

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:

  1. 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.
  2. 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 offiziellen Dokumentation. 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:

  1. Die Aktualisierung des Flows benötigt mehr Zeit. Es müssen alle Server aufgerufen werden.
  2. Es gibt Fehler bei der Aktualisierung der Templates. Hier wurde aktualisiert, aber dort wurde es vergessen.
  3. 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:

  1. MiNiFi anstelle von NiFi verwenden.
  2. NiFi CLI.
  3. NiPyAPI.

Einsatz von MiNiFi.

Apache 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 in diesem Artikel 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:

  1. Minifi unterstützt nicht alle Prozessoren von NiFi.
  2. 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 Beschreibung 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. hier herunter.

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:

Automatisierung der Datenübertragung in Apache NiFi

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.

Automatisierung der Datenübertragung in Apache NiFi

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. hier.

Automatisierung der Datenübertragung in Apache NiFi
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.

Automatisierung der Datenübertragung in Apache NiFi

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. Die Dokumentationsseite enthält die erforderlichen Informationen zur Verwendung der Bibliothek. Der schnelle Einstieg wird beschrieben in Projekt 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. hier.

Wir verbinden das Registry mit der nifi-Instanz über

nipyapi.versioning.create_registry_client

An 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_clients

Wir finden den Bucket für die weitere Suche nach Flow im Korb.

nipyapi.versioning.get_registry_bucket

Basierend auf dem gefundenen Bucket suchen wir den Flow.

nipyapi.versioning.get_flow_in_bucket

Es 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_groups

und 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_ver

Wir deployen die Prozessgruppe:

nipyapi.versioning.deploy_flow_version

Wir starten die Prozessoren:

nipyapi.canvas.schedule_process_group

Im 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_transmission

Also 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 — lassen Sie es mich wissen!

Quelle: habr.com

Zuverlässiges Webhosting mit DDoS-Schutz, VPS- und VDS-Server kaufen 🔥 Zuverlässiges Webhosting mit DDoS-Schutz, VPS- und VDS-Server kaufen | ProHoster