Автоматизация на доставката на flow в Apache NiFi

Здравейте на всички!

Автоматизация на доставката на flow в Apache NiFi

Задачата е следната — има flow, представен на изображението по-горе, който трябва да се разгърне на N сървъра с Apache NiFi. Flow е тестов — генерира се файл и се изпраща в друг инстанс NiFi. Предаването на данни се осъществява чрез протокола NiFi Site to Site.

NiFi Site to Site (S2S) — сигурен, лесно настраиваем начин за предаване на данни между инстанси NiFi. Как работи S2S можете да видите в документацията и е важно да не забравите да настроите инстанса NiFi да разреши S2S, вижте тук.

В случаите, когато става въпрос за предаване на данни с помощта на S2S — единият инстанс се нарича клиентски, а вторият - сървърен. Клиентският изпраща данни, а сървърният — приема. Има два начина да настроите предаването на данни между тях:

  1. Push. Данните от клиентския инстанс се изпращат с помощта на Remote Process Group (RPG). На сървърния инстанс данните се приемат с Input Port
  2. Pull. Сървърът приема данни с помощи на RPG, клиентът изпраща с Output port.


Flow за разгръщане го съхраняваме в Apache Registry.

Apache NiFi Registry — подпрожект на Apache NiFi, представляващ инструмент за съхраняване на flow и управление на версиите. Нещо подобно на GIT. Информация за инсталирането, настройката и работата с registry можете да намерите в официалната документация. Flow за съхранение се обединява в process group и така се съхранява в registry. По-късно в статията още ще се върнем на това.

В началото, когато N е малко число, flow се доставя и актуализира ръчно за разумно време.

Но с нарастващо N проблемите стават повече:

  1. за актуализиране на flow отнема повече време. Необходимо е да влезете на всички сървъри
  2. възникват грешки при актуализация на шаблони. Тук обновиха, а тук забравиха
  3. човешки грешки при извършване на голям брой еднотипни операции

Всичко това ни подтиква към автоматизиране на процеса. Опитах следните начини за решаване на тази задача:

  1. Използвайте MiNiFi вместо NiFi
  2. NiFi CLI
  3. NiPyAPI

Използване на MiNiFi

Apache MiNiFy — под проект Apache NiFi. MiNiFy — компактен агент, който използва същите процесори като NiFi, позволяващ създаването на същите потоци като в NiFi. Лекотата на агента се постига и чрез липсата на графичен интерфейс за конфигуриране на потоците. Липсата на графичен интерфейс в MiNiFy означава, че е необходимо да се реши проблема с доставката на потока до minifi. Поради активното използване на MiNiFy в IoT, компонентите са много и доставката на потоците до крайните инстанции на minifi трябва да бъде автоматизирана. Звучаща позната задача, нали?

Решението на такава задача ще бъде друг под проект — MiNiFi C2 Server. Този продукт е проектиран да бъде централна точка в архитектурата на разгръщане на конфигурации. Как да конфигурираме околната среда е описано в тази статия на Хабра и информацията е достатъчна за решаване на поставената задача. MiNiFi в комбинация с C2 server автоматично обновява конфигурацията си. Единственият недостатък на този подход е необходимостта от създаване на шаблони на C2 Server; простото добавяне в регистъра не е достатъчно.

Вариантът, описан в статията по-горе, е работещ и несложен за изпълнение, но не трябва да забравяме следното:

  1. В minifi не всички процесори от nifi са налични.
  2. Версиите на процесорите в Minifi са изостанали от версиите в NiFi.

Към момента на написване на публикацията последната версия на NiFi е 1.9.2. Версията на процесорите в последната версия на MiNiFi е 1.7.0. Процесорите могат да се добавят в MiNiFi, но поради несъответствие в версиите между процесорите NiFi и MiNiFi, това може да не сработи.

NiFi CLI

Съдя по описанието инструмента на официалния сайт, това е инструмент за автоматизация на взаимодействието между NiFI и NiFi Registry в областта на доставката на потоци или управление на процесите. За да започнете работа, е необходимо да изтеглите този инструмент. оттук.

Стартираме утилитата

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

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

За да заредим необходимия поток от регистъра, трябва да знаем идентификаторите на кошницата (bucket identifier) и самия поток (flow identifier). Тези данни могат да бъдат получени или чрез CLI, или в уеб интерфейса на NiFi Registry. В уеб интерфейса изглежда така:

Автоматизация на доставката на flow в Apache NiFi

С помощта на CLI става така:

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

Стартираме импорта на групата от процеси от регистъра:

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

Важно е — като хост, на който прилагаме групата от процеси, може да бъде посочен всеки инстанс на nifi.

Процесната група е добавена с неактивни процесори, които трябва да стартираме.

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

Чудесно, процесорите стартираха. Въпреки това, съгласно условията на задачата, трябва инстанците на NiFi да изпращат данни към други инстанси. Да предположим, че за предаване на данните на сървъра избрахме метода Push. За да организираме предаването на данни, трябва да активираме предаването на данни (Enable transmitting) на добавената Remote Process Group (RPG), която вече е включена в нашия поток.

Автоматизация на доставката на flow в Apache NiFi

В документацията на CLI и в други източници не намерих начин да се активира предаването на данни. Ако знаете как да го направите — моля, пишете в коментарите.

Тъй като имаме bash и сме готови да отидем до края — да намерим решение! Може да използваме NiFi API за решаване на този проблем. Ще използваме следния метод, ID взимаме от примери по-горе (в нашия случай това е 7f522a13-016e-1000-e504-d5b15587f2f3). Описание на методите на NiFi API. тук.

Автоматизация на доставката на flow в Apache NiFi
В тялото трябва да предадем JSON от следния вид:

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

Параметрите, които трябва да попълните, за да „заработи“:
state — статус на предаване на данни. Достъпен е TRANSMITTING за включване на предаването на данни и STOPPED за изключване.
version — версия на процесора.

версията по подразбиране ще е 0 при създаването, но тези параметри могат да се получат, използвайки метода.

Автоматизация на доставката на flow в Apache NiFi

За любителите на bash скриптове този метод може да изглежда полезен, но на мен ми е трудно — bash скриптовете не са точно моята любима част. Следващият метод ми се струва по-интересен и удобен.

NiPyAPI

NiPyAPI — библиотека за Python за взаимодействие с инстанците на NiFi. Страницата с документацията съдържа необходимата информация за работа с библиотеката. Бързото начало е описано в. проект в GitHub.

Нашият скрипт за внедряване на конфигурацията — програма на Python. Преминаваме към кодирането.
Настройваме конфигурациите за по-нататъшна работа. Нужни ще са ни следните параметри:

nipyapi.config.nifi_config.host = 'http://nifi:8080/nifi-api' #път до nifi-api инстанса, на който развиваме process group
nipyapi.config.registry_config.host = 'http://nifi-registry:18080/nifi-registry-api' #път до nifi-registry-api registry
nipyapi.config.registry_name = 'MyBeautifulRegistry' #име на registry, което ще се нарича в инстанса nifi
nipyapi.config.bucket_name = 'BucketName' #име на bucket, от който изтегляме потока
nipyapi.config.flow_name = 'FlowName' #име на потока, който изтегляме.

По-нататък ще посоча имената на методите на тази библиотека, които са описани. тук.

Свързваме registry с инстанса на nifi чрез.

nipyapi.versioning.create_registry_client

На тази стъпка можете да добавите и проверка дали регистрацията вече е добавена към инстанса, за което можете да използвате метода

nipyapi.versioning.list_registry_clients

Намираме bucket за по-нататъшно търсене на flow в кошницата

nipyapi.versioning.get_registry_bucket

По намерения bucket търсим flow

nipyapi.versioning.get_flow_in_bucket

След това е важно да разберем дали този process group вече е добавен. Process group се разполага по координати и може да се случи така, че втори компонент да се наложи върху първия. Проверих, може да се случи 🙂 За да получим всички добавени process group, използваме метода

nipyapi.canvas.list_all_process_groups

и после можем да търсим, например, по име.

Няма да описвам процеса по актуализиране на шаблона, ще кажа само, че ако в новата версия на шаблона се добавят процесори, проблеми с наличието на съобщения в опашките няма. Но ако процесорите се отстраняват, могат да възникнат проблеми (nifi не позволява да се изтрие процесор, ако пред него е натрупана опашка от съобщения). Ако ви интересува как реших този проблем — моля, напишете ми, ще обсъдим този момент. Контактите в края на статията. Преминаваме към стъпката по добавяне на process group.

При отстраняване на грешки в скрипта се сблъсках с особеност, че не винаги се извлича последната версия на flow, затова препоръчвам първо да уточните тази версия:

nipyapi.versioning.get_latest_flow_ver

Деплойваме process group:

nipyapi.versioning.deploy_flow_version

Стартираме процесорите:

nipyapi.canvas.schedule_process_group

В блока относно CLI беше посочено, че в remote process group автоматично не се включва предаване на данни? При реализирането на скрипта се сблъсках с този проблем също. В този момент не успях да стартирам предаването на данни чрез API и реших да пиша на разработчика на библиотеката NiPyAPI и да попитам за съвет/помощ. Разработчикът ми отговори, обсъдихме проблема и той написа, че му трябва време "да провери нещо". И ето, след няколко дни идва писмо, в което е написана функция на Python, решаваща моя проблем с пуска!!! В този момент версията на NiPyAPI беше 0.13.3 и в нея, разбира се, не е имало нищо такова. Но в версия 0.14.0, която излезе съвсем наскоро, тази функция вече беше включена в библиотеката. Поздравления,

nipyapi.canvas.set_remote_process_group_transmission

И така, с помощта на библиотеката NiPyAPI свързахме registry, инсталирахме flow и дори стартирахме процесори и предаване на данни. Следващата стъпка е да усъвършенстваме кода, да добавим различни проверки, логиране и всичко необходимо. Но това е съвсем друга история.

От разгледаните от мен варианти за автоматизация, последният ми се стори най-работещ. Първо, това все пак е код на python, в който можем да вграждаме помощен програмен код и да се възползваме от всички предимства на програмния език. Второ, проектът NiPyAPI активно се развива и в случай на проблеми можем да се обърнем към разработчика. Трето, NiPyAPI е по-гъвкав инструмент за взаимодействие с NiFi при решаването на сложни задачи. Например, при определяне дали опашките на съобщенията в момента са празни в flow и дали може да се актуализира групата от процеси.

Това е всичко. Описах 3 подхода за автоматизация на доставката на flow в NiFi, подводните камъни, с които могат да се сблъскат разработчиците, и предоставих работещ код за автоматизация на доставката. Ако темата ви интересува, както мен — пишете!

Източник: habr.com

Купете надежден хостинг за сайтове с защита от DDoS, VPS VDS сървъри 🔥 Купете надежден хостинг за сайтове с защита от DDoS, VPS VDS сървъри | ProHoster