lección 6
Kafka en profundidad: consumer groups, offsets y garantías
Domina consumer groups, gestión manual de offsets, rebalanceo, at-least-once vs exactly-once y patrones de commit.
⏱ 55 min
### Consumer Groups: el secreto del paralelismo en Kafka
En la lección anterior viste un consumidor simple leyendo de un topic. Pero en producción NUNCA tienes un solo consumidor. Tienes 3, 5, 10 instancias procesando en paralelo para aguantar el volumen. ¿Cómo se coordinan? Con CONSUMER GROUPS. Un consumer group es un conjunto de consumidores que COOPERAN para leer un topic. Kafka asigna cada partición a UN SOLO consumidor del grupo. Si el topic tiene 6 particiones y tu grupo tiene 3 consumidores, cada uno recibe 2 particiones.
La analogía: imagina un equipo de repartidores de pizza. El topic es "todas las pizzas que hay que repartir". Las particiones son zonas de la ciudad (norte, sur, centro). El consumer group es el equipo de repartidores. Kafka ASIGNA cada zona a un repartidor. Si un repartidor se va (fallo), su zona se reasigna a otro (rebalanceo). Ninguna pizza se reparte dos veces ni se queda sin repartir.
### Reglas de asignación de particiones
- Una partición solo puede ser asignada a UN consumidor del grupo a la vez.
- Un consumidor puede tener asignadas MÚLTIPLES particiones.
- Si hay más consumidores que particiones, los sobrantes están inactivos (standby).
- Si un consumidor muere, sus particiones se reasignan a los supervivientes (rebalanceo).
- Dos consumer groups DIFERENTES leen el topic de forma independiente (cada grupo ve TODOS los mensajes).
### Offsets: la memoria del consumidor
El offset es la POSICIÓN del consumidor en una partición. Es un número secuencial: offset 0 es el primer mensaje, offset 1 el segundo, etc. Kafka guarda el offset de cada consumer group en un topic interno llamado __consumer_offsets. Cuando un consumidor arranca, pregunta "¿por dónde iba?" y Kafka le responde con su último offset committed.
Es como un marcapáginas en un libro. Si dejas de leer y vuelves mañana, el marcapáginas te dice dónde continuar. Si pierdes el marcapáginas (offset no committed), tienes que decidir: ¿empiezo desde el principio (earliest) o salto al final (latest)?
### Auto-commit vs commit manual
Por defecto, kafka-python hace auto-commit cada 5 segundos: automáticamente guarda el offset del último mensaje LEÍDO (no procesado). El problema: si lees un mensaje, el auto-commit guarda el offset, y DESPUÉS tu procesamiento falla... el mensaje se pierde. Kafka cree que ya lo procesaste, pero en realidad falló.
La solución profesional: COMMIT MANUAL. Tú decides CUÁNDO commitear el offset, que debería ser DESPUÉS de procesar exitosamente el mensaje. Si el procesamiento falla, no commiteas, y al reiniciar el consumidor releerá ese mensaje. Nota importante: sin auto_offset_reset="earliest", un consumer group nuevo empieza por defecto en "latest" y solo ve mensajes que lleguen a partir de ese momento — por eso aquí siempre lo ponemos explícitamente.
1from kafka import KafkaConsumer2import json34# Consumidor con COMMIT MANUAL5consumer = KafkaConsumer(6 "pedidos",7 bootstrap_servers=["localhost:9092"],8 group_id="grupo-pagos-v2",9 auto_offset_reset="earliest",10 enable_auto_commit=False, # ← CLAVE: desactivar auto-commit11 value_deserializer=lambda m: json.loads(m.decode("utf-8")),12)1314print("Consumidor con commit manual escuchando...")15for message in consumer:16 pedido = message.value17 try:18 # 1. Procesar el mensaje19 resultado = procesar_pago(pedido)20 print(f" Procesado: {pedido['pedido_id']} -> {resultado}")2122 # 2. Commit DESPUÉS de procesar exitosamente23 consumer.commit()2425 except Exception as e:26 # NO commitear → al reiniciar, Kafka reenvía este mensaje27 print(f" ERROR procesando {pedido['pedido_id']}: {e}")28 print(f" Offset NO committed. Se reintentará al reiniciar.")29 break # O manejar de otra forma
Commit manual: solo confirmamos el offset DESPUÉS de procesar exitosamente
### Patrones de commit: por mensaje vs por lote
Commitear después de CADA mensaje es seguro pero lento (una petición de red por mensaje). Commitear por lotes es rápido pero arriesgado: si fallas a mitad del lote, al reiniciar reprocesas todo el lote. La decisión depende de tu tolerancia a duplicados vs tu necesidad de throughput.
1# Patrón: commit por lote (cada N mensajes)2BATCH_SIZE = 1003count = 045for message in consumer:6 procesar(message.value)7 count += 189 if count % BATCH_SIZE == 0:10 consumer.commit()11 print(f" Commit tras {count} mensajes")1213# Commit final por si quedaron mensajes sin commitear14consumer.commit()
Commit por lote: balance entre seguridad y rendimiento
### Rebalanceo: qué pasa cuando un consumidor muere
Cuando un consumidor del grupo se desconecta (crash, deploy, timeout), Kafka inicia un REBALANCEO: redistribuye las particiones del consumidor muerto entre los supervivientes. Durante el rebalanceo, NINGÚN consumidor del grupo puede leer — hay una pausa breve. Después, cada consumidor recibe sus nuevas particiones asignadas y continúa desde el último offset committed.
El rebalanceo también ocurre cuando AÑADES un consumidor al grupo. Si tenías 3 consumidores con 2 particiones cada uno y añades un cuarto, Kafka redistribuye: ahora 2 tienen 2 particiones y 2 tienen 1. Es elástico.
1# Configurar timeouts de heartbeat para rebalanceo rápido2consumer = KafkaConsumer(3 "pedidos",4 bootstrap_servers=["localhost:9092"],5 group_id="grupo-elastico",6 # Heartbeat: señal de "sigo vivo" cada 3s7 heartbeat_interval_ms=3000,8 # Si no hay heartbeat en 10s, se considera muerto9 session_timeout_ms=10000,10 # Timeout de poll: si no llamas a poll() en 5min, te echan11 max_poll_interval_ms=300000,12 enable_auto_commit=False,13 value_deserializer=lambda m: json.loads(m.decode("utf-8")),14)
Configurar heartbeats para que el rebalanceo sea rápido y predecible
### Garantías de entrega en Kafka
Kafka soporta tres niveles de garantía para el PRODUCTOR: acks=0 (fire and forget, puede perder datos), acks=1 (el líder confirma, pero si muere antes de replicar se pierde), acks=all (todos los replicas confirman, máxima durabilidad). Para el CONSUMIDOR, la garantía depende de cuándo commiteas: antes de procesar (at-most-once, puede perder), después de procesar (at-least-once, puede duplicar).
- AT-MOST-ONCE: commit antes de procesar. Si fallas al procesar, el mensaje se "pierde" (ya committed).
- AT-LEAST-ONCE: commit después de procesar. Si fallas al commitear, al reiniciar re-procesas (duplicado).
- EXACTLY-ONCE: Kafka Transactions (enable.idempotence=True + transactional.id). Complejo pero posible.
1# Productor con máxima garantía de durabilidad2from kafka import KafkaProducer34producer = KafkaProducer(5 bootstrap_servers=["localhost:9092"],6 acks="all", # Espera confirmación de TODAS las réplicas7 retries=5, # Reintenta hasta 5 veces si falla8 max_in_flight_requests_per_connection=1, # Mantiene orden en reintentos9 enable_idempotence=True, # Evita duplicados en reintentos del productor10 value_serializer=lambda v: json.dumps(v).encode("utf-8"),11)12# Con idempotence=True, Kafka deduplica automáticamente reintentos del productor.13# Combinado con acks=all, es la configuración más segura posible.
Productor idempotente con acks=all: máxima garantía contra pérdida de datos
### Exactly-once processing: transacciones en Kafka
Exactly-once es posible en Kafka desde la versión 0.11 usando TRANSACCIONES. El concepto: el productor escribe mensajes Y commitea offsets del consumidor en una TRANSACCIÓN ATÓMICA. O ambas cosas ocurren, o ninguna. Pero ojo: exactly-once en Kafka solo funciona cuando produces Y consumes DENTRO de Kafka. Si tu consumidor escribe en una base de datos externa, vuelves a at-least-once (necesitas idempotencia en la DB).
En la práctica, la mayoría de sistemas usan at-least-once + idempotencia en el consumidor. Es más simple, más robusto y cubre la inmensa mayoría de los casos. Exactly-once transaccional es para cuando procesas de un topic y escribes en OTRO topic de Kafka (transformaciones stream-to-stream).
### Monitorizar consumer groups
1# Ver el estado de un consumer group: lag, offset actual, offset final2docker exec kafka-broker kafka-consumer-groups \3 --bootstrap-server localhost:9092 \4 --describe --group grupo-pagos-v256# Salida esperada:7# GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG8# grupo-pagos-v2 pedidos 0 45 50 59# grupo-pagos-v2 pedidos 1 38 38 010# grupo-pagos-v2 pedidos 2 42 47 511#12# LAG = mensajes pendientes de procesar. Si crece, tu consumidor no da abasto.
Monitorizar el LAG: diferencia entre lo producido y lo consumido
Consejo de senior: MONITORIZA EL LAG. Es la métrica más importante de Kafka. Si el lag crece constantemente, tu consumidor es más lento que tu productor. Necesitas más consumidores (hasta el número de particiones) o consumidores más rápidos. Un lag estable y bajo es señal de salud.
Lo que le diría a mi yo de hace 5 años: usa enable_auto_commit=False SIEMPRE en producción. El auto-commit es cómodo para prototipos pero es una bomba de relojería en producción. La primera vez que pierdas datos por un auto-commit mal sincronizado, te arrepentirás de no haber usado commit manual desde el principio.
Si tienes más consumidores que particiones en un grupo, los sobrantes estarán INACTIVOS (no leen nada). No es un error, pero es un desperdicio de recursos. Regla: número de consumidores <= número de particiones.
## ejercicios
Consumidor con commit manual y manejo de errores
Implementa un consumidor que lea del topic "pedidos" con commit manual. Si el procesamiento falla (simula un 20% de fallos), NO debe commitear ese offset. Lleva la cuenta de procesados vs fallidos. Qué deberías ver: una mezcla de líneas OK y FAIL con partición:offset, y al final las stats con el ratio de éxito (~80%).
Monitorizar el lag de un consumer group
Escribe un script Python que use AdminClient de kafka-python para obtener el lag de cada partición del grupo "grupo-pagos-v2" del topic "pedidos". Imprime: partición, offset actual, offset final, lag. Qué deberías ver: una tabla con columnas Partition, Committed, End, Lag — y al final el lag total. Si el grupo nunca ha corrido, Committed será 0 y Lag = End.
Resetear offsets: volver a leer desde el inicio
A veces necesitas reprocesar datos (bug corregido, nuevo consumidor). Usa la CLI de Kafka para resetear los offsets del grupo "grupo-pagos-v2" en el topic "pedidos" al inicio (earliest). Luego verifica con --describe. Qué deberías ver: --reset-offsets muestra las particiones con NEW-OFFSET:0; el --describe posterior muestra CURRENT-OFFSET en 0 y LAG igual a LOG-END-OFFSET.
Dos consumer groups leyendo el mismo topic
Demuestra que dos consumer groups independientes leen TODOS los mensajes del topic. Crea grupo "analytics" y grupo "alertas" que lean del topic "pedidos". Ambos deben recibir los mismos mensajes. Qué deberías ver: las dos líneas con el mismo número de mensajes, y la confirmación "Ambos grupos leen los mismos N mensajes: True".
Regístrate para guardar tu progreso.
## comentarios
Reporta erratas, ayuda a otros o comparte tu opinión. Sé constructivo.
Inicia sesión para comentar y responder.
cargando comentarios...