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.idermö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.type—snappy,lz4oderzstd. Komprimierung reduziert die Netzwerk- und Festplattennutzung; zstd bietet in der Regel die beste Kompressionsrate.fetch.min.bytesundfetch.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 untermax.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=2und 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 untermax.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.recordszu hoch, wodurchmax.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.