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.

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.
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
});
}// ❌ 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.
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.
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.
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.
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.


