Event-Driven AI Agents with Apache Kafka

Build event-driven AI agents with Kafka: schemas, idempotent actions, approval workflows, retries, and safe replay.

Event-Driven AI Agents with Apache Kafka

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

Designing the AI Agent Architecture

Kafka AI Agents: From Events to Controlled Actions

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 in support.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 in support.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.opened payload with case priority and customer reference, an inference.requested payload with model, input reference, and requested deadline, and an agent.action.completed payload 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.

Related Blog Posts