Apache NIFI: Resumen breve de las capacidades en la práctica

Introducción

Resulta que en mi trabajo actual tuve que familiarizarme con esta tecnología. Comenzaré con una breve historia. un sistema conocido. Se entendía que este sistema conocido nos enviaría solicitudes a través de HTTP a un endpoint específico, y nosotros, por extraño que parezca, tendríamos que enviar respuestas en forma de mensajes SOAP. Parece todo simple y trivial. Por lo tanto, se requiere...

Tarea

Crear 3 servicios. El primero de ellos es el Servicio de Actualización de la Base de Datos. Este servicio, al recibir nuevos datos de un sistema externo, actualiza la base de datos y genera un archivo en formato CSV para pasarlo al siguiente sistema. Se invoca el endpoint del segundo servicio: el Servicio de Transporte a través de FTP, que recibe el archivo enviado, lo valida y lo coloca en el almacenamiento de archivos a través de FTP. El tercer servicio, el Servicio de Transferencia de Datos al Consumidor, funciona de manera asíncrona con los dos primeros. Este acepta solicitudes de un sistema externo para obtener el archivo mencionado anteriormente, toma el archivo de respuesta listo, lo modifica (actualiza los campos id, description, linkToFile) y envía la respuesta en forma de mensaje SOAP. En resumen, la situación es la siguiente: los primeros dos servicios inician su trabajo solo cuando llegan datos para su actualización. El tercer servicio funciona continuamente ya que hay muchos consumidores de información, alrededor de 1000 solicitudes para recibir datos por minuto. Los servicios están disponibles de manera continua y sus instancias se encuentran en diferentes entornos, como test, demo, preproducción y producción. A continuación se presenta el esquema de funcionamiento de estos servicios. Cabe aclarar que algunos detalles han sido simplificados para evitar complicaciones innecesarias.

Apache NIFI: Resumen breve de las capacidades en la práctica

Profundización técnica

Al planificar la solución del problema, inicialmente decidimos desarrollar aplicaciones en Java utilizando el framework Spring, con Nginx como balanceador, una base de datos Postgres y otras tecnologías, tanto técnicas como no tan técnicas. Dado que el tiempo para desarrollar la solución técnica permitía considerar otros enfoques, nos fijamos en la popular tecnología Apache NIFI en ciertos círculos. Debo decir que esta tecnología nos permitió identificar estos 3 servicios. Este artículo describirá el desarrollo del servicio de transporte de archivos y del servicio de transmisión de datos al consumidor; sin embargo, si el artículo resulta interesante, escribiré sobre el servicio de actualización de datos en la base de datos.

¿Qué es esto?

NIFI es una arquitectura distribuida para la carga y procesamiento rápido de datos, con una gran cantidad de complementos para fuentes y transformaciones, versionado de configuraciones y mucho más. Un agradable beneficio es que es muy fácil de usar. Procesos triviales, como getFile, sendHttpRequest y otros, se pueden representar en forma de cuadrados. Cada cuadrado representa un proceso, cuya interacción se puede visualizar en la imagen de abajo. En la documentación hay una excelente descripción de cómo interactuar y configurar los procesos. aquí , para aquellos que hablan ruso — aquí. La documentación explica muy bien cómo descomprimir y ejecutar NIFI, así como cómo crear procesos, los cuales son cuadrados.
La idea de escribir este artículo surgió después de una larga búsqueda y estructuración de la información obtenida en algo coherente, así como el deseo de facilitar un poco la vida a los futuros desarrolladores.

Ejemplo

Se ha considerado un ejemplo de cómo interactúan los cuadrados entre sí. El esquema general es bastante simple: recibimos una solicitud HTTP (teóricamente con un archivo en el cuerpo de la solicitud. Para demostrar las capacidades de NIFI, en este ejemplo la solicitud inicia el proceso de obtención de un archivo desde el FX local), luego enviamos de vuelta una respuesta indicando que se ha recibido la solicitud, mientras se inicia el proceso de obtención del archivo desde el FX y posteriormente se mueve a través de FTP al FX. Cabe aclarar que los procesos interactúan entre sí mediante lo que se denomina flowFile. Esta es la entidad básica en NIFI que almacena atributos y contenido. El contenido son los datos que se presentan en el archivo de flujo. Es decir, en términos simples, si recibes un archivo de un cuadrado y lo transmites a otro, el contenido será tu archivo.

Apache NIFI: Resumen breve de las capacidades en la práctica

Como puedes notar, en esta imagen se ilustra el proceso general. HandleHttpRequest recibe solicitudes, ReplaceText genera el cuerpo de respuesta, y HandleHttpResponse devuelve la respuesta. FetchFile obtiene un archivo del almacenamiento de archivos y lo pasa al cuadrado PutSftp, que coloca este archivo en FTP, en la dirección especificada. Ahora veamos este proceso en más detalle.

En este caso, la solicitud es el comienzo de todo. Revisemos sus parámetros de configuración.

Apache NIFI: Resumen breve de las capacidades en la práctica

Aquí todo es bastante trivial, excepto por StandartHttpContextMap, que es un servicio que permite enviar y recibir solicitudes. Puedes ver más detalles, incluso con ejemplos — aquí

A continuación, veamos los parámetros de configuración del cuadrado ReplaceText. Es importante prestar atención a ReplacementValue, que es lo que se devolverá al usuario como respuesta. En settings se puede ajustar el nivel de registro, los registros se pueden ver en {dónde descomprimimos nifi}/nifi-1.9.2/logs, donde también hay parámetros failure/success; basado en estos parámetros se puede regular el proceso en general. Es decir, en caso de un procesamiento exitoso del texto, se llamará al proceso de envío de respuesta al usuario; de lo contrario, simplemente registraremos el proceso no exitoso.

Apache NIFI: Resumen breve de las capacidades en la práctica

En las propiedades de HandleHttpResponse no hay nada particularmente interesante excepto el estado en caso de que la respuesta se genere con éxito.

Apache NIFI: Resumen breve de las capacidades en la práctica

Con la solicitud y respuesta aclaradas, pasemos a la obtención del archivo y su colocación en el servidor FTP. FetchFile recibe el archivo según la ruta especificada en la configuración y lo pasa al siguiente proceso.

Apache NIFI: Resumen breve de las capacidades en la práctica

Y luego el cuadrado PutSftp coloca el archivo en el almacenamiento de archivos. Los parámetros de configuración se pueden ver a continuación.

Apache NIFI: Resumen breve de las capacidades en la práctica

Es importante tener en cuenta que cada cuadrado es un proceso separado que debe ser iniciado. Hemos considerado el ejemplo más simple que no requiere ninguna personalización complicada. A continuación, examinaremos un proceso un poco más complejo, donde escribiremos un poco en Groovy.

Un ejemplo más complicado

El servicio de transmisión de datos al consumidor resultó ser un poco más complicado debido al proceso de modificación del mensaje SOAP. El proceso general se presenta en la figura a continuación.

Apache NIFI: Resumen breve de las capacidades en la práctica

Aquí la idea tampoco es muy compleja: recibimos una solicitud del consumidor, indicando que necesita datos, enviamos una respuesta confirmando que hemos recibido el mensaje, iniciamos el proceso de obtención del archivo de respuesta, luego lo editamos según una lógica determinada y, después, transferimos el archivo al consumidor en forma de mensaje SOAP al servidor.

Creo que no es necesario describir de nuevo los cuadrados que vimos arriba; pasemos directamente a los nuevos. Si necesitas editar algún archivo y los cuadrados comunes como ReplaceText no son adecuados, tendrás que escribir tu propio script. Esto se puede hacer con el cuadrado ExecuteGroovyScript. Sus configuraciones se presentan a continuación.

Apache NIFI: Resumen breve de las capacidades en la práctica

Hay dos opciones para cargar el script en este cuadrado. La primera es mediante la carga de un archivo con el script. La segunda es insertando el script en scriptBody. Hasta donde sé, el cuadrado executeScript admite varios lenguajes de programación, uno de ellos es groovy. Decepcionaré a los desarrolladores de Java: no se pueden escribir scripts en estos cuadrados en Java. Para aquellos que realmente lo desean, deberán crear su propio cuadrado personalizado y integrarlo en el sistema NIFI. Todo este proceso implica un considerable esfuerzo, del cual no nos ocuparemos en este artículo. Elegí el lenguaje groovy. A continuación, se presenta un script de prueba que simplemente actualiza incrementalmente el id en el mensaje SOAP. Es importante destacar que toma el archivo de flowFile, lo actualiza y no hay que olvidar devolverlo actualizado. También es importante mencionar que no todas las bibliotecas están conectadas. Puede suceder que aún así tenga que importar una de las librerías. Otro inconveniente es que depurar el script en este cuadrado es bastante complicado. Existe una forma de conectarse a la JVM de NIFI y comenzar el proceso de depuración. Personalmente, ejecuté una aplicación localmente e imité la recepción de un archivo de sesión. También realicé la depuración localmente. Los errores que surgen al cargar el script son bastante fáciles de buscar en Google y NIFI los escribe en el registro.

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 durant la procesión de validate.groovy", e)
    session.transfer(flowFile, REL_FAILURE)
}

De hecho, aquí termina la personalización del cuadrado. Luego, el archivo actualizado se envía al cuadrado que se encarga de enviar el archivo al servidor. A continuación, se presentan las configuraciones de este cuadrado.

Apache NIFI: Resumen breve de las capacidades en la práctica

Describimos el método por el cual se transmitirá el mensaje SOAP. Especificamos a dónde. A continuación, hay que indicar que se trata de un mensaje SOAP.

Apache NIFI: Resumen breve de las capacidades en la práctica

Agregamos varias propiedades como host y acción (soapAction). Guardamos, comprobamos. Para más detalles sobre cómo enviar solicitudes SOAP, puedes consultar. aquí

Hemos revisado varias opciones de uso de procesos NIFI. Cómo interactúan y cuál es su beneficio real. Los ejemplos considerados son pruebas y difieren un poco de lo que realmente se utiliza en producción. Espero que este artículo sea algo útil para los desarrolladores. Gracias por su atención. Si hay alguna pregunta, no duden en preguntar. Intentaré responder.

Fuente: habr.com

Compra un hosting fiable para sitios web con protección contra DDoS, servidores VPS VDS 🔥 Compra un hosting fiable para sitios web con protección contra DDoS, servidores VPS VDS | ProHoster