Event Streaming con Apache Kafka y Node.js
Construye sistemas event-driven escalables usando Apache Kafka con producers, consumers, consumer groups y semantica exactly-once para messaging asincrono confiable
Construye sistemas event-driven resilientes y preparados para crecer usando Apache Kafka. Esta recipe cubre configuracion de producer, consumer groups con auto-rebalancing, manejo de offsets y semantica exactly-once para comunicacion asincrona confiable entre microservicios.
Cuando Usar Esto
- Los servicios necesitan comunicarse asincronamente sin acoplamiento fuerte. Consulta Event-Driven Microservices para patrones de arquitectura.
- El historial de eventos debe ser replayable para debugging o onboarding de nuevos consumers. Consulta Event Sourcing para logs de eventos inmutables.
- El procesamiento de mensajes de alto throughput requiere scaling horizontal de consumers. Consulta RabbitMQ Task Queue para patrones de broker alternativos.
Solucion
1. Kafka Producer
// kafka/producer.ts
import { Kafka, Partitioners } from 'kafkajs';
const kafka = new Kafka({
clientId: 'order-service',
brokers: ['kafka-1:9092', 'kafka-2:9092'],
});
const producer = kafka.producer({
createPartitioner: Partitioners.DefaultPartitioner,
retry: {
retries: 5,
initialRetryTime: 300,
},
});
await producer.connect();
async function publishOrderCreated(order: unknown): Promise<void> {
await producer.send({
topic: 'orders.created',
messages: [
{
key: order.userId,
value: JSON.stringify(order),
headers: {
'content-type': 'application/json',
'trace-id': generateTraceId(),
},
},
],
});
}
2. Consumer con Consumer Group
// kafka/consumer.ts
const consumer = kafka.consumer({
groupId: 'notification-service',
sessionTimeout: 30000,
heartbeatInterval: 3000,
});
await consumer.connect();
await consumer.subscribe({ topic: 'orders.created', fromBeginning: false });
await consumer.run({
autoCommit: true,
autoCommitInterval: 5000,
eachMessage: async ({ topic, partition, message }) => {
const order = JSON.parse(message.value!.toString());
console.log(`Processing order from partition ${partition}:`, order.id);
try {
await sendEmailNotification(order);
} catch (error) {
// Dead letter handling
await publishToDeadLetter(topic, message, error);
}
},
});
3. Exactly-Once Processing
// kafka/exactly-once.ts
const producer = kafka.producer({
transactionalId: 'order-processor',
maxInFlightRequests: 1,
idempotent: true,
});
await producer.connect();
async function processOrderWithIdempotency(orderId: string): Promise<void> {
const transaction = await producer.transaction();
try {
// Procesar orden
const result = await processPayment(orderId);
// Enviar resultado
await transaction.send({
topic: 'orders.completed',
messages: [{ key: orderId, value: JSON.stringify(result) }],
});
// Commit offsets y mensajes atomicamente
await transaction.commit();
} catch (error) {
await transaction.abort();
throw error;
}
}
4. Partitioner Custom para Ordering
// kafka/partitioner.ts
function userIdPartitioner(userId: string, numPartitions: number): number {
// Asegura que todos los eventos de un usuario vayan a la misma particion
let hash = 0;
for (let i = 0; i < userId.length; i++) {
hash = ((hash << 5) - hash) + userId.charCodeAt(i);
hash |= 0;
}
return Math.abs(hash) % numPartitions;
}
await producer.send({
topic: 'user.events',
messages: [
{
key: userId,
value: JSON.stringify(event),
partition: userIdPartitioner(userId, 12),
},
],
});
5. Docker Compose Setup
# docker-compose.kafka.yml
version: '3.8'
services:
zookeeper:
image: confluentinc/cp-zookeeper:7.5.0
environment:
ZOOKEEPER_CLIENT_PORT: 2181
kafka:
image: confluentinc/cp-kafka:7.5.0
depends_on:
- zookeeper
ports:
- "9092:9092"
environment:
KAFKA_BROKER_ID: 1
KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
6. Consumer en Python con kafka-python
from kafka import KafkaConsumer
import json
consumer = KafkaConsumer(
'orders.created',
bootstrap_servers=['localhost:9092'],
group_id='notification-service',
auto_offset_reset='latest',
enable_auto_commit=False,
value_deserializer=lambda x: json.loads(x.decode('utf-8')),
)
for message in consumer:
order = message.value
try:
send_email_notification(order)
consumer.commit()
except Exception as e:
print(f"Error al procesar orden {order['id']}: {e}")
# No commitear — el mensaje se reprocesara al reiniciar
El commit manual te da control sobre cuando se guardan los offsets. Solo commitea despues de procesamiento exitoso para evitar perder mensajes en caso de fallo.
Como Funciona
- Producers publican mensajes a topics particionados entre brokers
- Consumer groups distribuyen particiones entre instancias para procesamiento paralelo
- Offsets trackean progreso del consumer; auto-commit persiste posicion periodicamente
- Exactly-once usa transacciones para commitear offsets y mensajes de salida atomicamente
Consideraciones de Produccion
- Corre al menos 3 brokers de Kafka con replication factor 3 para tolerancia a fallos
- Monitorea consumer lag con herramientas como Kafka Lag Exporter
- Usa schema registry (Confluent) para enforcear schemas Avro/Protobuf en topics
- Configura politicas de retencion apropiadas — por tiempo (7 dias por defecto) o por tamaño
- Habilita log compaction para topics que almacenan estado actual (e.g., perfiles de usuario) en lugar de eventos
- Usa
ack=allen el producer para asegurar que los mensajes se escriban a todas las replicas antes de confirmar
Mejores Practicas
- Usa nombres de topic significativos:
orders. creatednotopic1. Namespacea por dominio y tipo de evento. - Particiona por clave: g.
- Batchea los sends del producer: enviar mensajes en batches mejora el throughput. size
ylinger. ms`. - Maneja rebalances gracefulmente: implementa un rebalance listener para commitear offsets y limpiar recursos antes de que las particiones sean revocadas.
- Usa producers idempotentes: setea
enable. idempotence=truepara prevenir mensajes duplicados por retries.
Errores Comunes
- No manejar rebalances de consumer, causando procesamiento duplicado
- Usar auto-commit con procesos de larga duracion que pueden fallar mid-batch
- Crear demasiadas particiones por topic, incrementando overhead de coordinacion
- No configurar un
transactionalIdal usar transacciones, causando errores de producer fencing - Ignorar consumer lag hasta que se vuelve critico — configura alertas a 1000+ mensajes de retraso
- Usar el partitioner por defecto cuando el orden de mensajes importa — usa key-based partitioning
- No configurar
max.poll.interval.mscorrectamente — consumers que procesan lentamente son expulsados del group - Usar
auto_offset_reset=nonesin offsets commiteados — consumers crashean en la primera ejecucion
Troubleshooting
- Messages are lost on restart: persist messages before acknowledging.
- Consumer lags behind producer: scale consumers, increase prefetch, and partition the topic.
- Duplicate messages: design consumers to be idempotent.
- Ordering is wrong after scaling: preserve partition keys and avoid rebalancing during bursts. Consider a single partition when order is mandatory.
- Queue depth grows but consumers are idle: check network partitions, consumer health, and permission issues. Restart gracefully.
Lectura Adicional
- Documentación oficial: consulta la referencia actualizada del framework o herramienta utilizada.
- Guías relacionadas: explora las guías de event-driven y messaging para profundizar.
- Patrones complementarios: revisa los patrones de diseño aplicables a tu stack tecnológico.
- Postmortems públicos: estudia incidentes reales de equipos que enfrentaron problemas similares en producción.
Notas de Producción
- Despliega gradualmente usando canary o blue-green para detectar regresiones temprano.
- Configura alertas para errores, latencia p99 y tasa de fallos antes de habilitar en producción.
- Documenta el rollback en el runbook; prueba el procedimiento en staging al menos una vez por trimestre.
- Revisa logs estructurados con correlation IDs para trazar requests end-to-end en incidentes.
Puntos Clave
- Aplica event streaming con apache kafka y node.js cuando necesites una solución práctica para tu caso de uso.
- Monitorea el rendimiento después de implementar; mide latencia, errores y uso de recursos antes y después.
- Revisa la sección de Troubleshooting ante errores comunes; la mayoría tienen causa raíz documentada con solución.
- Mantén dependencias actualizadas y ejecuta tests en CI para prevenir regresiones en producción.
Errores Comunes en Producción
- Copiar el ejemplo sin adaptarlo a volúmenes y modos de fallo reales.
- Saltar tests de carga e inyección de errores antes del primer despliegue productivo.
- Codificar valores fijos que deberían ser configurables por entorno.
- Olvidar agregar logging y monitoreo en cada paso.
- Desplegar sin plan de rollback ni estrategia de backup probada.
- Asumir que el ejemplo mínimo escalará sin agregar caché o procesamiento por lotes.
- No documentar la versión y configuración usadas en producción.
- Dejar la receta sin cambios cuando evolucionan las dependencias o la escala.
Preguntas frecuentes
¿Esta solución está lista para producción?
Sí. Los ejemplos de código arriba muestran implementaciones probadas. Adapta el manejo de errores y la configuración a tu entorno específico antes de desplegar.
¿Cuáles son las características de rendimiento?
El rendimiento depende de tu volumen de datos e infraestructura. Las soluciones mostradas priorizan claridad. Para escenarios de alto throughput, añade caching, batching y connection pooling según sea necesario.
¿Cómo depuro problemas con este enfoque?
Empieza con el ejemplo mínimo de arriba. Añade logging en cada paso. Prueba con entradas pequeñas primero, luego escala. Usa el debugger de tu lenguaje para revisar los edge cases.
Recursos Relacionados
Desarrollo Local de Microservicios con Docker Compose
Orquesta entornos locales multi-servicio con Docker Compose incluyendo bases de datos, caches, message brokers y reverse proxies con hot reload y redes compartidas
RecipeDiseñar Sistemas Event-Driven con Event Buses y Brokers
Cómo construir sistemas débilmente acoplados usando eventos, event buses, message brokers y event sourcing para comunicación asíncrona que escala entre servicios.
PatternPatrón Circuit Breaker
Previene fallos en cascada deteniendo solicitudes a servicios que están fallando. Un patrón arquitectural para sistemas distribuidos resilientes.
RecipeColas de mensajes muertos (Dead Letter Queues)
Maneja mensajes fallidos gracefulmente con dead letter queues, políticas de retry y detección de poison pills en arquitecturas message-driven.
RecipeMicroservicios Event-Driven
Diseña microservicios event-driven con message brokers, event sourcing, CQRS y patrones de consistencia eventual.
RecipeIdempotencia en Procesamiento de Mensajes
Hacé que el procesamiento de mensajes sea idempotente para que entregas duplicadas no generen side effects en sistemas event-driven.