Cześć wszystkim!

Zadanie polega na tym, aby wdrożyć flow zaprezentowany na powyższym obrazku na N serwerach z . 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 i ważne, aby nie zapomnieć skonfigurować instansu NiFi, aby umożliwić S2S, patrz .
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:
- Push. Z instansu klienta dane są wysyłane za pomocą Remote Process Group (RPG). Na instansie serwerowym dane są odbierane za pomocą Input Port
- 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 . 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:
- aktualizacja flow zajmuje więcej czasu. Trzeba wejść na wszystkie serwery
- pojawiają się błędy aktualizacji szablonów. Tutaj zaktualizowano, a tutaj zapomniano
- 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:
- Użycie MiNiFi zamiast NiFi
- NiFi CLI
- NiPyAPI
Wykorzystanie MiNiFi
— 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 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:
- W minifi nie ma wszystkich procesorów z nifi.
- 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 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ć. .
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:

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

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 .

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

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. zawiera niezbędne informacje do pracy z biblioteką. Szybki start jest opisany w 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 .
Łączymy registry z instancją nifi za pomocą
nipyapi.versioning.create_registry_clientW 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_clientsZnalezienie bucketu do dalszego wyszukiwania flow w koszu
nipyapi.versioning.get_registry_bucketPo znalezieniu bucketu szukamy flow
nipyapi.versioning.get_flow_in_bucketKolejnym 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_groupsa 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_verDeployujemy grupę procesów:
nipyapi.versioning.deploy_flow_versionUruchamiamy procesory:
nipyapi.canvas.schedule_process_groupW 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_transmissionZa 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 —
Źródło: habr.com
