Event-Driven AI Agents with Apache Kafka
Build event-driven AI agents with Kafka: schemas, idempotent actions, approval workflows, retries, and safe replay.
I use Kafka to deliver events - not to give AI agents permission to act. I start with one workflow, let the agent propose a response, and route approved actions through a separate service with duplicate checks.
The key limit: Kafka’s exactly-once processing does not cover external actions. An eight-partition topic supports up to eight active consumers in a standard consumer group, but model delays and approval queues still affect completion time.
In this guide, I cover how to:
- Define event schemas, partition keys, and replay rules.
- Build Python workers with durable state, retries, and approval checks.
- Prevent repeated actions and confirm results with external services.
- Test failures, track decision quality, and control access and data retention.
- Plan capacity, recovery, costs, and team ownership.
My starting rule: <u>test in shadow mode before allowing actions</u>. If a scheduled job or direct API call meets your needs, choose that instead.
Beyond reactive systems: event-driven architecture for AI agents
sbb-itb-f123e37
Designing the AI Agent Architecture
Kafka AI Agents: From Events to Controlled Actions
Turning Business Events into Controlled Actions
Use topics, partitions, and offsets to keep raw events, enriched context, proposals, and outcomes in separate streams. Applications, database change streams, and devices send business signals to Kafka. Enrichment services add permitted context and publish an enriched event without changing the original fact.
Agent workers use that context, rules, models, and approved tools to create proposals, not direct actions. Keep proposals, action requests, execution results, and outcomes separate. The flow is:
Raw event → enriched event → proposal → approval → execution result → outcome
Action services manage credentials, authorization, retries, and rate limits. Outcome streams record confirmed results and human decisions.
Hypothetical small-business support workflow: A U.S. software company publishes
support.case.created. An enrichment worker adds subscription details, recent cases, documentation, and the service-level deadline. The agent recommends a queue and drafts a response insupport.case.proposed. Policy permits low-risk routing, while refunds, account changes, security-sensitive responses, and low-confidence proposals require supervisor approval. An approval deadline could be set to 4 business hours. After approval, the action service sends the reply and records the provider’s message ID insupport.response.sent. The case remains open until a customer reply, successful diagnostic, or agent resolution event supports closure. Store the original message separately from the draft, along with the exact proposal approved and the approver’s identity.
Agent Patterns, Event Ordering, and Scaling
| Pattern | Trigger | State | Scaling | Topic layout | Main risk |
|---|---|---|---|---|---|
| Reactive agent | Business event | Ephemeral | Partition-based workers | Input → proposal → outcome | Duplicate actions; weak context |
| Enrichment worker | Required lookup | Derived context | Partitions and dependency capacity | Raw → enriched | Stale data; slow lookups |
| Worker pool | Independent task | Durable task status | Shared consumer group | Work, retry, dead-letter | Uneven duration; retry storms |
| Multi-agent pipeline | Previous stage’s output | Workflow state and correlation IDs | Each stage scales separately | Stage-specific topics | Partial completion; schema drift |
| Feedback loop | Outcome or human feedback | Decision and evaluation history | Separate online work from evaluation | Outcome, feedback, evaluation | Self-reinforcing errors |
Partition choice controls both ordering and scale. Use a stable partition key, such as case_id or customer_id, for the entity whose events must stay in sequence. Check entity versions before applying changes.
Ordered input does not guarantee ordered completion. Asynchronous model calls can finish out of sequence. When side effects must happen in order, serialize work for each entity.
Within a consumer group, each partition has one active consumer. An eight-partition topic therefore supports eight active consumers. Separate groups let agents, audit, and reporting consume independently. Set concurrency limits based on model latency and downstream capacity - not event volume alone.
Managing Workflow State, Timers, and Approvals
Kafka stores event history; the application must store durable workflow state. Record the current stage, completed steps, pending approval, retry count, deadline, and last observed event version.
Timers can trigger escalations when approval expires or send reminders after 24 hours. Before acting, check the current state. Approval records must include the reviewer, exact proposal, scope, timestamp, and expiration.
If a CRM update succeeds but notification delivery fails, keep the completed update and retry only the notification - or assign manual correction. Carry a correlation_id through the workflow and use a causation_id to identify the event that triggered each step.
Offset commits are not business completion. Commit only after processing state or the next action request has been durably handled. Track external execution separately.
Give actions stable idempotency keys so duplicate requests can return an existing result. If a result is uncertain, reconcile it with the external system before retrying. Exactly-once processing does not cover external side effects, so close the workflow only after external confirmation.
For these controls to work reliably, define event schemas, topic design, event contracts, and delivery guarantees explicitly.
Defining Event Contracts and Processing Guarantees
With workflow state defined, set clear event contracts and delivery rules so replays stay safe.
Event Fields, Schemas, and Versioning
Give every event a consistent envelope and a domain-specific payload. Define event types for support cases, orders, payments, risk signals, telemetry, inference requests, and agent outcomes.
Examples include a
support.case.openedpayload with case priority and customer reference, aninference.requestedpayload with model, input reference, and requested deadline, and anagent.action.completedpayload with action ID, decision, result, and policy evaluation.
Publish tokenized references, not raw personal or payment data.
| Field | Type | Required? | Contract rule |
|---|---|---|---|
event_id |
String/UUID | Yes | Globally unique; used for tracing and deduplication |
event_type |
String | Yes | Stable domain meaning |
event_time |
RFC 3339 timestamp | Yes | Time the business event occurred; store in UTC with an explicit offset |
source |
String | Yes | Producing service or system |
entity_id |
String | Yes | Business entity identifier |
correlation_id |
String/UUID | Strongly recommended | Links related business events |
schema_version |
String/integer | Yes | Identifies the payload contract |
trace_id |
String | Optional | Connects distributed tracing records |
workflow_id |
String | Optional | Identifies the workflow instance |
payload |
Structured object | Yes | Domain-specific data validated against the schema |
Keep processing_time separate from event_time. Business windows use the time an event occurred; latency monitoring uses processing time. Require explicit timestamp offsets and normalize stored times to UTC. Validate events at publish time and again before execution.
Use a governed schema registry to enforce backward compatibility for old data, forward compatibility for new data, and full compatibility when consumers must read old and new records in both directions. Favor additive changes, such as an optional field with a safe default. Use transitive compatibility checks when all retained versions must remain readable. Changing a field’s meaning or unit calls for a migration, a new event type, or a parallel topic - not just a version bump.
| Format | Strengths | Limits | Best fit |
|---|---|---|---|
| JSON with JSON Schema | Human-readable and easy to inspect | Larger messages; validation must be enforced | Operational events, external integrations, and prototypes |
| Avro | Compact binary encoding with explicit schemas and evolution rules | Requires serialization tooling and schema-aware debugging | Kafka-centric data platforms and high-volume internal events |
| Protocol Buffers | Compact, efficient, strongly typed, with generated language bindings | Less human-readable; requires careful field-number management | High-throughput services and polyglot agent workers |
Topic Design, Retention, and Safe Replay
Once schemas are stable, put facts, commands, results, and errors in separate streams.
Use domain-oriented topics such as orders.events, payments.events, and risk.signals, alongside Separate topics for agent.commands, agent.results, and agent.errors as part of a robust AI automation strategy. Commands request work. Events record facts. Results report outcomes, while errors preserve failure context. Keep agent reasoning, execution requests, and audit outcomes separate so replay can reconstruct history without reissuing side effects.
Set retention based on recovery needs, privacy obligations, and storage cost. Compaction supports latest-state views, not complete workflow history. Keep immutable audit events on a retention-controlled topic.
Dead-letter records should retain the original topic, partition, offset, event ID, schema version, error category, diagnostic reference, and retry count.
For replay, select a bounded time or offset range, then use isolated topics with dry-run processing. Disable external actions by default and enable them only with approval. Carry replay_id, original_event_id, and replay_mode so replayed records remain separate from new business activity.
Delivery Guarantees and Idempotent Actions
Topic design needs an explicit processing boundary: what does the system guarantee, and where does that guarantee stop?
| Mode | Loss risk | Duplicate risk | Boundary and tradeoff |
|---|---|---|---|
| At-most-once | Possible if offsets are committed before processing | Low from redelivery | Simple offset handling, but work can be lost |
| At-least-once | Reduced with durable processing | Possible after crashes | Recoverable processing; actions need deduplication |
| Exactly-once | Avoided within a correctly configured Kafka transaction | Duplicate Kafka results avoided within that boundary | Coordinates Kafka writes and consumed offsets, not external effects; adds complexity |
Kafka producer idempotence requires acks=all, retries greater than zero, and max.in.flight.requests.per.connection no greater than 5. Kafka can deduplicate its own writes, not external side effects. Its guarantees stop at message delivery, Kafka writes, and offset handling. The action service still owns idempotency and reconciliation.
For actions, enforce a unique business-action key through conditional writes in a durable status store. Retain the request hash, input event ID, timestamps, and external reference. Use an inbox/outbox pattern when a database write and Kafka publish must succeed together. Treat timeouts as unknown, reconcile through the provider’s status endpoint or webhook, and reuse the original idempotency key.
Building and Running Agent Workers
Once contracts, guarantees, and workflow state are defined, the worker layer can process events safely.
Prerequisites and Development Setup
Assign a workflow owner. For organizations scaling these systems, AI consulting services can help design the underlying agentic foundation. Set event volume, latency targets, model quotas, timeout limits, and acceptable action risk. Provision the demo topics, separate dev/staging/production identities, and the worker’s model interface. Keep secrets outside source code.
Install confluent-kafka, create the demonstration topics below, and set KAFKA_BOOTSTRAP_SERVERS to your development broker. Use the event IDs, schema versions, and idempotency keys defined earlier.
Implementing Producers and Consumers in Python
Demonstration only: this worker reads one support case, classifies it with a replaceable stub, and publishes a proposed priority.
The example seeds a single demo event, then consumes it to show the full event-to-decision path.
import json
import os
from confluent_kafka import Consumer, Producer
BROKER = os.environ["KAFKA_BOOTSTRAP_SERVERS"]
INPUT_TOPIC = "support.cases.v1"
OUTPUT_TOPIC = "support.decisions.v1"
GROUP_ID = "support-agent-demo"
def validate_case(case):
required = {"case_id", "customer_id", "created_at", "text"}
missing = required - case.keys()
if missing:
raise ValueError(f"Missing fields: {sorted(missing)}")
if not isinstance(case["text"], str) or not case["text"].strip():
raise ValueError("text must be non-empty")
return case
def classify_case(case):
# Replace this deterministic stub with a model adapter.
text = case["text"].lower()
priority = "high" if "outage" in text or "fraud" in text else "normal"
return {
"case_id": case["case_id"],
"decision": {"priority": priority},
"decision_status": "proposed",
"model_version": "demo-rule-v1",
}
producer = Producer({
"bootstrap.servers": BROKER,
"client.id": "support-agent-producer",
})
consumer = Consumer({
"bootstrap.servers": BROKER,
"group.id": GROUP_ID,
"enable.auto.commit": False,
"auto.offset.reset": "earliest",
})
case = validate_case({
"case_id": "case-1001",
"customer_id": "customer-42",
"created_at": "2026-10-02T14:30:00Z",
"text": "Our service is experiencing an outage.",
})
producer.produce(
INPUT_TOPIC,
key=case["case_id"],
value=json.dumps(case).encode("utf-8"),
)
producer.flush()
consumer.subscribe([INPUT_TOPIC])
try:
while True:
message = consumer.poll(1.0)
if message is None:
continue
if message.error():
raise RuntimeError(message.error())
received = validate_case(json.loads(message.value()))
decision = classify_case(received)
producer.produce(
OUTPUT_TOPIC,
key=received["case_id"],
value=json.dumps(decision).encode("utf-8"),
)
producer.flush()
# Commit only after validation, inference, and output publication.
# For asynchronous work, commit only the next offset after the
# highest contiguous completed record in each partition.
# Offset 12 does not advance the commit if offset 11 is incomplete.
consumer.commit(message=message, asynchronous=False)
break
finally:
consumer.close()
Handling Failures and Protecting Actions
Keep polling responsive while inference runs in a bounded worker pool. When queues fill, pause partitions but continue polling. Resume them when capacity returns. On revocation, fence unfinished tasks and commit only completed work.
For temporary failures, use bounded retries with backoff and jitter. Isolate malformed records, and advance past them only after confirmed dead-letter publication. Require approval for high-impact actions. If dependencies fail, send uncertain proposals to review rather than executing them.
Testing and Staged Rollout
With failure handling in place, test rollout behavior under actual outage conditions.
Test schema compatibility, entity ordering, duplicates, replay, crashes, rebalances, model timeouts, tool failures, and peak load. Use mocked tools or sandbox services for destructive operations. After each failure, check external state - not just offsets.
Start in shadow mode, then add human approval and narrowly scoped permissions. Before increasing traffic, name the rollback owner, define safety and latency triggers, and test recovery procedures. Keep a kill switch that stops actions without dropping incoming events.
Monitoring Performance and AI Quality
Once the worker is stable, track both service health and decision quality.
Base alerts on documented service objectives, not generic thresholds. Evaluate normal, ambiguous, adversarial, and high-impact cases. Compare model, prompt, policy, and retrieval changes before releasing them to more traffic.
Set escalation thresholds, and review retrieval failures separately from model failures. Connect these signals to resolution time, manual effort, and measured cost per completed workflow in U.S. dollars.
| Area | Production-readiness checks | What operators must establish |
|---|---|---|
| Reliability | Lag, throughput, queue depth, rebalances, retries, dead-letter volume, commit failures | Processing keeps up without skipping unfinished records |
| Observability | End-to-end latency, poll-loop health, model and tool latency, correlation IDs, traces | A delayed workflow can be traced to its failing stage |
| AI quality | Output validity, decision quality, drift, retrieval failures, policy violations, human overrides | Decisions remain useful; unsafe or uncertain cases escalate |
| Recovery | Replay drills, idempotency conflicts, tool success, confirmed action completion, rollback and kill-switch tests | Recovery does not repeat actions; publication is not mistaken for completion |
Securing, Governing, and Deploying Agent Systems
Once agents can act safely, secure the access, data, and recovery paths that support their actions.
Securing Kafka Access and Agent Actions
Secure both Kafka access and external actions. Configure TLS for client and broker traffic, use SASL or mutual TLS authentication, and apply default-deny authorization. Give producers, consumers, and workers separate identities, with narrowly scoped permissions for topics and consumer groups. Restrict admin access, rotate credentials, and log access.
Treat events, retrieved documents, and model output as untrusted data. Enforce tool allowlists, argument validation, spending limits, and approval for sensitive actions outside the model. Keep credentials in a secrets manager - not in prompts or payloads.
Data Governance and Audit Records
With access controls in place, define what data the system may store, replay, and audit.
Before deployment, classify event fields, headers, topics, documents, prompts, model outputs, tool arguments, and audit records. Send models only the data they need, masking or redacting identifiers where possible. Document approved model-provider terms, training restrictions, retention rules, and data residency.
The deletion plan must cover topics, dead-letter topics, consumer offsets, indexes, logs, replicas, snapshots, and backups. Account for legal holds and retention expiry, too. Require approval for replay access, set a time limit, and log its use.
| Area | Required controls | Accountable owner | Evidence |
|---|---|---|---|
| Event data | Field minimization, classification, masking, retention and deletion rules | Data owner; privacy lead reviews | Data inventory and schema review |
| Kafka access | TLS, authentication, topic/group authorization, access reviews | Kafka platform owner | ACL exports and access-review records |
| Agent permissions | Tool allowlists, argument checks, bounded scopes | Agent owner | Permission manifest and security tests |
| Model operations | Approved models, provider terms, regional routing | AI product owner; privacy/legal reviews | Provider assessment and routing configuration |
| Decisions and prompts | Versioned prompt/model identifiers, decision reason codes, correlation IDs | Agent owner | Decision trace with redacted inputs |
| Business actions | Approval thresholds, spending limits, idempotency keys | Business-process owner | Approval record and confirmed action result |
| Audit records | Restricted access, tamper protection, redaction, deletion schedule | Compliance owner | Integrity checks and retention reports |
Link each event ID to its redacted inputs, schema version, model and prompt versions, policy version, decision, tool call, approval, and outcome. Use document references or hashes instead of copying content.
Keep operational telemetry separate from audit evidence. Engineers need error categories and latency data, not unrestricted customer records. Use these data-handling rules to plan deployment capacity and recovery controls.
Deployment Capacity and Disaster Recovery
Keep development, staging, and production isolated with separate clusters, accounts, credentials, topic namespaces, model keys, and data stores. Use production PII for testing only when approved and suitably masked.
Size partitions and worker concurrency for peak demand - not average load. Account for Kafka throughput, model quotas, tool latency, retry volume, and approval capacity. A consumer group cannot have more active assignments than partitions.
Define RPO for acceptable event/state loss and RTO for processing recovery. Test isolated restoration of schemas, offsets, ACLs, agent configuration, prompt and policy versions, approval state, and downstream safeguards. Replication alone does not prove that recovery works.
Coordinate schema and agent releases, retain rollback-compatible artifacts, and use compensating workflows when completed actions cannot be reversed. Assign named owners to these controls before production.
Implementation Roles and Consulting Support
Assign owners for schemas, Kafka controls, agent behavior, security review, privacy and compliance, approvals, operations, and incident response. The operations team should own capacity, deployment, recovery drills, and on-call support throughout the system’s life cycle.
Organizations that need implementation help should define what an outside partner will deliver while keeping internal ownership of data, approvals, and risk acceptance.
NAITIVE AI Consulting Agency can support AI consulting, autonomous-agent development, AI automation, and business-process automation for this type of implementation. Before engagement, verify Kafka-specific experience, delivery scope, measurable goals, source-code ownership, support expectations, and handoff requirements.
Conclusion: Building Reliable Event-Driven AI Agents
Kafka provides the event backbone; the agent decides, and a controlled service acts. Kafka preserves events for replay and recovery. The application still handles reasoning, persistent workflow state, policy checks, tool execution, and safe actions. This division works only when contracts and action controls are explicit.
Reliability rests on stable contracts, partition-safe keys, and idempotent actions that withstand replay and protect business outcomes. Kafka’s exactly-once processing does not make external actions exactly once. Keep historical replay separate from live actions. Measure decision quality alongside processing latency and cost.
After replay and side effects are under control, let measured results guide expansion. Start with one workflow that needs streaming - not one that merely can use it. Before building, define event volume, end-to-end latency, retention, cost, and a measurable baseline. Choose a simpler design when a scheduled job or direct API call is enough. Move from shadow evaluation to recommendations, then to low-risk execution. Add autonomy only when monitoring, security, human oversight, and recovery can support it.
FAQs
How do I know my workflow needs Kafka?
Your workflow needs Apache Kafka when it calls for large-scale parallel processing and reliable, real-time event streaming. That includes AI agents that depend on continuous data flows with sub-second latency, exactly-once processing, or durable audit trails for replaying historical data.
Kafka also fits high-throughput workflows that keep components separate. This lets services change independently, stay resilient, and handle traffic spikes through backpressure management.
What if an external service lacks idempotency support?
Build idempotency into your processing logic to keep results consistent. Give each event its own ID and track processed IDs to filter out duplicates. For external actions, use idempotency keys to prevent unintended side effects.
Use transactional messaging so event publishing and data updates succeed or fail together. Retry failed attempts with exponential backoff. Send persistent errors to a dead-letter queue for manual review.
When is my agent ready to act without approval?
Your agent is ready once it has passed thorough testing in a shadow environment, where it processes real-time inputs and produces outputs without changing live data. Before giving it full write access, use read-only middleware to protect data integrity and check performance. Set up rollback plans and train human oversight teams to handle high-risk actions.
NAITIVE AI Consulting Agency designs and manages autonomous AI agents and business process automation solutions to help integrate them into your business.