Automatyzacja dostarczania flow w Apache NiFi

Cześć wszystkim!

Automatyzacja dostarczania flow w Apache NiFi

Zadanie polega na tym, aby wdrożyć flow zaprezentowany na powyższym obrazku na N serwerach z Apache NiFi. Flow jest testowe — następuje generacja pliku i jego przesyłanie do innego instansu NiFi. Przesyłanie danych odbywa się za pomocą protokołu NiFi Site to Site.

NiFi Site to Site (S2S) to bezpieczny, łatwy w konfiguracji sposób przesyłania danych między instancjami NiFi. Jak działa S2S, można zobaczyć w dokumentacji i ważne, aby nie zapomnieć skonfigurować instansu NiFi, aby umożliwić S2S, patrz tutaj.

W przypadkach, gdy mowa jest o przesyłaniu danych za pomocą S2S — jeden instans nazywany jest klientem, drugi serwerem. Klient przesyła dane, serwer je przyjmuje. Istnieją dwa sposoby skonfigurowania przesyłania danych między nimi:

  1. Push. Z instansu klienta dane są wysyłane za pomocą Remote Process Group (RPG). Na instansie serwerowym dane są odbierane za pomocą Input Port
  2. Pobierz. Serwer przyjmuje dane za pomocą RPG, klient wysyła za pomocą Output port.


Flow do wdrożenia przechowujemy w Apache Registry.

Apache NiFi Registry to podprojekt Apache NiFi, który stanowi narzędzie do przechowywania flow i zarządzania wersjami. Taki GIT. Informacje o instalacji, konfiguracji i pracy z registry można znaleźć w oficjalnej dokumentacji. Flow do przechowywania łączy się w process group i w tej formie jest przechowywane w registry. Do tego jeszcze wrócimy w artykule.

Na początku, gdy N jest małą liczbą, flow jest dostarczane i aktualizowane ręcznie w akceptowalnym czasie.

Jednak w miarę wzrostu N, pojawia się coraz więcej problemów:

  1. aktualizacja flow zajmuje więcej czasu. Trzeba wejść na wszystkie serwery
  2. pojawiają się błędy aktualizacji szablonów. Tutaj zaktualizowano, a tutaj zapomniano
  3. błędy ludzkie przy wykonywaniu dużej ilości jednorodnych operacji

Wszystko to prowadzi nas do tego, że należy zautomatyzować proces. Próbowałem następujących sposobów rozwiązania tego problemu:

  1. Użycie MiNiFi zamiast NiFi
  2. NiFi CLI
  3. NiPyAPI

Wykorzystanie MiNiFi

Apache MiNiFy — podprojekt Apache NiFi. MiNiFy to kompaktowy agent, który wykorzystuje te same procesory co NiFi, umożliwiający tworzenie tych samych przepływów, co w NiFi. Lekkość agenta osiągana jest także dzięki temu, że MiNiFy nie ma interfejsu graficznego do konfiguracji przepływów. Brak interfejsu graficznego MiNiFy oznacza, że trzeba rozwiązać problem dostarczania przepływów do instancji minifi. Ponieważ MiNiFy jest aktywnie wykorzystywane w IoT, komponentów jest wiele, a proces dostarczania przepływów do końcowych instancji minifi należy zautomatyzować. Znana sprawa, prawda?

Rozwiązaniem tego problemu jest jeszcze jeden podprojekt — MiNiFi C2 Server. Ten produkt został zaprojektowany, aby być centralnym punktem w architekturze wdrażania konfiguracji. Jak skonfigurować środowisko — opisano w w tym artykule na Habra i informacji wystarczająco dużo, aby rozwiązać postawione zadanie. MiNiFi w połączeniu z serwerem C2 automatycznie aktualizuje swoją konfigurację. Jedyne wadą tego podejścia jest to, że trzeba tworzyć szablony na serwerze C2, proste zatwierdzenie w rejestrze nie wystarczy.

Opcja opisana w powyższym artykule jest działająca i niezbyt skomplikowana do wdrożenia, ale nie należy zapominać o następującym:

  1. W minifi nie ma wszystkich procesorów z nifi.
  2. Wersje procesorów w Minifi są opóźnione w stosunku do wersji procesorów w NiFi.

W momencie pisania publikacji ostatnia wersja NiFi to 1.9.2. Wersja procesorów ostatniej wersji MiNiFi to 1.7.0. Możliwe jest dodawanie procesorów do MiNiFi, ale z powodu rozbieżności wersji między procesorami NiFi i MiNiFi może to nie zadziałać.

NiFi CLI

Zgodnie z opisaniu narzędzia na oficjalnej stronie, jest to narzędzie do automatyzacji interakcji NiFi i NiFi Registry w zakresie dostarczania przepływów lub zarządzania procesami. Aby rozpocząć pracę, to narzędzie należy pobrać. stąd.

Uruchamiamy narzędzie

. /bin/cli.sh
           _     ___  _
 Apache   (_)  .' ..](_)   ,
 _ .--.   __  _| |_  __    )
[ `.-. | [  |'-| |-'[  |  /  
|  | | |  | |  | |   | | '    '
[___||__][___][___] [___]',  ,'
                           `'
          CLI v1.9.2

Wpisz 'help', aby zobaczyć listę dostępnych poleceń, użyj tab, aby autouzupełnić.

Aby załadować potrzebny przepływ z rejestru, musimy znać identyfikatory koszyka (bucket identifier) i samego przepływu (flow identifier). Te dane można uzyskać albo przez cli, albo w interfejsie webowym rejestru NiFi. W interfejsie webowym wygląda to tak:

Automatyzacja dostarczania flow w Apache NiFi

Za pomocą CLI robi się tak:

#> 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

Uruchamiamy import grupy procesów z rejestru:

#> 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

Ważny moment — jako host, na który nakładamy grupę procesów, może być wskazany dowolny instancja nifi.

Grupa procesów została dodana z zatrzymanymi procesorami, należy je uruchomić

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

Świetnie, procesory wystartowały. Jednakże, według warunków zadania, instancje NiFi muszą wysyłać dane do innych instancji. Załóżmy, że dla przesyłania danych na serwer wybrano metodę Push. Aby zorganizować przesyłanie danych, należy na dodanej Zdalnej Grupie Procesów (RPG), która już jest w naszym przepływie, włączyć przesyłanie danych (Enable transmitting).

Automatyzacja dostarczania flow w Apache NiFi

W dokumentacji w CLI i innych źródłach nie znalazłem sposobu na włączenie przesyłania danych. Jeśli wiesz, jak to zrobić — napisz proszę w komentarzach.

Skoro mamy bash i jesteśmy gotowi iść do końca — znajdźmy rozwiązanie! Można skorzystać z API NiFi, aby rozwiązać ten problem. Użyjemy następującej metody, ID bierzemy z powyższych przykładów (w naszym przypadku to 7f522a13-016e-1000-e504-d5b15587f2f3). Opis metod API NiFi tutaj.

Automatyzacja dostarczania flow w Apache NiFi
W ciele należy przesłać JSON w następującej formie:

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

Parametry, które należy wypełnić, aby „działało”:
stan — status przesyłania danych. Dostępne TRANSMITTING dla włączenia przesyłania danych, STOPPED dla wyłączenia
version — wersja procesora

wersja domyślnie będzie równa 0 przy tworzeniu, ale te parametry można uzyskać, korzystając z metody

Automatyzacja dostarczania flow w Apache NiFi

Dla miłośników skryptów bash ta metoda może wydawać się przydatna, ale ja mam z tym trudności — skrypty bash nie są moim ulubionym narzędziem. Następna metoda wydaje mi się ciekawsza i wygodniejsza.

NiPyAPI

NiPyAPI — biblioteka dla języka Python do interakcji z instancjami NiFi. Strona z dokumentacją zawiera niezbędne informacje do pracy z biblioteką. Szybki start jest opisany w projekcie na githubie.

Nasz skrypt do wdrażania konfiguracji — program w języku Python. Przechodzimy do kodowania.
Ustawiamy konfiguracje do dalszej pracy. Potrzebne będą nam następujące parametry:

nipyapi.config.nifi_config.host = 'http://nifi:8080/nifi-api' #ścieżka do instancji nifi-api, na której rozwijamy grupę procesów
nipyapi.config.registry_config.host = 'http://nifi-registry:18080/nifi-registry-api' #ścieżka do registry nifi-registry-api
nipyapi.config.registry_name = 'MyBeutifulRegistry' #nazwa registry, jak będzie się nazywać w instancji nifi
nipyapi.config.bucket_name = 'BucketName' #nazwa bucket, z którego pobieramy przepływ
nipyapi.config.flow_name = 'FlowName' #nazwa przepływu, który pobieramy

Będę teraz wstawiać nazwy metod tej biblioteki, które są opisane tutaj.

Łączymy registry z instancją nifi za pomocą

nipyapi.versioning.create_registry_client

W tym kroku można jeszcze dodać sprawdzenie, czy rejestr już został dodany do instancji, można to zrobić za pomocą metody

nipyapi.versioning.list_registry_clients

Znalezienie bucketu do dalszego wyszukiwania flow w koszu

nipyapi.versioning.get_registry_bucket

Po znalezieniu bucketu szukamy flow

nipyapi.versioning.get_flow_in_bucket

Kolejnym krokiem jest zrozumienie, czy ta grupa procesów już została dodana. Grupa procesów jest umieszczana według koordynatów i może zdarzyć się sytuacja, w której na jednym komponencie nałożony jest drugi. Sprawdzałem, taka sytuacja może wystąpić 🙂 Aby uzyskać wszystkie dodane grupy procesów, używamy metody

nipyapi.canvas.list_all_process_groups

a następnie możemy poszukać, na przykład według nazwy.

Nie będę opisywał procesu aktualizacji szablonu, powiem tylko, że jeśli w nowej wersji szablonu dodawane są procesory, nie ma problemów z obecnością wiadomości w kolejkach. Natomiast jeśli procesory są usuwane, mogą wystąpić problemy (nifi nie pozwala usunąć procesora, jeśli przed nim zgromadziła się kolejka wiadomości). Jeśli chcesz wiedzieć, jak rozwiązałem ten problem — napisz do mnie, proszę, porozmawiamy o tym. Dane kontaktowe na końcu artykułu. Przejdźmy do kroku dodawania grupy procesów.

Podczas debugowania skryptu natknąłem się na szczegół, że nie zawsze ściągana jest najnowsza wersja flow, dlatego zalecam najpierw ustalić tę wersję:

nipyapi.versioning.get_latest_flow_ver

Deployujemy grupę procesów:

nipyapi.versioning.deploy_flow_version

Uruchamiamy procesory:

nipyapi.canvas.schedule_process_group

W blokach dotyczących CLI napisano, że w remote process group automatycznie nie włącza się przesyłanie danych? Podczas realizacji skryptu napotkałem ten sam problem. W tym czasie nie udało mi się uruchomić przesyłania danych za pomocą API, więc postanowiłem napisać do twórcy biblioteki NiPyAPI i poprosić o radę/pomoc. Programista mi odpowiedział, omówiliśmy problem i napisał, że potrzebuje czasu, aby 'sprawdzić kilka rzeczy'. A oto, po kilku dniach, otrzymuję e-mail, w którym jest funkcja w Pythonie, rozwiązująca mój problem z uruchomieniem!!! W tym czasie wersja NiPyAPI wynosiła 0.13.3 i oczywiście nic takiego nie było. Natomiast w wersji 0.14.0, która została wydana całkiem niedawno, ta funkcja już weszła w skład biblioteki. Oto ją,

nipyapi.canvas.set_remote_process_group_transmission

Za pomocą biblioteki NiPyAPI podłączyliśmy registry, wdrożyliśmy flow i uruchomiliśmy procesory oraz przesył danych. Teraz można dopracować kod, dodać różne kontrole, logowanie i tym podobne. Ale to już zupełnie inna historia.

Spośród rozważonych przeze mnie opcji automatyzacji, ostatnia wydaje się najbardziej funkcjonalna. Po pierwsze, to wciąż kod w Pythonie, w który można wbudowywać pomocniczy kod programu i korzystać ze wszystkich zalet języka programowania. Po drugie, projekt NiPyAPI aktywnie się rozwija i w razie problemów można skontaktować się z deweloperem. Po trzecie, NiPyAPI jest bardziej elastycznym narzędziem do interakcji z NiFi w rozwiązywaniu skomplikowanych zadań. Na przykład w określaniu, czy kolejki wiadomości w flow są puste i czy można aktualizować grupy procesów.

Na tym kończymy. Opisałem 3 podejścia do automatyzacji dostarczania flow w NiFi, podniosłem trudności, z jakimi może się zmierzyć programista, oraz dostarczyłem działający kod do automatyzacji dostarczania. Jeśli interesuje Cię ten temat, podobnie jak mnie — napisz!

Źródło: habr.com

Kup solidny hosting stron z ochroną przed DDoS, serwery VPS VDS 🔥 Kup solidny hosting stron z ochroną przed DDoS, serwery VPS VDS | ProHoster