Einheit 4.3: Event-Streaming-Integration
Einführung
Sie verbinden den API-Dienst von TaskBoard mit dessen Hintergrundworker. Wenn ein Benutzer eine Aufgabe erstellt, verschiebt oder schließt, darf die API nicht darauf blockieren, eine E-Mail oder einen Webhook zu senden. Stattdessen erzeugt die API ein task-changed-Ereignis für Kafka und gibt sofort die Kontrolle zurück, während ein separater Worker-Dienst diese Ereignisse konsumiert und Benachrichtigungen auslöst. Dies entkoppelt den Anfragepfad von langsamen Nebeneffekten und ermöglicht es, den Worker unabhängig zu skalieren.
IONOS CLOUD Event Streams for Apache Kafka verwendet Apache Kafka 4.0.0, sodass Ihre vorhandenen Kafka-Client-Bibliotheken unverändert funktionieren. Im Vergleich zu einem selbst verwalteten Broker unterscheidet sich die Verbindung: Die Datenebene erfordert gegenseitiges TLS, nicht SASL oder Klartext, und die Broker-Adressen sowie Zertifikate stammen aus dem Cluster, das in Einheit 2.5 bereitgestellt wurde. Diese Einheit beginnt bei der Client-Verbindung und baut dann schrittweise das Topic-Design, die zuverlässige Zustellung, die Behandlung von Dead-Letter-Queues und die Skalierung der Konsumenten auf, alles im Kontext des TaskBoard-Ereignisflusses.
1. Verbindung von Clients mit mTLS
Die IONOS CLOUD Kafka Daten-Ebene ist mit TLS und gegenseitiger Authentifizierung gesichert. Beide Seiten authentifizieren einander: Ihr Client prüft das Broker-Zertifikat gegen die Zertifizierungsstelle des Clusters, und der Broker prüft Ihr Client-Zertifikat, das von der eigenen Zertifizierungsstelle des Clusters signiert wird. Es gibt keinen Klartext- oder SASL/PLAIN-Listener. Jeder Producer und Consumer präsentiert ein Client-Zertifikat.
Sie rufen die Zugangsdaten, die Zertifizierungsstelle, den privaten Client-Schlüssel und das Client-Zertifikat vom Cluster-Zugriffsendpunkt ab. Dies sind dieselben Werte, die das Cluster über seine Broker-Adressen und die Zugriffs-API bereitstellt. Jedes Client-Zertifikat ist 365 Tage gültig, daher gehört die Zertifikatsrotation in Ihr operatives Runbook, bevor das Jahr zu Ende geht.
1.1 Client-Verbindungs-Konfiguration
Die Konfiguration eines selbst verwalteten Kafka-Clients verwendet security.protocol=SSL mit einem Truststore (für die Cluster-Zertifizierungsstelle) und einem Keystore (für das Client-Zertifikat und den Schlüssel). Die IONOS CLOUD Dokumentation zeigt die kanonische Client-Eigenschaftsdatei:
security.protocol=SSL
ssl.truststore.type=PKCS12
ssl.truststore.location=ca-cert.p12
ssl.truststore.password=changeit
ssl.endpoint.identification.algorithm=
ssl.keystore.type=PKCS12
ssl.keystore.location=admin.p12
ssl.keystore.password=adminp12pass
Hinweis: ssl.endpoint.identification.algorithm bleibt leer. Da der Cluster eine private CA anstelle öffentlich signierter Zertifikate verwendet, ist die Hostname-Verifizierung gegenüber den Broker-Adressen hier deaktiviert. Der Truststore validiert stattdessen die Kette.
1.2 Producer-Verbindung in Python
Die confluent-kafka-Bibliothek akzeptiert dieselben SSL-Schlüssel. Verweisen Sie den Client auf die Broker-Adressen aus Ihrer Terraform-Ausgabe und übergeben Sie die CA, das Zertifikat und den Schlüssel als PEM-Dateien.
from confluent_kafka import Producer
# Broker addresses come from the ionoscloud_kafka_cluster Terraform output.
conf = {
"bootstrap.servers": "192.168.1.101:9093,192.168.1.102:9093,192.168.1.103:9093",
"security.protocol": "ssl",
"ssl.ca.location": "ca-cert.pem",
"ssl.certificate.location": "client-cert.pem",
"ssl.key.location": "client-key.pem",
"ssl.endpoint.identification.algorithm": "none",
"client.id": "taskboard-api",
}
producer = Producer(conf)
Der Cluster betreibt drei Broker über die Cluster-Topologie, daher sollten alle drei in bootstrap.servers für die Verbindungsresilienz aufgeführt werden. Der Client benötigt nur einen erreichbaren Broker zum Starten, aber das Auflisten aller drei vermeidet einen einzelnen Ausfallpunkt während des Startvorgangs.
2. Design von Topics und Partitions
Beim Topic-Design werden die IONOS CLOUD-spezifischen Einschränkungen erstmals spürbar. Verwenden Sie ein Topic pro Ereignistyp, anstatt ein einzelnes Firehose-Topic, da Verbraucher sich für die Ereignisse anmelden, die sie interessieren, und Sie Aufbewahrungsdauer und Partitionierung je nach Anliegen festlegen können. Für TaskBoard erstellen Sie ein task-changed-Topic für den Worker und später ein task-changed.dlq-Dead-Letter-Topic.
2.1 Erstellen von Topics über die API
Die Kafka-Verwaltungs-API ist regional und getrennt von der cloudapi. Der Host folgt dem Muster https://kafka.<region>.ionos.com, und im Gegensatz zur Datenebene authentifiziert sich die Verwaltungs-API mit einem Bearer-Token, nicht mit mTLS. Erstellen Sie ein Topic mit POST /clusters/{clusterId}/topics:
curl -X POST \
'https://kafka.de-txl.ionos.com/clusters/e69b22a5-8fee-56b1-b6fb-4a07e4205ead/topics' \
--header 'Content-Type: application/json' \
--header 'Authorization: Bearer '"$IONOS_TOKEN" \
--data '{
"properties": {
"name": "task-changed",
"replicationFactor": 3,
"numberOfPartitions": 6,
"logRetention": { "retentionTime": 604800000 }
}
}'
Die Body-Parameter sind name (das einzige Pflichtfeld), replicationFactor, numberOfPartitions und die logRetention-Einstellungen. Themen lassen sich mit GET /clusters/{clusterId}/topics und GET /clusters/{clusterId}/topics/{topicId} zurücklesen, welche name, replicationFactor, numberOfPartitions und logRetention zurückgeben. Die Löschung ist DELETE /clusters/{clusterId}/topics/{topicId} und gibt 202 Accepted zurück, da die Entfernung asynchron erfolgt.
2.2 Partition Count and Replication Constraints
Wählen Sie die Anzahl der Partitionen als Vielfaches von 3 (3, 6, 9, 12 usw.), um eine ungleichmäßige Verteilung der Partitionen auf die drei Broker zu vermeiden. Eine Partitionsanzahl, die kein Vielfaches der Brokeranzahl ist, führt dazu, dass ein Broker mehr Partitionen als die anderen trägt, was die Lastverteilung verzerrt.
Der empfohlene Replikationsfaktor ist 3, passend zur Topologie mit drei Brokern. Dies ist jedoch eine Empfehlung von IONOS CLOUD und kein vom System erzwungener Standardwert. Sie müssen replicationFactor beim Erstellen eines Themas explizit festlegen. Bei einem Replikationsfaktor von 3 behält der Cluster drei Kopien jeder Nachricht, sodass das Thema den Verlust eines Brokers ohne Datenverlust übersteht. Die standardmäßige Aufbewahrungszeit für Themen beträgt 604800000 ms (7 Tage). Die Einstellung von retentionTime auf -1 setzt keine zeitliche Grenze. In diesem Fall wird die Aufbewahrung nur durch den bereitgestellten Clusterspeicher begrenzt (beispielsweise 750 GB bei Größe S, 1200 GB bei Größe M).
Die folgende Tabelle zeigt die Zuordnung der Body-Parameter eines Themas zu ihrer Rolle bei der Themenkonzeption.
| Parameter | Rolle | TaskBoard-Wert |
|---|---|---|
name |
Themenkennung (Pflichtfeld) | task-changed |
replicationFactor |
Kopien pro Partition über die Broker | 3 |
numberOfPartitions |
Obergrenze für die Parallelität der Konsumenten | 6 (Vielfaches von 3) |
logRetention.retentionTime |
Lebensdauer der Nachrichten in ms (-1 = unbegrenzt) |
604800000 (7 Tage) |
Partitionen sind nicht nur für den Speicher relevant. Die Partitionsanzahl ist die feste Obergrenze für die Konsumentenparallelität innerhalb einer Konsumentengruppe. Legen Sie sie daher jetzt für den Spitzen-Durchsatz fest, da eine spätere Erhöhung der Partitionsanzahl die Zuordnung von Schlüsseln zu Partitionen ändert und die Reihenfolgegarantien stört.
3. Zuverlässigkeit des Produzenten
Ein Fire-and-forget-Produzent verwirft Nachrichten bei Störungen der Broker. Für die Task-Events von TaskBoard, bei denen ein verlorenes Event eine verpasste Benachrichtigung bedeutet, konfigurieren Sie den Produzenten für Persistenz, bevor Sie auf Durchsatz optimieren.
3.1 Acks, Idempotenz und Schlüssel
Setzen Sie acks=all, damit der Produzent darauf wartet, dass alle in-sync Replikate einen Schreibvorgang bestätigen, bevor er ihn als erfolgreich betrachtet. In Kombination mit einem Replikationsfaktor von 3 bedeutet dies, dass eine Nachricht über Broker hinweg persistent ist, bevor Ihr Code weiterverarbeitet wird. Aktivieren Sie den idempotenten Produzenten, damit Wiederholungen keine Duplikate auf der Broker-Seite erzeugen.
from confluent_kafka import Producer
import json
conf = {
"bootstrap.servers": "192.168.1.101:9093,192.168.1.102:9093,192.168.1.103:9093",
"security.protocol": "ssl",
"ssl.ca.location": "ca-cert.pem",
"ssl.certificate.location": "client-cert.pem",
"ssl.key.location": "client-key.pem",
"ssl.endpoint.identification.algorithm": "none",
"acks": "all",
"enable.idempotence": True,
"retries": 5,
"linger.ms": 20,
"batch.size": 65536,
}
producer = Producer(conf)
def on_delivery(err, msg):
if err:
# Persist for replay; never silently drop.
print(f"delivery failed: {err}")
else:
print(f"delivered to {msg.topic()} [{msg.partition()}] @ {msg.offset()}")
def emit_task_changed(task_id: str, event: dict):
producer.produce(
topic="task-changed",
key=task_id.encode("utf-8"), # same task -> same partition -> ordered
value=json.dumps(event).encode("utf-8"),
on_delivery=on_delivery,
)
producer.poll(0)
Die Verwendung von task_id als Schlüssel leitet alle Ereignisse für eine Aufgabe in dieselbe Partition, wodurch die Reihenfolge pro Aufgabe erhalten bleibt. Ohne einen Schlüssel verteilt Kafka Datensätze round-robin, sodass ein „task closed“-Ereignis vor einem „task created“-Ereignis verarbeitet werden könnte.
3.2 Batching und Flushing
linger.ms und batch.size tauschen einige Millisekunden Latenz gegen eine deutlich höhere Durchsatzrate aus, indem sie Datensätze pro Partition bündeln. Rufen Sie immer flush() auf, bevor der Prozess beendet wird, da sonst gepufferte Datensätze verloren gehen.
# At shutdown, block until all buffered messages are delivered or fail.
remaining = producer.flush(timeout=10)
if remaining > 0:
raise RuntimeError(f"{remaining} messages not delivered before shutdown")
4. Consumer Groups, Offsets und Liefersemantik
Der TaskBoard-Worker liest task-changed als Consumer Group. Jeder Consumer in derselben Gruppe teilt sich die Partitionen, sodass das Hinzufügen von Worker-Pods den Durchsatz bis zur Anzahl der Partitionen skaliert. Ein vierter Consumer auf einem Thema mit 3 Partitionen bleibt untätig.
4.1 At-Least-Once mit manuellen Commits
Für die Zustellung von Benachrichtigungen ist At-Least-Once der richtige Standardwert: Verarbeiten Sie die Nachricht und commiten Sie anschließend den Offset. Wenn der Worker nach der Verarbeitung, aber vor dem Commit abstürzt, wird die Nachricht erneut zugestellt und die Benachrichtigung kann zweimal ausgelöst werden, was Sie mit einem Idempotenzschlüssel sicherstellen. Das Gegenteil, das Commit vor der Verarbeitung (At-Most-Once), birgt das Risiko, bei einem Absturz Ereignisse zu verlieren.
from confluent_kafka import Consumer, KafkaException
import json
conf = {
"bootstrap.servers": "192.168.1.101:9093,192.168.1.102:9093,192.168.1.103:9093",
"security.protocol": "ssl",
"ssl.ca.location": "ca-cert.pem",
"ssl.certificate.location": "client-cert.pem",
"ssl.key.location": "client-key.pem",
"ssl.endpoint.identification.algorithm": "none",
"group.id": "taskboard-notifier",
"enable.auto.commit": False, # manual commit = control over semantics
"auto.offset.reset": "earliest",
}
consumer = Consumer(conf)
consumer.subscribe(["task-changed"])
while True:
msg = consumer.poll(1.0)
if msg is None:
continue
if msg.error():
raise KafkaException(msg.error())
event = json.loads(msg.value())
send_notification(event) # do the work first
consumer.commit(msg) # then commit only this offset
enable.auto.commit=False ist die entscheidende Wahl. Auto-commit verschiebt Offsets auf Basis eines Timers, unabhängig davon, ob die Verarbeitung erfolgreich war, wodurch Nachrichten, deren Verarbeitung fehlgeschlagen ist, stillschweigend verworfen werden.
4.2 Skalierung von Consumern und Rebalancing
Die Anzahl der Partitionen bildet die Obergrenze für Parallelität. Um den TaskBoard-Worker auf Managed Kubernetes (siehe Einheit 3.2) zu skalieren, erhöhen Sie die Replikatenanzahl des Deployments bis zur Anzahl der Partitionen. Jeder neue Pod tritt der taskboard-notifier-Gruppe bei und erhält durch ein Rebalancing einen Teil der Partitionen zugewiesen. Während eines Rebalancings werden Partitionen kurzzeitig nicht mehr konsumiert. Halten Sie die Verarbeitung daher idempotent und die Commits häufig, um erneute Verarbeitungen zu minimieren.
| Topic-Partitionen | Maximale nützliche Consumer | Auswirkung eines weiteren Consumers |
|---|---|---|
| 3 | 3 | 4. Consumer bleibt ungenutzt |
| 6 | 6 | lineare Skalierung bis 6 |
| 12 | 12 | Spielraum für zukünftiges Wachstum |
Dimensionieren Sie die Partitionen von vornherein für Ihre maximale Consumer-Anzahl. Sie können Partitionen später hinzufügen, was jedoch die Zuordnung von Schlüsseln zu Partitionen neu anordnet und die Reihenfolge für in Bearbeitung befindliche Schlüssel aufhebt.
5. Behandlung von Dead-Letter-Nachrichten und Schema-Disziplin
Nicht jede Nachricht kann verarbeitet werden. Ein fehlerhaftes Payload oder ein Ausfall in der nachgelagerten Verarbeitung wird bei erneuter Zustellung weiterhin fehlschlagen, und eine Gift-Nachricht, die Sie endlos erneut versuchen, blockiert ihre gesamte Partition.
5.1 Muster für Dead-Letter-Themen
Leiten Sie Nachrichten, die ihre Wiederholungsversuche erschöpft haben, an ein dediziertes Dead-Letter-Thema, task-changed.dlq, weiter und bestätigen Sie anschließend den ursprünglichen Offset, damit der Hauptdatenstrom in der Partition fortgesetzt werden kann. Sie untersuchen die DLQ später, beheben die Ursache und spielen die Nachrichten bei Bedarf erneut ab.
from confluent_kafka import Producer, Consumer
import json
dlq = Producer(conf_producer) # same mTLS settings as the main producer
def process_with_dlq(consumer, msg, max_attempts=3):
attempts = int(dict(msg.headers() or {}).get("x-attempts", b"0") or 0)
try:
send_notification(json.loads(msg.value()))
consumer.commit(msg)
except Exception as exc:
if attempts + 1 >= max_attempts:
dlq.produce(
"task-changed.dlq",
key=msg.key(),
value=msg.value(),
headers={"x-error": str(exc), "x-attempts": str(attempts + 1)},
)
dlq.flush(5)
consumer.commit(msg) # advance past the poison message
else:
raise # let it redeliver for a transient error
Unterscheiden Sie vorübergehende Fehler (z. B. ein Timeout im Benachrichtigungsdienst, der einen erneuten Versuch rechtfertigt) von dauerhaften Fehlern (z. B. ein fehlerhaftes Ereignis, das direkt in die DLQ geleitet wird). Ein blindes erneutes Versenden bei dauerhaften Fehlern blockiert die Partition; ein blindes Ablegen in die DLQ bei vorübergehenden Fehlern führt zum Verlust wiederherstellbarer Ereignisse.
5.2 Ereignisverträge mit JSON Schema
Produzenten und Konsumenten werden unabhängig voneinander bereitgestellt, daher stellt die Ereignispayload einen Vertrag zwischen beiden dar. Definieren Sie diese mit JSON Schema und validieren Sie sie auf beiden Seiten. Entwickeln Sie das Schema additiv: Fügen Sie optionale Felder hinzu, entfernen oder benennen Sie bestehende Felder niemals um, damit alte Konsumenten weiterhin funktionieren, während neue Produzenten ausgerollt werden.
import jsonschema
TASK_CHANGED_V1 = {
"type": "object",
"required": ["task_id", "action", "occurred_at"],
"properties": {
"task_id": {"type": "string"},
"action": {"enum": ["created", "moved", "closed"]},
"occurred_at": {"type": "string", "format": "date-time"},
"assignee": {"type": "string"}, # added in v1.1, optional = safe
},
"additionalProperties": False,
}
def validate_event(event: dict):
jsonschema.validate(event, TASK_CHANGED_V1) # raises ValidationError on breach
Führen Sie eine Schemaversion im Ereignis oder im Header mit, damit Verbraucher während eines Migrationszeitfensters darauf verzweigen können. Diese Disziplin ermöglicht es, die Ereignisstruktur von TaskBoard zu ändern, ohne ein koordiniertes, synchrones Neustartieren von Producer und Worker durchführen zu müssen.
API-Referenz Schnellkarte
Wichtige Endpunkte für die Kafka-Management-API. Der Host ist regional, beispielsweise https://kafka.de-txl.ionos.com:
| Methode | Endpunkt | Beschreibung |
|---|---|---|
POST |
/clusters |
Erstellen eines Kafka-Clusters |
GET |
/clusters/{clusterId} |
Abrufen von Clusterdetails und Broker-Adressen |
POST |
/clusters/{clusterId}/topics |
Erstellen eines Themas |
GET |
/clusters/{clusterId}/topics |
Auflisten aller Themen |
DELETE |
/clusters/{clusterId}/topics/{topicId} |
Löschen eines Themas (gibt 202 zurück) |
GET |
/clusters/{clusterId}/users/{userId}/access |
Abrufen der mTLS-Zugangsdaten (CA, Zertifikat, Schlüssel) |
Authentifizierung der Management-API: Authorization: Bearer <token>
Authentifizierung der Datenebene (Produzenten/Verbraucher): mutual TLS mit Client-Zertifikat
Code Lab
Ziel: Erzeugen und Verarbeiten von TaskBoard task-changed-Ereignissen auf einem IONOS CLOUD Kafka-Cluster mit einer Consumer Group, manuellen Offset-Commits und Dead-Letter-Routing bei Fehlern.
Voraussetzungen:
- IONOS CLOUD-Konto mit API-Token (
IONOS_TOKEN) - Ein in Einheit 2.5 bereitgestellter Kafka-Cluster mit verfügbaren Broker-Adressen
- Python 3.10 oder höher mit
confluent-kafkaundjsonschemainstalliert - Die Cluster-CA, das Client-Zertifikat und der Client-Schlüssel als PEM-Dateien gespeichert
Schritt 1: Broker-Adressen und Zugangsdaten abrufen
CLUSTER_ID="e69b22a5-8fee-56b1-b6fb-4a07e4205ead"
curl -s -X GET \
"https://kafka.de-txl.ionos.com/clusters/$CLUSTER_ID" \
-H "Authorization: Bearer $IONOS_TOKEN" | python3 -m json.tool
Erwartete Ausgabe:
{
"id": "e69b22a5-8fee-56b1-b6fb-4a07e4205ead",
"properties": { "name": "...", "version": "4.0.0", "size": "S",
"connections": [ { "brokerAddresses": ["192.168.1.101/24", ...] } ] }
}
Schritt 2: Erstellen der Haupt- und Dead-Letter-Themen
for T in task-changed task-changed.dlq; do
curl -s -X POST "https://kafka.de-txl.ionos.com/clusters/$CLUSTER_ID/topics" \
-H "Content-Type: application/json" -H "Authorization: Bearer $IONOS_TOKEN" \
--data "{\"properties\":{\"name\":\"$T\",\"replicationFactor\":3,\"numberOfPartitions\":6}}"
done
Erwartete Ausgabe:
{"id":"ae085c4c-...","type":"topic","properties":{"name":"task-changed",
"replicationFactor":3,"numberOfPartitions":6}}
Schritt 3: Erzeugen eines task-changed-Ereignisses
producer.produce("task-changed", key=b"task-42",
value=b'{"task_id":"task-42","action":"created","occurred_at":"2026-06-05T10:00:00Z"}')
producer.flush(10)
print("produced")
Erwartete Ausgabe:
delivered to task-changed [3] @ 0
produced
Schritt 4: Verwenden mit einer Consumer Group und manuellem Commit
consumer.subscribe(["task-changed"])
msg = consumer.poll(5.0)
print(msg.topic(), msg.partition(), msg.offset(), msg.value())
consumer.commit(msg)
Erwartete Ausgabe:
task-changed 3 0 b'{"task_id":"task-42",...}'
Schritt 5: Auslösen einer Dead-Letter-Meldung bei einem ungültigen Payload
producer.produce("task-changed", key=b"task-99", value=b'{not valid json')
producer.flush(10)
# Consumer's process_with_dlq routes it to task-changed.dlq after max_attempts
Erwartete Ausgabe:
delivered to task-changed.dlq [1] @ 0
Schritt 6: Prüfen, ob das Dead-Letter-Topic die Nachricht erhalten hat
dlq_consumer.subscribe(["task-changed.dlq"])
m = dlq_consumer.poll(5.0)
print(dict(m.headers()), m.value())
Erwartete Ausgabe:
{'x-error': b'...', 'x-attempts': b'3'} b'{not valid json'
Prüfliste:
- [ ] Themen
task-changedundtask-changed.dlqmit Replikationsfaktor 3 erstellt - [ ] Producer liefert mit
acks=allund aktivierter Idempotenz - [ ] Consumer bestätigt Offsets manuell nach der Verarbeitung
- [ ] Eine fehlerhafte Nachricht landet im Dead-Letter-Topic, nicht in einer endlosen Wiederholungsschleife
Aufräumen:
for ID in $(curl -s "https://kafka.de-txl.ionos.com/clusters/$CLUSTER_ID/topics" \
-H "Authorization: Bearer $IONOS_TOKEN" | python3 -c \
"import sys,json;[print(t['id']) for t in json.load(sys.stdin)['items']]"); do
curl -s -X DELETE "https://kafka.de-txl.ionos.com/clusters/$CLUSTER_ID/topics/$ID" \
-H "Authorization: Bearer $IONOS_TOKEN"
done
Häufige Fehler
-
Verwendung von SASL oder Klartext statt mTLS
- Problem: Der Client hängt bei der Verbindung oder schlägt mit
SSL handshake failedfehl, obwohl die Anmeldeinformationen korrekt erscheinen. - Ursache: Die IONOS CLOUD Daten-Ebene akzeptiert ausschließlich gegenseitiges TLS. Es gibt keinen SASL/PLAIN- oder Klartext-Listener, und die Broker verwenden eine private CA, sodass ein Standard-Truststore diese ablehnt.
- Lösung: Setzen Sie
security.protocol=ssl, stellen Sie die Cluster-CA alsssl.ca.locationbereit, das Clientzertifikat und den Schlüssel, und deaktivieren Sie die Hostname-Verifizierung mitssl.endpoint.identification.algorithm=none.
- Problem: Der Client hängt bei der Verbindung oder schlägt mit
-
Auto-Commit verwirft fehlgeschlagene Nachrichten stillschweigend
- Problem: Einige Aufgabenbenachrichtigungen werden nie ausgelöst, aber es erscheinen keine Fehler in den Protokollen und die Konsumentenverzögerung ist null.
- Ursache: Mit
enable.auto.commit=Truewerden Offsets anhand eines Timers voranbewegt, unabhängig davon, obsend_notificationerfolgreich war, sodass eine Nachricht, die eine Ausnahme ausgelöst hat, als konsumiert markiert wird. - Lösung: Setzen Sie
enable.auto.commit=Falseund rufen Sieconsumer.commit(msg)erst nach erfolgreicher Verarbeitung auf und leiten Sie dauerhafte Fehler an das Dead-Letter-Thema weiter.
-
Partitionsanzahl kein Vielfaches der Brokeranzahl
- Problem: Ein Broker zeigt eine höhere CPU- und Festplattenauslastung als die anderen beiden, und die Durchsatzrate stagniert unter den Erwartungen.
- Ursache: Bei 3 Brokern verteilt ein Thema, das mit 4 oder 5 Partitionen erstellt wurde, die Last ungleichmäßig, sodass ein Broker einen zusätzlichen Partitionsleiter hostet.
- Lösung: Erstellen Sie Themen mit einer Partitionsanzahl, die ein Vielfaches von 3 ist (3, 6, 9, 12), damit sich die Partitionen gleichmäßig auf die drei Broker verteilen.
Zusammenfassung
Sie können IONOS CLOUD Event Streams for Apache Kafka nun als gleichberechtigten Bestandteil der TaskBoard-Architektur in Anwendungscode integrieren. Sie verbinden Produzenten und Konsumenten über gegenseitiges TLS mit Broker-Adressen und Zertifikaten aus Ihrer bereitgestellten Cluster-Instanz, erstellen Themen über die regionale Management-API mit Bearer-Authentifizierung und gestalten Partitionierungs- und Replikationseinstellungen, die die Topologie mit drei Brokern berücksichtigen. Darüber hinaus verfügen Sie über einen zuverlässigen Produzenten mit Idempotenz und acks=all, eine Konsumentengruppe mit manuellen at-least-once-Commits, ein Dead-Letter-Topic für fehlerhafte Nachrichten sowie JSON-Schema-Verträge, die es Produzent und Konsument ermöglichen, unabhängig voneinander weiterzuentwickelt zu werden.
Wichtige Punkte:
- Die Kafka-Datenebene erfordert gegenseitiges TLS mit einem Client-Zertifikat; die Management-API verwendet Bearer-Tokens. Es handelt sich um zwei unterschiedliche Authentifizierungsmodelle auf zwei verschiedenen Endpunkten.
- Die Anzahl der Topic-Partitionen sollte ein Vielfaches von 3 sein, um sie gleichmäßig auf die drei Broker zu verteilen; der empfohlene Replikationsfaktor ist 3, muss aber explizit gesetzt werden, da er kein systemseitig erzwungener Standardwert ist.
- At-least-once-Zustellung entsteht durch Deaktivierung des Auto-Commit und das Commiten von Offsets erst nach erfolgreicher Verarbeitung; machen Sie die Verarbeitung idempotent, um erneute Zustellungen zu tolerieren.
- Die Partitionanzahl ist die feste Obergrenze für die Konsumentenparallelität; dimensionieren Sie Partitionen für die Spitzenlast vor der Bereitstellung, da spätere Hinzufügungen die Schlüsselreihenfolge stören.
- Dead-Letter-Topics isolieren fehlerhafte Nachrichten, damit eine einzelne fehlerhafte Nutzlast ihre Partition nicht blockiert, und JSON-Schema erzwingt den Vertrag zwischen Produzent und Konsument.
Wichtige Begriffe:
- mTLS (mutual TLS): Client und Broker authentifizieren sich gegenseitig mit Zertifikaten; die einzige Authentifizierungsmethode der IONOS CLOUD-Datenebene, mit Client-Zertifikaten, die 365 Tage gültig sind.
- Konsumentengruppe: Eine Gruppe von Konsumenten, die sich die Partitionen eines Themas teilen; die Einheit für horizontale Skalierung des TaskBoard-Workers.
- Offset-Commit: Die Aufzeichnung der Position, bis zu der ein Konsument verarbeitet hat; das Commiten nach der Verarbeitung ergibt at-least-once-Zustellung.
- Dead-Letter-Topic (DLQ): Ein separates Topic für Nachrichten, die nach Wiederholungsversuchen fehlschlagen, um die Hauptpartition freizuhalten.
- Replikationsfaktor: Die Anzahl der Kopien jeder Partition über die Broker hinweg; 3 auf IONOS CLOUD, passend zum Cluster mit drei Brokern, sodass das Topic den Verlust eines Brokers übersteht.
Nächste Schritte
Weiterlernen: Einheit 4.4: AI Model Hub Integration
Verwandte Themen: