Saltar al contenido

Pipelines de Datos en Tiempo Real con Apache Kafka y TypeScript

Crea pipelines de datos en tiempo real con Kafka y TypeScript: particionado, consumer groups, semántica exactly-once y dead letter queues.

4 min de lectura
Un clúster de Kafka con producers que envían eventos a través de las partitions de un topic hacia consumer groups que procesan los datos en paralelo

Por qué Kafka para pipelines en tiempo real

El procesamiento por lotes tiene su lugar, pero las aplicaciones modernas necesitan reaccionar a los eventos a medida que ocurren: detección de fraude, analítica en tiempo real, actualización de inventario, sistemas de notificaciones. Kafka ofrece una plataforma de streaming de eventos distribuida, duradera y de alto rendimiento que desacopla los producers de los consumers y garantiza el orden de los mensajes dentro de cada partition.

Producción de eventos

Los producers envían eventos a los topics de Kafka. Cada evento tiene una key que determina a qué partition va. Los eventos con la misma key siempre van a la misma partition, lo que garantiza el orden para esa key.

tstypescript
import { Kafka, Partitioners } from "kafkajs";
 
const kafka = new Kafka({
  clientId: "order-service",
  brokers: ["kafka-1:9092", "kafka-2:9092", "kafka-3:9092"],
});
 
const producer = kafka.producer({
  createPartitioner: Partitioners.DefaultPartitioner,
  idempotent: true,    // Prevent duplicate messages on retry
  maxInFlightRequests: 5,
  retry: { retries: 5 },
});
 
interface OrderEvent {
  eventType: "order.created" | "order.updated" | "order.cancelled";
  orderId: string;
  userId: string;
  payload: Record<string, unknown>;
  timestamp: string;
}
 
async function publishOrderEvent(event: OrderEvent): Promise<void> {
  await producer.send({
    topic: "orders",
    messages: [
      {
        // Key = orderId ensures all events for an order go to same partition
        key: event.orderId,
        value: JSON.stringify(event),
        headers: {
          "event-type": event.eventType,
          "correlation-id": crypto.randomUUID(),
        },
      },
    ],
    acks: -1,  // Wait for all replicas to acknowledge
  });
}
tstypescript
// ❌ No key — events for same order scattered across partitions
await producer.send({
  topic: "orders",
  messages: [{ value: JSON.stringify(event) }], // Random partition
});
 
// ✅ Key-based partitioning preserves ordering per entity
await producer.send({
  topic: "orders",
  messages: [{
    key: event.orderId,  // Same order always same partition
    value: JSON.stringify(event),
  }],
});

Consumer groups y procesamiento en paralelo

Los consumers de un mismo consumer group se reparten las partitions entre sí. Si un topic tiene 12 partitions y el consumer group tiene 4 consumers, cada consumer procesa 3 partitions. Al añadir consumers, la carga de trabajo se rebalancea automáticamente.

tstypescript
const consumer = kafka.consumer({
  groupId: "analytics-pipeline",
  sessionTimeout: 30_000,
  heartbeatInterval: 3_000,
  maxWaitTimeInMs: 100,
});
 
async function startConsumer(): Promise<void> {
  await consumer.connect();
  await consumer.subscribe({
    topics: ["orders"],
    fromBeginning: false,
  });
 
  await consumer.run({
    autoCommit: false, // Manual commit for exactly-once processing
    eachMessage: async ({ topic, partition, message }) => {
      const event: OrderEvent = JSON.parse(message.value!.toString());
 
      try {
        await processEvent(event);
 
        // Commit only after successful processing
        await consumer.commitOffsets([
          {
            topic,
            partition,
            offset: (Number(message.offset) + 1).toString(),
          },
        ]);
      } catch (error) {
        await handleProcessingError(event, error as Error);
      }
    },
  });
}

Dead letter queues para mensajes fallidos

Algunos mensajes no se pueden procesar: datos inválidos, referencias faltantes, errores transitorios. En lugar de bloquear la partition o perder el mensaje, los fallos se enrutan a una dead letter queue para su investigación y reprocesamiento.

tstypescript
interface DeadLetterMessage {
  originalTopic: string;
  originalPartition: number;
  originalOffset: string;
  originalKey: string | null;
  originalValue: string;
  errorMessage: string;
  errorStack: string;
  failedAt: string;
  retryCount: number;
}
 
async function handleProcessingError(
  event: OrderEvent,
  error: Error,
  context: { topic: string; partition: number; offset: string }
): Promise<void> {
  const deadLetter: DeadLetterMessage = {
    originalTopic: context.topic,
    originalPartition: context.partition,
    originalOffset: context.offset,
    originalKey: event.orderId,
    originalValue: JSON.stringify(event),
    errorMessage: error.message,
    errorStack: error.stack ?? "",
    failedAt: new Date().toISOString(),
    retryCount: 0,
  };
 
  await producer.send({
    topic: `${context.topic}.dead-letter`,
    messages: [
      {
        key: event.orderId,
        value: JSON.stringify(deadLetter),
      },
    ],
  });
 
  console.error(
    `Event ${event.orderId} sent to dead letter queue: ${error.message}`
  );
}

Procesamiento exactly-once con transacciones

Kafka admite transacciones que producen mensajes y confirman offsets de consumer de forma atómica. Esto logra una semántica exactly-once: cada mensaje se procesa exactamente una vez, incluso ante fallos.

tstypescript
const transactionalProducer = kafka.producer({
  idempotent: true,
  transactionalId: "analytics-transformer",
  maxInFlightRequests: 1,
});
 
async function processWithTransaction(
  messages: Array<{ topic: string; partition: number; message: KafkaMessage }>
): Promise<void> {
  const transaction = await transactionalProducer.transaction();
 
  try {
    for (const { topic, partition, message } of messages) {
      const event: OrderEvent = JSON.parse(message.value!.toString());
 
      // Transform and produce to downstream topic
      const enriched = await enrichEvent(event);
      await transaction.send({
        topic: "orders-enriched",
        messages: [
          {
            key: event.orderId,
            value: JSON.stringify(enriched),
          },
        ],
      });
 
      // Commit consumer offset within the transaction
      await transaction.sendOffsets({
        consumerGroupId: "analytics-pipeline",
        topics: [
          {
            topic,
            partitions: [
              {
                partition,
                offset: (Number(message.offset) + 1).toString(),
              },
            ],
          },
        ],
      });
    }
 
    await transaction.commit();
  } catch (error) {
    await transaction.abort();
    throw error;
  }
}

Evolución del esquema con eventos versionados

A medida que tu dominio evoluciona, los esquemas de los eventos cambian. Usa un schema registry o un versionado embebido para gestionar la compatibilidad hacia atrás y hacia adelante.

tstypescript
interface VersionedEvent<T> {
  schemaVersion: number;
  eventType: string;
  data: T;
}
 
// Version-aware deserializer
type EventHandler<T> = (data: T) => Promise<void>;
 
class EventDeserializer {
  private handlers: Map<string, Map<number, EventHandler<unknown>>> = new Map();
 
  register<T>(
    eventType: string,
    version: number,
    handler: EventHandler<T>
  ): void {
    if (!this.handlers.has(eventType)) {
      this.handlers.set(eventType, new Map());
    }
    this.handlers.get(eventType)!.set(version, handler as EventHandler<unknown>);
  }
 
  async handle(raw: string): Promise<void> {
    const event: VersionedEvent<unknown> = JSON.parse(raw);
    const versionHandlers = this.handlers.get(event.eventType);
 
    if (!versionHandlers) {
      throw new Error(`Unknown event type: ${event.eventType}`);
    }
 
    const handler = versionHandlers.get(event.schemaVersion);
    if (!handler) {
      // Try to find the latest version handler that can upcast
      const latest = Math.max(...versionHandlers.keys());
      const latestHandler = versionHandlers.get(latest);
      if (latestHandler) {
        const upcasted = upcastEvent(event, latest);
        await latestHandler(upcasted.data);
        return;
      }
      throw new Error(
        `No handler for ${event.eventType} v${event.schemaVersion}`
      );
    }
 
    await handler(event.data);
  }
}

Puntos clave

Los pipelines de Kafka desacoplan a los producers de los consumers, lo que permite un procesamiento en tiempo real a gran escala. Usa keys de mensaje para garantizar el orden por entidad: todos los eventos de un pedido van a la misma partition. Los consumer groups paralelizan el procesamiento automáticamente; añade consumers para escalar horizontalmente.

Desactiva el auto-commit y gestiona los offsets manualmente para un procesamiento confiable. Enruta los mensajes fallidos a dead letter queues en lugar de bloquear partitions o perder datos. Usa transacciones de Kafka para lograr semántica exactly-once al transformar eventos entre topics. Versiona los esquemas de tus eventos desde el primer día: los cambios incompatibles en el formato de los eventos son dolorosos de corregir después. Empieza simple, con un producer, un consumer group y una dead letter queue; añade transacciones y schema registries a medida que tu pipeline madure.

Wilfredo Rujel

Wilfredo Rujel

Ingeniero de Software Full Stack

Compartir esta publicaciónX