Apache NiFi
Introducció
Apache NiFi és un projecte d'Apache (desenvolupat inicialment en Java per la NSA) que planteja un sistema distribuït dedicat a ingestar i transformar dades mitjançant el paradigma streaming, amb una interfície gràfica de disseny visual de fluxos.
Característiques principals:
- Projecte open source, multiplataforma.
- Fluxos de dades escalables, amb més de 300 connectors i processadors personalitzats.
- Ingesta de dades en streaming, amb el paradigma de programació basat en fluxos.
- Defineix les aplicacions com a grafs de processos dirigits (DAG) mitjançant connexions que transporten missatges.
- Lliurament garantit, sense pèrdua de dades.
- Escalabilitat horitzontal mitjançant un clúster de màquines, política d'usuaris (LDAP), validador de configuracions i llinatge i procedència de la dada integrats.
Casos d'ús: transferències de dades entre sistemes (JSON a una base de dades, FTP a Hadoop...), preparació i enriquiment de dades, encaminament segons característiques i prioritats, i conversió de formats. No és apropiat per a processament i còmput distribuït complex, operacions de streaming amb joins/agregacions elaborades, ni processament complex d'esdeveniments (per a això, vegeu Apache Kafka i Apache Spark).
Posada en marxa amb Docker
docker run --name nifi -p 8443:8443 -d \
-e SINGLE_USER_CREDENTIALS_USERNAME=nifi \
-e SINGLE_USER_CREDENTIALS_PASSWORD=nifinifinifi \
-e NIFI_JVM_HEAP_MAX=2g \
apache/nifi:latest
Un cop iniciat (cal esperar un parell de minuts), s'accedeix a https://localhost:8443/nifi. La interfície té quatre zones principals: el menú superior (processadors, ports d'entrada/sortida...), la barra inferior amb l'estat d'execució, el quadre Navigate per fer zoom, i l'àrea de treball drag&drop.
Components
Flowfile (FF): la unitat de dada que viatja pel flux, persistida en disc en crear-se (realment un punter a la dada, per accelerar-ne el rendiment). Es compon de contingut (la dada) i atributs (metadades clau/valor).
Processador: executa una transformació o regla sobre les dades, generant un nou FF a partir de l'FF d'entrada. Tots els processadors s'executen en paral·lel i es poden distribuir entre nodes d'un clúster. NiFi n'ofereix més de 300 predefinits, agrupats per família: transformació (ReplaceText, JoltTransformJSON), enrutat (RouteOnAttribute, RouteOnContent), accés a BD (ExecuteSQL, PutSQL), extracció d'atributs (EvaluateJsonPath, UpdateAttribute), ingestió (GetFile, GetFTP, GetHTTP, GetHDFS), enviament (PutFile, PutFTP, PutKafka, PutMongo) i divisió/agregació (SplitText, SplitJson, MergeContent).
Connector: cua dirigida que uneix la sortida d'un processador amb l'entrada d'un altre (o d'ell mateix, per a reintents), transportant els FF encara no processats amb una prioritat configurable (FIFO, LIFO...). Es distingeixen les relacions success (l'FF que torna un processador en acabar correctament) i failure (quan la tasca ha fallat).
Autoterminar les relacions
Si s'obliden connexions sense connectar o relacions sense autoterminar, els processadors implicats no es podran iniciar.
Cas 1 — Movent dades
Flux mínim per moure un fitxer d'un directori a un altre: GetFile (propietat Input Directory) → connector success → PutFile (propietat directory de sortida, amb ambdues relacions success/failure autoterminades).
Gestió d'errors i reanomenat: per no perdre fitxers repetits, es canvia la propietat Conflict Resolution Strategy de PutFile a fail (en lloc d'ignore o replace), i es redirigeix la relació failure cap a un segon PutFile que emmagatzemi els fitxers conflictius en una carpeta a part. Per no perdre'n l'històric, s'afegeix abans un UpdateAttribute que reanomena el fitxer amb un prefix temporal usant el NiFi Expression Language: ${now():toNumber()}-${filename}.
La icona Data Provenance permet consultar, per a cada FF, el llinatge complet: en quins passos ha estat, amb quin contingut i en quin instant.
Cas 2 — Treballant amb atributs
GenerateFlowFile crea FF amb dades aleatòries o personalitzades (útil per a testejar fluxos): es configura File Size, Batch Size i una planificació (Run Schedule, per exemple 3s). Connectat a ReplaceText (canvia el contingut) i a LogAttribute (escriu els atributs al log), es pot inspeccionar la cua (list queue) per veure el contingut de cada FF en cada pas.
Amb ExtractText es pot copiar el contingut d'un FF a un nou atribut (propietat amb expressió * per capturar-ho tot), per després consultar-lo des del Data Provenance, que ofereix el llinatge complet de la dada: l'origen, els moviments, les transformacions i la qualitat, facilitant la documentació i la governança.
Cas 3 — Filtratge de dades
A partir d'un CSV (ProductID;Date;Zip;Units;Revenue;Country), es crea un nou fitxer amb només les vendes de França amb més d'una unitat:
GetFilellegeix el fitxer (keep source file atrue).SplitRecordsepara cada fila en un FF, configurant unCSVReaderi unCSVRecordSetWriter(separador;, primera fila com a capçalera) i indicant1a Records per Split.QueryRecordexecuta una consulta SQL contra el FF (amb els mateixos Reader/Writer), per exemple:
UpdateAttribute+PutFilepersisteixen el resultat.
Amb els processadors de tipus Record (ConvertRecord, LookupRecord, QueryRecord, ConsumeKafkaRecord_N_M, PublishKafkaRecord_N_M) es tracten els FF com a conjunts de registres en lloc de contingut pla, cosa que simplifica els fluxos i en millora el rendiment (FF més grans, amb múltiples registres cadascun).
Cas 4 — Fusionar dades i persistir a MongoDB
Flux més complex: escoltar peticions HTTP, distingir si contenen la paraula "ERROR", fusionar els missatges en fitxers periòdics segons el seu estat, i desar el resultat a MongoDB.
ListenHTTPescolta al port i endpoint configurats (per exemple, port8081, Base Pathiabd).RouteOnContentsepara el flux en dos segons si el contingut conté la cadenaERROR(Match Requirement:content must contain match).MergeContentagrupa els FF de cada branca (Merge Strategy: Bin-Packing,Correlation Attribute Name: el nom de la ruta,Maximum Bin Age: per exemple300sperquè el fitxer fusionat surti com a molt tard cada 5 minuts).ExtractText+UpdateAttribute+AttributesToJSONpreparen el contingut en format JSON (contingut + estat).PutMongoinsereix el document a la col·lecció configurada (Mongo URI, Database Name, Collection Name).
# docker-compose.yml — NiFi + MongoDB al mateix contenidor/xarxa
services:
nifi:
image: apache/nifi:latest
ports:
- "8443:8443"
environment:
SINGLE_USER_CREDENTIALS_USERNAME: nifi
SINGLE_USER_CREDENTIALS_PASSWORD: nifinifinifi
NIFI_JVM_HEAP_MAX: 2g
links:
- mongodb
mongodb:
image: mongo:latest
ports:
- "27017:27017"
Grups, ports i funnels
Els grups de processos permeten encapsular un tros del flux (amb un Input port i un Output port propis) com una caixa negra reutilitzable, millorant la llegibilitat d'un flux complex. Els funnels permeten centralitzar en un únic punt diverses connexions en paral·lel, de manera que canviar el processador de destinació no obligui a refer totes les connexions una per una.
Plantilles
NiFi permet exportar i importar plantilles (fitxers XML) per a reutilitzar fluxos entre instal·lacions. Hi ha col·leccions públiques de plantilles mantingudes per la comunitat (per exemple, a hortonworks-gallery/nifi-templates), com la plantilla csv-to-json-flow.xml, que converteix directament un flux CSV a JSON sense haver de configurar processador per processador.
Integració amb Elasticsearch
Amb PutElasticSearchHttp es poden indexar documents directament (Elasticsearch URL, Index). Si el document d'origen és un array (per exemple, un movies.json amb una llista de pel·lícules), cal separar-lo prèviament amb SplitJson i l'expressió JSONPath corresponent ($.movies.*), o Elasticsearch indexarà un únic document amb tot l'array en lloc d'un document per pel·lícula.
API REST
NiFi ofereix una API REST equivalent a la interfície gràfica, accessible per exemple a https://localhost:8443/nifi-api/flow/about, que permet automatitzar la gestió del flux des de scripts externs.
Miniactivitat — AC5075/02/01
Dissenya (en un diagrama o descripció textual) un flux NiFi que llegeixi fitxers CSV d'un directori, en filtri les files amb un camp numèric per sobre d'un llindar, i publiqui el resultat a un topic de Kafka mitjançant PublishKafkaRecord. Indica quins processadors Record (Reader/Writer) faries servir i per què.
RA2 | Mòdul 5075 Big Data Aplicat | Institut Sa Palomera (Blanes) | Curs IABD 2026-2027