Kafka remplace vos middlewares — une supply chain de 200 magasins pilotée par 3 agents IA


TL;DR — Ce qu’un distributeur met 4 heures à traiter — détecter une rupture, décider un réappro, exécuter le transfert, tracer la décision — trois agents IA sur Kafka le font en 30 secondes. MCP Confluent pour interroger les topics en langage naturel, KIP-932 pour le scaling élastique, un skill métier pour la logique de décision. Zéro middleware, code exécutable.

Au sommaire


Architecture supply chain en cartoon — magasins, pipeline Kafka, trois agents

1. Pourquoi cet article, pourquoi maintenant

L'évolution de Kafka et des agents IA a convergé en une décennie. Deux trajectoires distinctes — l'une côté infrastructure de streaming, l'autre côté frameworks d'agents — qui se rejoignent en 2026 au point de ne plus nécessiter la moindre couche de glue. 2015 Bus de logs LinkedIn 2020 Stream processing KSQL · Streams · EOS 2024 Script Python + prompt system Mi-2025 LangChain · CrewAI Kafka opaque 2026 Convergence Apache Kafka Agents IA

Deux briques récentes changent la donne :

Brique Date De quoi il s’agit Référence
MCP Confluent GA juillet 2026 Un serveur MCP open-source qui expose 50+ primitives Kafka — topics, consumer groups, schémas Avro, Flink SQL, connectors — comme des tools appelables en langage naturel par n’importe quel agent compatible (Claude, Cursor, ADK). L’agent n’apprend pas le protocole Kafka ; il appelle des tools, le MCP traduit. Docs Confluent · GitHub
KIP-932 Share Groups GA avril 2026 (broker/Java) Introduit les share groups : une consommation coopérative où N consommateurs lisent les mêmes partitions sans la contrainte 1 consommateur = 1 partition. Chaque message suit un cycle fetch → lock → ack/release, avec retry natif (delivery_count). Kafka devient une file de tâches, sans RabbitMQ. Note : le client Python confluent-kafka est encore en Preview — seule l’API Java est GA. KIP-932 (Apache)

Ce que cet article démontre : ces deux briques en action sur un cas concret.

2. Le use case — réapprovisionnement retail

Un distributeur exploite 200 magasins, chacun référençant 50 produits. Chaque produit a un seuil minimum de stock par magasin. Dès qu’un stock passe sous ce seuil, une anomalie est détectée, une décision de réapprovisionnement est prise — transfert depuis un magasin voisin ou commande fournisseur — la tâche est exécutée, et l’opération est journalisée.

flowchart TD
    A["Réseau de 200 magasins × 50 produits<br/>chaque produit a un seuil minimum par magasin"] --> B{"Stock d'un produit<br/>< seuil minimum ?"}
    B -- Non --> A
    B -- Oui --> C["Anomalie détectée<br/>(rupture de stock imminente)"]
    C --> D{"Un magasin voisin<br/>a-t-il du stock disponible ?"}
    D -- Oui --> E["Décision : transfert interne<br/>depuis le magasin voisin"]
    D -- Non --> F["Décision : commande fournisseur"]
    E --> G["Exécution de la tâche<br/>(transfert ou commande)"]
    F --> G
    G --> H["Traçabilité : journal d'audit"]

Le circuit humain actuel — remontée d’information, analyse, validation, exécution — prend environ 4 heures par anomalie. Avec 200 magasins et 50 produits, le volume d’anomalies quotidien peut atteindre plusieurs centaines. L’automatisation n’est pas un luxe, c’est une nécessité opérationnelle.

Le use case est volontairement simplifié — 200 magasins, 50 produits, une seule règle de décision. L’objectif n’est pas de modéliser un réseau de distribution réel avec ses entrepôts, ses contraintes logistiques et ses fournisseurs multiples, mais d’isoler le pattern architectural : comment des agents coopèrent sur Kafka sans middleware de glue. La complexité métier s’ajoute en enrichissant le SKILL.md — l’architecture, elle, ne change pas.

Le code source complet, les règles métier détaillées et les instructions pour exécuter le PoC sont sur github.com/arabaaoui/kafka-for-agents.


3. Pourquoi Kafka plutôt qu’un orchestrateur

On aurait pu coder ça avec un workflow engine, une file de tâches, un orchestrateur Python. Le problème, c’est que chaque brique aurait nécessité un middleware différent :

Brique Approche classique Problème
Interroger les flux de données API REST custom sur Kafka Une couche de glue à maintenir
Distribuer les tâches RabbitMQ ou SQS Deux stacks à opérer
Éviter les doubles exécutions Lock distribué (Redis) Un composant de plus à surveiller
Tracer les décisions Base de données d’audit Synchro Kafka → DB à coder

L’architecture retenue fait le pari inverse : tout sur Kafka. Pas de glue. Pas de middlewares. Trois agents, cinq topics, et les primitives natives de la plateforme.

graph TB
    SIM["Simulateur<br/>200 magasins × 50 produits"] -->|stocks| KAFKA

    subgraph KAFKA["Apache Kafka 4.2.1 (KRaft)"]
        STOCKS[("stocks")]
        ANOMALIES[("anomalies")]
        TASKS[("tasks")]
        AUDIT[("audit")]
    end

    STOCKS -->|poll 5s| DET["Agent Détection<br/>MCP Confluent"]
    DET -->|produce_anomaly| ANOMALIES
    ANOMALIES -->|poll 3s| DEC["Agent Décision<br/>SKILL.md métier"]
    DEC -->|produce_task| TASKS
    DEC -->|log_audit| AUDIT
    TASKS -->|ShareGroup KIP-932| EXE["Agent Exécution ×N<br/>scale horizontal"]

4. MCP Confluent — l’agent qui voit Kafka comme une API

L’Agent Détection n’a pas de KafkaConsumer codé en dur. Il utilise le serveur MCP Confluent — un bridge qui expose les primitives Kafka (consume-messages, list-topics, get-topic-config) comme des tools que l’agent ADK peut appeler en langage naturel.

sequenceDiagram
    participant K as Kafka (stocks)
    participant P as Agent Détection (ADK)
    participant M as MCP Confluent
    participant A as Kafka (anomalies)

    K->>P: stock P-42 = 5, seuil_min = 20
    P->>P: pré-filtre déterministe : 5 < 20 → candidat
    P->>M: call_mcp("consume-messages", {topic:"stocks", max_messages:5})
    M-->>P: historique récent du produit P-42
    P->>P: qualification LLM : sévérité, tendance
    P->>A: produce_anomaly("rupture_stock", CRITIQUE)

L’agent ne « sait » pas interroger Kafka. Il sait appeler des tools. Le bridge MCP traduit. Si demain on change de cluster Kafka, l’agent ne change pas — seul le endpoint MCP change.

5. Skill métier — la logique de décision sans code

L’Agent Décision ne contient pas de if stock_voisin > besoin: transfert() else: commande(). Sa logique est dans un fichier SKILL.md, lu une fois au démarrage et injecté dans l’instruction de l’agent :

SKILL: supply-chain-replenishment
1. Identifier le produit et le magasin concernés
2. Vérifier le stock des magasins voisins (même région)
3. Si transfert possible → générer ordre de transfert
4. Sinon → générer commande fournisseur prioritaire
5. Si produit périssable → ajuster la quantité (buffer +10%)
6. Logger la décision dans le topic audit
sequenceDiagram
    participant A as Kafka (anomalies)
    participant P as Agent Décision (ADK)
    participant S as Kafka (stocks voisins)
    participant T as Kafka (tasks)
    participant AU as Kafka (audit)

    Note over P: instruction = prompt + SKILL.md<br/>(chargé au démarrage)

    A->>P: anomalie: P-42, Lyon-Part-Dieu, qty=5
    P->>S: get_neighbor_stocks("Lyon", "P-42")
    S-->>P: Bellecour: 200u, Confluence: 80u
    P->>P: applique SKILL.md → transfert possible
    P->>T: produce_task(transfert, 50u, Bellecour→Part-Dieu)
    P->>AU: log_audit(décision, raison=voisin avec surplus)

Le métier peut modifier la règle sans toucher une ligne de Python. Le skill est dans Git — versionné, reviewable, auditable. Un docker compose restart decision-agent et la nouvelle règle est en production.

6. KIP-932 Share Groups — le scaling élastique natif

Le matin à 8h, 200 anomalies sont détectées (pic de la nuit). L’Agent Exécution doit traiter 200 tâches en quelques minutes.

Sans share groups, le scaling est plafonné par le nombre de partitions : 3 partitions = 3 agents max. Avec KIP-932, tous les agents coopèrent sur les mêmes partitions :

stateDiagram-v2
    [*] --> AVAILABLE: tâche produite dans tasks
    AVAILABLE --> ACQUIRED: fetch par Agent-7 (lock 30s)
    ACQUIRED --> ACKNOWLEDGED: ACK après exécution réussie
    ACQUIRED --> AVAILABLE: crash agent → lock expire → retry
    ACQUIRED --> DEAD: delivery_count > 5 → abandon

    note right of ACQUIRED
        Les N-1 autres agents continuent
        de fetcher sur les mêmes partitions.
        Pas de rebalance. Pas de limite.
    end note
# Pic de charge à 8h
docker compose up -d --scale execution-agent=20
# → 200 tâches traitées en 2 minutes

# Creux à 14h
docker compose up -d --scale execution-agent=2
# → redescente sans rebalance, sans perte de message

Kafka devient une file de tâches native. Plus besoin de RabbitMQ à côté. Le delivery_count natif gère les retries. Si un agent crashe en plein traitement, le lock expire après 30 secondes et la tâche est automatiquement remise en circulation.

7. Le pipeline complet en 30 secondes

sequenceDiagram
    participant S as Simulateur
    participant D as Détection (ADK + MCP)
    participant DEC as Décision (ADK + SKILL.md)
    participant E as Exécution ×N (KIP-932)
    participant A as Audit

    S->>D: stock P-42 Lyon-Part-Dieu = 5 (seuil=20)
    D->>D: pré-filtre → MCP → qualification CRITIQUE
    D->>DEC: anomalie: rupture_stock, P-42, Part-Dieu

    DEC->>DEC: vérifie voisins → Bellecour a 200u
    DEC->>DEC: SKILL.md → transfert interne 50u
    DEC->>E: task: transfert(Bellecour→Part-Dieu, 50u)
    DEC->>A: audit: décision=transfert, raison=surplus voisin

    E->>E: exécute transfert (ACK après succès)
    E->>A: audit: task completed, 2.0s
Étape Durée Ce qui se passe
t=0s — Le simulateur publie un stock sous le seuil
t=5s +5s L’Agent Détection qualifie l’anomalie via MCP
t=8s +3s L’Agent Décision choisit le transfert en appliquant SKILL.md
t=10s +2s L’Agent Exécution traite la tâche et acquitte
t=10s — La décision est journalisée dans l’audit

30 secondes pour ce qui prenait 4 heures en circuit humain. Sans un seul if codé en dur dans la logique métier. Sans middleware de glue entre les briques.

8. Ce qui tourne vraiment

Le PoC est exécutable — pas une maquette, pas des print() qui font semblant.

Chaque agent est un vrai google.adk.Agent avec un vrai LLM derrière :

Agent Stack Mode déterministe
Détection ADK + LiteLLM (OpenAI/Anthropic/Gemini) + MCP Confluent Anomalie basique sans qualification LLM
Décision ADK + SKILL.md injecté + tools Kafka réels Commande fournisseur par défaut
Exécution ShareConsumer natif (ACK/RELEASE/REJECT) + scaling natif Délai fixe 2s, 100% succès

Et si vous n’avez pas de clé LLM sous la main, chaque agent bascule en mode déterministe. Le pipeline tourne de bout en bout sans rien payer.

Note sur le ShareConsumer natif — le client Python confluent-kafka expose ShareConsumer en Preview depuis la version 2.15.0. Le PoC l’utilise directement pour le cycle de vie KIP-932 — AVAILABLE → ACQUIRED → ACKNOWLEDGED, expiration des locks, compteur de tentatives, dead-letter. Reste une limite du client Preview : pas d’action RENEW, ce qui empêche de prolonger un verrou pendant un appel LLM long — un agent lent sur un traitement risque de perdre son lock avant d’avoir acquitté. Détail complet dans le README du repo.

git clone https://github.com/arabaaoui/kafka-for-agents.git
cd kafka-for-agents
cp .env.example .env   # laisser vide = full déterministe
make all               # cluster Kafka test + agents + simulateur
make check             # vérifier que tout tourne

Kafka UI sur http://localhost:8081 permet de voir les messages circuler en temps réel à travers les 5 topics. make check affiche l’état des consumer groups, le nombre de messages par topic, et les dernières anomalies/tâches.

9. Démonstration et prochaine étape

Ce que cet article a démontré :

Cette architecture fait disparaître deux middlewares : la glue d’accès aux données (MCP remplace les scripts AKHQ/API REST) et la file de tâches externe (KIP-932 remplace RabbitMQ/SQS). Mais elle repose sur un skill métier chargé localement — un fichier markdown monté dans un container.

Le prochain article change de terrain : les agents ne sont plus dans le flux métier, ils sont autour — l’ops. Un poison message bloque un consumer group facturation. L’agent diagnostique la cause exacte via MCP en 2 minutes — là où un ops mettrait 30 à 90 minutes à la main. Le fix, lui, reste manuel : une commande CLI de 30 secondes. Pas d’écriture sur le cluster, pas d’autonomie déléguée — l’agent élimine le temps de diagnostic, pas le jugement humain. Un use case concret, toujours.


Références


Les agents utilisent Google ADK, LiteLLM pour le multi-provider, et Kafka 4.2.1 en KRaft. Testable localement avec Docker en deux commandes.

Commentaires