Most agent demos run behind an HTTP endpoint: a request arrives, the agent thinks, the response goes back on the same connection. That breaks down when work arrives from many systems at once, when a single turn can take a minute of model and tool calls, or when you need to replay last Tuesday's traffic against a new prompt. Apache Kafka fits those cases well. It gives you a durable, partitioned log of requests, consumer groups that spread load across workers, and per-key ordering so one conversation's messages are handled in sequence.

This article shows how to put Google's Agent Development Kit for Java behind Kafka properly. It covers the topic and key design, a worker loop that survives minute-long agent turns without being evicted from its group, where ADK sessions must live when partitions move between workers, which delivery guarantees Kafka actually gives you and which it cannot, and how to handle poison messages, rate limits and replays. Code uses the Kafka Java client and the ADK Java Runner API as documented at the time of writing; check both against the versions you deploy.

Why Kafka, and what it actually guarantees

Kafka is a log, not a queue. Producers append records to a topic, which is split into partitions; each record has a key, and records with the same key always land in the same partition, in order. A consumer group shares the partitions among its members, so each partition is read by exactly one worker at a time, and each worker tracks its progress as a committed offset per partition. If a worker dies, its partitions are reassigned and the new owner resumes from the last committed offset.

Mapped onto agents, that gives three useful properties. Load spreads across workers up to the partition count. Ordering holds per key, so if you key by conversation, a user's second message is never processed before their first. And because the log is retained, you can rewind a consumer group to re-run traffic, which is how you evaluate a prompt change against real inputs. The cost is that you now own delivery semantics, and an LLM call is the slowest, least deterministic thing that has ever sat inside a Kafka consumer.

Reference architecture

Kafka-driven ADK Java agent workersProducersweb, email, jobsagent.requestskey = conversationIdWorker groupconsumer + ADK Runneragent.resultskey preservedpollsendDedup storerequestId seen?Session serviceshared, durableOutbox tabletool side effectsagent.requests.dlqpoison messagesafter N triesCommit offsets only after the result is acknowledged; everything outside Kafka must be idempotent.
Requests are keyed by conversation, consumed by a group of ADK workers and answered on a results topic. Dedup, sessions and side effects live outside Kafka.

The topology has two topics and one dead-letter topic. agent.requests carries a JSON envelope with a unique requestId, the userId, a conversationId used as the record key, the user text and a reply-to hint. agent.results carries the final answer under the same key so downstream consumers can stitch it back to the conversation. agent.requests.dlq receives requests that failed a fixed number of times, with the error attached as headers.

Partition count bounds parallelism, because one partition is processed by one worker at a time and, in the design below, sequentially within that worker. If a turn averages 20 seconds and you need 30 turns per second, you need at least 600 partitions in flight, either as 600 partitions or as fewer partitions with per-key parallelism inside a worker. Kafka does not let you reduce a topic's partition count, and adding partitions changes which partition a key maps to, so size the topic for growth up front.

Sessions must outlive partition ownership

ADK keeps conversation history in a session, addressed by app name, user id and session id, and the runner appends every event of a turn to it. The quick-start uses InMemoryRunner, whose session service lives in the worker's heap. Behind Kafka that is a bug waiting for the first rebalance: the conversation's partition moves to another worker, which has never seen the session, and the agent silently starts the conversation from scratch. Partition affinity does not give you session affinity.

So the session service must be shared and durable: an implementation of ADK's BaseSessionService backed by a database every worker can reach. Use the conversation id as the session id, so a worker can find or create the session from the record alone. The service's getSession returns an RxJava Maybe<Session> and createSession a Single<Session>, which compose neatly:

import com.google.adk.events.Event;
import com.google.adk.runner.Runner;
import com.google.adk.sessions.BaseSessionService;
import com.google.adk.sessions.Session;
import com.google.genai.types.Content;
import com.google.genai.types.Part;
import io.reactivex.rxjava3.core.Flowable;
import io.reactivex.rxjava3.core.Single;
import java.util.Optional;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;

final class AgentTurn {
    private final Runner runner;            // built with a shared, durable session service
    private final BaseSessionService sessions;

    AgentTurn(Runner runner) {
        this.runner = runner;
        this.sessions = runner.sessionService();
    }

    String run(AgentRequest req) {
        Session s = sessions
            .getSession(runner.appName(), req.userId(), req.conversationId(), Optional.empty())
            .switchIfEmpty(Single.defer(() -> sessions.createSession(   // lazy: only if missing
                runner.appName(), req.userId(), new ConcurrentHashMap<>(), req.conversationId())))
            .blockingGet();

        Content msg = Content.fromParts(Part.fromText(req.text()));
        StringBuilder answer = new StringBuilder();
        runner.runAsync(s.userId(), s.id(), msg)
            .takeUntil(Flowable.timer(90, TimeUnit.SECONDS))   // whole-turn deadline
            .blockingForEach(ev -> {
                if (ev.finalResponse()) {
                    answer.append(ev.stringifyContent());
                }
            });
        if (answer.length() == 0) {
            throw new TurnTimeoutException(req.requestId());  // retried, then DLQ
        }
        return answer.toString();
    }
}

Three details matter. The deadline bounds the whole turn, which the consumer loop depends on; RxJava's timeout(long, TimeUnit) would not, because it measures the gap between events, so a turn that emits something every minute could run forever. The create call is wrapped in Single.defer so it only runs when the session is missing. And there is a race between two workers creating the same session during a rebalance overlap; the store should treat session id as unique and the losing create should fall back to a read.

A consumer loop that survives slow turns

The classic consumer loop calls poll(), processes the batch and commits. With agents that loop gets your worker thrown out of the group. Kafka expects poll() to be called at least every max.poll.interval.ms (default 300,000 ms); with the default max.poll.records of 500 and 20-second turns, one batch takes nearly three hours. The group coordinator assumes the worker is dead, revokes its partitions and hands them to someone else, who re-processes the same records.

The robust pattern is to decouple polling from processing. The polling thread hands each partition's records to an executor, pauses that partition so further polls return nothing for it, keeps polling to stay alive, and commits and resumes the partition when its work completes.

Map<TopicPartition, Future<Long>> inFlight = new HashMap<>();

while (running) {
    ConsumerRecords<String, String> batch = consumer.poll(Duration.ofMillis(200));
    for (TopicPartition tp : batch.partitions()) {
        List<ConsumerRecord<String, String>> recs = batch.records(tp);
        consumer.pause(List.of(tp));                    // no more records until done
        inFlight.put(tp, pool.submit(() -> processInOrder(recs)));  // returns last offset
    }
    Map<TopicPartition, OffsetAndMetadata> done = new HashMap<>();
    for (Iterator<Map.Entry<TopicPartition, Future<Long>>> it = inFlight.entrySet().iterator(); it.hasNext(); ) {
        Map.Entry<TopicPartition, Future<Long>> e = it.next();
        if (e.getValue().isDone()) {
            done.put(e.getKey(), new OffsetAndMetadata(e.getValue().get() + 1));
            it.remove();
        }
    }
    if (!done.isEmpty()) {
        consumer.commitSync(done);                      // commit the NEXT offset to read
        consumer.resume(done.keySet());
    }
}

Records within a partition are processed in order by one task, so per-conversation ordering survives. The committed offset is the last processed offset plus one, because Kafka's offset means the next record to read. A rebalance listener completes the picture: in onPartitionsRevoked, wait briefly for in-flight work on revoked partitions, commit what finished and abandon the rest, which the new owner will redo. Use the CooperativeStickyAssignor so a rebalance only moves the partitions that need to move instead of stopping every worker. Disable auto-commit, which would otherwise commit offsets for records you have not finished.

Delivery semantics and duplicate turns

With the loop above, delivery is at least once: a crash between producing the result and committing the offset re-runs the turn. Kafka's exactly-once support does not make that go away. A transactional producer can send the result record and the consumer's offsets atomically with sendOffsetsToTransaction(offsets, consumer.groupMetadata()), and downstream consumers reading with isolation.level=read_committed will never see a result whose offset was not committed. But that guarantee covers only Kafka-to-Kafka. The model call, the session writes and every tool that charges a card or sends an email sit outside the transaction and will run twice on a retry.

So design for duplicates explicitly:

  • Dedup on request id. Before running a turn, insert the requestId into a table with a unique constraint. If it already exists with a stored result, republish that result and skip the model call.
  • Idempotent tools. Pass a key derived from the request id and tool call into every side-effecting API, or write the intent to a transactional outbox and let a relay deliver it; see the outbox pattern for ADK Java.
  • Idempotent producer. Keep enable.idempotence=true (the default in current clients) with acks=all so producer retries do not duplicate results within a partition.
  • Idempotent readers. Consumers of agent.results should key on request id and ignore repeats, which also makes replays safe.

Worked example: a support inbox

A support platform routes email, chat and web-form tickets into agent.requests with 48 partitions and eight workers, each processing up to six partitions concurrently. A customer sends two emails a minute apart. Both carry conversation id c-7731, hash to partition 17 and are processed in order by worker 3, which loads the session, runs a refund-lookup tool and publishes an answer.

Midway through the second turn, worker 3 is redeployed. Its partitions are revoked, the in-flight task is abandoned, and partition 17 moves to worker 5. Worker 5 resumes from the committed offset, which still points at the second email, and re-runs it. The dedup table shows the request id as started but not finished, so the turn runs. The refund tool call carries an idempotency key, so the payment service returns the first attempt's outcome rather than issuing a second refund. The shared session store already holds the first turn, so the agent answers in context. The customer sees one reply.

Now suppose an email contains a 6 MB pasted log that makes the model call fail every time. Without a limit that record blocks partition 17 forever, and every conversation hashed there stalls behind it. The worker counts attempts in the dedup row; after three, it publishes the record to agent.requests.dlq with the error in headers, commits past it and moves on. An operator triages the DLQ and can re-publish fixed records.

Operating the worker group

Setting or metricRecommendationWhy
enable.auto.commitfalsecommit only finished work
max.poll.recordssmall, for example 10 to 50bounds work handed out per poll
max.poll.interval.mskeep default, keep polling while pausedliveness, not turn length
partition.assignment.strategyCooperativeStickyAssignorincremental rebalances
Consumer lag per partitionalert on growth, not absolute valueshows a stuck or slow key
Turn latency p50 and p99track per agent versionsizes partitions and timeouts
DLQ ratealert on any sustained ratepoison input or a broken tool
Model 429 responsesback off by pausing all partitionskeeps retries off the hot path

Rate limits deserve their own rule. When the model provider returns throttling errors, do not retry inside the task with tight loops; pause every assigned partition, keep polling, and resume after a backoff. Kafka is doing exactly what it is good at here: absorbing the backlog durably until capacity returns. Timeouts belong at two levels, per tool call as in ADK Java tool timeout handling and per turn as in the runner code above. Emit the request id, conversation id, partition and offset on every log line and trace span so a single turn can be followed across retries, and map ADK events to telemetry as described in the ADK Java events article.

Failure modes

  • Rebalance storms. Processing inside the poll loop exceeds the poll interval; workers are evicted, rejoin and evict each other. Fix with pause and resume, smaller batches and cooperative assignment.
  • Amnesiac agents. In-memory sessions lose context when partitions move. Fix with a shared session service keyed by conversation id.
  • Head-of-line blocking. One poison or slow record stalls its partition. Fix with a per-turn timeout, an attempt counter and a DLQ.
  • Double side effects. Retries re-run tools outside Kafka's transaction. Fix with dedup, idempotency keys and an outbox.
  • Hot partitions. A single very active key, such as a bulk import keyed by one tenant, pins one worker. Key by conversation rather than tenant, or split bulk work onto its own topic.
  • Replays with side effects. Rewinding a group to evaluate a new prompt re-executes real tools. Replay into a separate consumer group with tools stubbed or pointed at a sandbox.

Trade-offs

Kafka is the right front door when traffic is bursty, work is asynchronous and ordering per entity matters. It is the wrong one for interactive chat that needs token streaming to a browser; there, serve the turn over HTTP or server-sent events and publish the finished turn to Kafka for audit and analytics. A simpler queue such as a managed task queue gives per-message acknowledgement without partition head-of-line blocking, at the cost of weaker ordering and no cheap replay. Sequential processing per partition is the simplest way to keep order; per-key parallelism inside a partition raises throughput but forces you to track offsets per key and commit only the lowest completed offset, which is real complexity. Start sequential, measure lag, and add partitions before adding cleverness. For how sessions and context grow across long-running conversations, see session context depth in the ADK Java runtime.

What to do next

  • Define the request and result envelopes with a unique request id and the conversation id as the record key.
  • Replace InMemoryRunner with a runner backed by a shared, durable session service, and test a forced rebalance mid-conversation.
  • Implement the pause, process, commit, resume loop with auto-commit off and a rebalance listener.
  • Add a dedup table, idempotency keys on every side-effecting tool, and a DLQ after a fixed attempt count.
  • Set a per-turn timeout below your worst acceptable lag, and alert on consumer lag growth and DLQ rate.
  • Before your first prompt change, practise a replay into a separate consumer group with sandboxed tools.
Key takeaway: Put ADK Java behind Kafka by keying requests by conversation, decoupling polling from minute-long turns with pause and resume, storing sessions in a shared durable service, and treating every turn as at-least-once: dedup on request id and make every side effect idempotent, because Kafka's transactions stop at the edge of Kafka.