Programmierung und Softwareentwicklung

Wie man eine Kafka-Verarbeitungspipeline baut, die die Reihenfolge von Sitzungen mit Go wahrt

Joshua Oluikpe erläutert den praktischen Entwurf zur Verarbeitung von Kafka-Nachrichten in strikter Reihenfolge innerhalb jeder Sitzung, während unabhängige Sitzungen weiterhin parallel arbeiten können. Der Entwurf basiert auf konsistentem Hashing, einem Go-Worker pro Sitzung, Wiederholungsversuchen innerhalb desselben Pfads und der Verwaltung zusammenhängender Commit-Punkte, um Nachrichtenverluste nach einem Ausfall zu vermeiden.

2026-10-07
5 Min. Lesezeit
1 Aufrufe
certi.news Editorial Team
Wie man eine Kafka-Verarbeitungspipeline baut, die die Reihenfolge von Sitzungen mit Go wahrt

Es reicht nicht aus, sich auf die Kafka-Reihenfolge innerhalb einer Partition zu verlassen, wenn sich Tausende Sitzungen dieselbe Partition teilen. In Chat- und sequenziellen Aufgabensystemen kann beispielsweise eine Korrekturanweisung wie „Setze das Ziel auf Paris“ nach der Nachricht „Buche einen Flug nach London“ eintreffen. Wird sie vor der ersten Nachricht ausgeführt, wird der Kontext verfälscht. Die von Joshua Oluikpe bei InfoQ beschriebene Lösung verlagert einen Teil der Verantwortung für die Reihenfolge in die Anwendungsschicht, die in Go geschrieben ist.

Reihenfolge innerhalb einer Sitzung und Parallelität zwischen Sitzungen

Der Entwurf basiert auf zwei Ebenen von Workern. Die Verteilungsebene empfängt Kafka-Datensätze und leitet sie anhand der Sitzungs-ID mithilfe von konsistentem Hashing weiter. Dadurch gelangen die Nachrichten derselben Sitzung zum selben Worker, während verschiedene Sitzungen auf mehrere Worker verteilt werden können. Anschließend erhält jede aktive Sitzung eine eigene Goroutine, die jeweils eine Nachricht verarbeitet und dies in der Reihenfolge ihres Eingangs tut.

Dadurch wird eine langsame Sitzung nicht zum Hindernis für andere Sitzungen, wie es in einem flachen Worker-Pool geschehen kann, wenn der Wiederholungsversuch für eine Sitzung den gesamten Worker blockiert. Goroutines werden bei Bedarf erstellt und nach einer gewissen Inaktivitätsdauer entfernt. Dadurch hängt der Ressourcenverbrauch von der Zahl der aktiven Sitzungen statt von der Gesamtzahl möglicher Sitzungen ab.

Wiederholungsversuche, ohne Nachrichten zu überspringen

Der Wiederholungsversuch erfolgt innerhalb der Goroutine derselben Sitzung. Tritt ein vorübergehender Fehler auf, etwa das Ablaufen des Timeouts eines abhängigen Dienstes, wartet der Worker mit exponentiellem Backoff und einem gewissen Zufallsanteil und verarbeitet anschließend dieselbe Nachricht erneut, bevor er mit den nachfolgenden Nachrichten fortfährt. Da die folgenden Nachrichten hinter der fehlgeschlagenen Nachricht im Kanal verbleiben, kann die vierte Nachricht die dritte nicht überholen.

Nicht behebbare Fehler, etwa ungültige Daten, werden direkt an eine Dead-Letter-Queue (DLQ) weitergeleitet. Diese Trennung zwischen vorübergehenden und endgültigen Fehlern macht eine separate Zustandsmaschine zur Verwaltung von Wiederholungsversuchen überflüssig. Sie bedeutet jedoch auch, dass sich die Nachrichten einer bestimmten Sitzung während der Backoff-Phase ansammeln können.

Sicheres Commit und Wiederherstellung nach einem Ausfall

Wenn mehrere Sitzungen ihre Nachrichten parallel verarbeiten, darf nicht einfach der höchste abgeschlossene Offset committed werden. Wird beispielsweise Offset 104 vor Offset 102 abgeschlossen und anschließend 104 committed, kann ein Ausfall des Consumers dazu führen, dass 102 endgültig übersprungen wird. Daher verfolgt das System für jede Partition die Offsets, die sich in Verarbeitung befinden, sowie diejenigen, die abgeschlossen wurden, und verschiebt den Commit-Punkt nur bis zum höchsten lückenlos abgeschlossenen Offset.

Ein Neustart kann infolgedessen bei einem älteren Offset beginnen und einige Nachrichten erneut verarbeiten. Der Entwurf verwendet eine unveränderliche Ereignis-ID, um bei der Deduplizierung zu helfen, behauptet jedoch keine „Exactly-once“-Verarbeitung. Externe Auswirkungen hängen weiterhin davon ab, ob die abhängigen Dienste mit Duplikaten idempotent umgehen können. In seltenen Fällen kann eine festgefahrene Lücke nach Ablauf eines festgelegten Timeouts übersprungen werden, um den Fortschritt der Partition wiederherzustellen. Dies wird als operativer Fehler protokolliert und die Nachricht an die DLQ weitergeleitet. Dabei handelt es sich um einen ausdrücklichen Kompromiss zwischen Verfügbarkeit und Vollständigkeit der Verarbeitung.

Was ist für den tatsächlichen Produktionseinsatz erforderlich?

Die Reihenfolgegarantien sind nicht allein durch die grundlegende Logik vollständig. Während einer Neuzuweisung von Partitionen stoppt der Consumer den Empfang neuer Datensätze, wartet auf den Abschluss der laufenden Arbeit und committed anschließend den letzten zusammenhängenden Watermark, bevor er die Besitzrechte überträgt. Wenn die Worker-Kanäle voll sind, pausiert er den Abruf der Nachrichten vorübergehend, anstatt Nachrichten zu verwerfen. Bereits abgerufene Datensätze werden in partitionseigenen Puffern gespeichert, damit neuere Nachrichten nicht an ihnen vorbeiziehen.

Das System benötigt außerdem eine individuelle Erkennung festgefahrener Offsets, eine durchgängige Weitergabe des Tracings über die asynchronen Kafka-Grenzen hinweg, die Aufrechterhaltung der für die automatische Skalierung erforderlichen Rückstandssignale während des pausierten Abrufs sowie das Abarbeiten der Aufgaben und das Schreiben in die DLQ vor dem Herunterfahren. Das Produkt wendet konsistentes Hashing auch auf der Senderseite an, um die Reihenfolge der Nachrichten bei asynchroner Produktion sicherzustellen.

Ergebnisse und Einschränkungen

In einem synthetischen Test mit 50.000 Nachrichten und 10 Sitzungen, bei dem bei 10 % der Nachrichten vorübergehende Fehler injiziert wurden, verarbeitete das System 14.027 Nachrichten pro Sekunde ohne Verletzungen der Reihenfolge bei einer medianen Latenz von 11,7 Millisekunden. In einem Lasttest mit 1.000 Sitzungen lag der tatsächliche Durchsatz bei 48.805 Nachrichten pro Sekunde gegenüber einem Ziel von 50.000, ohne Übertragungsfehler. Die Tail-Latenz überschritt jedoch das Budget von 200 Millisekunden aufgrund des Nachrichtenratenlimits der API Gateway und nicht wegen des Mechanismus zur Sitzungsreihenfolge.

Der Autor berichtet, dass das System mehr als 40 Millionen Nachrichten aus der Produktion ohne erkennbare Verletzungen der Reihenfolge verarbeitet habe und lediglich 460 Nachrichten an die DLQ weitergeleitet worden seien, also etwa 0,002 %. Diese Zahl beruht jedoch auf der operativen Überwachung und nicht auf einer umfassenden Prüfung der Abfolge jeder einzelnen Nachricht. Außerdem überprüfte der umfangreichere Lasttest die Reihenfolge für jede Sitzung nicht unabhängig. Der praktische Wert des Entwurfs besteht daher darin, die Reihenfolge von einer Annahme auf Kafka-Ebene in eine von der Anwendung erzwungene Eigenschaft umzuwandeln. Gleichzeitig bleiben explizite Sequenztests und Idempotenzgarantien in den abhängigen Diensten erforderlich.

Nachrichtenquelle
InfoQ - Architecture Articles
Originalquelle öffnen ↗
c
Autor

certi.news Editorial Team

Aus derselben Kategorie

Das könnte Sie interessieren

Alle Nachrichten anzeigen