Apache Kafka
Introducció
Apache Kafka és, en poques paraules, un middleware de missatgeria entre sistemes heterogenis que, mitjançant un sistema de cues (topics), facilita la comunicació asíncrona i desacobla els fluxos de dades dels sistemes que els produeixen o consumeixen. Funciona com un broker de missatges que enruta els missatges entre clients de manera molt ràpida.
Sense Kafka, si múltiples generadors de dades (servidors web, bases de dades, xat) han d'emmagatzemar les seves dades en múltiples destinacions (logs, mètriques, carretó de compra, fallades), es genera una xarxa de dependències directes entre productors i consumidors. Kafka hi posa remei connectant tots els productors a Kafka, i tots els consumidors a Kafka:
graph LR
P1[Productor A] --> K[(Apache Kafka)]
P2[Productor B] --> K
P3[Productor C] --> K
K --> C1[Consumidor X]
K --> C2[Consumidor Y]
K --> C3[Consumidor Z]
Es tracta d'una plataforma open source distribuïda de transmissió d'esdeveniments en temps real amb emmagatzematge durador, que ofereix alt rendiment (bilions de peticions al dia, latència inferior a 10 ms), tolerància a fallades, disponibilitat i escalabilitat horitzontal. Més del 80% de les 100 companyies més grans dels EUA (Uber, Netflix, Spotify, LinkedIn, PayPal) processen els seus missatges amb Kafka.
L'arquitectura de Kafka té dues directives clares: no bloquejar els productors (gestionant la back pressure, quan un publicador produeix més ràpid del que un subscriptor pot consumir) i aïllar productors i consumidors entre si (no es coneixen mútuament).
Amazon Kinesis i Confluent
Amazon Kinesis és l'equivalent gestionat dins d'AWS (no open source), amb gran facilitat d'escalabilitat i integració amb la resta de serveis AWS. Confluent és una solució PaaS que ofereix desplegament i monitoratge de Kafka com a producte, amb una versió Community provable en local (encara que amb requisits alts de RAM).
Model publicador/subscriptor
Tres elements clau: el publicador (productor) genera una dada i la col·loca en un topic com a missatge; el topic és el magatzem temporal/durador que actua com a cua; el subscriptor (consumidor) rep el missatge. Un productor mai es comunica directament amb un consumidor: sempre ho fa a través d'un topic.
Primer contacte: creant, produint i consumint
# Crear un topic
kafka-topics.sh --create --topic iabd-topic --bootstrap-server iabdkafka:9092
# Llistar els topics existents
kafka-topics.sh --list --bootstrap-server iabdkafka:9092
# Descripció d'un topic (particions, líders, rèpliques)
kafka-topics.sh --describe --topic iabd-topic --bootstrap-server iabdkafka:9092
# Produir missatges (cada línia és un event)
kafka-console-producer.sh --topic iabd-topic --bootstrap-server iabdkafka:9092
# Consumir missatges des de l'inici
kafka-console-consumer.sh --topic iabd-topic --from-beginning --bootstrap-server iabdkafka:9092
Elements de l'arquitectura Kafka
Topics i particions
Un topic és un flux particular de dades, amb un nom unívoc, que funciona com una cua que emmagatzema dades temporalment o de forma duradora. Cada topic es divideix en particions numerades (la primera és la 0); en crear-lo es pot indicar la quantitat inicial amb --partitions. Cada partició és un commit log ordenat, on cada missatge té un identificador incremental anomenat offset: l'ordre només es garanteix dins d'una partició, mai entre particions.
Les dades d'una partició tenen un retention period (per defecte, una setmana) i, un cop escrites, són immutables. Per defecte les dades s'assignen aleatòriament a una partició, tot i que es pot indicar una clau de particionat.
kafka-topics.sh --create --topic iabd-topic-p3 --partitions 3 --bootstrap-server iabdkafka:9092
kafka-topics.sh --delete --topic iabd-topic --bootstrap-server iabdkafka:9092
Esborrar un topic
En esborrar un topic, els seus índexs amb l'històric no s'eliminen: tornar a crear un topic amb el mateix nom pot produir dades incongruents.
Brokers i factor de replicació
Un clúster de Kafka es compon de múltiples nodes anomenats brokers (cada un, un servidor Kafka amb un id enter). Cada broker conté un subconjunt de particions —mai totes les dades, ja que Kafka és distribuït. En connectar-se a un broker qualsevol (bootstrap broker), el client es connecta automàticament a tot el clúster. Es recomana començar amb una arquitectura de 3 brokers.
Per a la tolerància a fallades, els topics han de tenir un factor de replicació superior a 1 (habitualment entre 2 i 3):
kafka-topics.sh --create --topic TopicA --partitions 2 --replication-factor 2 --bootstrap-server iabdkafka:9092
En qualsevol instant, cada partició té una única rèplica líder, l'única que rep i serveix les lectures/escriptures; la resta són ISR (in-sync replicas) que se sincronitzen amb el líder. Si el broker líder cau, una altra rèplica en pren el relleu automàticament.
Productors
Els productors escriuen als topics coneixent automàticament el broker i la partició de destinació, i es recuperen automàticament si un broker falla. Sense clau, Kafka reparteix els missatges amb Round Robin entre particions. Configuració d'ACK:
| Valor | Significat |
|---|---|
ack=0 |
El productor no espera confirmació (possible pèrdua de dades, enviament asíncron) |
ack=1 |
Espera la confirmació del líder (enviament síncron, limita la pèrdua de dades) |
ack=all |
Espera la confirmació del líder i de totes les rèpliques (sense pèrdua de dades) |
Clau de missatge: si s'envia una clau, tots els missatges amb la mateixa clau van sempre a la mateixa partició (útil quan cal ordenar per un camp concret, com un identificador d'operació).
Consumidors, grups i offsets
Els consumidors llegeixen en ordre dins de cada partició (mai entre particions), i poden llegir de diverses particions en paral·lel. Un consumidor pertany a un consumer group: cada partició del topic s'assigna a un únic consumidor del grup, permetent processament paral·lel. Diferents grups de consumidors reben, cadascun, la mateixa dada de cada partició (útil quan dues aplicacions diferents —per exemple, ML i BI— han de rebre les mateixes dades).
Kafka emmagatzema l'offset de cada consumer group en un topic intern (__consumer_offsets), a mode de checkpoint. En fer commit d'un offset després de processar un missatge, si el consumidor cau, en reprendre llegirà des de l'últim offset confirmat. Semàntiques d'entrega:
- At most once: commit just en rebre el missatge — si falla el processament, es perd.
- At least once (la més equilibrada): commit després de processar — pot duplicar-se, per la qual cosa cal idempotència.
- Exactly once: només amb Kafka Streams de punta a punta, o amb un consumidor idempotent si hi ha un sistema extern (base de dades) implicat.
# Llistar grups de consumidors
kafka-consumer-groups.sh --list --bootstrap-server iabdkafka:9092
# Detall d'un grup (CURRENT-OFFSET, LOG-END-OFFSET, LAG)
kafka-consumer-groups.sh --describe --group iabd-app1 --bootstrap-server iabdkafka:9092
El LAG (missatges pendents de llegir) és la mètrica de monitorització més important de Kafka — vegeu Monitorització de sistemes Big Data.
Zookeeper i Kraft
ZooKeeper gestiona la llista de brokers, ajuda en l'elecció de la rèplica líder i notifica canvis a Kafka (creació/eliminació de topics, caiguda o recuperació d'un broker). En un entorn real, s'instal·la un nombre senar de servidors (3, 5, 7). Des de Kafka 3.x, Kraft ofereix un nou protocol de consens que permet prescindir de ZooKeeper (encara no és l'opció més recomanable en producció el 2026, però guanya adopció).
Kafka i Python
# producer.py
from kafka import KafkaProducer
from json import dumps
producer = KafkaProducer(
value_serializer=lambda m: dumps(m).encode('utf-8'),
bootstrap_servers=['iabdkafka:9092'])
for i in range(10):
producer.send("iabd-topic", value={"nom": f"producer {i}"})
producer.flush()
# consumer.py
from kafka import KafkaConsumer
from json import loads
consumer = KafkaConsumer(
'iabd-topic',
auto_offset_reset='earliest',
enable_auto_commit=True,
group_id='iabd-grup-1',
value_deserializer=lambda m: loads(m.decode('utf-8')),
bootstrap_servers=['iabdkafka:9092'])
for m in consumer:
print(m.value)
Clúster de Kafka: creant múltiples brokers
# server101.properties
broker.id=101
listeners=PLAINTEXT://:9092
log.dirs=/opt/kafka/logs/broker_101
zookeeper.connect=localhost:2181
Amb tres fitxers de configuració anàlegs (server101/102/103.properties, canviant broker.id, listeners i log.dirs), s'arrenca cada broker en un terminal separat, i es crea un topic ja repartit:
kafka-topics.sh --create --topic iabd-topic-3p2r \
--bootstrap-server iabd-virtualbox:9092 \
--partitions 3 --replication-factor 2
Guia de rendiment: en un clúster petit (< 6 brokers), crear el doble de particions que brokers; en un de gran (> 12 brokers), la mateixa quantitat. El factor de replicació hauria de ser com a mínim 2, recomanablement 3 (amb un màxim de 4): més replicació implica millor tolerància a fallades, però més latència (si acks=all) i més espai en disc. Un broker no hauria de superar les 2.000-4.000 particions, ni el clúster sencer les 20.000.
Kafka Connect
Kafka Connect permet importar/exportar dades des de/cap a Kafka mitjançant connectors ja construïts, sense escriure codi:
- Connectors font (source): obtenen dades des de les fonts (la "E" d'ETL).
- Connectors destinació (sink): publiquen les dades als magatzems de destí (la "L" d'ETL).
# mysql-source-connector.properties
name=mysql-source
connector.class=io.confluent.connect.jdbc.JdbcSourceConnector
tasks.max=1
connection.url=jdbc:mysql://localhost/retail_db
connection.user=iabd
connection.password=iabd
table.whitelist=categories
mode=incrementing
incrementing.column.name=category_id
topic.prefix=iabd-retail_db-
connect-standalone.sh config/connect-standalone.properties config/mysql-source-connector.properties
# API REST de Kafka Connect (port 8083 per defecte)
curl http://iabdkafka:8083/connectors
Qualsevol INSERT posterior a la taula d'origen apareix automàticament al topic corresponent, sense necessitat de codi addicional.
Kafka Streams i el paper de Kafka en el Big Data
Kafka Streams és la tercera pota de l'ecosistema Kafka: permet processar i transformar dades directament dins de Kafka mitjançant llibreries Java/Scala. En aquest mòdul, aquest tipus de processament es fa amb Spark Structured Streaming, que ofereix una API més accessible en Python.
Kafka està pensat principalment per al tractament en streaming, però amb Kafka Connect també dona suport al processament batch, cosa que el converteix en l'element central de l'arquitectura Kappa: les fonts de dades se substitueixen per cues de Kafka, unificant la ingesta batch i en streaming.
Miniactivitat — AC5075/02/02
Crea un topic amb 3 particions i factor de replicació 1 en un únic broker local. Escriu un productor Python que enviï 20 missatges amb una clau basada en un id_client simulat, i un consumidor que els mostri indicant partició, offset i clau. Explica per què, sense clau, l'ordre de recepció no coincideix amb l'ordre d'enviament.
RA2 | Mòdul 5075 Big Data Aplicat | Institut Sa Palomera (Blanes) | Curs IABD 2026-2027