Kafka replaces your middlewares — a 200-store supply chain run by 3 AI agents


TL;DR — What a distributor takes 4 hours to process — detect a stockout, decide on replenishment, execute the transfer, trace the decision — three AI agents on Kafka do it in 30 seconds. MCP Confluent to query topics in natural language, KIP-932 for elastic scaling, a business skill for decision logic. Zero middleware, executable code.

In this article


Cartoon architecture — stores, Kafka pipeline, three agents

1. Why this article, why now


The evolution of Kafka and AI agents converged over a decade. Two distinct trajectories — one on the streaming infrastructure side, the other on the agent frameworks side — that join up in 2026 to the point of no longer requiring the slightest glue layer. 2015 Log bus LinkedIn 2020 Stream processing KSQL · Streams · EOS 2024 Script Python + prompt system Mi-2025 LangChain · CrewAI Kafka opaque 2026 Convergence Apache Kafka AI agents

Two recent building blocks change the game:

Building block Date What it is Reference
MCP Confluent GA July 2026 An open-source MCP server that exposes 50+ Kafka primitives — topics, consumer groups, Avro schemas, Flink SQL, connectors — as callable tools in natural language by any compatible agent (Claude, Cursor, ADK). The agent doesn’t learn the Kafka protocol; it calls tools, the MCP translates. Confluent Docs · GitHub
KIP-932 Share Groups GA April 2026 (broker/Java) Introduces share groups: cooperative consumption where N consumers read the same partitions without the 1 consumer = 1 partition constraint. Each message follows a fetch → lock → ack/release cycle, with native retry (delivery_count). Kafka becomes a task queue, without RabbitMQ. Note: the confluent-kafka Python client is still in Preview — only the Java API is GA. KIP-932 (Apache)

What this article demonstrates: these two building blocks in action on a concrete use case.

2. The use case — retail replenishment

A distributor operates 200 stores, each carrying 50 products. Each product has a minimum stock threshold per store. As soon as stock drops below that threshold, an anomaly is detected, a replenishment decision is made — transfer from a neighboring store or supplier order — the task is executed, and the operation is logged.

flowchart TD
    A["Network of 200 stores × 50 products<br/>each product has a minimum threshold per store"] --> B{"Product stock<br/>< minimum threshold?"}
    B -- No --> A
    B -- Yes --> C["Anomaly detected<br/>(imminent stockout)"]
    C --> D{"Does a neighboring store<br/>have stock available?"}
    D -- Yes --> E["Decision: internal transfer<br/>from neighboring store"]
    D -- No --> F["Decision: supplier order"]
    E --> G["Task execution<br/>(transfer or order)"]
    F --> G
    G --> H["Traceability: audit log"]

The current manual workflow — information gathering, analysis, validation, execution — takes about 4 hours per anomaly. With 200 stores and 50 products, the daily volume of anomalies can reach several hundred. Automation isn’t a luxury; it’s an operational necessity.

The use case is deliberately simplified — 200 stores, 50 products, a single decision rule. The goal isn’t to model a real distribution network with its warehouses, logistical constraints, and multiple suppliers, but to isolate the architectural pattern: how agents cooperate on Kafka without glue middleware. Business complexity is added by enriching the SKILL.md — the architecture doesn’t change.

The complete source code, detailed business rules, and instructions for running the PoC are on github.com/arabaaoui/kafka-for-agents.


3. Why Kafka over an orchestrator

We could have built this with a workflow engine, a task queue, a Python orchestrator. The problem is that each component would have required a different middleware:

Component Classic approach Problem
Querying data streams Custom REST API on Kafka A glue layer to maintain
Distributing tasks RabbitMQ or SQS Two stacks to operate
Preventing duplicate executions Distributed lock (Redis) Yet another component to monitor
Tracing decisions Audit database Kafka → DB sync to build

The chosen architecture takes the opposite bet: everything on Kafka. No glue. No middlewares. Three agents, five topics, and the platform’s native primitives.


graph TB
    SIM["Simulator<br/>200 stores × 50 products"] -->|stocks| KAFKA

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

    STOCKS -->|poll 5s| DET["Detection Agent<br/>MCP Confluent"]
    DET -->|produce_anomaly| ANOMALIES
    ANOMALIES -->|poll 3s| DEC["Decision Agent<br/>Business SKILL.md"]
    DEC -->|produce_task| TASKS
    DEC -->|log_audit| AUDIT
    TASKS -->|ShareGroup KIP-932| EXE["Execution Agent ×N<br/>horizontal scaling"]

4. MCP Confluent — the agent that sees Kafka as an API

The Detection Agent has no hardcoded KafkaConsumer. It uses the MCP Confluent server — a bridge that exposes Kafka primitives (consume-messages, list-topics, get-topic-config) as tools that the ADK agent can call in natural language.


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

    K->>P: stock P-42 = 5, min_threshold = 20
    P->>P: deterministic pre-filter: 5 < 20 → candidate
    P->>M: call_mcp("consume-messages", {topic:"stocks", max_messages:5})
    M-->>P: recent history of product P-42
    P->>P: LLM qualification: severity, trend
    P->>A: produce_anomaly("stockout", CRITICAL)

The agent doesn’t “know” how to query Kafka. It knows how to call tools. The MCP bridge translates. If tomorrow we switch Kafka clusters, the agent doesn’t change — only the MCP endpoint changes.

5. Business Skill — decision logic without code

The Decision Agent doesn’t contain if neighbor_stock > needed: transfer() else: supplier_order(). Its logic lives in a SKILL.md file, read once at startup and injected into the agent’s instruction:


SKILL: supply-chain-replenishment
1. Identify the product and store concerned
2. Check stock of neighboring stores (same region)
3. If transfer possible → generate transfer order
4. Otherwise → generate priority supplier order
5. If perishable product → adjust quantity (buffer +10%)
6. Log the decision in the audit topic
sequenceDiagram
    participant A as Kafka (anomalies)
    participant P as Decision Agent (ADK)
    participant S as Kafka (neighbor stocks)
    participant T as Kafka (tasks)
    participant AU as Kafka (audit)

    Note over P: instruction = prompt + SKILL.md<br/>(loaded at startup)

    A->>P: anomaly: P-42, Lyon-Part-Dieu, qty=5
    P->>S: get_neighbor_stocks("Lyon", "P-42")
    S-->>P: Bellecour: 200u, Confluence: 80u
    P->>P: applies SKILL.md → transfer possible
    P->>T: produce_task(transfer, 50u, Bellecour→Part-Dieu)
    P->>AU: log_audit(decision, reason=neighbor with surplus)

The business team can modify the rule without touching a single line of Python. The skill is in Git — versioned, reviewable, auditable. A docker compose restart decision-agent and the new rule is in production.

6. KIP-932 Share Groups — native elastic scaling

At 8 AM, 200 anomalies are detected (overnight peak). The Execution Agent must process 200 tasks within minutes.

Without share groups, scaling is capped by the partition count: 3 partitions = 3 agents max. With KIP-932, all agents cooperate on the same partitions:


stateDiagram-v2
    [*] --> AVAILABLE: task produced in tasks
    AVAILABLE --> ACQUIRED: fetch by Agent-7 (lock 30s)
    ACQUIRED --> ACKNOWLEDGED: ACK after successful execution
    ACQUIRED --> AVAILABLE: agent crash → lock expires → retry
    ACQUIRED --> DEAD: delivery_count > 5 → dead letter

    note right of ACQUIRED
        The N-1 other agents continue
        to fetch on the same partitions.
        No rebalance. No limit.
    end note
# Load spike at 8 AM
docker compose up -d --scale execution-agent=20
# → 200 tasks processed in 2 minutes

# Lull at 2 PM
docker compose up -d --scale execution-agent=2
# → scale down without rebalance, without message loss

Kafka becomes a native task queue. No need for RabbitMQ alongside it. The native delivery_count handles retries. If an agent crashes mid-processing, the lock expires after 30 seconds and the task is automatically put back into circulation.

7. The complete pipeline in 30 seconds

sequenceDiagram
    participant S as Simulator
    participant D as Detection (ADK + MCP)
    participant DEC as Decision (ADK + SKILL.md)
    participant E as Execution ×N (KIP-932)
    participant A as Audit

    S->>D: stock P-42 Lyon-Part-Dieu = 5 (threshold=20)
    D->>D: pre-filter → MCP → CRITICAL qualification
    D->>DEC: anomaly: stockout, P-42, Part-Dieu

    DEC->>DEC: checks neighbors → Bellecour has 200u
    DEC->>DEC: SKILL.md → internal transfer 50u
    DEC->>E: task: transfer(Bellecour→Part-Dieu, 50u)
    DEC->>A: audit: decision=transfer, reason=neighbor surplus

    E->>E: executes transfer (ACK after success)
    E->>A: audit: task completed, 2.0s
Step Duration What happens
t=0s — The simulator publishes a stock below threshold
t=5s +5s The Detection Agent qualifies the anomaly via MCP
t=8s +3s The Decision Agent selects the transfer by applying SKILL.md
t=10s +2s The Execution Agent processes the task and acknowledges
t=10s — The decision is logged in the audit trail

30 seconds for what took 4 hours in a human circuit. Without a single hardcoded if in the business logic. No glue middleware between the components.

8. What actually runs

The PoC is executable — not a mockup, not print() statements pretending to work.

Each agent is a real google.adk.Agent with a real LLM behind it:

Agent Stack Deterministic mode
Detection ADK + LiteLLM (OpenAI/Anthropic/Gemini) + MCP Confluent Basic anomaly without LLM qualification
Decision ADK + SKILL.md injected + real Kafka tools Default supplier order
Execution Native ShareConsumer (ACK/RELEASE/REJECT) + native scaling Fixed 2s delay, 100% success

And if you don’t have an LLM key handy, each agent falls back to deterministic mode. The pipeline runs end-to-end without paying anything.

Note on the native ShareConsumer — the confluent-kafka Python client exposes ShareConsumer in Preview since version 2.15.0. The PoC uses it directly for the KIP-932 lifecycle — AVAILABLE → ACQUIRED → ACKNOWLEDGED, lock expiration, retry counter, dead-letter. One limitation of the Preview client: no RENEW action, which prevents extending a lock during a long LLM call — a slow agent risks losing its lock before acknowledging. Full details in the repo README.


git clone https://github.com/arabaaoui/kafka-for-agents.git
cd kafka-for-agents
cp .env.example .env   # leave empty = fully deterministic
make all               # Kafka test cluster + agents + simulator
make check             # verify everything is running

Kafka UI on http://localhost:8081 lets you see messages flowing in real time across the 5 topics. make check displays consumer group states, message counts per topic, and the latest anomalies/tasks.

9. Demonstration and Next Step

What this article demonstrated:

This architecture eliminates two middlewares: the data access glue (MCP replaces AKHQ/REST API scripts) and the external task queue (KIP-932 replaces RabbitMQ/SQS). But it relies on a business skill loaded locally — a markdown file mounted in a container.

The next article changes scope: the agents are no longer in the business flow — they’re around it: ops. A poison message blocks a billing consumer group. The agent diagnoses the exact cause via MCP in 2 minutes — where an ops engineer would take 30 to 90 minutes by hand. The fix remains manual: a 30-second CLI command. No cluster writes, no delegated autonomy — the agent eliminates diagnostic time, not human judgment. A concrete use case, as always.


References


The agents use Google ADK, LiteLLM for multi-provider, and Kafka 4.2.1 in KRaft. Testable locally with Docker in two commands.

Commentaires