← 🌊 Streaming
intermedio🧪 Lab práctico

4.4 · Lab 4: Un pipeline de streaming con Kafka

⏱ 60 minMódulo 4: Streaming e Ingesta de Datos

🎯 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, con netstat -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"]
Dentro de un mismo grupo, cada particion la lee un solo consumidor: asi se escala horizontalmente el consumo.

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.

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é?

🧹 Limpieza

# Ctrl+C en las tres terminales para detener productor y consumidores
docker stop redpanda-lab
docker rm redpanda-lab

📚 Lecturas y fuentes

RecursoTipoPor qué leerlo
Kafka DesignDocumentación oficialExplica 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 StartGuía prácticaAmplí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 groupsReferenciaDetalla 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.