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.

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.
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 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.
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.
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.
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.
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.


