Where AI Agents Belong in a Kafka Architecture

Stéphane Derosiaux October 1, 2026 8 min read
A dense stream of small wireframe squares flows left to right through a vertical grid sieve; one square, glowing lime, leaves the stream and travels along a dashed line to a single spherical node.

When does an LLM agent belong on a Kafka stream?

I often see demos with an agent plugged straight into a topic, calling an LLM on every event, each order, each payment. This makes no sense:

  • It's slow and expensive. An agent costs seconds and cents per call, on a stream that moves in milliseconds.
  • Most events carry no question. Telemetry or business, they need no interpreting, searching or deciding, which is what LLMs are good at.
  • Per-record enrichment is a classic ML job. When an LLM does show up on a live stream, it's usually to enrich each record (think sentiment scoring). A simple ML model does that better.

From adtech. I used to work in adtech, where we ran Vowpal Wabbit to predict clicks on a high-throughput Kafka topic of ad impressions. An LLM there would have been absurd: low latency and a clear probabilistic model were what drove business performance.

LLMs make that kind of capability trivial to add. They're generalist, they seem to work, and it's tempting to skip building an ML model altogether. But a mundane task doesn't need a generalist, and an LLM gives you no consistent prediction model you can evaluate and forecast against.

The real use case I see is narrower: events that represent a business anomaly worth investigating, where something has to interpret it, pull in more context, and make a decision a downstream process can act on. Everything else should be deterministic processing and classic ML.

telemetry · ordersevery event?LLM agent?

"Kafka agent" covers three different architectures:

  • Model 1, the agent that operates Kafka. An assistant for platform engineers and developers: why is this consumer group lagging, who owns this topic, create a topic that follows our policies. Most of the demand we hear is here: it's closer to an AI SRE than to a data pipeline.
  • Model 2, the agent that consumes topics. A stream processing app detects something worth investigating (an account that looks compromised, a shipment that missed three checkpoints) and hands the agent a case, not every record.
  • Model 3, the event-driven app with an LLM inside. A model is one step of a pipeline: ticket triage, document summarization, multi-agent task queues. Low volume, latency-tolerant work, not the high-volume business events.

Model 1: the agent that operates Kafka

An agent that helps engineers build applications on Kafka and operate it.

EngineerAI assistantMCPtoken: read / writeKafka metadatatopics · groups · ACLs · lagConduktor metadataowners · applications · policies · labelstopic records

Platform engineers often already live in Claude Code or Cursor. Their assistant should answer the questions product teams keep asking them:

  • Why is this consumer group lagging since 9am?
  • Where is this data, who owns this topic, and who do I ask for access?
  • Create a topic for this new service, compliant with our naming and retention policies.
  • Which topics contain PII and haven't been read in 90 days?

The agent here is a client of the control plane. It reads metadata (topics, consumer groups, schemas, ACLs, lag) and correlates it. When a change is needed, it runs with the engineer's permissions. This kind of agent never subscribes to a topic. The tricky part is making it secure by construction:

  • Delegated identity: no shared "ai-bot" service account. The agent sees only what the engineer's RBAC allows, and every action is attributable to the engineer who triggered it. The AI is never the accountable party.
  • Short-lived tokens: agent tokens expire quickly and are scoped to the task at hand.
  • Data masking: if the agent has to search data, apply field-level masking to sensitive fields, because the agent will see payloads.

"[An] MCP server where you can search data would also mean that this data get to the LLM. And that's probably not exactly what security would like us to do." — An engineer, on enabling MCP access to Kafka

The agent is a new user type of your control plane. If your Kafka platform already has RBAC, audit, and masking for humans, you're extending it. If you don't have that capability yet (Conduktor provides it), thinking about agents will force two conversations: ownership (why ownership comes first) and delegated identity with bounded permissions.

Model 2: the agent that consumes topics

This is where the FOMO is: "real-time context for agents", the agent drinking fresh events straight from the stream. We covered the useful side of that idea in Kafka and Flink Are the Infrastructure for AI Agents.

Start with throughput. At 10,000 records/s, one model call per record is 864 million model calls a day. Latency, quotas and cost all break long before anything else. Nice for demos, useless in real life.

stream10,000 msg/s1 call per recordLLM agent864Mmodel calls/day

An agent answers a request. To build the answer, it runs an agentic loop: it calls MCP tools, does external lookups, queries APIs and databases, and reasons over what it finds. Most events don't need that machinery. Only some do.

Example: anomaly detection with an agent to investigate

In production, agents work from a topic dedicated to things worth investigating. The goal of each stage is to move fewer events to more expensive stages:

eventsevery record10,000/sdeterministicrules · windows · MLcasesanomaliesa few/minLLM agentinvestigateshumanif neededrareattention is expensive
  • Deterministic first. Stream processing (Kafka Streams, Flink) does the joins, windows, aggregations and rules, or runs an anomaly detection model. This layer reads every event.
  • An investigation request event. The output is a situation somebody would investigate (an account issue, a shipment issue), with the issue and some information attached.
  • The agent handles the ambiguous tail. Cases where the resolution path isn't known in advance and the data to check is spread across systems, the kind of work a human would otherwise spend an hour gathering and reasoning about.
  • A human only when needed. The agent resolves the case or escalates it. Human attention is the most expensive step, so it comes last.

To classify events, you have options: a simple algorithm, an ML model, a small fine-tuned model, or one of the new fast classifiers like TypeSafe's Jev or OpenAI's Decisions API, which return a typed decision (escalate or not, category, urgency) with a probability.

Agents are for ambiguity, not volume.

Would Postgres or ClickHouse + MCP be enough?

Once Kafka is in place, it's tempting to make everything streaming. Often it isn't necessary: replace Kafka + Flink + a materialized context store with Postgres or ClickHouse + MCP, and see if that answers your needs.

Kafka pays off when detection must be continuous.

  • A time-series database recomputes any window when queried, late events included, but only when someone asks.
  • A stream evaluates the condition on every event and fires the moment it's true.

Polling the database every few seconds to get the same reaction means rebuilding, badly, a pipeline Kafka stream processing gives you directly. And once detection runs on the stream, replay comes for free: the same code re-reads history, with no separate batch job to build.

swap the stream forPostgres / ClickHouse + MCPsame answerquery systems of recordover MCP · ships sooneranswer changesyou need the streamcontinuous detection · replay for free

Model 3: the event-driven app with an LLM inside

An application consumes Kafka events, calls an LLM somewhere in its processing, and produces results back to Kafka.

It works when the stream is low-volume and latency isn't critical: support tickets triaged and routed, documents or emails summarized, notes drafted for a human.

tickets · docslow volumeLLM stepinside the appschema checkoutput topicdead letterorders · paymentshigh volumeLLM stepseconds per record

But it doesn't work on high-volume business streams. Enriching every order or payment with an LLM adds seconds of latency and a per-record cost to a pipeline built for milliseconds. Nobody does fraud detection with an LLM: it runs on deterministic rules, specialized ML models and scoring services.

So what is a "Kafka agent" here? Just an application that consumes from and produces to Kafka. It has a service account, ACLs, quotas and schema contracts, like any other consumer and producer. Kafka doesn't know there's an LLM inside, and doesn't care.

I expect more multi-agent systems to be built on Kafka, now that share groups give it queue semantics and make it a strong backbone for asynchronous agentic workflows:

  • an orchestrator publishes tasks or events to a topic;
  • a pool of agent workers consumes them;
  • results come back on another topic, asynchronously.

That's ordinary queue work between applications, with what Kafka already brings: horizontal scaling, proven resilience, and a large ecosystem.

An agent job can run for minutes. Since Kafka 4.2, a worker can extend its lock on a record with a RENEW acknowledgement. Without it, Kafka treats the worker as stuck and hands the record to another worker, which runs the side effects a second time unless your consumers are idempotent.

Conclusion

LLMs on raw event streams make good demos and no sense in production. On Kafka, they help engineers operate it, investigate anomalies, and handle low-volume steps inside applications.

Whichever you use, you have to answer three questions: who is the agent acting for, what can it read, and what did it do? Identity, scope, audit.

We built Conduktor MCP around those questions. It runs with the permissions of the token behind it, inheriting the user's own Console permissions. A read-only token lets an agent investigate with no risk of changing anything. Grant write permissions and it can automate low-risk tasks, such as updating topic labels and metadata.


Want to see what governed agent access to Kafka looks like? Book a demo, or read how we use MCP and skills for AI agents.