Event Streaming

Apache Kafka

Kafka est un journal de commits distribué en ajout seul (append-only), et non une file d'attente. Les enregistrements sont écrits une fois et lus plusieurs fois, ce qui lui confère des garanties de relecture, de fan-out et d'ordonnancement qu'un broker classique ne peut offrir.

advanced15 min readUpdated 16 sept. 2026
producer.ts
ts
// producer.ts
import { Kafka } from "kafkajs";

const kafka = new Kafka({
  clientId: "orders",
  brokers: ["localhost:9092"],
});

const producer = kafka.producer();

await producer.connect();
await producer.send({
  topic: "orders.created",
  messages: [
    // The key decides the partition, so all events for one
    // customer stay in order.
    { key: order.customerId, value: JSON.stringify(order) },
  ],
});
await producer.disconnect();
Sortie
2011
Origine
LinkedIn
Modèle
Append-only commit log
Ordonnancement
Par partition
Livraison par défaut
At-least-once
Rétention
Basée sur le temps ou la taille

Pourquoi c'est important

Pourquoi les équipes choisissent Kafka

Un log durable et rejouable

Les enregistrements sont ajoutés et ne sont jamais supprimés lors de la lecture. N'importe quel consumer peut revenir en arrière et retraiter l'historique, ce qu'une file d'attente classique ne permet tout simplement pas.

Passage à l'échelle et ordre via les partitions

Un topic est divisé en partitions réparties sur plusieurs brokers. Le parallélisme provient des partitions, et l'ordre est garanti au sein de chacune d'elles.

Plusieurs lecteurs indépendants

Les consumer groups suivent chacun leurs propres offsets ; ainsi, un job de facturation, un indexeur de recherche et un pipeline d'analyse peuvent lire les mêmes événements sans interférer.

Le tableau complet

Les trois concepts fondamentaux de Kafka

Un topic est un log en ajout seul, les producers y écrivent des enregistrements avec une clé, et les consumers suivent leur propre position via des offsets.

Topic

Ajout

Un log nommé et partitionné. Les producers ajoutent des enregistrements à la fin et chaque enregistrement conserve son offset jusqu'à l'expiration de la rétention.

Consumer group

Partage

Un groupe de consumers se répartit les partitions et suit ses propres offsets commités, permettant ainsi une mise à l'échelle et une progression indépendantes.

Offset

Suivi

La position d'un groupe dans une partition. Le commiter est l'accusé de réception qui détermine si l'on est en at-least-once ou at-most-once.

HTML5 en un coup d'oeil

Les briques de base de Kafka

Topics

Logs nommés en ajout seul qui conservent les enregistrements jusqu'à l'expiration de la rétention.

Partitions

L'unité de parallélisme et la limite de l'ordonnancement.

Producers

Sérialisent un enregistrement, choisissent une clé et l'ajoutent à une partition.

Consumer groups

Répartissent les partitions entre les membres et commitent les offsets.

Offsets

La position d'un groupe dans une partition et son point d'acquittement.

Rétention

Suppression par temps ou taille, ou compaction pour ne garder que la dernière valeur par clé.

Flux

Le voyage d'un enregistrement

Un enregistrement est écrit une fois et lu plusieurs fois. Rien n'est supprimé lorsqu'un consumer le lit, ce qui rend la relecture possible.

  1. 1

    Produire

    Le producer sérialise un enregistrement et choisit une partition en hachant sa clé, afin que la même clé arrive toujours au même endroit.

  2. 2

    Ajouter

    Le broker leader de la partition ajoute l'enregistrement à la fin de son log et renvoie le nouvel offset au producer.

  3. 3

    Consommer

    Chaque consumer d'un groupe possède un sous-ensemble de partitions et lit ses enregistrements dans l'ordre des offsets.

  4. 4

    Committer l'offset

    Après traitement, le consumer commit sa position pour que le groupe puisse reprendre à partir de là après un redémarrage.

  5. 5

    Conserver ou rejouer

    L'enregistrement reste dans le log jusqu'à l'expiration de la rétention, permettant à un autre groupe ou un job ultérieur de revenir en arrière et de le relire.

Un bref aperçu

Des logs de LinkedIn au standard du streaming

  1. 2011

    Kafka devient open-source

    LinkedIn publie Kafka en tant que log de commits distribué pour ses flux d'activité.

    11
  2. 2012

    Incubation Apache

    Kafka devient un projet Apache et son adoption s'étend bien au-delà de LinkedIn.

    12
  3. 2016

    Kafka Streams

    La version 0.10 introduit une bibliothèque de traitement de flux, suivie plus tard par ksqlDB.

    16
  4. 2017

    Sémantique Exactly-once

    La version 0.11 ajoute des producers idempotents et des transactions pour des pipelines exactly-once.

    17
  5. 2021

    Début de KRaft

    Le KIP-500 commence à remplacer ZooKeeper par un quorum de métadonnées intégré aux brokers.

    21
  6. 2025

    Kafka 4.0 abandonne ZooKeeper

    Les nouveaux clusters fonctionnent uniquement en mode KRaft, marquant la fin de l'ère ZooKeeper.

    25

Le guide complet

Apache Kafka: Tout ce que vous devez savoir

Ce qu’est réellement Kafka

Apache Kafka est un journal de commit distribué, en ajout seul (append-only). Cette seule phrase explique presque tout le reste. Les enregistrements sont ajoutés à la fin d’un journal, chacun reçoit un numéro croissant appelé offset, et rien n’est jamais modifié sur place. Les lecteurs ne suppriment pas les enregistrements lorsqu’ils les consomment ; ils déplacent simplement un curseur vers l’avant.

C’est là que réside la différence fondamentale avec une file d’attente de messages classique. Dans une queue, un consommateur récupère un message et celui-ci disparaît. Dans Kafka, un consommateur lit un enregistrement, mémorise son offset, et l’enregistrement reste là où il se trouve tant que la politique de rétention du topic le permet. Dix consommateurs différents — et dix applications différentes — peuvent lire le même enregistrement indépendamment, à leur propre rythme, sans avoir à se coordonner entre eux.

Kafka a été conçu chez LinkedIn pour gérer des flux d’activité : vues de pages, clics, lignes de logs, le tout à des millions d’événements par seconde. Il a été open-sourcé en 2011 et est devenu un projet Apache en 2012. Aujourd’hui, c’est l’épine dorsale par défaut pour le streaming d’événements : change data capture, pipelines de métriques, event sourcing, agrégation de logs et stream processing.

Si vous venez de RabbitMQ ou Redis, le changement de paradigme consiste à ne plus penser en termes de « message à livrer », mais en termes de « fait qui a été enregistré ».

Topics, partitions et ordonnancement

Un topic est un journal (log) durable et nommé. Les producteurs y écrivent et les consommateurs y lisent. Les topics sont divisés en partitions, et la partition est l’unité de base tant pour le parallélisme que pour l’ordonnancement.

  • Les enregistrements au sein d’une même partition sont strictement ordonnés par offset.
  • Il n’y a aucune garantie d’ordre entre les partitions.
  • Plus il y a de partitions, plus on peut avoir de consommateurs en parallèle, mais cela implique davantage de fichiers, plus de réplication et des rebalances plus lents.

C’est la garantie la plus importante de Kafka : l’ordre est géré par partition. Si deux événements doivent être traités dans l’ordre, ils doivent atterrir dans la même partition. Dans le cas contraire, Kafka pourrait les transmettre à des consommateurs différents s’exécutant simultanément, et l’ordre serait alors perdu.

Les partitions déterminent également le parallélisme maximal des consommateurs au sein d’un groupe : un groupe peut avoir au maximum un seul consommateur actif par partition. Six partitions signifient donc au maximum six consommateurs utiles ; un septième resterait inactif.

# three partitions, replicated across three brokers
kafka-topics.sh --create \
  --topic orders.created \
  --partitions 6 \
  --replication-factor 3 \
  --bootstrap-server localhost:9092

Il est possible d’ajouter des partitions ultérieurement, mais cela modifie le mappage entre la clé et la partition pour les clés existantes. Par conséquent, les événements d’une même entité peuvent se retrouver répartis sur deux partitions et perdre leur ordre relatif. Déterminez le nombre de partitions lors de la création du topic, en prévoyant une marge pour la croissance.

Producteurs, clés et partitionnement

Un producteur sérialise un enregistrement et détermine à quelle partition il appartient. Le partitionneur par défaut hache la clé de l’enregistrement et l’associe à une partition. Une même clé est toujours envoyée vers la même partition, ce qui permet de préserver l’ordre pour une entité donnée tout en répartissant les différentes entités sur plusieurs partitions.

await producer.send({
  topic: "orders.created",
  messages: [
    { key: order.customerId, value: JSON.stringify(order) },
  ],
});

Le choix de la clé est une décision de conception, et non un simple détail. Utilisez customerId comme clé et tous les événements d’un client seront ordonnés ensemble. Utilisez orderId et vous obtiendrez une répartition maximale, mais sans ordre entre les événements. Les enregistrements avec une clé nulle sont distribués pour l’équilibrage — le partitionneur “sticky” de Kafka remplit une partition avant de passer à la suivante — mais ils ne bénéficient d’aucune garantie d’ordre.

Le producteur contrôle également la durabilité et le débit via le batching et les acquittements. acks détermine combien de répliques doivent confirmer une écriture, linger.ms et batch.size définissent le temps d’attente pour remplir un batch, et compression.type détermine la quantité de CPU à sacrifier pour optimiser le réseau et le disque. Plus de détails sur ces points ci-dessous.

Brokers, réplication et le leader

Un broker est un serveur Kafka. Un cluster est composé de plusieurs brokers travaillant ensemble. Chaque partition possède un broker leader et zéro ou plusieurs followers. Les producers et consumers communiquent avec le leader ; les followers répliquent le log.

La réplication est ce qui rend Kafka durable. Si un leader tombe, l’un des replicas synchronisés (ISR) est promu et le cluster continue de fonctionner. Le paramètre acks détermine le niveau d’attente du producer :

  • acks=0 — “fire and forget” ; le plus rapide et le moins sûr.
  • acks=1 — le leader a écrit la donnée ; perdue si le leader tombe avant la réplication.
  • acks=all — tous les replicas synchronisés ont acquitté la réception ; le plus sûr.

Associez acks=all avec min.insync.replicas=2 et un facteur de réplication de trois pour une configuration de production : une écriture n’est acquittée que lorsque au moins deux replicas la possèdent, ainsi la perte d’un broker n’entraîne aucune perte de données.

const producer = kafka.producer({
  idempotent: true,
  maxInFlightRequests: 1,
  transactionalId: "orders-producer",
});

Le producer idempotent ajoute un numéro de séquence à chaque batch afin que le broker puisse rejeter les doublons, ce qui élimine les doublons accidentels que les tentatives de réessai (retries) pourraient autrement introduire. Un cluster possède également un controller qui gère le leadership des partitions et les métadonnées. Les anciennes versions de Kafka utilisaient ZooKeeper pour cela ; Kafka moderne utilise KRaft, où les brokers forment leur propre quorum de métadonnées et ZooKeeper a totalement disparu.

Groupes de consommateurs et affectation des partitions

Les consommateurs appartiennent à un groupe de consommateurs, identifié par groupId. Kafka affecte chaque partition à un seul consommateur au sein du groupe. Cela vous apporte deux avantages simultanément : le passage à l’échelle horizontal, car les partitions sont partagées, et l’équilibrage de charge, car aucun consommateur d’un même groupe ne traite la même partition.

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);
  },
});

Lorsque des consommateurs rejoignent ou quittent le groupe, Kafka effectue un rebalance : il révoque les affectations et en distribue de nouvelles. Les rebalances interrompent la consommation, ils sont donc coûteux ; les protocoles de rebalancing coopératif réduisent cette perturbation en ne déplaçant que les partitions qui doivent réellement être déplacées. Un consommateur qui cesse d’envoyer des heartbeats pendant session.timeout.ms est considéré comme mort et ses partitions sont réaffectées.

Le groupe est également le moyen utilisé par Kafka pour mémoriser la progression. Chaque groupe possède ses propres offsets commités, stockés dans le topic interne __consumer_offsets, ainsi deux groupes lisant le même topic peuvent se trouver à des positions complètement différentes sans aucune coordination.

Rééquilibrage et boucle de poll

Sous le capot, un consommateur est une boucle de poll. Il récupère des enregistrements, les transmet à votre handler, puis recommence l’opération. Le broker suit l’état de présence séparément via des heartbeats ; ainsi, un handler lent n’est pas immédiatement considéré comme inactif. Cependant, il existe une limite stricte : si un seul poll prend plus de max.poll.interval.ms, le broker considère que le consommateur est bloqué et déclenche un rééquilibrage.

const consumer = kafka.consumer({
  groupId: "billing",
  sessionTimeout: 45_000,
  heartbeatInterval: 3_000,
  rebalanceTimeout: 60_000,
});

Un rééquilibrage est déclenché dès qu’un consommateur rejoint le groupe, le quitte ou en est expulsé, lorsque les partitions ou les topics changent, ou encore lorsqu’une souscription est modifiée. Pendant le rééquilibrage, Kafka révoque les assignations et met la consommation en pause ; des rééquilibrages fréquents nuisent donc gravement au débit. Trois bonnes pratiques permettent de les limiter :

  • Bornez le travail par enregistrement. Un handler qui s’exécute parfois pendant plusieurs minutes finira par dépasser max.poll.interval.ms.
  • Utilisez le rééquilibrage coopératif. Le protocole coopératif incrémental ne déplace que les partitions qui doivent impérativement bouger, au lieu d’arrêter l’ensemble du groupe.
  • Utilisez l’appartenance statique (static membership). Le fait de définir un group.instance.id stable permet à un consommateur redémarré de récupérer ses anciennes partitions sans déclencher de rééquilibrage.

Kafka vous permet également d’utiliser consumer.pause() et consumer.resume() pour appliquer une backpressure lorsqu’une dépendance en aval est en difficulté. Mettre en pause est préférable à l’idée de bloquer la boucle de poll : le consommateur continue d’envoyer des heartbeats, reste dans le groupe et s’arrête simplement de récupérer des données jusqu’à ce que vous soyez prêt.

Offsets et garanties de livraison

Un offset représente la position d’un groupe de consommateurs dans une partition. Lorsqu’un consommateur traite un enregistrement, il peut committer l’offset, ce qui indique à Kafka : « ce groupe a tout terminé jusqu’ici ». Committer trop tôt et un crash entraîne une perte de données ; committer trop tard et un crash entraîne un retraitement des données.

  • At-most-once (au plus une fois) — commit avant le traitement. En cas de crash, l’enregistrement est perdu.
  • At-least-once (au moins une fois) — traitement, puis commit. En cas de crash, l’enregistrement est retraité. C’est le comportement par défaut le plus courant.
  • Exactly-once (exactement une fois) — les transactions et les producteurs idempotents rendent le cycle lecture-traitement-écriture atomique au sein de Kafka.

L’auto-commit s’effectue en arrière-plan via un minuteur. C’est pratique mais dangereux, car cela peut committer des offsets pour des enregistrements encore en cours de traitement. Le commit manuel après la fin du travail est le choix par défaut le plus fiable.

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) },
    ]);
  },
});

Comme le mode at-least-once est la norme, les handlers doivent être idempotents. Le guide sur le Batch Processing couvre le même contrat pour les files d’attente de tâches : partez du principe que le travail peut être exécuté deux fois et faites en sorte que la seconde exécution n’ait aucun effet (no-op). Une clé de déduplication dérivée du travail — et non de l’offset Kafka — rend le handler sûr, même si le même événement logique est produit deux fois.

Le mode exactly-once existe réellement, mais son champ d’application est plus restreint qu’il n’y paraît. Il fonctionne pour les pipelines Kafka-vers-Kafka avec des producteurs transactionnels et isolation.level=read_committed, mais dès que vous écrivez dans une base de données externe, vous devez à nouveau mettre en œuvre des écritures idempotentes.

Rétention, replay et compaction de logs

Kafka conserve les enregistrements selon une politique de rétention, et non selon un accusé de réception de livraison. retention.ms (sept jours par défaut) et retention.bytes limitent chaque partition ; dès que l’une de ces limites est dépassée, les anciens segments de logs sont supprimés. Comme rien n’est supprimé lors de la lecture, un consommateur peut effectuer un replay de l’historique en se repositionnant sur un offset antérieur ou en réinitialisant le groupe.

# 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

Le replay est le super-pouvoir de Kafka. Un nouveau service peut être initialisé en lisant l’intégralité de l’historique d’un topic. Un correctif de bug peut être déployé et la dernière journée de données retraitée. Un pipeline d’analyse peut être reconstruit à partir du log d’événements bruts. Une file d’attente classique ne peut rien de tout cela, car les données ont déjà disparu.

Pour l’état avec clé (keyed state), l’autre mode de rétention est la compaction de logs. Avec cleanup.policy=compact, Kafka conserve au moins la dernière valeur pour chaque clé et rejette les valeurs plus anciennes. Une valeur null est une tombstone qui supprime la clé. Le log devient alors un changelog capable de reconstruire une table — c’est précisément ce que Kafka Streams utilise pour ses state stores et ce sur quoi s’appuient les pipelines de change-data-capture.

# 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

Schémas et Schema Registry

Comme le log survit à n’importe quelle application individuelle, la structure d’un enregistrement devient un contrat entre les équipes. Un producteur et un consommateur sont couplés par les octets qu’ils échangent, et ce couplage persiste malgré les déploiements. Un schema registry transforme cela en un problème de compatibilité géré plutôt qu’en une mauvaise surprise.

Le registry stocke des schémas versionnés — généralement Avro, Protobuf ou JSON Schema — et assigne un identifiant numérique à chacun. Les producteurs enregistrent un schéma et écrivent l’id aux côtés du payload ; les consommateurs récupèrent le schéma via son id et procèdent à la désérialisation. Lorsqu’un schéma change, le registry impose un mode de compatibilité (comme backward ou forward), rejetant tout changement qui casserait les lecteurs existants.

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 }],
});

Même si vous n’utilisez pas de registry, traitez vos payloads comme une API : ajoutez des champs plutôt que de les renommer, donnez une valeur par défaut aux nouveaux champs, et versionnez lorsque la sémantique change. Un consommateur utilisant la version précédente doit être capable de lire un enregistrement écrit par la version suivante.

Kafka Connect et Kafka Streams

Deux composants de la plateforme vous évitent d’écrire deux fois la même plomberie technique.

Kafka Connect est un framework permettant de déplacer des données vers et depuis Kafka via de la configuration plutôt que du code. Les connecteurs source extraient des données depuis des bases de données, du stockage objet ou des API SaaS ; les connecteurs sink poussent les données vers des warehouses, des index de recherche ou une autre base de données. Debezium, par exemple, transforme le write-ahead log de Postgres en un flux d’événements de modification. Connect fonctionne sous forme de cluster, suit les offsets et gère les tentatives en cas d’échec, ce qui en fait la solution standard pour « injecter des données dans Kafka » et « extraire des données de Kafka ».

Kafka Streams est une bibliothèque client pour le traitement des données au sein de Kafka. Elle vous offre un DSL de streaming — map, filter, groupByKey, join, window — basé sur les abstractions KStream et KTable, avec des state stores s’appuyant sur des topics compactés. Elle s’exécute à l’intérieur de votre application, s’adapte en ajoutant des instances et gère la tolérance aux pannes via des topics de changelog. Si vous préférez le SQL, ksqlDB propose une couche de requête basée sur les mêmes concepts.

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

Il est utile de connaître ces deux outils même si vous commencez par de simples producers et consumers, car ils définissent ce qu’est « la méthode Kafka » à grande échelle.

Un exemple minimal de bout en bout

Il est utile de voir l’ensemble dans un seul fichier : connecter un producteur et un consommateur, s’abonner, exécuter et s’arrêter proprement.

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);
});

Trois éléments de ce fichier sont cruciaux en production. Le producteur est connecté une seule fois et réutilisé, et non créé pour chaque message. Le consommateur s’abonne avant de s’exécuter. Enfin, le processus gère SIGTERM, car un arrêt brutal pendant un rééquilibrage ou un commit laisse le groupe dans un état bien pire qu’une fermeture propre.

L’Event Sourcing et le pattern Outbox

Kafka est souvent décrit comme la colonne vertébrale de l’event sourcing, où le log est la source de vérité et l’état actuel est une projection d’événements. Au lieu de stocker uniquement la dernière ligne, vous stockez la séquence des faits — order.created, order.paid, order.shipped — et reconstruisez n’importe quelle vue en les rejouant. Les topics compactés transforment la projection en table ; la rétention la rend auditable.

La partie difficile n’est jamais l’écriture de l’événement ; c’est l’écriture de l’événement et de la ligne en base de données de manière atomique. Un processus peut planter après avoir validé la ligne et avant la publication, ou après la publication et avant la validation. Publier à l’intérieur d’une transaction de base de données est impossible, donc la réponse standard est le outbox pattern : écrivez l’événement dans une table outbox au sein de la même transaction que le changement d’état, puis laissez un relais ou un connecteur CDC publier ces lignes vers Kafka.

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 ou un petit relais de polling suit ensuite l’outbox et produit vers le topic, en supprimant les lignes une fois qu’elles sont publiées. Comme l’écriture dans l’outbox partage la transaction, l’événement existe exactement au moment où l’état existe. Les consommateurs doivent tout de même être idempotents, car le relais peut publier une ligne deux fois après un crash.

Développement et tests locaux

Vous n’avez pas besoin d’un cluster Kafka complet pour développer. Un broker à nœud unique via Docker Compose, ou Redpanda en mode compatibilité, vous permet de disposer de topics, de consumer groups et d’offsets directement sur votre ordinateur.

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"

Pour les tests, le schéma est le même que pour toute autre dépendance d’intégration : lancez le broker dans un conteneur, créez les topics nécessaires au test, et utilisez un group id unique par exécution de test afin que les offsets commités ne fuitent jamais d’une exécution à l’autre. Réinitialisez explicitement les offsets lorsqu’un test doit lire l’historique, et effectuez vos assertions sur les records reçus par votre consumer plutôt que sur le timing.

const groupId = `test-${crypto.randomUUID()}`;
const consumer = kafka.consumer({ groupId });

Gardez vos consumer handlers purs et légers — parsez, validez, appelez un service — afin que la majeure partie de la logique puisse être testée via des tests unitaires sans aucun broker. Réservez les tests d’intégration pour le câblage : est-ce qu’un record produit atteint le bon handler, et l’offset est-il commité ensuite ?

Débit, latence et batching

Kafka est rapide car il utilise le batching et écrit les données de manière séquentielle. Ces deux aspects sont configurables, et le réglage consiste à trouver l’équilibre entre latence et débit.

  • linger.ms — le temps pendant lequel un producer attend pour accumuler un batch. Une valeur plus élevée signifie des batchs plus volumineux et un débit accru, au prix d’une latence supplémentaire.
  • batch.size — le nombre maximum d’octets par batch et par partition.
  • compression.typesnappy, lz4 ou zstd. La compression réduit l’utilisation du réseau et du disque ; zstd l’emporte généralement sur le ratio de compression.
  • fetch.min.bytes et fetch.max.wait.ms — le temps pendant lequel les consumers attendent pour remplir une réponse de fetch.
  • max.poll.records — le nombre d’enregistrements qu’un consumer traite par poll. Augmentez cette valeur pour résorber les backlogs plus rapidement, mais veillez à ce que le traitement reste inférieur à max.poll.interval.ms, sinon le consumer est expulsé du groupe.
const producer = kafka.producer({
  linger: { ms: 20 },
  compression: CompressionTypes.GZIP,
});

const consumer = kafka.consumer({
  groupId: "billing",
  maxBytesPerPartition: 1_048_576,
  maxWaitTimeInMs: 500,
});

Un producer optimisé pour le débit pourrait utiliser linger.ms=20, un batch d’un mégaoctet et zstd. Un producer qui doit publier en quelques millisecondes utilisera linger.ms=0. Il n’y a pas de réponse unique ; il n’y a que le compromis que vous choisissez et le p99 que vous pouvez accepter.

Choisir entre Kafka, RabbitMQ et Redis

Ces trois outils sont souvent comparés comme s’ils étaient interchangeables. Ce n’est pas le cas.

  • Kafka est un log rejouable. Choisissez-le lorsque vous avez besoin d’un débit élevé, d’un historique durable, d’un fan-out vers de nombreux consommateurs indépendants, d’event sourcing, de stream processing ou de la possibilité de retraiter des données. C’est l’outil le plus lourd à exploiter.
  • RabbitMQ est un message broker. Choisissez-le lorsque le routage est primordial : exchanges, routing keys, accusés de réception par message, priorités et dead-lettering. Il est parfaitement adapté à la distribution de tâches et au routage complexe, et il est plus léger que Kafka pour des volumes modestes. Consultez le guide RabbitMQ.
  • Redis est un store de structures de données en mémoire qui fait également office de file d’attente de jobs. Choisissez-le lorsque le travail est simple, le volume modéré et que vous utilisez déjà Redis. BullMQ sur Redis vous offre des tentatives de rejeu (retries), de la planification et une interface utilisateur avec presque aucun coût opérationnel. Consultez Redis Queues.

Une règle simple : si les consommateurs doivent rejouer des messages, diffuser vers de nombreux groupes ou lire l’historique, utilisez Kafka. Si un message doit être routé vers un worker spécifique puis oublié, utilisez RabbitMQ. S’il s’agit d’un job d’arrière-plan avec une politique de retry, utilisez Redis.

Sécurité et contrôle d’accès

Kafka contient généralement les données les plus précieuses de votre infrastructure, il est donc primordial de le sécuriser.

  • Chiffrement — activez TLS pour le trafic entre les brokers ainsi que pour le trafic client-broker. Un Kafka en texte clair à l’intérieur d’un VPC reste du texte clair.
  • Authentification — utilisez SASL/SCRAM ou mTLS pour les clients. Évitez les listeners non authentifiés, sauf pour une installation locale jetable.
  • Autorisation — utilisez des ACL pour accorder les droits de lecture, d’écriture ou de création sur des topics et des groupes spécifiques. Un service ne doit pas pouvoir lire tous les topics simplement parce qu’il peut se connecter.
  • Quotas — limitez la bande passante des producers et des consumers par client afin qu’un service incontrôlé ne puisse pas saturer le cluster.
  • Secrets — n’intégrez jamais d’identifiants en dur dans un producer ; injectez-les via les variables d’environnement ou un gestionnaire de secrets.

Monitoring et consumer lag

Kafka peut échouer silencieusement si vous n’y prenez pas garde. Quatre indicateurs couvrent l’essentiel de la réalité opérationnelle.

  • Consumer lag — l’écart entre le dernier offset et l’offset commit d’un groupe, par partition. Une augmentation du lag est le premier signe que les consommateurs ne parviennent plus à suivre la cadence.
  • Under-replicated partitions — un nombre non nul signifie qu’un broker est hors service ou lent, et que la durabilité des données est menacée.
  • Latence réseau et des requêtes — les percentiles côté broker révèlent quand les disques ou le réseau deviennent le goulot d’étranglement.
  • Utilisation du disque et nombre de segments — Kafka est limité par les performances disque, et la réplication multiplie chaque octet par le facteur de réplication.

Affichez le lag sur le même dashboard que vos services et configurez vos alertes sur la tendance, et non sur un pic isolé. Un lag qui croît régulièrement tout au long de la journée est un problème de capacité ; un lag qui pique puis redescend est généralement dû à un rebalance ou à un déploiement lent.

Bonnes pratiques

  • Concevez vos clés en fonction de l’ordre dont vous avez réellement besoin ; n’oubliez pas que l’ordre est garanti par partition.
  • Utilisez acks=all, min.insync.replicas=2 et un facteur de réplication de trois pour les topics importants.
  • Activez le producteur idempotent et privilégiez les commits d’offset manuels.
  • Rendez vos consommateurs idempotents, car le contrat par défaut est le “at-least-once” (au moins une fois).
  • Maintenez le traitement de max.poll.records en dessous de max.poll.interval.ms pour éviter les tempêtes de rebalancement (rebalance storms).
  • Utilisez la compaction de log pour les états avec clés et les tombstones pour les suppressions.
  • Suivez le retard (lag) des consommateurs par partition, et pas seulement le total du cluster.
  • Traitez les schémas d’enregistrements comme un contrat versionné, avec un registre lorsque plusieurs équipes sont impliquées.
  • Séparez les topics par cycle de vie et rétention plutôt que de tout regrouper dans un seul.
  • Définissez la rétention en fonction d’un besoin réel ; sept jours est une valeur par défaut, pas une politique.
  • Privilégiez un Kafka managé, à moins que l’exploitation de l’infrastructure ne soit le cœur de votre métier.
  • Utilisez Connect pour les intégrations standards et Streams pour le traitement au sein du cluster avant d’écrire du code d’intégration personnalisé.

Erreurs courantes

  • Supposer un ordonnancement global alors que Kafka ne garantit l’ordre qu’au sein d’une partition.
  • Utiliser une clé aléatoire ou nulle pour des événements qui doivent rester ordonnés.
  • Valider automatiquement les offsets (auto-commit) avant que le travail ne soit terminé, entraînant une perte d’enregistrements en cas de crash.
  • Créer des handlers non-idempotents, ce qui duplique les effets de bord lors d’un rebalance.
  • Créer un topic avec une seule partition et s’étonner que les consumers ne puissent pas monter en charge.
  • Augmenter le nombre de partitions ultérieurement, modifiant ainsi silencieusement le routage des clés.
  • Traiter Kafka comme une API de type requête/réponse avec des réponses par message.
  • Régler max.poll.records trop haut et dépasser max.poll.interval.ms.
  • Oublier que la réplication multiplie les coûts de stockage et de réseau.
  • Ne pas surveiller le consumer lag sur le dashboard jusqu’à ce que le backlog s’accumule sur plusieurs heures.
  • Faire tourner ZooKeeper sur un nouveau cluster en 2026.

Et après ?

Kafka est le journal (log) au centre d’un système piloté par les événements. La roadmap backend aborde l’architecture event-driven et les microservices, expliquant comment s’articulent les événements, les producteurs et les consommateurs, et comment un log devient la colonne vertébrale entre les services. Si vous comparez différents brokers, les guides sur RabbitMQ et les Redis Queues présentent des alternatives axées respectivement sur le routage et sur la gestion de tâches (jobs). Enfin, comme les consommateurs ne sont que des workers d’arrière-plan avec des offsets, le guide sur le Batch Processing traite des tentatives de réexécution (retries), de l’idempotence et de l’observabilité, des concepts qui s’appliquent également ici.

En pratique

Produire, consommer, commiter, administrer

Les quatre opérations qui constituent presque toutes les applications Kafka.

producer.ts
import { Kafka } from "kafkajs";

const kafka = new Kafka({
  clientId: "orders",
  brokers: ["localhost:9092"],
});

const producer = kafka.producer({
  idempotent: true,
  maxInFlightRequests: 1,
});

await producer.connect();

await producer.send({
  topic: "orders.created",
  acks: -1, // wait for all in-sync replicas
  messages: [
    { key: order.customerId, value: JSON.stringify(order) },
  ],
});

await producer.disconnect();

Log Kafka vs file d'attente classique

Une file d'attente supprime un message une fois acquitté. Un log Kafka le conserve, permettant à tout groupe de rejouer l'historique et à plusieurs lecteurs de partager le même flux.

Log Kafka
// A second group can read the same history
// from the beginning, months later.
await consumer.subscribe({
  topic: "orders.created",
  fromBeginning: true,
});
File d'attente classique
// The message was acknowledged by the first
// worker and removed; it cannot be replayed
// or read by a second independent consumer.
await channel.ack(msg);

Auto-commit vs commit manuel

L'auto-commit s'exécute via un timer et peut dépasser des enregistrements encore en cours de traitement. Commiter après le travail terminé rend le mode de défaillance explicite.

Manuel
await consumer.run({
  autoCommit: false,
  eachMessage: async ({ topic, partition, message }) => {
    await handle(message);
    await consumer.commitOffsets([
      {
        topic,
        partition,
        offset: String(Number(message.offset) + 1),
      },
    ]);
  },
});
Auto
await consumer.run({
  // Offsets are committed on a timer, so a
  // crash can skip records that were never
  // processed.
  eachMessage: async ({ message }) => {
    await handle(message);
  },
});

Compromis

Kafka vaut-il son poids opérationnel ?

Kafka résout des problèmes que les files d'attente ne peuvent pas gérer, mais il exige une infrastructure réelle et un modèle mental différent. Choisissez-le pour ses garanties, pas pour le hype.

Strengths

  • La relecture change tout

    L'historique étant conservé, vous pouvez initialiser un nouveau service à partir du log, reconstruire une projection après la correction d'un bug et analyser des événements bruts longtemps après leur occurrence.

  • Fan-out sans coordination

    N'importe quel nombre de consumer groups indépendants peut lire le même topic à son propre rythme. Ajouter un lecteur ne coûte rien au producer et ne perturbe jamais les consumers existants.

  • Débit scalable horizontalement

    Les écritures séquentielles, le batching et les partitions permettent à un petit cluster d'absorber des millions d'enregistrements par seconde ; on scale en ajoutant des brokers et des partitions.

Trade-offs

  • C'est une plateforme, pas une bibliothèque

    Brokers, réplication, rebalances, dimensionnement des disques et lag des consumers sont sous votre responsabilité. Les services managés aident, mais Kafka n'est jamais une simple dépendance qu'on oublie.

  • Le routage n'est pas son rôle

    Kafka n'a ni exchanges ni clés de routage. Les consumers lisent des topics entiers, donc le routage sélectif par message doit être intégré à la conception des topics ou au code applicatif.

  • Ordre uniquement au sein d'une partition

    L'ordre global n'est pas disponible. Garantir l'ordre implique de choisir soigneusement les clés et le nombre de partitions, et vous ne pouvez pas ajouter de partitions plus tard sans casser le routage par clé.

FAQ

Foire aux questions

Keep learning

Related topics from the roadmap.

$ commencer à apprendre

Prêt à apprendre Apache Kafka ?

Notre tutoriel interactif vous guide à travers Apache Kafka pas à pas — avec des quiz et du vrai code que vous pouvez exécuter dans le navigateur.