Automatizarea livrării flow-ului în Apache NiFi

Salut tuturor!

Automatizarea livrării flow-ului în Apache NiFi

Sarcina constă în următoarele: există un flux, prezentat în imaginea de mai sus, care trebuie implementat pe N servere cu Apache NiFi. Fluxul este de test — se generează un fișier și se trimite într-o altă instanță NiFi. Transferul de date se realizează prin protocolul NiFi Site to Site.

NiFi Site to Site (S2S) este o metodă sigură și ușor de configurat pentru transferul de date între instanțele NiFi. Cum funcționează S2S puteți vedea în documentation și este important să nu uitați să configurați instanța NiFi pentru a permite S2S, consultați aici.

În cazurile în care este vorba despre transferul de date prin S2S, o instanță este denumită client, iar cealaltă server. Clientul trimite datele, serverul le primește. Există două modalități de a configura transferul de date între ele:

  1. Push. Din instanța client, datele sunt trimise prin Remote Process Group (RPG). Pe instanța server, datele sunt primite prin Input Port
  2. Pull. Serverul primește datele prin RPG, clientul trimite prin Output port.


Fluxul pentru implementare este stocat în Apache Registry.

Apache NiFi Registry este un subproiect Apache NiFi, care reprezintă un instrument pentru stocarea fluxurilor și gestionarea versiunilor. Este ca un GIT. Informațiile despre instalare, configurare și utilizarea registry-ului pot fi găsite în documentația oficială. Fluxul pentru stocare este combinat într-un process group și este stocat în acest mod în registry. Vom reveni la acest subiect în continuare în articol.

La început, când N este un număr mic, fluxul este livrat și actualizat manual într-un timp rezonabil.

Dar odată cu creșterea lui N, problemele devin mai multe:

  1. actualizarea fluxului durează mai mult timp. Trebuie să vă conectați la toate serverele
  2. apar erori la actualizarea șabloanelor. Aici s-a actualizat, iar aici a fost uitat
  3. erori umane în executarea unui număr mare de operațiuni repetate

Toate acestea ne conduc spre necesitatea de a automatiza procesul. Am încercat următoarele metode de a rezolva această problemă:

  1. Utilizați MiNiFi în loc de NiFi
  2. NiFi CLI
  3. NiPyAPI

Utilizarea MiNiFi

Apache MiNiFi — subproiectul Apache NiFi. MiNiFy — un agent compact care utilizează același procesor ca și NiFi, permițând crearea acelorași fluxuri ca în NiFi. Ușoritatea agentului este, de asemenea, datorată faptului că MiNiFy nu are o interfață grafică pentru configurarea fluxului. Absența interfeței grafice în MiNiFy înseamnă că este necesară soluționarea problemei livrării fluxului în minifi. Având în vedere că MiNiFy este utilizat activ în IOT, există multe componente și procesul de livrare a fluxului către instanțele finale minifi trebuie automatizat. O sarcină cunoscută, nu-i așa?

Rezolvarea unei astfel de sarcini va fi ajutată de un alt subproiect — MiNiFi C2 Server. Acest produs este destinat să fie punctul central în arhitectura desfășurării configurațiilor. Cum să configurați mediul este descris în această articole pe Habr și informațiile sunt suficiente pentru a rezolva sarcina. MiNiFi, împreună cu serverul C2, actualizează automat configurația pe el. Singurul dezavantaj al acestui abord este că trebuie să creați șabloane pe C2 Server, un simplu commit în registry nu este suficient.

Varianta descrisă în articolul de mai sus este funcțională și ușor de implementat, dar trebuie să nu uităm următoarele:

  1. În minifi nu sunt toate procesoarele din nifi
  2. Versiunile procesoarelor în Minifi sunt mai vechi decât versiunile procesoarelor din NiFi.

La momentul scrierii publicației, ultima versiune NiFi era 1.9.2. Versiunea procesoarelor ultimei versiuni MiNiFi era 1.7.0. Procesoarele pot fi adăugate în MiNiFi, dar din cauza discrepanțelor de versiune între procesoarele NiFi și MiNiFi, acest lucru poate să nu funcționeze.

NiFi CLI

Judecând după descrierea instrumentului de pe site-ul oficial, acesta este un instrument pentru automatizarea interacțiunii între NiFi și NiFi Registry în domeniul livrării fluxului sau gestionării proceselor. Pentru a începe utilizarea, acest instrument trebuie descărcat de aici.

Pornim utilitarul

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

Type 'help' to see a list of available commands, use tab to auto-complete.

Pentru a încărca fluxul necesar din registry, trebuie să cunoaștem identificatorii coșului (bucket identifier) și ai fluxului (flow identifier). Aceste date pot fi obținute fie prin cli, fie în interfața web a NiFi registry. În interfața web arată astfel:

Automatizarea livrării flow-ului în Apache NiFi

Folosind CLI, se face așa:

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

Pornim importarea grupului de procese din 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

Un aspect important este că orice instanță NiFi poate fi specificată ca gazdă pe care să aplicăm grupul de procese.

Grupul de procese a fost adăugat cu procesele oprite, trebuie să le pornim.

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

Excelent, procesele au început. Totuși, conform cerințelor, trebuie să ne asigurăm că instanțele NiFi trimit date către alte instanțe. Să presupunem că am ales metoda Push pentru a transmite datele către server. Pentru a organiza transmiterea datelor, trebuie să activăm transmiterea pe Remote Process Group (RPG) adăugat, care este deja inclus în fluxul nostru.

Automatizarea livrării flow-ului în Apache NiFi

În documentația CLI și în alte surse nu am găsit o modalitate de a activa transmiterea datelor. Dacă știți cum să faceți acest lucru, vă rog să scrieți în comentarii.

Deoarece avem bash și suntem pregătiți să mergem până la capăt, să găsim o soluție! Putem folosi API-ul NiFi pentru a rezolva această problemă. Vom folosi următoarea metodă, ID-ul luat din exemplele de mai sus (în cazul nostru este 7f522a13-016e-1000-e504-d5b15587f2f3). Descrierea metodelor NiFi API. aici.

Automatizarea livrării flow-ului în Apache NiFi
În body trebuie să trimitem JSON de forma următoare:

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

Parametrii care trebuie completati pentru a „funcționa”:
state — statusul transmiterii datelor. Disponibil TRANSMITTING pentru a activa transmiterea, STOPPED pentru a o dezactiva.
version — versiunea procesorului.

versiunea va fi 0 în mod implicit la creare, dar acești parametri pot fi obținuți utilizând metoda.

Automatizarea livrării flow-ului în Apache NiFi

Pentru iubitorii de scripturi bash, această metodă ar putea părea utilă, dar nu este preferata mea — scripturile bash nu sunt cel mai plăcut lucru pentru mine. Următoarea metodă este mai interesantă și mai convenabilă, din punctul meu de vedere.

NiPyAPI

NiPyAPI — o bibliotecă pentru limbajul Python pentru interacțiunea cu instanțele NiFi. Pagina cu documentația conține informațiile necesare pentru a lucra cu biblioteca. Quick start este descris în la care lucrez acum, trebuie să ajung la host-uri din afară, prin NAT. Folosind pentru aceasta protocoale cu criptografie avansată, nu m-a părăsit niciodată senzația că e ca și cum ai folosi un tun de război pentru vrăbii. Deoarece tunelul este folosit în mare parte doar pentru a pătrunde în NAT, traficul intern este de obicei și el criptat, toată lumea susține HTTPS. pe github.

Scriptul nostru pentru desfășurarea configurației - un program scris în Python. Să trecem la codare.
Configurăm fișierele pentru lucrul ulterior. Avem nevoie de următorii parametri:

nipyapi.config.nifi_config.host = 'http://nifi:8080/nifi-api' #calea către instanța nifi-api pe care desfășurăm grupul de procese
nipyapi.config.registry_config.host = 'http://nifi-registry:18080/nifi-registry-api' #calea către registry-ul nifi-registry-api
nipyapi.config.registry_name = 'MyBeautifulRegistry' #numele registry-ului, așa cum va apărea în instanța NiFi
nipyapi.config.bucket_name = 'BucketName' #numele bucket-ului din care preluăm fluxul
nipyapi.config.flow_name = 'FlowName' #numele fluxului pe care îl preluăm

În continuare voi insera denumirile metodelor acestei biblioteci, care sunt descrise aici.

Conectăm registry la instanța nifi folosind

nipyapi.versioning.create_registry_client

În acest pas, putem adăuga încă o verificare pentru a vedea dacă registry este deja adăugat la instanță, pentru aceasta putem folosi metoda

nipyapi.versioning.list_registry_clients

Găsim bucket-ul pentru a căuta flow-ul în coș

nipyapi.versioning.get_registry_bucket

Pe bucket-ul găsit căutăm flow-ul

nipyapi.versioning.get_flow_in_bucket

Apoi, este important să înțelegem dacă acest process group a fost deja adăugat. Process group-ul este plasat în funcție de coordonate și poate apărea situația în care un component este suprapus de un altul. Am verificat, asta se poate întâmpla 🙂. Pentru a obține toate process group-urile adăugate, folosim metoda

nipyapi.canvas.list_all_process_groups

și ulterior putem căuta, de exemplu, după nume.

Nu voi descrie procesul de actualizare a șablonului, voi spune doar că dacă în noua versiune a șablonului sunt adăugate procesoare, nu sunt probleme cu prezența mesajelor în cozi. Dar, dacă procesoarele sunt șterse, problemele pot apărea (nifi nu permite ștergerea procesorului dacă înaintea lui s-au acumulat mesaje în coadă). Dacă te interesează cum am rezolvat această problemă — scrie-mi, te rog, să discutăm acest aspect. Contactele sunt la sfârșitul articolului. Să trecem la pasul de adăugare a process group-ului.

În timpul depanării scriptului, m-am confruntat cu o caracteristică, că nu întotdeauna se trage ultima versiune a flow-ului, de aceea recomand să clarifici mai întâi această versiune:

nipyapi.versioning.get_latest_flow_ver

Deploiem process group-ul:

nipyapi.versioning.deploy_flow_version

Pornim procesoarele:

nipyapi.canvas.schedule_process_group

În secțiunea despre CLI s-a menționat că în remote process group transmiterea datelor nu se activează automat? Când am implementat scriptul, m-am confruntat și eu cu această problemă. La acel moment, nu am reușit să pornesc transmiterea datelor folosind API-ul și am decis să scriu dezvoltatorului bibliotecii NiPyAPI și să cer sfatul/ajutorul. Dezvoltatorul mi-a răspuns, am discutat problema și mi-a scris că are nevoie de timp „să verifice câteva lucruri”. Și iată, după câteva zile vine un e-mail în care este scrisă o funcție în Python, care rezolvă problema mea de pornire!!! La acel moment, versiunea NiPyAPI era 0.13.3 și, desigur, nu conținea nimic de acest gen. Însă, în versiunea 0.14.0, care a fost lansată foarte recent, această funcție a fost inclusă în biblioteca. Întâmpinați,

nipyapi.canvas.set_remote_process_group_transmission

Așadar, cu ajutorul bibliotecii NiPyAPI am conectat registry-ul, am implementat flow-ul și am pornit procesoarele și transferul de date. Acum putem îmbunătăți codul, adăugând diverse verificări, logare și alte lucruri. Dar aceasta este o cu totul altă poveste.

Dintre opțiunile de automatizare pe care le-am analizat, ultima mi s-a părut cea mai eficientă. În primul rând, este un cod scris în Python, în care putem integra cod auxiliar și beneficia de toate avantajele acestui limbaj de programare. În al doilea rând, proiectul NiPyAPI este în continuă dezvoltare și, în caz de probleme, putem contacta dezvoltatorul. În al treilea rând, NiPyAPI este un instrument mai flexibil pentru interacțiunea cu NiFi în rezolvarea sarcinilor complexe. De exemplu, în determinarea dacă coada de mesaje este acum goală în flow și dacă putem actualiza group-ul de procese.

Asta a fost tot. Am descris 3 abordări pentru automatizarea livrării flow-ului în NiFi, capcanele cu care se poate confrunta un dezvoltator și am furnizat cod funcțional pentru automatizarea livrării. Dacă v-ați interesat de acest subiect, la fel ca și mine — scrieți!

Sursa: habr.com

Cumpără un hosting fiabil pentru site-uri cu protecție DDoS, servere VPS VDS 🔥 Cumpără un hosting fiabil pentru site-uri cu protecție DDoS, servere VPS VDS | ProHoster