19 min de lecture

Objectifs d'apprentissage

À la fin de ce module, vous serez en mesure de:

  • Connecter les producteurs et consommateurs Kafka à un cluster IONOS CLOUD Event Streams for Apache Kafka en utilisant l'authentification mTLS et les adresses des brokers récupérées à partir de la sortie Terraform
  • Concevoir les sujets, partitions et groupes de consommateurs pour le flux d'événements TaskBoard, en appliquant les contraintes de partitionnement et de réplication d'IONOS CLOUD
  • Mettre en œuvre une livraison au moins une fois avec des engagements manuels des décalages et un sujet de lettre morte pour l'isolation des messages toxiques
  • Configurer la fiabilité du producteur avec l'idempotence, `acks=all` et le groupement, et ajuster le parallélisme du consommateur par rapport à la limite du nombre de partitions
  • Construire des contrats d'événements avec JSON Schema et les faire évoluer sans casser les consommateurs existants

Unité 4.3 : Intégration du streaming d'événements

Introduction

Vous reliez le service API de TaskBoard à son travailleur d'arrière-plan. Lorsqu'un utilisateur crée, déplace ou ferme une tâche, l'API ne doit pas bloquer l'envoi d'un e-mail ou d'un webhook. À la place, l'API produit un événement task-changed vers Kafka et renvoie immédiatement, tandis qu'un service travailleur distinct consomme ces événements et déclenche les notifications. Cela découple le chemin de requête des effets secondaires lents et vous permet de mettre à l'échelle le travailleur de manière indépendante.

IONOS CLOUD Event Streams for Apache Kafka exécute Apache Kafka 4.0.0, de sorte que vos bibliothèques de clients Kafka existantes fonctionnent sans modification. Ce qui diffère par rapport à un courtier auto-géré, c'est la connexion : le plan de données exige une authentification TLS mutuelle, et non SASL ou texte en clair, et les adresses des courtiers et les certificats proviennent du cluster que vous avez provisionné dans l'Unité 2.5. Cette unité commence par la connexion du client, puis aborde la conception des sujets, la livraison fiable, le traitement des lettres mortes et la mise à l'échelle des consommateurs, le tout dans le contexte du flux d'événements de TaskBoard.

1. Connexion des clients avec mTLS

Le plan de données Kafka d'IONOS CLOUD est sécurisé par TLS avec authentification mutuelle. Les deux parties s'authentifient mutuellement : votre client vérifie le certificat du broker auprès de l'autorité de certification du cluster, et le broker vérifie votre certificat client, qui est signé par l'autorité de certification propre au cluster. Il n'existe aucun écouteur en texte brut ou SASL/PLAIN. Chaque producteur et chaque consommateur présente un certificat client.

Vous récupérez les identifiants, l'autorité de certification, la clé privée du client et le certificat client depuis le point d'accès du cluster. Il s'agit des mêmes valeurs que le cluster expose via ses adresses de brokers et son API d'accès. Chaque certificat client est valide pendant 365 jours, il est donc recommandé d'intégrer la rotation des certificats à votre guide opérationnel avant la fin de l'année.

1.1 Configuration de la connexion du client

La configuration d'un client Kafka auto-géré utilise security.protocol=SSL avec un magasin de confiance (pour l'autorité de certification du cluster) et un magasin de clés (pour le certificat et la clé du client). La documentation IONOS CLOUD présente le fichier de propriétés client de référence :

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

Remarque : ssl.endpoint.identification.algorithm est laissé vide. Étant donné que le cluster utilise une CA privée plutôt que des certificats signés publiquement, la vérification du nom d'hôte par rapport aux adresses des brokers est désactivée ici. Le truststore valide la chaîne à la place.

1.2 Connexion du producteur en Python

La bibliothèque confluent-kafka accepte les mêmes clés SSL. Pointez le client vers les adresses des brokers issues de votre sortie Terraform et fournissez la CA, le certificat et la clé sous forme de fichiers PEM.

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)

Le cluster exécute trois brokers sur l'ensemble de la topologie du cluster, il est donc nécessaire d'énumérer les trois dans bootstrap.servers afin d'assurer la résilience de la connexion. Le client n'a besoin que d'un broker accessible pour l'initialisation, mais l'énumération des trois évite un point de défaillance unique lors de la démarrage.

2. Conception des sujets et des partitions

La conception des sujets est le point où les contraintes spécifiques à IONOS CLOUD se font d'abord sentir. Utilisez un sujet par type d'événement plutôt qu'un unique sujet de type « tuyau d'évacuation », car les consommateurs s'abonnent aux événements qui les intéressent et vous pouvez définir la rétention et le partitionnement pour chaque préoccupation. Pour TaskBoard, vous créez un sujet task-changed pour le travailleur, puis plus tard un sujet de lettre morte task-changed.dlq.

2.1 Création de sujets via l'API

L'API de gestion Kafka est régionale et distincte de cloudapi. L'hôte suit le modèle https://kafka.<region>.ionos.com, et contrairement au plan de données, l'API de gestion s'authentifie avec un jeton Bearer, et non avec mTLS. Créez un sujet avec 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 }
    }
  }'

Les paramètres du corps sont name (le seul champ obligatoire), replicationFactor, numberOfPartitions, et les paramètres logRetention. Relisez les sujets avec GET /clusters/{clusterId}/topics et GET /clusters/{clusterId}/topics/{topicId}, qui renvoient name, replicationFactor, numberOfPartitions, et logRetention. La suppression est DELETE /clusters/{clusterId}/topics/{topicId} et renvoie 202 Accepted, car le retrait est asynchrone.

2.2 Contraintes de nombre de partitions et de réplication

Choisissez le nombre de partitions comme un multiple de 3 (3, 6, 9, 12, etc.) pour éviter une répartition inégale des partitions entre les trois brokers. Un nombre de partitions qui n'est pas un multiple du nombre de brokers laisse un broker supporter plus de partitions que les autres, ce qui déséquilibre la charge.

Le facteur de réplication recommandé est 3, correspondant à la topologie à trois brokers, mais il s'agit d'une recommandation IONOS CLOUD, et non d'une valeur par défaut imposée par le système ; vous devez définir replicationFactor explicitement lors de la création d'un sujet. Avec un facteur de réplication de 3, le cluster conserve trois copies de chaque message, de sorte que le sujet survit à la perte d'un broker sans perte de données. La rétention par défaut des sujets est de 604800000 ms (7 jours). Définir retentionTime à -1 n'applique aucune limite de temps, auquel cas la rétention est limitée uniquement par le stockage du cluster que vous avez provisionné (par exemple 750 Go pour la taille S, 1200 Go pour M).

Le tableau suivant associe les paramètres du corps du sujet à leur rôle lors de la conception d'un sujet.

Paramètre Rôle Valeur TaskBoard
name Identifiant du sujet (obligatoire) task-changed
replicationFactor Copies par partition entre les brokers 3
numberOfPartitions Plafond de parallélisme pour les consommateurs 6 (multiple de 3)
logRetention.retentionTime Durée de vie des messages en ms (-1 = illimité) 604800000 (7 jours)

Les partitions comptent au-delà du stockage. Le nombre de partitions est la limite stricte du parallélisme des consommateurs au sein d'un groupe de consommateurs, il convient donc de le dimensionner pour le débit de pointe dès maintenant, car l'augmentation des partitions ultérieurement modifie la correspondance clé-partition et perturbe les garanties d'ordonnancement.

3. Fiabilité du producteur

Un producteur en mode « fire-and-forget » abandonne les messages en cas de dysfonctionnement du courtier. Pour les événements de tâches de TaskBoard, où la perte d'un événement signifie une notification manquée, configurez le producteur pour la durabilité avant d'optimiser le débit.

3.1 Accusés de réception, idempotence et clés

Définissez acks=all afin que le producteur attende que tous les répliques synchrones accusent réception d'une écriture avant de la considérer comme réussie. Combiné à un facteur de réplication de 3, cela signifie qu'un message est durable sur les courtiers avant que votre code ne passe à la suite. Activez le producteur idempotent afin que les nouvelles tentatives ne créent pas de doublés côté courtier.

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)

La clé de routage par task_id achemine tous les événements d'une même tâche vers la même partition, ce qui préserve l'ordre des événements par tâche. Sans clé, Kafka distribue les enregistrements en mode round-robin et un événement « tâche fermée » pourrait être traité avant « tâche créée ».

3.2 Regroupement et vidage

linger.ms et batch.size sacrifient quelques millisecondes de latence au profit d'un débit nettement supérieur en regroupant les enregistrements par partition. Il est impératif d'effectuer flush() avant la sortie du processus, sous peine de perdre les enregistrements en mémoire tampon.

# 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. Groupes de consommateurs, décalages et sémantiques de livraison

Le worker TaskBoard lit task-changed en tant que groupe de consommateurs. Chaque consommateur au sein du même groupe partage les partitions, de sorte que l'ajout de pods de worker augmente le débit jusqu'au nombre de partitions. Un quatrième consommateur sur un sujet à 3 partitions reste inactif.

4.1 Au moins une fois avec des confirmations manuelles

Pour la livraison des notifications, la sémantique « au moins une fois » est le choix par défaut approprié : traiter le message, puis confirmer le décalage. Si le worker échoue après le traitement mais avant la confirmation, le message est redéclaré et la notification peut se déclencher deux fois, ce que vous sécurisez à l'aide d'une clé d'idempotence. À l'inverse, confirmer avant le traitement (au plus une fois) présente un risque de perte d'événements en cas d'échec.

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 est le choix clé. L'auto-commit fait avancer les décalages sur un minuterie, indépendamment de la réussite du traitement, ce qui entraîne la perte silencieuse des messages dont le traitement a échoué.

4.2 Mise à l'échelle des consommateurs et rééquilibrage

Le nombre de partitions constitue la limite de parallélisme. Pour mettre à l'échelle le worker TaskBoard sur Managed Kubernetes (voir l'unité 3.2), augmentez le nombre de répliques du Deployment jusqu'au nombre de partitions, et chaque nouveau pod rejoint le groupe taskboard-notifier et se voit attribuer un sous-ensemble de partitions lors d'un rééquilibrage. Pendant un rééquilibrage, les partitions cessent brièvement d'être consommées, il est donc essentiel de maintenir le traitement idempotent et les confirmations fréquentes afin de minimiser le retraitement.

Partitions du sujet Nombre maximal de consommateurs utiles Effet d'un consommateur supplémentaire
3 3 4e consommateur inactif
6 6 mise à l'échelle linéaire jusqu'à 6
12 12 marge pour une croissance future

Dimensionnez les partitions en fonction de votre nombre maximal de consommateurs dès le départ. Vous pouvez ajouter des partitions ultérieurement, mais cela réorganise la correspondance entre les clés et les partitions et rompt l'ordre pour les clés en cours de traitement.

5. Gestion des messages en échec et discipline des schémas

Tous les messages ne peuvent pas être traités. Une charge utile mal formée ou une panne en aval continuera d'échouer lors de la nouvelle livraison, et un message toxique que vous réessayez indéfiniment bloque toute sa partition.

5.1 Motif de sujet de messages en échec

Routez les messages qui épuisent leurs tentatives vers un sujet de messages en échec dédié, task-changed.dlq, puis validez le décalage d'origine afin que la partition principale continue de fonctionner. Vous inspectez la file d'attente d'échec ultérieurement, corrigez la cause et, le cas échéant, rejouez les messages.

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

Distinguer les pannes transitoires (un dépassement de délai sur le service de notification, qui justifie une nouvelle tentative) des pannes permanentes (un événement mal formé, qui est envoyé directement à la file d'attente des messages en échec). Retenter aveuglément les pannes permanentes bloque la partition ; envoyer aveuglément les pannes transitoires à la file d'attente des messages en échec entraîne la perte d'événements récupérables.

5.2 Contrats d'événements avec JSON Schema

Les producteurs et les consommateurs sont déployés de manière indépendante, de sorte que la charge utile de l'événement constitue un contrat entre eux. La définir avec JSON Schema et la valider des deux côtés. Faire évoluer le schéma de manière additive : ajouter des champs facultatifs, ne jamais supprimer ni renommer les champs existants, afin que les anciens consommateurs continuent de fonctionner pendant que les nouveaux producteurs sont déployés.

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

Incluez une version de schéma dans l'événement ou dans un en-tête afin que les consommateurs puissent en tenir compte lors d'une fenêtre de migration. Cette rigueur permet de modifier la structure des événements de TaskBoard sans nécessiter un déploiement synchronisé et coordonné du producteur et du travailleur.

Fiche de référence rapide de l'API

Principales terminaisons de l'API de gestion Kafka. L'hôte est régional, par exemple https://kafka.de-txl.ionos.com :

Méthode Terminaison Description
POST /clusters Créer un cluster Kafka
GET /clusters/{clusterId} Récupérer les détails du cluster et les adresses des brokers
POST /clusters/{clusterId}/topics Créer un sujet
GET /clusters/{clusterId}/topics Lister tous les sujets
DELETE /clusters/{clusterId}/topics/{topicId} Supprimer un sujet (renvoie 202)
GET /clusters/{clusterId}/users/{userId}/access Récupérer les identifiants d'accès mTLS (CA, certificat, clé)

Authentification de l'API de gestion : Authorization: Bearer <token> Authentification du plan de données (producteurs/consommateurs) : TLS mutuel avec certificat client

Atelier de code

Objectif : Produire et consommer des événements TaskBoard task-changed sur un cluster Kafka IONOS CLOUD avec un groupe de consommateurs, des engagements manuels des décalages et un routage vers un journal des lettres mortes en cas d'échec.

Prérequis :

  • Compte IONOS CLOUD avec jeton API (IONOS_TOKEN)
  • Un cluster Kafka provisionné depuis l'Unité 2.5 avec les adresses des brokers disponibles
  • Python 3.10 ou supérieur avec confluent-kafka et jsonschema installés
  • Le certificat autorité du cluster, le certificat client et la clé client enregistrés sous forme de fichiers PEM

Étape 1 : Récupérer les adresses des brokers et les identifiants d'accès

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

Sortie attendue :

{
  "id": "e69b22a5-8fee-56b1-b6fb-4a07e4205ead",
  "properties": { "name": "...", "version": "4.0.0", "size": "S",
    "connections": [ { "brokerAddresses": ["192.168.1.101/24", ...] } ] }
}

Étape 2 : Créer les sujets principaux et de file d'attente des messages en échec

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

Sortie attendue :

{"id":"ae085c4c-...","type":"topic","properties":{"name":"task-changed",
 "replicationFactor":3,"numberOfPartitions":6}}

Étape 3 : Produire un événement de modification de tâche

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")

Sortie attendue :

delivered to task-changed [3] @ 0
produced

Étape 4 : Consommer avec un groupe de consommateurs et une confirmation manuelle

consumer.subscribe(["task-changed"])
msg = consumer.poll(5.0)
print(msg.topic(), msg.partition(), msg.offset(), msg.value())
consumer.commit(msg)

Sortie attendue :

task-changed 3 0 b'{"task_id":"task-42",...}'

Étape 5 : Déclencher une lettre morte sur une charge utile invalide

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

Sortie attendue :

delivered to task-changed.dlq [1] @ 0

Étape 6 : Vérifier que le sujet de messages en échec a reçu le message

dlq_consumer.subscribe(["task-changed.dlq"])
m = dlq_consumer.poll(5.0)
print(dict(m.headers()), m.value())

Sortie attendue :

{'x-error': b'...', 'x-attempts': b'3'} b'{not valid json'

Liste de contrôle de validation :

  • [ ] Sujets task-changed et task-changed.dlq créés avec un facteur de réplication de 3
  • [ ] Le producteur délivre avec acks=all et l'idempotence activée
  • [ ] Le consommateur engage les décalages manuellement après le traitement
  • [ ] Un message mal formé atterrit dans le sujet de messages morts, et non dans une boucle de réessai infinie

Nettoyage :

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

Pièges courants

  1. Utilisation de SASL ou de texte en clair au lieu de mTLS

    • Problème : Le client reste bloqué lors de la connexion ou échoue avec SSL handshake failed, bien que les identifiants semblent corrects.
    • Cause : Le plan de données IONOS CLOUD n'accepte que le TLS mutuel. Il n'existe aucun écouteur SASL/PLAIN ou texte en clair, et les brokers utilisent une CA privée, de sorte qu'un magasin de confiance par défaut les rejette.
    • Correction : Définir security.protocol=ssl, fournir la CA du cluster en tant que ssl.ca.location, le certificat et la clé du client, et désactiver la vérification du nom d'hôte avec ssl.endpoint.identification.algorithm=none.
  2. La validation automatique supprime silencieusement les messages en échec

    • Problème : Certaines notifications de tâches ne se déclenchent jamais, mais aucune erreur n'apparaît dans les journaux et le retard du consommateur est nul.
    • Cause : Avec enable.auto.commit=True, les décalages avancent selon un minuteur, indépendamment de la réussite de send_notification, de sorte qu'un message ayant levé une exception est marqué comme consommé.
    • Correction : Définir enable.auto.commit=False et appeler consumer.commit(msg) uniquement après un traitement réussi, en acheminant les échecs permanents vers le sujet de messages morts.
  3. Le nombre de partitions n'est pas un multiple du nombre de brokers

    • Problème : Un broker affiche une charge CPU et disque plus élevée que les deux autres, et le débit plafonne en dessous des attentes.
    • Cause : Avec 3 brokers, un sujet créé avec 4 ou 5 partitions est réparti de manière inégale, de sorte qu'un broker héberge un leader de partition supplémentaire.
    • Correction : Créer des sujets avec un nombre de partitions qui est un multiple de 3 (3, 6, 9, 12) afin que les partitions soient réparties uniformément sur les trois brokers.

Résumé

Vous pouvez désormais intégrer IONOS CLOUD Event Streams for Apache Kafka dans le code applicatif en tant que composant à part entière de l'architecture TaskBoard. Vous connectez les producteurs et les consommateurs via TLS mutuel en utilisant les adresses des brokers et les certificats de votre cluster provisionné, vous créez des sujets via l'API de gestion régionale avec une authentification Bearer, et vous concevez les paramètres de partitionnement et de réplication qui respectent la topologie à trois brokers. En plus de cela, vous disposez d'un producteur fiable avec idempotence et acks=all, d'un groupe de consommateurs avec des engagements manuels au moins une fois, d'un sujet de messages morts pour les messages toxiques, et de contrats JSON Schema qui permettent au producteur et au consommateur d'évoluer indépendamment.

Points clés :

  • Le plan de données Kafka nécessite un TLS mutuel avec un certificat client ; l'API de gestion utilise des jetons Bearer. Il s'agit de deux modèles d'authentification distincts sur deux points d'accès différents.
  • Le nombre de partitions d'un sujet doit être un multiple de 3 pour une répartition uniforme sur les trois brokers ; le facteur de réplication recommandé est 3, mais vous devez le définir explicitement, car il ne s'agit pas d'une valeur par défaut imposée par le système.
  • La livraison au moins une fois est obtenue en désactivant l'engagement automatique et en engageant les décalages uniquement après un traitement réussi ; assurez-vous que le traitement est idempotent pour tolérer la redélivrance.
  • Le nombre de partitions constitue la limite supérieure stricte de la parallélisation des consommateurs ; dimensionnez les partitions pour la charge maximale avant le déploiement, car leur ajout ultérieur perturbe l'ordre des clés.
  • Les sujets de messages morts isolent les messages toxiques afin qu'une charge utile défectueuse ne bloque pas sa partition, et JSON Schema impose le contrat entre le producteur et le consommateur.

Terminologie importante :

  • mTLS (TLS mutuel) : Le client et le broker s'authentifient mutuellement à l'aide de certificats ; c'est la seule méthode d'authentification du plan de données IONOS CLOUD, avec des certificats clients valides pendant 365 jours.
  • Groupe de consommateurs : Un ensemble de consommateurs partageant le travail des partitions d'un sujet, l'unité de mise à l'échelle horizontale pour le travailleur TaskBoard.
  • Engagement de décalage : Enregistrement de la position jusqu'à laquelle un consommateur a traité les messages ; l'engagement après le traitement garantit une livraison au moins une fois.
  • Sujet de messages morts (DLQ) : Un sujet distinct pour les messages qui échouent au traitement après des tentatives de reprise, maintenant la partition principale débloquée.
  • Facteur de réplication : Le nombre de copies de chaque partition sur les brokers ; 3 sur IONOS CLOUD, correspondant au cluster à trois brokers, afin que le sujet survive à la perte d'un broker.

Prochaines étapes

Poursuivre l'apprentissage : Unité 4.4 : Intégration AI Model Hub

Sujets connexes :