Event Streaming

Apache Kafka

Kafka ist ein verteiltes, Append-only Commit Log und keine Queue. Records werden einmal geschrieben und mehrfach gelesen, was Replay, Fan-out und Ordnungsgarantien ermöglicht, die ein klassischer Broker nicht bieten kann.

advanced15 min readUpdated 16. Sept. 2026
producer.ts
ts
// producer.ts
import { Kafka } from "kafkajs";

const kafka = new Kafka({
  clientId: "orders",
  brokers: ["localhost:9092"],
});

const producer = kafka.producer();

await producer.connect();
await producer.send({
  topic: "orders.created",
  messages: [
    // The key decides the partition, so all events for one
    // customer stay in order.
    { key: order.customerId, value: JSON.stringify(order) },
  ],
});
await producer.disconnect();
Veröffentlicht
2011
Ursprung
LinkedIn
Modell
Append-only commit log
Reihenfolge
Pro Partition
Standard-Zustellung
At-least-once
Retention
Zeit- oder größenbasiert

Warum es wichtig ist

Warum Teams zu Kafka greifen

Ein dauerhaftes, replaybares Log

Records werden angehängt und beim Lesen niemals entfernt. Jeder Consumer kann die Historie zurückspulen und erneut verarbeiten, was eine klassische Queue schlichtweg nicht leisten kann.

Partitions für Skalierung und Ordnung

Ein Topic wird in Partitions aufgeteilt, die über Broker verteilt sind. Parallelität ergibt sich aus den Partitions, und die Reihenfolge ist innerhalb jeder einzelnen Partition garantiert.

Viele unabhängige Leser

Consumer Groups tracken jeweils ihre eigenen Offsets. So können ein Billing-Job, ein Suchindex-Indexer und eine Analytics-Pipeline dieselben Events lesen, ohne sich gegenseitig zu beeinflussen.

Das Gesamtbild

Die drei Kernkonzepte hinter Kafka

Ein Topic ist ein Append-only Log, Producers schreiben Keyed Records hinein und Consumers tracken ihre eigene Position mittels Offsets.

Topic

Append

Ein benanntes, partitioniertes Log. Producers hängen Records am Ende an, und jeder Record behält seinen Offset, bis die Retention abläuft.

Consumer group

Share

Eine Gruppe von Consumern teilt die Partitions unter sich auf und trackt ihre eigenen committed Offsets, sodass Skalierung und Fortschritt unabhängig bleiben.

Offset

Track

Die Position einer Gruppe innerhalb einer Partition. Das Committen des Offsets ist die Bestätigung, die über At-least-once versus At-most-once entscheidet.

HTML5 auf einen Blick

Kafkas Bausteine

Topics

Benannte Append-only Logs, die Records speichern, bis die Retention abläuft.

Partitions

Die Einheit für Parallelität und die Grenze für die Reihenfolgegarantie.

Producers

Serialisieren einen Record, wählen einen Key und hängen ihn an eine Partition an.

Consumer groups

Teilen Partitions unter ihren Mitgliedern auf und committen Offsets.

Offsets

Die Position einer Gruppe in einer Partition und ihr Bestätigungspunkt.

Retention

Löschen nach Zeit oder Größe, oder Log Compaction, um den neuesten Wert pro Key zu behalten.

Ablauf

Die Reise eines Records

Ein Record wird einmal geschrieben und mehrfach gelesen. Es wird nichts entfernt, wenn ein Consumer ihn liest, was genau das Replay ermöglicht.

  1. 1

    Produce

    Der Producer serialisiert einen Record und wählt eine Partition durch Hashing seines Keys aus, sodass derselbe Key immer an derselben Stelle landet.

  2. 2

    Append

    Der Leader-Broker der Partition hängt den Record an das Ende seines Logs an und gibt den neuen Offset an den Producer zurück.

  3. 3

    Consume

    Jeder Consumer in einer Gruppe besitzt eine Teilmenge der Partitions und liest seine Records in Offset-Reihenfolge.

  4. 4

    Commit the offset

    Nach der Verarbeitung committet der Consumer seine Position, damit die Gruppe nach einem Neustart dort fortfahren kann.

  5. 5

    Retain or replay

    Der Record bleibt im Log, bis die Retention abläuft, sodass eine andere Gruppe oder ein späterer Job zurückspulen und ihn erneut lesen kann.

Eine kurze Geschichte

Von LinkedIn-Logs zum Streaming-Standard

  1. 2011

    Kafka wird Open Source

    LinkedIn veröffentlicht Kafka als verteiltes Commit Log für seine Activity Streams.

    11
  2. 2012

    Apache Inkubation

    Kafka wird ein Apache-Projekt und die Verbreitung reicht weit über LinkedIn hinaus.

    12
  3. 2016

    Kafka Streams

    Das 0.10 Release liefert eine Stream-Processing-Library, später folgt ksqlDB.

    16
  4. 2017

    Exactly-once Semantik

    Das 0.11 Release fügt idempotente Producers und Transaktionen für Exactly-once Pipelines hinzu.

    17
  5. 2021

    KRaft beginnt

    KIP-500 beginnt, ZooKeeper durch ein in die Broker integriertes Metadata Quorum zu ersetzen.

    21
  6. 2025

    Kafka 4.0 entfernt ZooKeeper

    Neue Cluster laufen nur noch im KRaft-Modus, die ZooKeeper-Ära endet.

    25

Der vollständige Leitfaden

Apache Kafka: Alles was Sie wissen müssen

Was Kafka eigentlich ist

Apache Kafka ist ein verteiltes, Append-only Commit Log. Dieser eine Satz erklärt fast alles Weitere. Datensätze werden am Ende eines Logs angehängt, jeder erhält eine monoton steigende Nummer, den sogenannten Offset, und es wird niemals etwas direkt an Ort und Stelle verändert. Leser entfernen Datensätze nicht, wenn sie diese konsumieren; sie verschieben lediglich einen Cursor nach vorne.

Dies ist der grundlegende Unterschied zu einer klassischen Message Queue. In einer Queue nimmt ein Consumer eine Nachricht entgegen und die Nachricht ist anschließend weg. In Kafka liest ein Consumer einen Datensatz, merkt sich dessen Offset, und der Datensatz bleibt dort gespeichert, solange es die Retention Policy des Topics erlaubt. Zehn verschiedene Consumer – und damit zehn verschiedene Anwendungen – können denselben Datensatz unabhängig voneinander und in ihrem eigenen Tempo lesen, ohne sich untereinander abstimmen zu müssen.

Kafka wurde bei LinkedIn entwickelt, um Activity Streams zu verarbeiten: Seitenaufrufe, Klicks, Log-Zeilen – alles mit Millionen von Events pro Sekunde. Es wurde 2011 als Open-Source-Projekt veröffentlicht und wurde 2012 ein Apache-Projekt. Heute ist es das Standard-Backbone für Event Streaming: Change Data Capture, Metrics Pipelines, Event Sourcing, Log Aggregation und Stream Processing.

Wenn Sie von RabbitMQ oder Redis kommen, besteht der mentale Wechsel darin, nicht mehr in Kategorien von „einer zuzustellenden Nachricht“ zu denken, sondern in „einer aufgezeichneten Tatsache“.

Topics, Partitions und Reihenfolge

Ein Topic ist ein benannter, dauerhafter Log. Producer schreiben hinein, Consumer lesen daraus. Topics sind in Partitions unterteilt, wobei die Partition die Einheit für sowohl Parallelität als auch Reihenfolge ist.

  • Records innerhalb einer Partition sind strikt nach ihrem Offset geordnet.
  • Es gibt keine Garantie für die Reihenfolge über Partitions hinweg.
  • Mehr Partitions ermöglichen mehr parallele Consumer, allerdings auf Kosten von mehr Dateien, mehr Replikation und langsameren Rebalances.

Dies ist die wichtigste Garantie in Kafka: Die Reihenfolge gilt pro Partition. Wenn zwei Events in einer bestimmten Reihenfolge verarbeitet werden müssen, müssen sie in derselben Partition landen. Wenn dies nicht der Fall ist, kann Kafka sie an verschiedene, gleichzeitig laufende Consumer übergeben, wodurch die Reihenfolge verloren geht.

Partitions bestimmen zudem die maximale Consumer-Parallelität innerhalb einer Gruppe: Eine Gruppe kann maximal einen Consumer pro Partition haben, der diese aktiv liest. Sechs Partitions bedeuten also maximal sechs nützliche Consumer; ein siebter würde untätig bleiben.

# three partitions, replicated across three brokers
kafka-topics.sh --create \
  --topic orders.created \
  --partitions 6 \
  --replication-factor 3 \
  --bootstrap-server localhost:9092

Das nachträgliche Hinzufügen von Partitions ist möglich, verändert jedoch das Mapping vom Key zur Partition für bestehende Keys. Dadurch können Events einer Entität auf zwei verschiedene Partitions verteilt werden und ihre relative Reihenfolge verlieren. Legen Sie die Anzahl der Partitions bereits bei der Erstellung des Topics fest und planen Sie dabei ausreichend Puffer für zukünftiges Wachstum ein.

Producer, Keys und Partitionierung

Ein Producer serialisiert einen Record und entscheidet, welcher Partition dieser zugeordnet wird. Der Standard-Partitioner hasht den Key des Records und ordnet ihn einer Partition zu. Derselbe Key landet immer in derselben Partition; so wird die Reihenfolge für eine einzelne Entität gewahrt, während verschiedene Entitäten über die Partitionen verteilt werden.

await producer.send({
  topic: "orders.created",
  messages: [
    { key: order.customerId, value: JSON.stringify(order) },
  ],
});

Die Wahl des Keys ist eine Design-Entscheidung und kein bloßes Detail. Nutzt man customerId als Key, sind alle Events für einen Kunden gemeinsam geordnet. Nutzt man orderId, erreicht man eine maximale Verteilung, aber keine Reihenfolge über verschiedene Events hinweg. Records mit einem null-Key werden zur Balance verteilt – Kafkas Sticky Partitioner füllt eine Partition, bevor er zur nächsten übergeht –, bieten jedoch keinerlei Garantie bezüglich der Reihenfolge.

Der Producer steuert zudem die Durability und den Durchsatz über Batching und Acknowledgements. acks legt fest, wie viele Replicas einen Schreibvorgang bestätigen müssen, linger.ms und batch.size bestimmen, wie lange gewartet wird, bis ein Batch gefüllt ist, und compression.type entscheidet, wie viel CPU-Leistung zugunsten von Netzwerk und Disk aufgewendet wird. Mehr dazu weiter unten.

Broker, Replikation und der Leader

Ein Broker ist ein Kafka-Server. Ein Cluster besteht aus mehreren zusammenarbeitenden Brokern. Jede Partition hat einen Leader-Broker und null oder mehr Follower. Producer und Consumer kommunizieren mit dem Leader; die Follower replizieren das Log.

Die Replikation ist das, was Kafka ausfallsicher macht. Wenn ein Leader ausfällt, wird eine der In-Sync Replicas (ISR) befördert und der Cluster arbeitet ohne Unterbrechung weiter. Die Einstellung acks bestimmt, worauf der Producer wartet:

  • acks=0 — Fire-and-Forget; am schnellsten und am wenigsten sicher.
  • acks=1 — Der Leader hat die Daten geschrieben; Daten gehen verloren, wenn der Leader vor der Replikation ausfällt.
  • acks=all — Alle In-Sync Replicas haben den Empfang bestätigt; am sichersten.

Kombinieren Sie acks=all mit min.insync.replicas=2 und einem Replikationsfaktor von drei für ein Production-Setup: Ein Schreibvorgang wird erst bestätigt, wenn mindestens zwei Replicas die Daten haben, sodass der Ausfall eines einzelnen Brokers nicht zu Datenverlust führt.

const producer = kafka.producer({
  idempotent: true,
  maxInFlightRequests: 1,
  transactionalId: "orders-producer",
});

Der idempotente Producer fügt jedem Batch eine Sequenznummer hinzu, sodass der Broker Duplikate verwerfen kann. Dies verhindert versehentliche Duplikate, die sonst durch Retries entstehen könnten. Ein Cluster besitzt zudem einen Controller, der die Partition-Leadership und die Metadaten verwaltet. Ältere Kafka-Versionen nutzten ZooKeeper dafür; modernes Kafka verwendet KRaft, wobei die Broker ihr eigenes Metadaten-Quorum bilden und ZooKeeper vollständig entfällt.

Consumer-Gruppen und Partition-Zuweisung

Consumer gehören zu einer Consumer-Gruppe, die durch groupId identifiziert wird. Kafka weist jede Partition genau einem Consumer innerhalb der Gruppe zu. Das bietet zwei Vorteile gleichzeitig: horizontale Skalierung, da Partitionen aufgeteilt werden, und Load Balancing, da keine zwei Consumer in einer Gruppe dieselbe Partition verarbeiten.

const consumer = kafka.consumer({ groupId: "billing" });

await consumer.connect();
await consumer.subscribe({ topic: "orders.created", fromBeginning: true });

await consumer.run({
  eachMessage: async ({ partition, message }) => {
    const order = JSON.parse(message.value!.toString());
    console.log(`p${partition} @ ${message.offset}`, order.id);
  },
});

Wenn Consumer beitreten oder die Gruppe verlassen, führt Kafka ein Rebalancing durch: Zuweisungen werden entzogen und neue vergeben. Rebalances pausieren den Konsum, weshalb sie kostspielig sind; kooperative Rebalancing-Protokolle reduzieren diese Unterbrechungen, indem nur die Partitionen verschoben werden, die tatsächlich verschoben werden müssen. Ein Consumer, der innerhalb von session.timeout.ms keine Heartbeats mehr sendet, wird als tot betrachtet und seine Partitionen werden neu zugewiesen.

Über die Gruppe merkt sich Kafka zudem den Fortschritt. Jede Gruppe hat ihre eigenen committed Offsets, die im internen __consumer_offsets Topic gespeichert werden. So können zwei Gruppen, die dasselbe Topic lesen, an völlig unterschiedlichen Positionen stehen, ohne dass eine Koordination erforderlich ist.

Rebalancing und der Poll-Loop

Im Hintergrund ist ein Consumer ein Poll-Loop. Er ruft Records ab, übergibt sie an Ihren Handler und ruft anschließend erneut Daten ab. Der Broker überwacht die Liveness separat über Heartbeats, sodass ein langsamer Handler nicht sofort als “tot” gilt. Es gibt jedoch ein hartes Limit: Wenn ein einzelner poll länger als max.poll.interval.ms dauert, geht der Broker davon aus, dass der Consumer feststeckt, und löst ein Rebalance aus.

const consumer = kafka.consumer({
  groupId: "billing",
  sessionTimeout: 45_000,
  heartbeatInterval: 3_000,
  rebalanceTimeout: 60_000,
});

Ein Rebalance wird ausgelöst, wenn ein Consumer beitritt, die Gruppe verlässt oder entfernt wird, wenn sich Partitions oder Topics ändern oder wenn sich eine Subscription ändert. Während des Rebalancing entzieht Kafka die Zuweisungen und pausiert den Consumption, weshalb häufige Rebalances den Durchsatz massiv beeinträchtigen. Drei Gewohnheiten helfen dabei, sie selten zu halten:

  • Begrenzen Sie die Arbeit pro Record. Ein Handler, der gelegentlich minutenlang läuft, wird früher oder später max.poll.interval.ms überschreiten.
  • Nutzen Sie Cooperative Rebalancing. Das inkrementelle Cooperative-Protokoll verschiebt nur die Partitions, die tatsächlich verschoben werden müssen, anstatt die gesamte Gruppe zu stoppen.
  • Nutzen Sie Static Membership. Die Definition einer stabilen group.instance.id ermöglicht es einem neu gestarteten Consumer, seine alten Partitions zurückzufordern, ohne überhaupt ein Rebalance auszulösen.

Kafka bietet Ihnen zudem consumer.pause() und consumer.resume(), um Backpressure anzuwenden, wenn eine Downstream-Dependency Probleme bereitet. Das Pausieren ist besser, als den Poll-Loop zu blockieren: Der Consumer sendet weiterhin Heartbeats, bleibt in der Gruppe und hört einfach auf, Daten abzurufen, bis Sie wieder bereit sind.

Offsets und Delivery-Garantien

Ein Offset ist die Position einer Consumer-Gruppe innerhalb einer Partition. Wenn ein Consumer einen Record verarbeitet, kann er den Offset committen, was Kafka signalisiert: „Diese Gruppe hat alles bis zu diesem Punkt abgeschlossen“. Committet man zu früh, geht bei einem Absturz Arbeit verloren; committet man zu spät, wird Arbeit nach einem Absturz erneut verarbeitet.

  • At-most-once — Commit vor der Verarbeitung. Bei einem Absturz geht der Record verloren.
  • At-least-once — Verarbeitung, dann Commit. Bei einem Absturz wird der Record erneut verarbeitet. Dies ist der gängige Standard.
  • Exactly-once — Transaktionen und idempotente Producer machen den Read-Process-Write-Zyklus innerhalb von Kafka atomar.

Auto-Commit erfolgt im Hintergrund über einen Timer. Das ist bequem, aber gefährlich, da Offsets für Records committet werden können, die sich noch in der Verarbeitung befinden. Ein manueller Commit nach Abschluss der Arbeit ist der zuverlässige Standard.

await consumer.run({
  autoCommit: false,
  eachMessage: async ({ topic, partition, message }) => {
    await chargeOrder(JSON.parse(message.value!.toString()));
    await consumer.commitOffsets([
      { topic, partition, offset: String(Number(message.offset) + 1) },
    ]);
  },
});

Da At-least-once die Norm ist, sollten Handler idempotent sein. Der Guide zu Batch Processing beschreibt denselben Vertrag für Job-Queues: Gehen Sie davon aus, dass die Arbeit zweimal ausgeführt werden kann, und sorgen Sie dafür, dass der zweite Durchlauf ein No-Op ist. Ein Deduplizierungs-Key, der aus der Arbeit selbst – nicht aus dem Kafka-Offset – abgeleitet wird, macht den Handler sicher, selbst wenn dasselbe logische Event zweimal produziert wird.

Exactly-once ist real, aber spezifischer, als es klingt. Es funktioniert für Kafka-zu-Kafka-Pipelines mit transaktionalen Producern und isolation.level=read_committed, aber in dem Moment, in dem Sie in eine externe Datenbank schreiben, benötigen Sie dort wieder idempotente Writes.

Retention, Replay und Log Compaction

Kafka bewahrt Datensätze gemäß einer Retention Policy auf, nicht basierend auf einer Zustellbestätigung. retention.ms (standardmäßig sieben Tage) und retention.bytes begrenzen jede Partition; sobald einer dieser Werte überschritten wird, werden alte Log-Segmente gelöscht. Da beim Lesen nichts entfernt wird, kann ein Consumer die Historie replayen, indem er zu einem früheren Offset springt oder die Gruppe zurücksetzt.

# rewind a group to the beginning of a topic
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --group billing --topic orders.created \
  --reset-offsets --to-earliest --execute

Replay ist die Superkraft von Kafka. Ein neuer Service kann gestartet werden, indem die gesamte Historie eines Topics gelesen wird. Ein Bugfix kann bereitgestellt und der letzte Tag erneut verarbeitet werden. Eine Analytics-Pipeline kann aus dem rohen Event-Log neu aufgebaut werden. Eine Queue kann nichts davon leisten, da die Daten bereits gelöscht sind.

Für keyed state gibt es einen weiteren Retention-Modus: Log Compaction. Mit cleanup.policy=compact behält Kafka mindestens den neuesten Wert für jeden Key bei und verwirft ältere Werte. Ein null-Wert ist ein Tombstone, der den Key löscht. Das Log wird so zu einem Changelog, mit dem eine Tabelle wiederhergestellt werden kann – genau das, was Kafka Streams für seine State Stores nutzt und worauf Change-Data-Capture-Pipelines basieren.

# keep the latest value per key instead of deleting by age
kafka-configs.sh --bootstrap-server localhost:9092 \
  --alter --entity-type topics --entity-name user.profiles \
  --add-config cleanup.policy=compact,min.cleanable.dirty.ratio=0.1

Schemas und die Schema Registry

Da das Log länger existiert als jede einzelne Anwendung, wird die Struktur eines Records zu einem Vertrag zwischen den Teams. Ein Producer und ein Consumer sind durch die Bytes, die sie austauschen, gekoppelt – und diese Kopplung bleibt auch über Deployments hinweg bestehen. Eine Schema Registry verwandelt dies in ein verwaltetes Kompatibilitätsproblem, anstatt in eine böse Überraschung.

Die Registry speichert versionierte Schemas – üblicherweise Avro, Protobuf oder JSON Schema – und weist jedem eine numerische ID zu. Producer registrieren ein Schema und schreiben die ID zusammen mit dem Payload; Consumer rufen das Schema per ID ab und deserialisieren es. Wenn sich ein Schema ändert, erzwingt die Registry einen Kompatibilitätsmodus (z. B. backward oder forward) und lehnt Änderungen ab, die bestehende Reader beeinträchtigen würden.

import { SchemaRegistry } from "@kafkajs/confluent-schema-registry";

const registry = new SchemaRegistry({ host: "http://localhost:8081" });

const encoded = await registry.encode(schemaId, {
  orderId: order.id,
  totalCents: order.totalCents,
});

await producer.send({
  topic: "orders.created",
  messages: [{ key: order.customerId, value: encoded }],
});

Auch wenn Sie keine Registry betreiben, sollten Sie Payloads wie eine API behandeln: Fügen Sie Felder hinzu, anstatt sie umzubenennen, geben Sie neuen Feldern einen Default-Wert und führen Sie eine Versionierung ein, wenn sich die Bedeutung ändert. Ein Consumer, der die vorherige Version ausführt, muss in der Lage sein, einen Record zu lesen, der von der nächsten Version geschrieben wurde.

Kafka Connect und Kafka Streams

Zwei Teile der Plattform ersparen es Ihnen, dasselbe “Plumbing” doppelt zu schreiben.

Kafka Connect ist ein Framework, um Daten mittels Konfiguration statt Code in und aus Kafka zu bewegen. Source Connectors ziehen Daten aus Datenbanken, Object Storage oder SaaS APIs; Sink Connectors pushen diese in Warehouses, Suchindizes oder eine andere Datenbank. Debezium beispielsweise verwandelt ein Postgres Write-Ahead Log in einen Stream von Change-Events. Connect läuft als Cluster, trackt Offsets und führt Retries bei Fehlern durch – es ist daher die Standardantwort auf die Fragen “Wie bringe ich Daten in Kafka?” und “Wie bringe ich Daten aus Kafka heraus?”.

Kafka Streams ist eine Client-Library für die Verarbeitung von Daten in Kafka. Sie bietet Ihnen eine Stream-DSL — map, filter, groupByKey, join, window — basierend auf KStream und KTable Abstraktionen, mit State Stores, die durch compacted topics abgesichert sind. Sie läuft innerhalb Ihrer Anwendung, skaliert durch das Hinzufügen von Instanzen und bewältigt Fehlertoleranz über Changelog-Topics. Wenn Sie SQL bevorzugen, bietet ksqlDB eine Query-Layer über denselben Konzepten an.

const stream = builder.stream("orders.created");

stream
  .filter((key, order) => order.totalCents > 10_000)
  .groupBy((key) => order.customerId)
  .windowedBy(tumblingWindow({ size: 60 * 60 * 1000 }))
  .count()
  .toStream()
  .to("customer.hourly_orders");

Beides ist es wert, bekannt zu sein, selbst wenn Sie mit einfachen Producern und Consumern starten, da sie definieren, wie “der Kafka-Weg” in großem Maßstab aussieht.

Ein minimales End-to-End-Beispiel

Es ist hilfreich, das Ganze in einer einzigen Datei zu sehen: einen Producer und einen Consumer verbinden, abonnieren, ausführen und sauber beenden.

import { Kafka } from "kafkajs";

const kafka = new Kafka({
  clientId: "orders-app",
  brokers: ["localhost:9092"],
});

const producer = kafka.producer();
const consumer = kafka.consumer({ groupId: "orders-app" });

async function main() {
  await producer.connect();
  await producer.send({
    topic: "orders.created",
    messages: [{ key: "customer_1", value: JSON.stringify({ id: "order_1" }) }],
  });

  await consumer.connect();
  await consumer.subscribe({ topic: "orders.created", fromBeginning: true });
  await consumer.run({
    eachMessage: async ({ partition, message }) => {
      const order = JSON.parse(message.value!.toString());
      console.log(`p${partition} @ ${message.offset}`, order.id);
    },
  });
}

async function shutdown() {
  await consumer.disconnect();
  await producer.disconnect();
  process.exit(0);
}

process.on("SIGTERM", shutdown);
process.on("SIGINT", shutdown);

main().catch((err) => {
  console.error(err);
  process.exit(1);
});

Drei Dinge in dieser Datei sind für die Produktion entscheidend. Der Producer wird einmal verbunden und wiederverwendet, anstatt ihn pro Nachricht neu zu erstellen. Der Consumer abonniert das Topic, bevor er startet. Und der Prozess verarbeitet SIGTERM, da ein hartes Beenden während eines Rebalance oder eines Commits die Gruppe in einem schlechteren Zustand hinterlässt als ein Graceful Shutdown.

Event Sourcing und das Outbox-Pattern

Kafka wird oft als das Rückgrat von Event Sourcing beschrieben, wobei das Log die „Source of Truth“ ist und der aktuelle Zustand eine Projektion von Events darstellt. Anstatt nur die aktuellste Zeile zu speichern, speichern Sie die Sequenz von Fakten — order.created, order.paid, order.shipped — und stellen jede Ansicht wieder her, indem Sie diese erneut abspielen. Compacted Topics machen die Projektion zu einer Tabelle; die Retention macht sie prüfbar.

Die Schwierigkeit liegt nicht im Schreiben des Events, sondern darin, das Event und die Datenbankzeile atomar zu schreiben. Ein Prozess kann nach dem Commit der Zeile, aber vor dem Publishing abstürzen, oder nach dem Publishing, aber vor dem Commit. Das Publishing innerhalb einer Datenbanktransaktion ist unmöglich, daher ist die Standardlösung das Outbox-Pattern: Schreiben Sie das Event in eine outbox-Tabelle in derselben Transaktion wie die Zustandsänderung und lassen Sie dann ein Relay oder einen CDC-Connector diese Zeilen an Kafka übertragen.

BEGIN;
INSERT INTO orders (id, status, total_cents)
VALUES ($1, 'created', $2);

INSERT INTO outbox (id, topic, payload)
VALUES ($1, 'orders.created', $2);
COMMIT;

Debezium oder ein kleines Polling-Relay liest dann die Outbox aus und produziert die Daten in das Topic, wobei die Zeilen nach dem Publishing gelöscht werden. Da der Outbox-Schreibvorgang Teil der Transaktion ist, existiert das Event genau dann, wenn auch der Zustand existiert. Consumer müssen dennoch idempotent sein, da das Relay nach einem Absturz eine Zeile doppelt veröffentlichen kann.

Lokale Entwicklung und Testing

Sie benötigen keinen vollständigen Kafka-Cluster für die Entwicklung. Ein Single-Node-Broker in Docker Compose oder Redpanda im Kompatibilitätsmodus bietet Ihnen Topics, Consumer-Gruppen und Offsets direkt auf Ihrem Laptop.

services:
  kafka:
    image: redpandadata/redpanda:latest
    command: >
      redpanda start --overprovisioned --smp 1
      --kafka-addr PLAINTEXT://0.0.0.0:9092
      --advertise-kafka-addr PLAINTEXT://localhost:9092
    ports:
      - "9092:9092"

Für Tests gilt dasselbe Muster wie bei jeder anderen Integrationsabhängigkeit: Starten Sie den Broker in einem Container, erstellen Sie die für den Test benötigten Topics und verwenden Sie eine eindeutige Group ID pro Testlauf, damit committete Offsets niemals zwischen den Läufen durchsickern. Setzen Sie Offsets explizit zurück, wenn ein Test den Verlauf lesen muss, und prüfen Sie die vom Consumer empfangenen Records anstatt auf das Timing zu setzen.

const groupId = `test-${crypto.randomUUID()}`;
const consumer = kafka.consumer({ groupId });

Halten Sie die Consumer-Handler rein und schlank – Parsen, Validieren, Service-Aufruf –, sodass der Großteil der Logik ohne Broker mittels Unit-Tests geprüft werden kann. Reservieren Sie die Integrationstests für das Wiring: Erreicht ein produzierter Record den richtigen Handler und wird der Offset danach committet?

Durchsatz, Latenz und Batching

Kafka ist schnell, weil es Batching nutzt und sequenziell schreibt. Beides lässt sich konfigurieren, wobei das Tuning im Grunde ein Regler zwischen Latenz und Durchsatz ist.

  • linger.ms — wie lange ein Producer wartet, um einen Batch zu sammeln. Ein höherer Wert bedeutet größere Batches und mehr Durchsatz, allerdings auf Kosten einer höheren Latenz.
  • batch.size — die maximale Byte-Zahl pro Partition-Batch.
  • compression.typesnappy, lz4 oder zstd. Komprimierung reduziert die Netzwerk- und Festplattennutzung; zstd bietet in der Regel die beste Kompressionsrate.
  • fetch.min.bytes und fetch.max.wait.ms — wie lange Consumer warten, bis eine Fetch-Antwort gefüllt ist.
  • max.poll.records — wie viele Records ein Consumer pro Poll verarbeitet. Erhöhen Sie diesen Wert, um Backlogs schneller abzuarbeiten, aber halten Sie die Verarbeitungszeit unter max.poll.interval.ms, da der Consumer sonst aus der Gruppe geworfen wird.
const producer = kafka.producer({
  linger: { ms: 20 },
  compression: CompressionTypes.GZIP,
});

const consumer = kafka.consumer({
  groupId: "billing",
  maxBytesPerPartition: 1_048_576,
  maxWaitTimeInMs: 500,
});

Ein auf Durchsatz optimierter Producer könnte linger.ms=20, einen Ein-Megabyte-Batch und zstd verwenden. Ein Producer, der innerhalb von einstelligen Millisekunden publizieren muss, nutzt linger.ms=0. Es gibt keine einzelne richtige Antwort; es geht nur um den Trade-off, den Sie wählen, und den p99-Wert, mit dem Sie leben können.

Die Wahl zwischen Kafka, RabbitMQ und Redis

Diese drei werden oft so verglichen, als wären sie austauschbar. Das sind sie nicht.

  • Kafka ist ein replayable Log. Wähle es, wenn du einen hohen Durchsatz, eine dauerhafte Historie, Fan-out an viele unabhängige Consumer, Event Sourcing, Stream Processing oder die Möglichkeit zur erneuten Verarbeitung benötigst. Es ist in der Bedienung am aufwendigsten.
  • RabbitMQ ist ein Message Broker. Wähle es, wenn das Routing entscheidend ist: Exchanges, Routing Keys, Bestätigungen pro Nachricht (Acknowledgements), Prioritäten und Dead-Lettering. Es eignet sich hervorragend für die Aufgabenverteilung und komplexes Routing und ist bei moderaten Volumina leichter zu betreiben als Kafka. Siehe den RabbitMQ-Guide.
  • Redis ist ein In-Memory-Datenspeicher, der gleichzeitig als Job-Queue fungiert. Wähle es, wenn die Aufgaben simpel sind, das Volumen moderat ist und du Redis bereits im Einsatz hast. BullMQ auf Redis bietet dir Retries, Scheduling und ein UI bei fast keinem operationalen Overhead. Siehe Redis Queues.

Eine nützliche Faustregel: Wenn Consumer Nachrichten erneut abspielen (replay), an viele Gruppen verteilen (fan out) oder die Historie lesen müssen, nutze Kafka. Wenn eine Nachricht an einen spezifischen Worker geroutet und danach vergessen werden soll, nutze RabbitMQ. Wenn es sich um einen Hintergrundjob mit einer Retry-Policy handelt, nutze Redis.

Sicherheit und Zugriffskontrolle

Kafka enthält in der Regel die wertvollsten Daten im gesamten System, daher ist eine strikte Absicherung unerlässlich.

  • Verschlüsselung — aktivieren Sie TLS für den Traffic zwischen den Brokern sowie zwischen Client und Broker. Plaintext Kafka innerhalb einer VPC bleibt weiterhin Plaintext.
  • Authentifizierung — nutzen Sie SASL/SCRAM oder mTLS für Clients. Vermeiden Sie nicht authentifizierte Listener, außer in kurzlebigen lokalen Setups.
  • Autorisierung — verwenden Sie ACLs, um Lese-, Schreib- oder Erstellungsrechte für spezifische Topics und Groups zu vergeben. Ein Service sollte nicht in der Lage sein, jedes Topic zu lesen, nur weil er eine Verbindung herstellen kann.
  • Quotas — begrenzen Sie die Bandbreite für Producer und Consumer pro Client, damit ein außer Kontrolle geratener Service nicht den gesamten Cluster blockiert.
  • Secrets — betten Sie Zugangsdaten niemals direkt in einen Producer ein; injizieren Sie diese über die Umgebungsvariablen oder einen Secret Manager.

Monitoring und Consumer Lag

Kafka scheitert still und leise, wenn man es zulässt. Vier Signale decken den Großteil der betrieblichen Realität ab.

  • Consumer lag — die Differenz zwischen dem neuesten Offset und dem committed Offset einer Gruppe pro Partition. Ein steigender Lag ist das früheste Anzeichen dafür, dass die Consumer nicht hinterherkommen.
  • Under-replicated partitions — ein Wert ungleich Null bedeutet, dass ein Broker down oder langsam ist und die Durability gefährdet ist.
  • Request- und Netzwerk-Latenz — Perzentile auf Broker-Seite zeigen auf, wann Festplatten oder das Netzwerk den Flaschenhals bilden.
  • Festplattenbelegung und Segment-Anzahl — Kafka ist disk-bound, und die Replikation multipliziert jedes Byte mit dem Replikationsfaktor.

Stellen Sie den Lag im selben Dashboard wie Ihre Services dar und richten Sie Alarme basierend auf dem Trend ein, nicht auf einzelne Spitzen. Ein Lag, der über den Tag stetig wächst, ist ein Kapazitätsproblem; ein Lag, der kurzzeitig ansteigt und sich dann wieder erholt, ist in der Regel ein Rebalance oder ein langsames Deployment.

Best Practices

  • Entwerfen Sie Keys basierend auf der Sortierung, die Sie tatsächlich benötigen; bedenken Sie, dass die Sortierung pro Partition erfolgt.
  • Verwenden Sie acks=all, min.insync.replicas=2 und einen Replikationsfaktor von drei für wichtige Topics.
  • Aktivieren Sie den idempotent producer und bevorzugen Sie manuelle offset commits.
  • Gestalten Sie Consumer idempotent, da at-least-once der Standard-Contract ist.
  • Halten Sie die max.poll.records-Verarbeitung unter max.poll.interval.ms, um rebalance storms zu vermeiden.
  • Nutzen Sie log compaction für keyed state und tombstones für Löschvorgänge.
  • Überwachen Sie den consumer lag pro Partition und nicht nur die Gesamtsumme des Clusters.
  • Behandeln Sie record schemas als versionierten Contract, idealerweise mit einer registry, wenn mehr als ein Team involviert ist.
  • Trennen Sie Topics nach Lebenszyklus und retention, anstatt alles in ein einziges Topic zu werfen.
  • Legen Sie die retention basierend auf realen Anforderungen fest; sieben Tage sind ein Standardwert, keine Richtlinie.
  • Bevorzugen Sie managed Kafka, es sei denn, der Betrieb von Kafka ist tatsächlich Ihr Kerngeschäft.
  • Nutzen Sie Connect für Standard-Integrationen und Streams für die Verarbeitung innerhalb des Clusters, bevor Sie eigenen glue-code schreiben.

Häufige Fehler

  • Die Annahme einer globalen Reihenfolge, obwohl Kafka die Reihenfolge nur innerhalb einer Partition garantiert.
  • Die Verwendung eines zufälligen oder null-Keys für Events, die zwingend in der richtigen Reihenfolge bleiben müssen.
  • Das automatische Committen von Offsets, bevor die Arbeit abgeschlossen ist, was bei einem Absturz zum Datenverlust führt.
  • Die Implementierung von nicht-idempotenten Handlern, wodurch bei einem Rebalance Seiteneffekte dupliziert werden.
  • Das Erstellen eines Topics mit nur einer Partition und die anschließende Verwunderung darüber, dass Consumer nicht skalieren können.
  • Das nachträgliche Erhöhen der Partitionsanzahl, was die Key-Routing-Logik stillschweigend verändert.
  • Die Behandlung von Kafka als Request/Response-API mit Antworten pro Nachricht.
  • Das Setzen von max.poll.records zu hoch, wodurch max.poll.interval.ms überschritten wird.
  • Das Vergessen, dass die Replikation die Speicher- und Netzwerkkosten vervielfacht.
  • Das Auslassen des Consumer-Lags im Dashboard, bis der Backlog bereits mehrere Stunden beträgt.
  • Der Betrieb von ZooKeeper in einem neuen Cluster im Jahr 2026.

Wie geht es weiter?

Kafka ist das Log im Zentrum eines event-gesteuerten Systems. Die Backend-Roadmap behandelt event-gesteuerte Architekturen und Microservices – also wie Events, Producer und Consumer ineinandergreifen und ein Log zum Rückgrat zwischen den Services wird. Wenn Sie verschiedene Broker vergleichen, zeigen die Guides zu RabbitMQ und Redis Queues die Alternativen mit Fokus auf Routing bzw. Job-Verarbeitung. Und da Consumer im Grunde Hintergrundworker mit Offsets sind, behandelt der Guide zu Batch Processing die Themen Retries, Idempotenz und Observability, die auch hier relevant sind.

In der Praxis

Produce, consume, commit, administer

Die vier Operationen, aus denen fast jede Kafka-Anwendung besteht.

producer.ts
import { Kafka } from "kafkajs";

const kafka = new Kafka({
  clientId: "orders",
  brokers: ["localhost:9092"],
});

const producer = kafka.producer({
  idempotent: true,
  maxInFlightRequests: 1,
});

await producer.connect();

await producer.send({
  topic: "orders.created",
  acks: -1, // wait for all in-sync replicas
  messages: [
    { key: order.customerId, value: JSON.stringify(order) },
  ],
});

await producer.disconnect();

Kafka Log vs. klassische Queue

Eine Queue löscht eine Nachricht, sobald sie bestätigt wurde. Ein Kafka Log behält sie, sodass jede Gruppe die Historie replayen kann und viele Leser denselben Stream nutzen können.

Kafka Log
// A second group can read the same history
// from the beginning, months later.
await consumer.subscribe({
  topic: "orders.created",
  fromBeginning: true,
});
Klassische Queue
// The message was acknowledged by the first
// worker and removed; it cannot be replayed
// or read by a second independent consumer.
await channel.ack(msg);

Auto-Commit vs. manueller Commit

Auto-Commit läuft über einen Timer und kann über Records hinausspringen, die noch verarbeitet werden. Ein Commit nach Abschluss der Arbeit macht das Fehlerverhalten explizit.

Manuell
await consumer.run({
  autoCommit: false,
  eachMessage: async ({ topic, partition, message }) => {
    await handle(message);
    await consumer.commitOffsets([
      {
        topic,
        partition,
        offset: String(Number(message.offset) + 1),
      },
    ]);
  },
});
Auto
await consumer.run({
  // Offsets are committed on a timer, so a
  // crash can skip records that were never
  // processed.
  eachMessage: async ({ message }) => {
    await handle(message);
  },
});

Abwägungen

Ist Kafka den operationalen Aufwand wert?

Kafka löst Probleme, die Queues nicht bewältigen können, erfordert aber echte Infrastruktur und ein anderes mentales Modell. Wähle es wegen der Garantien, nicht wegen des Hypes.

Strengths

  • Replay ändert alles

    Da die Historie erhalten bleibt, kannst du einen neuen Service aus dem Log bootstrappen, eine Projektion nach einem Bugfix neu aufbauen und Analytics auf Raw Events ausführen, lange nachdem sie passiert sind.

  • Fan-out ohne Koordination

    Beliebig viele unabhängige Consumer Groups lesen dasselbe Topic in ihrem eigenen Tempo. Ein zusätzlicher Leser kostet den Producer nichts und stört bestehende Consumer nicht.

  • Horizontal skalierender Durchsatz

    Sequenzielle Schreibvorgänge, Batching und Partitions ermöglichen es einem kleinen Cluster, Millionen von Records pro Sekunde zu absorbieren. Skalierung erfolgt durch Hinzufügen von Brokern und Partitions.

Trade-offs

  • Es ist eine Plattform, keine Library

    Broker, Replikation, Rebalances, Disk-Sizing und Consumer Lag müssen alle überwacht werden. Managed Services helfen, aber Kafka ist niemals eine Abhängigkeit, die man einfach vergisst.

  • Routing ist nicht seine Aufgabe

    Kafka hat keine Exchanges oder Routing Keys. Consumer lesen ganze Topics, daher muss selektives Per-Message-Routing im Topic-Design oder im Anwendungscode implementiert werden.

  • Reihenfolge nur innerhalb einer Partition

    Eine globale Reihenfolge gibt es nicht. Um die Ordnung korrekt zu gewährleisten, müssen Keys und Partition-Counts sorgfältig gewählt werden; Partitions können später nicht hinzugefügt werden, ohne das Key-Routing zu brechen.

Häufig gestellte Fragen

Häufig gestellte Fragen

Keep learning

Related topics from the roadmap.

$ Lernen Sie jetzt

Bereit, Apache Kafka zu lernen?

Unser interaktives Tutorial führt Sie Schritt für Schritt durch Apache Kafka — mit Quizzen und echtem Code, den Sie im Browser ausführen können.