Aller au contenu
Riadh Mnasri
← Retour au blog
8 min de lecture

Kafka « au moins une fois » : rendre un consommateur idempotent

Kafka « au moins une fois » : rendre un consommateur idempotent

Dans l'article sur l'event-driven avec Kafka, une ligne du tableau résume un compromis en quelques mots : « au moins une fois (le consommateur doit gérer les doublons) ». Cette parenthèse cache la moitié du travail quand on construit un système sur Kafka. Cet article explique d'où viennent les doublons, pourquoi l'option « exactly-once » ne les supprime pas tous, et comment écrire un consommateur qui les absorbe sans effet de bord.

D'où viennent les doublons#

Un consommateur Kafka fait deux choses distinctes : il traite un message, puis il enregistre (commit) l'offset de ce message pour ne pas le relire. Ces deux actions ne sont pas atomiques. Tout se joue dans l'intervalle entre les deux.

 lire le message ──▶ traiter (écrire en base, appeler une API) ──▶ commit de l'offset
                                         ▲
                        crash, rebalance ou timeout ici :
                        le traitement est fait, l'offset n'est pas commité,
                        le message sera relu par le prochain consommateur

Si on inverse l'ordre (commit d'abord, traitement ensuite), on obtient l'inverse : un crash entre les deux fait perdre le message. C'est le mode « au plus une fois ». Entre perdre un message et le traiter deux fois, la plupart des systèmes choisissent le doublon, parce qu'un doublon se détecte et une perte ne se voit pas.

Les causes concrètes d'une relecture sont banales et fréquentes :

  • un redéploiement qui arrête le pod entre le traitement et le commit ;
  • un rebalance du groupe de consommateurs, par exemple parce qu'un traitement a dépassé max.poll.interval.ms (5 minutes par défaut) ;
  • un producteur qui renvoie un message après un timeout réseau alors que le broker l'avait bien reçu (côté producteur, ce cas est réglé par l'idempotence du producteur, voir plus bas).

Autrement dit, les doublons ne sont pas un cas limite. Sur un système qui tourne depuis quelques mois avec des déploiements réguliers, ils sont certains.

Ce que « exactly-once » garantit vraiment#

Kafka propose deux mécanismes qu'on présente souvent comme la solution.

Le producteur idempotent (enable.idempotence=true, activé par défaut depuis Kafka 3.0). Le broker attribue un identifiant au producteur et un numéro de séquence à chaque message. Si le producteur renvoie un message après un timeout, le broker reconnaît le doublon et l'ignore. Ça règle les doublons à l'écriture dans Kafka, pas à la lecture.

Les transactions (transactional.id côté producteur, isolation.level=read_committed côté consommateur). Elles permettent d'écrire des messages dans plusieurs topics et de commiter les offsets consommés dans une même transaction atomique. C'est sur ce mécanisme que Kafka Streams s'appuie pour offrir un vrai exactly-once.

La limite est dans la définition : ces garanties couvrent les échanges entre topics Kafka. Le schéma « lire un topic, transformer, écrire dans un autre topic » peut être exactly-once. Le schéma « lire un topic, écrire en base PostgreSQL, envoyer un email » ne l'est pas : ni la base ni le serveur mail ne participent à la transaction Kafka.

ScénarioExactly-once possible avec Kafka seul ?
Topic → transformation → topicOui, avec les transactions (ou Kafka Streams)
Topic → écriture en baseNon, la base n'est pas dans la transaction
Topic → appel d'une API externeNon
Topic → envoi d'email ou de SMSNon, et c'est le cas le plus visible pour l'utilisateur

Pour tous les cas de la partie basse du tableau, qui sont la majorité des consommateurs réels, l'idempotence est à la charge du consommateur.

Idempotent, ça veut dire quoi exactement#

Un traitement est idempotent quand l'appliquer deux fois produit le même état que l'appliquer une fois. solde = 100 est idempotent. `solde = solde

  • 100` ne l'est pas. Toute la technique consiste à transformer le second type d'opération en premier type, ou à détecter qu'on a déjà traité le message.

Il y a trois façons de faire, de la plus simple à la plus générale.

1. Rendre l'opération naturellement idempotente#

Parfois, il suffit de changer la forme de l'écriture. Un événement StatutCommandeChange(commandeId=42, statut=EXPEDIEE) peut être appliqué autant de fois qu'on veut : on écrit le statut, on ne l'incrémente pas. Un upsert sur une clé métier a la même propriété :

sql
INSERT INTO commande_statut (commande_id, statut, mis_a_jour_le)
VALUES (:commandeId, :statut, :horodatage)
ON CONFLICT (commande_id) DO UPDATE
  SET statut = EXCLUDED.statut, mis_a_jour_le = EXCLUDED.mis_a_jour_le
  WHERE commande_statut.mis_a_jour_le < EXCLUDED.mis_a_jour_le;

La clause WHERE sur l'horodatage protège aussi contre un message ancien relu après un plus récent. Quand c'est possible, c'est la meilleure solution : aucune table technique, aucune logique de déduplication.

2. Une table de messages traités, dans la même transaction#

Quand l'opération n'est pas idempotente par nature (créditer un compte, créer une facture), on enregistre l'identifiant de chaque message traité, dans la même transaction de base que l'effet métier. Si le message revient, l'insertion échoue sur la contrainte d'unicité et on sait qu'il est déjà traité.

sql
CREATE TABLE message_traite (
  message_id  VARCHAR(64) PRIMARY KEY,
  traite_le   TIMESTAMP NOT NULL
);
kotlin
@Component
class PaiementRecuConsumer(
    private val jdbc: JdbcTemplate,
    private val comptes: CompteRepository,
    private val tx: TransactionTemplate,
) {
 
    @KafkaListener(topics = ["paiements"], groupId = "comptabilite")
    fun onMessage(event: PaiementRecu) {
        tx.executeWithoutResult {
            val nouveau = jdbc.update(
                """
                INSERT INTO message_traite (message_id, traite_le)
                VALUES (?, now())
                ON CONFLICT (message_id) DO NOTHING
                """.trimIndent(),
                event.paiementId,
            ) == 1
 
            if (nouveau) {
                comptes.crediter(event.compteId, event.montant)
            }
        }
    }
}

Trois détails comptent ici.

L'identifiant vient du producteur, pas de Kafka. On pourrait utiliser le triplet topic, partition, offset, mais il change si le message est republié (rejeu depuis une autre source, migration de cluster). Un identifiant métier (paiementId) ou un UUID généré à la création de l'événement reste stable.

L'insertion et l'effet sont dans la même transaction. Si on insère l'identifiant, qu'on commite, puis qu'on crédite le compte dans une seconde transaction, un crash entre les deux marque le message comme traité alors que le crédit n'a jamais eu lieu. On a remplacé un doublon par une perte.

La table grossit. Il faut la purger, en gardant une fenêtre plus longue que la rétention du topic : un message ne peut pas revenir après avoir été supprimé de Kafka.

3. Les effets externes : clé d'idempotence ou outbox#

Pour un appel d'API ou un envoi d'email, la base locale ne suffit plus : on ne peut pas annuler un email envoyé si la transaction échoue ensuite.

Deux options, selon ce que l'API appelée permet.

L'API accepte une clé d'idempotence. Beaucoup d'API de paiement ou d'envoi le font (un en-tête Idempotency-Key). On transmet l'identifiant du message, et c'est le service appelé qui ignore le second appel. C'est la solution la plus simple, il faut juste vérifier la durée pendant laquelle le service conserve les clés.

L'API ne l'accepte pas. On sépare la décision de l'exécution : le consommateur écrit dans une table « à envoyer », dans la même transaction que le marquage du message (méthode 2). Un second processus lit cette table, envoie, et marque la ligne comme envoyée. Le doublon Kafka est absorbé par la table, et il ne reste qu'un risque de double envoi en cas de crash pendant l'envoi lui-même, qu'on réduit au minimum sans pouvoir l'éliminer complètement sans l'aide du service appelé.

Le problème inverse : publier de façon fiable#

Les doublons ont une cause symétrique côté producteur. Un service qui écrit en base puis publie un événement a le même problème d'atomicité :

kotlin
// Fragile : si la publication échoue, la commande existe
// mais personne n'est prévenu. Si la base échoue après
// la publication, on annonce une commande qui n'existe pas.
commandes.save(commande)
kafkaTemplate.send("commandes", CommandeCreee(commande.id))

Le pattern transactional outbox règle ça : on écrit l'événement dans une table outbox dans la même transaction que la commande, et un processus séparé (un job qui lit la table, ou Debezium qui lit le journal de la base) publie ensuite dans Kafka. La publication peut être rejouée, donc elle produit parfois des doublons, ce qui ramène au consommateur idempotent. Les deux patterns vont ensemble : l'outbox garantit qu'aucun événement n'est perdu, l'idempotence garantit qu'un événement en double ne fait pas de dégât.

Tester l'idempotence#

Un test suffit à prouver l'essentiel : envoyer le même message deux fois et vérifier l'état final.

kotlin
@Test
fun `un paiement reçu deux fois ne crédite le compte qu'une fois`() {
    val event = PaiementRecu(paiementId = "p-123", compteId = "c-1", montant = 50.euros)
 
    consumer.onMessage(event)
    consumer.onMessage(event)
 
    assertThat(comptes.solde("c-1")).isEqualTo(50.euros)
}

Ce test est bon marché et attrape la majorité des régressions. Le cas plus difficile, un crash au milieu du traitement, se teste en injectant une exception après l'effet métier et avant la fin de la transaction, puis en rejouant le message : l'état doit être le même que s'il avait été traité une seule fois.

Ce qu'il faut retenir#

« Au moins une fois » n'est pas un défaut de configuration qu'on corrige en activant une option. C'est le contrat normal d'un consommateur qui a des effets hors de Kafka. L'exactly-once de Kafka est réel, mais il s'arrête à la frontière du broker. Au-delà, trois outils suffisent dans la grande majorité des cas : une écriture idempotente par nature quand c'est possible, une table de messages traités dans la même transaction sinon, et une clé d'idempotence ou une outbox pour les effets externes. La question à poser pour chaque nouveau consommateur tient en une ligne : que se passe-t-il si ce message arrive deux fois ?