Въведение
Така се случи, че на текущото ми работно място ми се наложи да се запозная с тази технология. Ще започна с кратка предистория. На поредната среща, на нашия екип беше казано, че трябва да създадем интеграция с известна система. Под интеграцията се подразбираше, че тази известна система ще ни изпраща заявки през HTTP на определен ендпойнт, а ние, колкото и странно да звучи, ще изпращаме обратно отговори под формата на SOAP съобщения. Изглежда всичко е просто и тривиално. Оттук следва, че трябва да...
Задача
Създадем 3 услуги. Първата от тях — Услуга за обновление на БД. Тази услуга, при получаване на нови данни от външната система, обновява данните в базата данни и генерира определен файл в CSV формат, за да го предаде на следващата система. Извиква се ендпойнтът на втората услуга — Услуга за транспортировка през FTP, която получава предадения файл, валидира го и го съхранява във файлово хранилище през FTP. Третата услуга — Услуга за предаване на данни на потребителя, работи асинхронно с първите две. Тя приема заявка от външната система за получаване на файла, за който ставаше дума по-горе, взема готовия файл с отговор, модифицира го (обновява полетата id, description, linkToFile) и изпраща отговор под формата на SOAP съобщение. Тоест, общата картина е следната: първите две услуги започват работа единствено когато получат данни за обновление. Третата услуга работи постоянно, тъй като потребителите на информация са многобройни, с около 1000 заявки за получаване на данни в минута. Услугите са постоянно достъпни и техните инстанции се намират в различни среди, като тест, демо, препрод и прод. По-долу е представена схема на работата на тези услуги. Веднага уточнявам, че някои детайли са опростени, за да се избегне ненужна сложност.

Техническо задълбочаване
Когато планирахме решението на задачата, първо решихме да направим приложения на Java с помощта на Spring framework, балансировщик Nginx, база данни Postgres и други технически и не толкова технически неща. Понеже времето за разработване на техническото решение позволяваше да разгледаме и други подходи, вниманието ни бе насочено към модната в определени кръгове технология Apache NIFI. Веднага ще кажа, че тази технология ни позволи да забележим тези 3 услуги. В тази статия ще опишем разработката на услугата за пренос на файл и услугата за предаване на данни на потребителя, но ако статията се хареса, ще напиша и за услугата за актуализиране на данни в БД.
Какво е това
NIFI е разпределена архитектура за бързо паралелно зареждане и обработка на данни, предлагаща множество плъгини за източници и трансформации, версиониране на конфигурации и много други. Приятен бонус е, че е много лесен за употреба. Тривиалните процеси, като getFile, sendHttpRequest и други, могат да се представят под формата на квадратчета. Всяко квадратче представлява определен процес, взаимодействието на който може да се види на изображението по-долу. По-подробна документация за взаимодействието при настройването на процесите е написана , за тези, които говорят руски — . В документацията е чудесно описано как да разархивирате и стартирате NIFI, както и как да създавате процеси, те също са квадратчета.
Идеята за написването на статията възникна след продължителни търсения и структуриране на получената информация в нещо осъзнато, а също и желанието да улесня бъдещите разработчици.
Пример
Разгледан е пример как взаимодействат квадратите помежду си. Общата схема е доста проста: Първо получаваме HTTP запитване (теоретично с файл в тялото на запитването. За демонстрация на възможностите на NIFI, в този пример запитването стартира процеса на получаване на файл от локалното FP), след което обратно изпращаме отговор, че запитването е получено, като паралелно се стартира процес на получаване на файл от FP и след това процес на прехвърляне на файла чрез FTP в FP. Струва си да се поясни, че процесите взаимодействат помежду си чрез така наречения flowFile. Това е основната същност в NIFI, която съхранява атрибути и съдържание. Съдържанието — данни, представени във файлов поток. Тоест, грубо казано, ако получите файл от един квадрат и го предадете на друг, съдържанието ще бъде вашият файл.

Както може да се забележи — на тази схема е показан общият процес. HandleHttpRequest — приема запитвания, ReplaceText — генерира тялото на отговора, HandleHttpResponse — изпраща отговора. FetchFile — получава файл от файловото хранилище и го предава на квадрата PutSftp — поставя този файл на FTP, на указан адрес. Сега по-подробно за този процес.
В този случай — request е всичко начало. Нека да разгледаме параметрите му за конфигурация.

Тук всичко е доста тривиално с изключение на StandartHttpContextMap — това е вид услуга, която позволява изпращането и получаването на запитвания. По-подробно, дори с примери, може да се види тук —
Следващата стъпка е да разгледаме параметрите за конфигурация на квадрата ReplaceText. Тук е важно да се обърне внимание на ReplacementValue — това е това, което ще се върне на потребителя под формата на отговор. В настройките може да се регулира нивото на логване, логовете може да се видят {къде разпакувахме nifi}/nifi-1.9.2/logs, там има и параметри failure/success — въз основа на тези параметри може да се регулира процесът като цяло. Тоест, в случай на успешно обработване на текста — ще бъде извикан процес на изпращане на отговор на потребителя, а в противен случай просто ще логнем неуспешния процес.

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

С разглеждането на запитването и отговора приключихме — да преминем напред към получаване на файла и поставянето му на FTP сървър. FetchFile — получава файл по указан в настройките път и го предава на следващия процес.

И по-нататък квадратът PutSftp поставя файла в хранилището. Конфигурационните параметри могат да бъдат видяни по-долу.

Важно е да се обърне внимание на факта, че всеки квадрат представлява отделен процес, който трябва да бъде стартиран. Разгледахме най-простия пример, който не изисква сложна персонализация. Сега ще разгледаме малко по-сложен процес, при който ще напишем малко по грозно.
По-сложен пример
Сервизът за предаване на данни на потребителя е малко по-сложен поради процеса на модификация на SOAP съобщението. Общият процес е представен на изображението по-долу.

Идеята тук също не е особено сложна: получихме запитване от потребителя, че му трябват данни, изпратихме отговор, че сме получили съобщението, стартирахме процеса на получаване на файла с отговора, след което го редактирахме с определена логика, след което предадохме файла на потребителя под формата на SOAP съобщение на сървъра.
Мисля, че няма нужда да описвам отново квадратите, които видяхме по-горе — ще преминем направо към новите. Ако ви е нужно да редактирате някой файл и обикновените квадрати тип ReplaceText не са подходящи, ще ви се наложи да напишете свой собствен скрипт. Можете да направите това с помощта на квадрата ExecuteGroogyScript. Настройките му са представени по-долу.

Има два варианта за зареждане на скрипта в този квадрат. Първият е чрез зареждане на файл със скрипта. Вторият е чрез вмъкване на скрипта в scriptBody. Доколкото знам, квадратът executeScript поддържа няколко езика за програмиране - един от тях е groovy. За съжаление на java разработчиците - не можете да пишете скриптове на java в такива квадрати. За тези, които много искат - трябва да създадете свой собствен персонализиран квадрат и да го добавите в системата NIFI. Цялата тази операция се придружена от доста дълги танци с бубна, с които няма да се занимаваме в рамките на тази статия. Избрах езика groovy. По-долу е представен тестов скрипт, който просто инкрементално обновява id в SOAP съобщението. Важно е да се отбележи. Вие взимате файла от flowFile, обновявате го, и не трябва да забравяте, че трябва да го върнете обратно там, обновен. Също така е важно да се отбележи, че не всички библиотеки са свързани. Може да се случи така, че все пак ще трябва да импортирате една от библиотеките. Друг минус е, че скриптът в този квадрат е доста труден за отстраняване на грешки. Има начин да се свържете с JVM NIFI и да започнете процеса на отстраняване на грешки. Лично аз стартирах локално приложение и симулирах получаване на файл от сесия. Отстраняването на грешки също извършвах локално. Грешките, които се появяват при зареждане на скрипта, са доста лесни за търсене в Google и NIFI ги записва в логовете.
import org.apache.commons.io.IOUtils
import groovy.xml.XmlUtil
import java.nio.charset.*
import groovy.xml.StreamingMarkupBuilder
def flowFile = session.get()
if (!flowFile) return
try {
flowFile = session.write(flowFile, { inputStream, outputStream ->
String result = IOUtils.toString(inputStream, "UTF-8");
def recordIn = new XmlSlurper().parseText(result)
def element = recordIn.depthFirst().find {
it.name() == 'id'
}
def newId = Integer.parseInt(element.toString()) + 1
def recordOut = new XmlSlurper().parseText(result)
recordOut.Body.ClientMessage.RequestMessage.RequestContent.content.MessagePrimaryContent.ResponseBody.id = newId
def res = new StreamingMarkupBuilder().bind { mkp.yield recordOut }.toString()
outputStream.write(res.getBytes(StandardCharsets.UTF_8))
} as StreamCallback)
session.transfer(flowFile, REL_SUCCESS)
}
catch(Exception e) {
log.error("Error during processing of validate.groovy", e)
session.transfer(flowFile, REL_FAILURE)
}Същност, тук приключва персонализацията на квадрата. След това обновеният файл се предава на квадрата, който се занимава с изпращането на файла на сървъра. По-долу са представени настройките на този квадрат.

Описваме метода, по който ще се предава SOAP съобщението. Пишем накъде. След това трябва да укажем, че това е именно SOAP.

Добавяме няколко свойства като хост и действие (soapAction). Запазваме, проверяваме. За повече подробности относно изпращането на SOAP заявки можете да видите.
Разгледахме няколко варианта за използване на процесите NIFI. Как взаимодействат помежду си и каква е реалната полза от тях. Обсъдените примери са тестови и леко се различават от реалните в продукция. Надявам се, тази статия да е малко полезна за разработчиците. Благодаря за вниманието. Ако имате въпроси, не се колебайте да пишете. Ще се постарая да отговоря.
Източник: habr.com
