Zum Inhalt springen

Echtzeit-Datenpipelines mit Apache Kafka und TypeScript

Baue Echtzeit-Datenpipelines mit Apache Kafka und TypeScript: Partitionierung, Consumer Groups, Exactly-once-Semantik und Dead Letter Queues.

4 Min. Lesezeit
Ein Kafka-Cluster, in dem Producer Events über Topic-Partitionen an Consumer Groups senden, die die Daten parallel verarbeiten

Warum Kafka für Echtzeit-Pipelines

Batch-Verarbeitung hat ihre Daseinsberechtigung, aber moderne Anwendungen müssen auf Events reagieren, sobald sie eintreten – Betrugserkennung, Echtzeit-Analysen, Bestandsaktualisierungen, Benachrichtigungssysteme. Kafka bietet eine verteilte, dauerhafte Event-Streaming-Plattform mit hohem Durchsatz, die Producer und Consumer voneinander entkoppelt und die Reihenfolge der Nachrichten innerhalb einzelner Partitionen garantiert.

Events produzieren

Producer senden Events an Kafka-Topics. Jedes Event hat einen Key, der bestimmt, in welcher Partition es landet. Events mit demselben Key gehen immer in dieselbe Partition, wodurch die Reihenfolge für diesen Key garantiert wird.

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 und parallele Verarbeitung

Consumer in derselben Consumer Group teilen sich die Partitionen untereinander auf. Hat ein Topic 12 Partitionen und die Gruppe 4 Consumer, verarbeitet jeder Consumer 3 Partitionen. Werden weitere Consumer hinzugefügt, wird die Arbeitslast automatisch neu verteilt.

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 für fehlgeschlagene Nachrichten

Manche Nachrichten lassen sich nicht verarbeiten – ungültige Daten, fehlende Referenzen, vorübergehende Bugs. Anstatt die Partition zu blockieren oder die Nachricht zu verlieren, werden fehlgeschlagene Nachrichten in eine Dead Letter Queue geleitet, wo sie untersucht und erneut verarbeitet werden können.

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}`
  );
}

Exactly-once-Verarbeitung mit Transaktionen

Kafka unterstützt Transaktionen, die Nachrichten produzieren und Consumer-Offsets atomar committen. Dadurch entsteht Exactly-once-Semantik: Jede Nachricht wird genau einmal verarbeitet, selbst wenn Fehler auftreten.

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;
  }
}

Schemaentwicklung mit versionierten Events

Mit der Weiterentwicklung deiner Domäne ändern sich auch die Event-Schemas. Nutze eine Schema Registry oder eine eingebettete Versionierung, um Abwärts- und Vorwärtskompatibilität sicherzustellen.

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);
  }
}

Wichtigste Erkenntnisse

Kafka-Pipelines entkoppeln Producer von Consumern und ermöglichen so Echtzeitverarbeitung im großen Maßstab. Nutze Nachrichten-Keys, um die Reihenfolge pro Entität zu garantieren – alle Events einer Bestellung landen in derselben Partition. Consumer Groups parallelisieren die Verarbeitung automatisch; füge weitere Consumer hinzu, um horizontal zu skalieren.

Deaktiviere Auto-Commit und verwalte Offsets für zuverlässige Verarbeitung manuell. Leite fehlgeschlagene Nachrichten in Dead Letter Queues um, statt Partitionen zu blockieren oder Daten zu verlieren. Nutze Kafka-Transaktionen für Exactly-once-Semantik, wenn du Events zwischen Topics transformierst. Versioniere deine Event-Schemas von Anfang an – nachträgliche Breaking Changes im Event-Format sind mühsam zu beheben. Starte einfach mit einem Producer, einer Consumer Group und einer Dead Letter Queue; ergänze Transaktionen und Schema Registries, sobald deine Pipeline reift.

Wilfredo Rujel

Wilfredo Rujel

Full-Stack-Softwareentwickler

Diesen Beitrag teilenX