Automatisierung der Flow-Zustellung in Apache NiFi

Hallo zusammen!

Automatisierung der Flow-Zustellung in Apache NiFi

Die Aufgabe besteht darin, einen Flow, der oben im Bild dargestellt ist, auf N Servern zu verteilen, Apache NiFi. Der Flow ist ein Test – es findet eine Dateigenerierung statt und diese wird an eine andere NiFi-Instanz gesendet. Der Datenaustausch erfolgt über das NiFi Site-to-Site-Protokoll.

NiFi Site-to-Site (S2S) ist eine sichere, einfach konfigurierbare Methode zur Datenübertragung zwischen NiFi-Instanzen. Wie S2S funktioniert, erfahren Sie in Dokumentation und es ist wichtig, die NiFi-Instanz so zu konfigurieren, dass S2S erlaubt ist. Siehe hier.

In Fällen, in denen Daten über S2S übertragen werden, wird eine Instanz als Client und die andere als Server bezeichnet. Der Client sendet Daten, der Server empfängt sie. Es gibt zwei Möglichkeiten, die Datenübertragung zwischen ihnen zu konfigurieren:

  1. Push. Vom Client-Instanz aus werden die Daten über eine Remote Process Group (RPG) gesendet. Auf der Server-Instanz werden die Daten über einen Input Port empfangen.
  2. Pull. Der Server empfängt die Daten über RPG, der Client sendet über den Output Port.


Den Flow zur Verteilung speichern wir im Apache Registry.

Apache NiFi Registry ist ein Unterprojekt von Apache NiFi und bietet ein Werkzeug zur Speicherung von Flows und zur Versionsverwaltung. So etwas wie GIT. Informationen zur Installation, Konfiguration und zur Arbeit mit dem Registry finden Sie in offiziellen Dokumentation. Die Flows zur Speicherung werden in einer Process Group zusammengefasst und in dieser Form im Registry gespeichert. Darauf werden wir im weiteren Verlauf des Artikels zurückkommen.

Zu Beginn, wenn N eine kleine Zahl ist, wird der Flow von Hand in akzeptabler Zeit bereitgestellt und aktualisiert.

Mit dem Anstieg von N werden jedoch mehr Probleme sichtbar:

  1. Die Aktualisierung des Flows benötigt mehr Zeit. Es müssen alle Server aufgerufen werden.
  2. Es treten Fehler bei der Aktualisierung der Vorlagen auf. Hier wurde aktualisiert, aber dort wurde vergessen.
  3. Menschliche Fehler beim Ausführen einer großen Anzahl ähnlicher Operationen.

All dies führt uns zu der Erkenntnis, dass der Prozess automatisiert werden muss. Ich habe folgende Ansätze zur Lösung dieses Problems ausprobiert:

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

Verwendung von MiNiFi

Apache MiNiFi — Teilprojekt Apache NiFi. MiNiFy ist ein kompakter Agent, der dieselben Prozessoren wie NiFi verwendet und es ermöglicht, dieselben Flows wie in NiFi zu erstellen. Die Leichtgewichtigkeit des Agents wird unter anderem dadurch erreicht, dass MiNiFy keine grafische Benutzeroberfläche zur Konfiguration von Flows hat. Das Fehlen einer grafischen Benutzeroberfläche bei MiNiFy bedeutet, dass das Problem der Bereitstellung von Flows in Minifi gelöst werden muss. Da MiNiFy aktiv im IoT eingesetzt wird, gibt es viele Komponenten, und der Prozess der Bereitstellung von Flows zu den Endinstanzen von Minifi muss automatisiert werden. Eine vertraute Aufgabe, nicht wahr?

Um eine solche Aufgabe zu lösen, hilft ein weiteres Teilprojekt — MiNiFi C2 Server. Dieses Produkt ist dafür gedacht, der zentrale Punkt in der Architektur der Rollout-Konfigurationen zu sein. Wie man die Umgebung konfiguriert, wird in diesem Artikel auf Habré beschrieben, und die Informationen sind ausreichend, um die gestellte Aufgabe zu lösen. MiNiFi aktualisiert im Zusammenspiel mit dem C2 Server automatisch die Konfiguration bei sich. Der einzige Nachteil dieses Ansatzes ist, dass man Vorlagen auf dem C2 Server erstellen muss, ein einfacher Commit im Registry reicht nicht aus.

Die in dem obigen Artikel beschriebene Variante ist funktional und nicht schwer umzusetzen, aber man sollte Folgendes nicht vergessen:

  1. In Minifi sind nicht alle Prozessoren aus NiFi vorhanden.
  2. Die Versionen der Prozessoren in Minifi hinken den Versionen der Prozessoren in NiFi hinterher.

Zum Zeitpunkt der Veröffentlichung war die neueste Version von NiFi — 1.9.2. Die Version der Prozessoren der neuesten MiNiFi-Version ist 1.7.0. Prozessoren können in MiNiFi hinzugefügt werden, aber aufgrund der Versionsabweichungen zwischen den Prozessoren von NiFi und MiNiFi kann dies möglicherweise nicht funktionieren.

NiFi CLI

einem Bericht Beschreibung des Werkzeugs auf der offiziellen Website, es ist ein Werkzeug zur Automatisierung der Interaktion zwischen NiFi und NiFi Registry im Bereich der Bereitstellung von Flows oder der Verwaltung von Prozessen. Um mit der Arbeit zu beginnen, muss dieses Werkzeug heruntergeladen werden. von hier.

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 erforderlichen Flow aus dem Registry zu laden, müssen wir die Identifikatoren des Buckets (bucket identifier) und des Flows (flow identifier) kennen. Diese Informationen können entweder über die CLI oder im Webinterface der NiFi Registry abgerufen werden. Im Webinterface sieht es so aus:

Automatisierung der Flow-Zustellung in Apache NiFi

So wird es mit der CLI gemacht:

#> 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 des Process Groups 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 wichtiger Punkt — als Host, auf den wir das Process Group anwenden, kann jede Instanz von NiFi angegeben werden.

Die Prozessgruppe wurde mit gestoppten Prozessoren hinzugefügt, sie müssen gestartet werden.

#> nifi pg-start -pgid 7f522a13-016e-1000-e504-d5b15587f2f3 -u http://nifi:8080

Ausgezeichnet, die Prozessoren haben gestartet. Allerdings müssen wir gemäß den Bedingungen der Aufgabe sicherstellen, dass die NiFi-Instanzen Daten an andere Instanzen senden. Angenommen, wir haben die Push-Methode für die Datenübertragung auf den Server gewählt. Um die Datenübertragung zu organisieren, müssen wir die Übertragung auf der hinzugefügten Remote Process Group (RPG), die bereits in unseren Flow integriert ist, aktivieren (Enable transmitting).

Automatisierung der Flow-Zustellung in Apache NiFi

In der Dokumentation, im CLI und in anderen Quellen habe ich keinen Weg gefunden, um die Datenübertragung zu aktivieren. Wenn Sie wissen, wie das geht - bitte schreiben Sie es in die Kommentare.

Da wir Bash haben und bereit sind, bis zum Ende zu gehen - lassen Sie uns einen Ausweg finden! Wir können die NiFi API nutzen, um dieses Problem zu lösen. Wir verwenden die folgende Methode, ID entnehmen wir den obigen Beispielen (in unserem Fall ist es 7f522a13-016e-1000-e504-d5b15587f2f3). Beschreibung der Methoden der NiFi API. hier.

Automatisierung der Flow-Zustellung in Apache NiFi
Im Body müssen wir JSON im folgenden Format übergeben:

{
    "revision": {
	    "clientId": "value",
	    "version": 0,
	    "lastModifier": "value"
	},
    "state": "value",
    "disconnectedNodeAcknowledged": true
}

Die Parameter, die ausgefüllt werden müssen, damit es "funktioniert":
state — Status der Datenübertragung. Verfügbar ist TRANSMITTING zum Aktivieren der Datenübertragung, STOPPED zum Deaktivieren.
version — Version des Prozessors.

Die Version ist standardmäßig 0 bei der Erstellung, aber diese Parameter können mit der Methode abgerufen werden.

Automatisierung der Flow-Zustellung in Apache NiFi

Für Fans von Bash-Skripten mag diese Methode geeignet erscheinen, aber ich finde es schwierig - Bash-Skripte sind nicht mein Favorit. Die nächste Methode ist interessanter und praktischer aus meiner Sicht.

NiPyAPI

NiPyAPI - eine Bibliothek für die Programmiersprache Python zur Interaktion mit NiFi-Instanzen. Die Seite mit der Dokumentation enthält die notwendigen Informationen zur Arbeit mit der Bibliothek. Schneller Einstieg ist beschrieben in dem Projekt auf github.

Unser Skript zur Bereitstellung der Konfiguration ist ein Programm in der Sprache Python. Lassen Sie uns mit dem Codieren beginnen.
Wir konfigurieren die Konfigurationen für die weitere Arbeit. Wir benötigen die folgenden Parameter:

nipyapi.config.nifi_config.host = 'http://nifi:8080/nifi-api' #Pfad zur nifi-api Instanz, auf der wir die Prozessgruppe bereitstellen
nipyapi.config.registry_config.host = 'http://nifi-registry:18080/nifi-registry-api' #Pfad zur nifi-registry-api Registry
nipyapi.config.registry_name = 'MyBeutifulRegistry' #Name der Registry, wie sie in der NiFi-Instanz genannt 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 Methoden dieser Bibliothek einfügen, die beschrieben sind hier.

Verbinden Sie das Registry mit der NiFi-Instanz mithilfe von

nipyapi.versioning.create_registry_client

An dieser Stelle kann auch überprüft werden, ob das Registry bereits zur Instanz hinzugefügt wurde. Dazu kann die Methode verwendet werden

nipyapi.versioning.list_registry_clients

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

nipyapi.versioning.get_registry_bucket

Wir suchen den Flow im gefundenen Bucket

nipyapi.versioning.get_flow_in_bucket

Darüber hinaus ist es wichtig zu verstehen, ob diese Prozessgruppe bereits hinzugefügt wurde. Eine Prozessgruppe wird an den Koordinaten platziert und es kann passieren, dass sich eine zweite über eine vorhandene Komponente legt. Ich habe das überprüft, das kann vorkommen 🙂 Um alle hinzugefügten Prozessgruppen zu erhalten, verwenden wir die Methode

nipyapi.canvas.list_all_process_groups

und können dann zum Beispiel nach dem Namen suchen.

Ich werde den Prozess der Aktualisierung des Templates nicht beschreiben; ich sage nur, dass es keine Probleme mit der Anwesenheit von Nachrichten in den Warteschlangen gibt, wenn in der neuen Version des Templates Prozessoren hinzugefügt werden. Wenn Prozessoren jedoch entfernt werden, können Probleme auftreten (nifi erlaubt es nicht, einen Prozessor zu löschen, wenn sich davor eine Warteschlange von Nachrichten angesammelt hat). Wenn Sie interessiert sind, wie ich dieses Problem gelöst habe - schreiben Sie mir bitte, wir werden diesen Punkt besprechen. Kontaktdaten am Ende des Artikels. Lassen Sie uns zum Schritt des Hinzufügens einer Prozessgruppe übergehen.

Während der Fehlersuche des Skripts stieß ich auf das Problem, dass nicht immer die letzte Version des Flows geladen wird. Daher empfehle ich, diese Version zunächst 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 in der Remote-Prozessgruppe automatisch nicht aktiviert ist? Bei der Umsetzung des Skripts stieß ich auch auf dieses Problem. Zu diesem Zeitpunkt gelang es mir nicht, die Datenübertragung über die API zu starten, und ich beschloss, den Entwickler der NiPyAPI-Bibliothek zu kontaktieren und um Rat/Hilfe zu bitten. Der Entwickler antwortete mir, wir diskutierten das Problem, und er sagte, dass er Zeit brauche, um "einige Dinge zu überprüfen". Und bald nach ein paar Tagen kam eine E-Mail, in der eine Python-Funktion beschrieben wurde, die mein Problem beim Starten löst!!! Zu diesem Zeitpunkt war die Version von NiPyAPI 0.13.3 und hatte natürlich nichts dergleichen. Aber in der Version 0.14.0, die vor kurzem herauskam, war diese Funktion bereits Bestandteil der Bibliothek. Willkommen,

nipyapi.canvas.set_remote_process_group_transmission

Mit der NiPyAPI-Bibliothek haben wir das Registry verbunden, den Flow implementiert und sogar die Prozessoren gestartet sowie die Datenübertragung durchgeführt. Jetzt können wir den Code verbessern, verschiedene Überprüfungen und Protokollierungen hinzufügen und all das. Aber das ist eine ganz andere Geschichte.

Von den von mir betrachteten Automatisierungsoptionen erschien mir die letzte als die funktionalste. Erstens handelt es sich immer noch um Code in Python, in den unterstützender Programmcode integriert werden kann und von dem alle Vorteile der Programmiersprache genutzt werden können. Zweitens entwickelt sich das NiPyAPI-Projekt aktiv weiter, und im Falle von Problemen kann der Entwickler kontaktiert werden. Drittens ist NiPyAPI ein flexibleres Werkzeug für die Interaktion mit NiFi bei komplexen Aufgaben. Zum Beispiel bei der Feststellung, 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 Fallstricke, mit denen Entwickler konfrontiert werden können, angesprochen und funktionierenden Code zur Automatisierung der Bereitstellung bereitgestellt. Wenn Sie sich für dieses Thema interessieren, wie ich — schreiben Sie mir!

Quelle: habr.com

60GB SSD 8Gb DDR4