StackPractices
intermediate Por Mathias Paulenko

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=all en el producer para asegurar que los mensajes se escriban a todas las replicas antes de confirmar

Mejores Practicas

  • Usa nombres de topic significativos: orders. created no topic1. Namespacea por dominio y tipo de evento.
  • Particiona por clave: g.
  • Batchea los sends del producer: enviar mensajes en batches mejora el throughput. sizeylinger. 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=true para 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 transactionalId al 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.ms correctamente — consumers que procesan lentamente son expulsados del group
  • Usar auto_offset_reset=none sin 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.