🎯 Objetivo
Montar un pipeline de streaming completo en tu máquina —broker, productor y consumidores— para ver con tus propios ojos cómo se reparten las particiones, cómo avanzan los offsets y qué ocurre durante un rebalanceo.
✅ Antes de empezar
Comprueba los requisitos antes de tocar nada. Si algún comando falla, resuélvelo ahora: a mitad del lab es mucho más molesto.
docker --version # Docker version 24.x o superior
python --version # Python 3.10 o superior
pip install kafka-python-ng
Sobre el cliente de Python: instalamos kafka-python-ng, el fork mantenido de kafka-python. El paquete original lleva tiempo sin mantenimiento y falla con Python 3.12 o superior. El import es idéntico: from kafka import KafkaProducer, KafkaConsumer.
Además necesitas:
- ~1.5 GB de RAM libre para el contenedor del broker.
- El puerto 9092 libre. Verifícalo con
docker ps(ningún contenedor debe exponerlo) o, en Windows, connetstat -ano | findstr 9092— si no devuelve nada, está libre.
⚠️ Este lab necesita DOS terminales abiertas a la vez (y TRES en el Paso 5). El productor se queda corriendo mientras tú trabajas en otra ventana. Rotula mentalmente cada terminal desde ya: Terminal 1 = productor, Terminal 2 = consumidor A, Terminal 3 = consumidor B.
Paso 1: Levantar el broker
Usaremos Redpanda en lugar de Kafka. Redpanda habla exactamente el mismo protocolo que Kafka, así que el código Python que escribas aquí funciona sin cambiar una sola línea contra un clúster de Kafka real. La ventaja: un único contenedor, sin ZooKeeper y con un arranque de segundos.
Terminal 1:
docker run -d --name redpanda-lab -p 9092:9092 redpandadata/redpanda:latest redpanda start --overprovisioned --smp 1 --memory 1G --reserve-memory 0M --node-id 0 --check=false --kafka-addr PLAINTEXT://0.0.0.0:9092 --advertise-kafka-addr PLAINTEXT://localhost:9092
docker ps
docker exec -it redpanda-lab rpk cluster info
Paso 2: Crear el topic con 3 particiones
Un topic con una sola partición no permite paralelismo: da igual cuántos consumidores tengas, solo uno trabajará. Creamos clicks con 3 particiones para poder repartirlas después entre varios consumidores.
Terminal 1:
docker exec -it redpanda-lab rpk topic create clicks -p 3
docker exec -it redpanda-lab rpk topic list
Paso 3: El productor (Terminal 1)
Ahora generamos el flujo. Cada segundo emitimos un evento JSON de compra con usuario, producto, precio y timestamp. Lo importante: pasamos el usuario como key, así el hash de la clave decide la partición y todos los eventos del mismo usuario caen siempre en la misma partición.
Crea productor.py:
import json
import random
import time
from kafka import KafkaProducer
producer = KafkaProducer(
bootstrap_servers="localhost:9092",
key_serializer=lambda k: k.encode("utf-8"),
value_serializer=lambda v: json.dumps(v).encode("utf-8"),
)
USUARIOS = ["ana", "beto", "carla", "diego", "elena", "fabio"]
PRECIOS = {
"laptop": 1200.0,
"mouse": 25.0,
"teclado": 80.0,
"monitor": 300.0,
"webcam": 45.0,
}
print("Produciendo eventos en el topic 'clicks'. Ctrl+C para parar.")
n = 0
while True:
producto = random.choice(list(PRECIOS))
evento = {
"usuario": random.choice(USUARIOS),
"producto": producto,
"precio": PRECIOS[producto],
"timestamp": time.time(),
}
# la key fuerza el reparto: mismo usuario -> misma particion, orden garantizado
futuro = producer.send("clicks", key=evento["usuario"], value=evento)
meta = futuro.get(timeout=10)
n += 1
print(f"[{n}] {evento['usuario']:>6} -> particion {meta.partition} offset {meta.offset}")
time.sleep(1)
Ejecútalo en la Terminal 1 y déjalo corriendo durante todo el lab:
python productor.py
Paso 4: El consumidor (Terminal 2)
Abre una segunda terminal y déjala junto a la primera. El consumidor se une a un consumer group (group_id) e imprime la partición y el offset de cada mensaje: esos dos números son el corazón de este lab.
Crea consumidor.py:
import json
import sys
from kafka import KafkaConsumer
grupo = sys.argv[1] if len(sys.argv) > 1 else "analitica"
consumer = KafkaConsumer(
"clicks",
bootstrap_servers="localhost:9092",
group_id=grupo,
auto_offset_reset="latest", # solo mensajes nuevos; usa "earliest" para releer todo
value_deserializer=lambda v: json.loads(v.decode("utf-8")),
)
print(f"Consumidor del grupo '{grupo}' escuchando el topic 'clicks'...")
for msg in consumer:
e = msg.value
print(f"p{msg.partition} off={msg.offset} | {e['usuario']:>6} compro {e['producto']:<8} {e['precio']:>7.2f}")
Terminal 2:
python consumidor.py analitica
graph LR P["Productor<br/>1 evento/s"] --> T["Topic clicks"] T --> P0["Particion 0"] T --> P1["Particion 1"] T --> P2["Particion 2"] P0 --> C1["Consumidor A<br/>grupo analitica"] P1 --> C1 P2 --> C2["Consumidor B<br/>grupo analitica"]
Paso 5: El experimento clave — rebalanceo
Aquí está toda la lección del lab. Vas a arrancar un segundo consumidor en el mismo grupo y verás cómo Kafka reparte las particiones entre ambos. Luego repetirás con un grupo distinto y el comportamiento cambiará por completo.
Terminal 3 (deja las otras dos corriendo):
# 5a. MISMO grupo que la Terminal 2 -> se reparten las particiones
python consumidor.py analitica
Mira ahora las Terminales 2 y 3 a la vez durante unos 15 segundos. Después detén la Terminal 3 con Ctrl+C y arráncala con otro grupo:
# 5b. Grupo DISTINTO -> recibe una copia completa del stream
python consumidor.py auditoria
Es el patrón doble de Kafka: dentro de un grupo funciona como cola de trabajo (reparto), entre grupos funciona como pub/sub (difusión). Si arrancas un cuarto consumidor en
analitica, quedará ocioso: solo hay 3 particiones.
Paso 6: Agregación en ventana
Consumir mensajes de uno en uno es fácil; lo interesante es agregar. Vamos a acumular ingresos por producto en ventanas tumbling de 10 segundos, usando un simple diccionario en memoria y poll() para no bloquearnos.
Crea ventanas.py y ejecútalo en la Terminal 3 (Ctrl+C primero al consumidor anterior):
import time
from collections import defaultdict
import json
from kafka import KafkaConsumer
VENTANA = 10 # segundos
consumer = KafkaConsumer(
"clicks",
bootstrap_servers="localhost:9092",
group_id="ventanas",
auto_offset_reset="latest",
value_deserializer=lambda v: json.loads(v.decode("utf-8")),
)
ingresos = defaultdict(float)
inicio = time.time()
print(f"Ingresos por producto, ventanas de {VENTANA}s. Ctrl+C para parar.")
while True:
lotes = consumer.poll(timeout_ms=500)
for _tp, mensajes in lotes.items():
for msg in mensajes:
e = msg.value
ingresos[e["producto"]] += e["precio"]
if time.time() - inicio >= VENTANA:
print(f"--- ventana cerrada {time.strftime('%H:%M:%S')} ---")
if not ingresos:
print(" (sin eventos)")
for producto, total in sorted(ingresos.items(), key=lambda x: -x[1]):
print(f" {producto:<8} {total:>9.2f}")
ingresos.clear()
inicio = time.time()
Ojo con lo que acabas de construir: si este proceso se cae, el diccionario en memoria se pierde y la ventana a medias desaparece. Además agrupamos por processing time, no por event time. Exactamente eso es lo que Spark Structured Streaming o Flink te dan resuelto: estado con checkpoints, tolerancia a fallos, ventanas por event time y watermarks para los eventos rezagados.
🧪 Reto final
🧪 Ejercicio
Alerta en tiempo real: comprador compulsivo
Escribe alerta.py: consume el topic 'clicks' y detecta, en tiempo real, a cualquier usuario que haga MAS de 3 compras en una ventana deslizante de 30 segundos. Imprime una alerta la primera vez que un usuario cruza el umbral y no la repitas hasta que vuelva a bajar. Pista de diseno: mantén por usuario una cola con los timestamps de sus compras y descarta los que ya salieron de la ventana.
🔍 Usa collections.deque por usuario. Al llegar un evento, añade su timestamp al final y saca por la izquierda todos los que sean más antiguos que 30 segundos. La longitud de la cola es el conteo dentro de la ventana. Guarda en un set qué usuarios ya alertaste para no spamear.
import json
import time
from collections import defaultdict, deque
from kafka import KafkaConsumer
VENTANA = 30 # segundos
UMBRAL = 3 # mas de 3 compras dispara la alerta
consumer = KafkaConsumer(
"clicks",
bootstrap_servers="localhost:9092",
group_id="alertas",
auto_offset_reset="latest",
value_deserializer=lambda v: json.loads(v.decode("utf-8")),
)
compras = defaultdict(deque) # usuario -> timestamps dentro de la ventana
alertados = set() # usuarios ya alertados
print("Vigilando compras. Ctrl+C para parar.")
for msg in consumer:
evento = msg.value
usuario = evento["usuario"]
ahora = evento["timestamp"]
cola = compras[usuario]
cola.append(ahora)
# descartar lo que ya salio de la ventana deslizante
while cola and ahora - cola[0] > VENTANA:
cola.popleft()
conteo = len(cola)
if conteo > UMBRAL and usuario not in alertados:
print("ALERTA:", usuario, "hizo", conteo, "compras en", VENTANA, "segundos",
"| particion", msg.partition, "offset", msg.offset)
alertados.add(usuario)
elif conteo <= UMBRAL:
# vuelve a estar por debajo del umbral: puede volver a alertarse
alertados.discard(usuario)Comprueba lo aprendido
Comprueba que lo pillaste
En el Paso 5a arrancaste un segundo consumidor con el MISMO group_id sobre el topic clicks (3 particiones). ¿Qué observaste y por qué?
Dentro de un consumer group cada partición se asigna a exactamente un consumidor, así que al entrar un nuevo miembro se dispara un rebalanceo que reparte las 3 particiones (por ejemplo 2 y 1). No hay duplicados. La difusión completa solo ocurre entre grupos distintos, como viste en el Paso 5b.
🧹 Limpieza
# Ctrl+C en las tres terminales para detener productor y consumidores
docker stop redpanda-lab
docker rm redpanda-lab
📚 Lecturas y fuentes
| Recurso | Tipo | Por qué leerlo |
|---|---|---|
| Kafka Design | Documentación oficial | Explica por qué el log append-only y el sistema de ficheros secuencial hacen a Kafka tan rápido: la teoría detrás de los offsets que acabas de ver moverse. |
| Redpanda Quick Start | Guía práctica | Amplía el broker de este lab a un clúster de 3 nodos con rpk y muestra más comandos de administración de topics. |
| Consumer configuration and groups | Referencia | Detalla el protocolo de consumer groups, el commit de offsets (automático vs. manual) y los parámetros que controlan cuándo se dispara un rebalanceo. |