Saltar al contenido

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).
Consumer Groups: cada grupo lee todo el topic, pero internamente reparte particiones entre sus miembros

### 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 KafkaConsumer
2import json
3
4# Consumidor con COMMIT MANUAL
5consumer = 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-commit
11 value_deserializer=lambda m: json.loads(m.decode("utf-8")),
12)
13
14print("Consumidor con commit manual escuchando...")
15for message in consumer:
16 pedido = message.value
17 try:
18 # 1. Procesar el mensaje
19 resultado = procesar_pago(pedido)
20 print(f" Procesado: {pedido['pedido_id']} -> {resultado}")
21
22 # 2. Commit DESPUÉS de procesar exitosamente
23 consumer.commit()
24
25 except Exception as e:
26 # NO commitear → al reiniciar, Kafka reenvía este mensaje
27 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 = 100
3count = 0
4
5for message in consumer:
6 procesar(message.value)
7 count += 1
8
9 if count % BATCH_SIZE == 0:
10 consumer.commit()
11 print(f" Commit tras {count} mensajes")
12
13# Commit final por si quedaron mensajes sin commitear
14consumer.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ápido
2consumer = KafkaConsumer(
3 "pedidos",
4 bootstrap_servers=["localhost:9092"],
5 group_id="grupo-elastico",
6 # Heartbeat: señal de "sigo vivo" cada 3s
7 heartbeat_interval_ms=3000,
8 # Si no hay heartbeat en 10s, se considera muerto
9 session_timeout_ms=10000,
10 # Timeout de poll: si no llamas a poll() en 5min, te echan
11 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 durabilidad
2from kafka import KafkaProducer
3
4producer = KafkaProducer(
5 bootstrap_servers=["localhost:9092"],
6 acks="all", # Espera confirmación de TODAS las réplicas
7 retries=5, # Reintenta hasta 5 veces si falla
8 max_in_flight_requests_per_connection=1, # Mantiene orden en reintentos
9 enable_idempotence=True, # Evita duplicados en reintentos del productor
10 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 final
2docker exec kafka-broker kafka-consumer-groups \
3 --bootstrap-server localhost:9092 \
4 --describe --group grupo-pagos-v2
5
6# Salida esperada:
7# GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG
8# grupo-pagos-v2 pedidos 0 45 50 5
9# grupo-pagos-v2 pedidos 1 38 38 0
10# grupo-pagos-v2 pedidos 2 42 47 5
11#
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

[01]

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%).

Cargando editor...
[02]

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.

Cargando editor...
[03]

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.

Cargando editor...
[04]

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".

Cargando editor...

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...