On this page
Kafka for agents in one paragraph
Use topics as the agent's context and history: change data capture from your databases keeps retrieval indexes current, a compacted topic holds the latest state per entity, and a retained topic is a replayable record of every decision. Key messages by session or conversation so each one stays ordered on one partition. Treat tool calls as request and reply over topics with correlation IDs, run agents as consumer groups or, on Kafka 4.2 and later, share groups, and make every external side effect idempotent, because Kafka transactions stop at the cluster boundary. Put a schema on agent messages. Off Confluent Cloud you lose Streaming Agents and managed Flink, and you replace them with self-run or cloud-managed Flink, Kafka Streams, the preview Apache Flink Agents project, and an open-source schema registry. Any MCP server that touches the cluster gets its own principal and an allow-list of tools.
Event streams as agent context and memory
An agent is only as good as what it knows when it acts. Kafka gives it two useful forms of memory. A retained topic is history: every event in order, replayable from any offset, which is how you rebuild an agent's view after a bug or test a new prompt against last week's traffic. A compacted topic is current state: Kafka keeps at least the last value for every key, so a consumer reading from the start gets the latest record per customer, order or ticket without a separate database.
Records with the same key go to the same partition, and Kafka guarantees consumers read a partition in the order it was written. That is the property agents need most: if the key is the session or conversation ID, every message in that conversation arrives in sequence at whichever agent instance owns the partition.
Kafka Connect and CDC for RAG pipelines
Retrieval is only useful if the index matches the source. The common pattern is change data capture: Debezium runs as Kafka Connect source connectors, reads the database transaction log, and writes one change stream per table to Kafka. A processing step turns each change into chunks and embeddings, and a sink writes them to the vector store, whether that is Redis, OpenSearch, PostgreSQL with a vector extension or a dedicated database. Updates re-embed the row; deletes arrive as delete events and tombstones, and the pipeline must remove the matching vectors, or the agent will keep citing records that no longer exist.
Keep the embedding step separate from the connector so you can re-embed from the topic when you change models, which you will. Retain the change topics long enough to make that replay possible.
Stream processing for features and embeddings
Two engines do most of this work. Kafka Streams is a library inside your application: it keeps fault-tolerant local state stores, supports exactly-once processing between Kafka topics, and suits enrichment, joins and per-session aggregates that feed an agent's context. Apache Flink is a separate cluster with SQL, the Table API and asynchronous I/O, which is the usual way to call an embedding model or inference endpoint from a stream without blocking it.
The Flink community has also started Apache Flink Agents, a framework for event-driven agents on Flink in Java and Python, with tool calling, short-term memory with TTLs, long-term memory and vector store integrations. Its 0.3 release in June 2026 describes the APIs as experimental and subject to incompatible change, so treat it as something to evaluate, not a foundation for a production system yet.
MCP servers for Kafka, and what they should not be allowed to do
Confluent publishes an open-source MCP server, mcp-confluent, under the MIT license. It works against Confluent Cloud, Confluent Platform and standalone Apache Kafka through bootstrap servers, exposes more than 50 tools across topics, produce and consume, consumer groups and lag, Flink SQL, connectors, Schema Registry and Tableflow, and runs over stdio, HTTP or SSE. Tools that depend on Confluent services only work where those services exist. It also supports allow-lists and block-lists of tools, which is the feature that matters most. Community Kafka MCP servers exist as well; review them like any third-party code you hand cluster credentials.
On a production cluster, give the MCP server its own principal with ACLs scoped to named topics and groups. Allow describe, read and lag inspection. Block topic deletion, configuration changes, partition reassignment and offset resets outright, and allow produce only to a staging or request topic that a human-reviewed process consumes. Set client quotas so a looping agent cannot saturate a broker, and keep audit logs of every call.
Patterns: tool calls, agent pools and ordering
Request and reply for tool calls. The agent writes a request to a tool's topic with a correlation ID and a reply topic in the headers; the tool service consumes it, does the work and writes the result to the reply topic. The agent matches replies by correlation ID and times out the ones that never come. This decouples agents from tools, gives every call a durable record, and lets a slow tool queue rather than fail.
Consumer groups as agent pools. Run agent instances as one consumer group and Kafka spreads partitions across them, but parallelism stops at the partition count. Queues for Kafka (KIP-932), production-ready since Kafka 4.2, adds share groups: consumers can outnumber partitions, records are acknowledged individually and delivery attempts are counted, which fits independent tasks such as document enrichment. Share groups give up per-partition ordering, so keep conversations on ordinary consumer groups.
Ordering per conversation. Key every message by session or conversation ID. One key, one partition, one ordered sequence. If a single conversation can outgrow one consumer, the design problem is the conversation, not the partition count.
Exactly-once stops at the cluster boundary
Idempotent producers and transactions give exactly-once processing from Kafka topic to Kafka topic, and Kafka Streams uses them to do so. An agent's side effects are rarely Kafka writes: it sends an email, charges a card, files a ticket. Kafka's own design documentation is clear that writing to an external system needs the consumer's position coordinated with the output, either by storing offsets alongside the result or by making the write idempotent. For agents, the practical rule is an idempotency key on every tool call, derived from the conversation and step, checked by the tool before it acts.
Retention, compaction and schemas for agent messages
Set retention per topic by what the data is for. Conversation and decision logs need long retention for replay and audit. Current agent state belongs in a compacted topic, where a record with a null value is a tombstone that removes the key. Tombstones are themselves cleaned up after delete.retention.ms, 24 hours by default, so a consumer that lags longer than that can miss a deletion; size that setting to your slowest rebuild. For a compacted topic that should also age out, use the combined compact and delete cleanup policy.
Agent messages change shape often, so give them schemas. Confluent Schema Registry is under the Confluent Community License, with its client and serializer modules under Apache 2.0. Apicurio Registry and Karapace are Apache 2.0 registries that implement the Confluent-compatible API, so most existing serializers work against them. Use Avro, Protobuf or JSON Schema, and enforce compatibility rules, so a prompt change that renames a field cannot break every downstream consumer.
Without Confluent Cloud: what you give up and what replaces it
Confluent Streaming Agents run on Confluent Cloud for Apache Flink. They let you define agents in SQL, call tools through MCP, invoke models from several providers, and replay agent runs against historical events for testing. That is a real convenience, and it does not exist on open-source Kafka, Amazon MSK or Aiven. Neither does fully managed Flink tied to the same cluster.
The replacements are components, not a product. Run Flink yourself, or on AWS use Amazon Managed Service for Apache Flink, which reads from MSK. Write agents in Kafka Streams or a consumer with your own agent framework, and keep an eye on Apache Flink Agents as it matures. Swap Schema Registry for Apicurio or Karapace if licensing matters. mcp-confluent still works against standalone Kafka for the core topic and consumer tools. What you are really giving up is integration and someone else's pager, which is the question to answer before choosing.
Who runs it
An agent platform on Kafka is a Kafka estate with more consumers, more connectors and stricter audit needs. The work is the same as any production cluster: partition and retention planning, ACLs and quotas, connector operations, upgrades to the release that has the features you need. Decide early whether your team holds that pager or hands it to someone who does.
Frequently asked questions
Can I use Confluent's MCP server with open-source Kafka?
Yes. mcp-confluent connects to standalone Apache Kafka through bootstrap servers, as well as to Confluent Cloud and Confluent Platform. Tools that depend on Confluent services, such as Flink SQL or Tableflow, need those services. Use its allow-list to expose only the tools an agent needs.
Does Kafka give exactly-once delivery for agent tool calls?
No. Transactions give exactly-once processing between Kafka topics. A tool call that writes to an external system needs an idempotency key or offsets stored with the result, so a retried message does not repeat the action.
Are Kafka share groups ready for production?
Yes, from Apache Kafka 4.2, where Queues for Kafka (KIP-932) was declared production-ready after a preview in 4.1. Share groups let consumers outnumber partitions and acknowledge records individually, at the cost of per-partition ordering.
Is Apache Flink Agents production-ready?
Not yet. The 0.3 release notes describe the APIs and configuration as experimental and subject to non-backward-compatible changes. Evaluate it, but build production agents on Flink or Kafka Streams with your own agent framework for now.
What do I lose by running agents on MSK or Aiven instead of Confluent Cloud?
Confluent Streaming Agents and the managed Flink they run on. The underlying patterns work on any Kafka; you assemble Flink, a schema registry and agent code yourself, or use a cloud-managed Flink service such as Amazon Managed Service for Apache Flink.
Related
Where this gets done
The work behind this page, run by the same engineers who wrote it.
- 24/7 Kafka supportSelf-managed, MSK or Confluent Platform
- Managed Kafka servicesWe run the clusters, on your infrastructure
- Debezium CDC supportChange data capture from Oracle, SQL Server and PostgreSQL
- Kafka consultingPartition strategy, sizing, security and migration
- RabbitMQ supportIf the estate runs both brokers
- Kubernetes and container servicesKafka on Kubernetes, operated with your team
- Enterprise MQ supportOne contract across Kafka, RabbitMQ and IBM MQ
- Enterprise support plansSLA tiers and what each covers
- Enterprise MQ consultingMulti-broker architecture and migration
Other Kafka guides, comparisons and research
Recent Kafka articles
- How Much Does Kafka Enterprise Support Cost?Oct 2026
- Kafka Disaster Recovery: Architectures, RPO/RTO, TestingOct 2026
- Apache Kafka End of Life: Supported Versions in 2026Oct 2026
- Confluent vs Independent Kafka SupportOct 2026
- Kafka Consumer Lag: Causes and FixesSep 2026
- Kafka Rebalance: Triggers and How to Stop ItSep 2026
Need this done on your Kafka estate?
Named senior Kafka engineers, 24/7, with a 15-minute emergency SLA — self-managed, MSK or Confluent Platform.