Pub/Sub and agents fit together naturally on paper. Work arrives as messages, a pool of workers picks them up, and each message becomes one agent turn. In practice the two have opposite temperaments. Pub/Sub delivers at least once, expects fast acknowledgements and redelivers anything it is unsure about. An ADK agent turn is slow, varies from two seconds to two minutes, costs money on every attempt and may call tools that change the world: issuing a refund, sending an email, opening a ticket. Run a refund agent naively behind a subscription and a redelivered message refunds the customer twice.

This article is about the consumer side: an ADK Java worker that pulls from a subscription, runs an LlmAgent per message and publishes the result. It covers the message contract, the subscriber configuration that turns flow control into a concurrency cap on model calls, deduplication before any side effect, mapping messages onto ADK sessions, the split between retrying and dead-lettering, and a worked sizing example. For the producer side, where an agent publishes events reliably alongside a state change, see the transactional outbox with ADK Java.

Architecture

A Pub/Sub-driven ADK Java worker: pull, dedup, run the agent, publish, then ackProducerstickets, webhooksagent-requeststopicpublishagent-worker subretry policy, DLQdead-letter topicafter N attemptsforwardWorker process (Cloud Run / GKE)Subscriberflow control = capstreaming pullDedup storeclaim, completeADK RunnerLlmAgent + toolsSession servicekeyed by conversationagent-resultstopicpublishack / nackonly after publishAck only after the result is durable. A crash before the ack means redelivery, so every step before it must be safe to repeat.
Figure 1. Producers publish requests; the worker's subscriber pulls under flow control, checks a dedup store, runs the agent, publishes the result and only then acks. Messages that keep failing are forwarded to a dead-letter topic.

The order inside the worker is the whole design: claim the message in the dedup store, run the agent, publish the result, record completion, ack. Anything can crash between any two steps, and Pub/Sub will then redeliver. The question for every step is therefore what happens if it runs twice, and the answers below make each step either idempotent or guarded.

The message contract

Put the payload in the message data as versioned JSON and the routing facts in attributes, where the worker can read them without parsing the body. A minimal contract:

attributes:
  schema          = "support-triage/v2"
  idempotency_key = "ticket-48113:rev-3"      # business key, set by the producer
  conversation_id = "ticket-48113"            # maps to the ADK session
  user_id         = "cust-9921"
data (JSON):
  { "text": "My order arrived damaged, can I get a refund?",
    "order_id": "A-77120", "channel": "email" }

Why a business idempotency key instead of the Pub/Sub message ID? The message ID identifies one publish. If the producer itself retries a publish after a timeout, the same request arrives with two different message IDs, and dedup on message ID lets both through. A key derived from the business event (ticket and revision) survives producer retries as well as redelivery.

Use conversation_id as the ordering key if turns in one conversation must be processed in order. With ordering enabled on the subscription, messages that share an ordering key, published in the same region, are received in publish order. The cost is head-of-line blocking: if one message for a key keeps failing, later messages for that key wait behind it, which is one more reason the dead-letter path below matters.

Configuring the subscriber

The Java client's Subscriber uses streaming pull and calls your MessageReceiver on its executor threads. Two settings shape everything. Flow control limits how many messages the client holds unacknowledged at once; since each outstanding message is an in-flight agent turn, that number is your concurrency cap on model calls. And the maximum ack extension period bounds how long the client keeps extending a message's lease while your handler works on it; it must exceed your slowest acceptable turn.

ProjectSubscriptionName sub = ProjectSubscriptionName.of(PROJECT, "agent-worker");

Subscriber subscriber = Subscriber.newBuilder(sub, handler::receive)
    .setFlowControlSettings(FlowControlSettings.newBuilder()
        .setMaxOutstandingElementCount(48L)        // = max concurrent agent turns per instance
        .setMaxOutstandingRequestBytes(16L * 1024 * 1024)
        .build())
    .setMaxAckExtensionPeriod(Duration.ofMinutes(10))   // longest turn we will wait for
    .setParallelPullCount(1)
    .build();

subscriber.startAsync().awaitRunning();
Runtime.getRuntime().addShutdownHook(new Thread(() ->
    subscriber.stopAsync().awaitTerminated()));    // stop pulling, let in-flight turns finish

The subscription's own ack deadline (10 to 600 seconds, default 10) matters less than you might expect, because the client extends leases automatically while the handler runs. It matters a great deal if the process freezes: a long ack deadline then delays redelivery to a healthy worker. Leave it short and rely on the client's extension.

Flow control is per subscriber instance, so the fleet-wide cap is the per-instance number times the instance count. If the model quota is the real constraint, size from the quota down, and put a shared limiter in front of the model as well; rate limiting in ADK Java covers token-bucket designs that work across instances.

The handler: claim, run, publish, ack

The handler claims the key, runs the agent, publishes, records completion and acks. The ADK calls here are the documented Java API: LlmAgent.builder(), InMemoryRunner (swap in a persistent session service for production), runAsync(userId, sessionId, content, runConfig) returning a Flowable<Event>, and event.finalResponse() to pick out the answer.

void receive(PubsubMessage msg, AckReplyConsumer reply) {
    String key = msg.getAttributesOrThrow("idempotency_key");
    int attempt = Optional.ofNullable(Subscriber.getDeliveryAttempt(msg)).orElse(0);

    Claim claim = dedup.claim(key, msg.getMessageId(), Duration.ofMinutes(15));
    switch (claim.state()) {
        case DONE -> { reply.ack(); return; }            // already processed: just ack
        case HELD -> {                                   // a turn for this key is running somewhere
            if (claim.messageId().equals(msg.getMessageId())) reply.nack();  // our own redelivery: retry later
            else reply.ack();                            // producer duplicate: the running turn covers it
            return;
        }
        case CLAIMED -> { }                              // ours, continue
    }
    try {
        Request req = Request.parse(msg);                // throws PoisonMessage on bad schema
        Session session = sessions.forConversation(req.conversationId(), req.userId(), key);
        Content input = Content.fromParts(Part.fromText(req.text()));

        String answer = runner.runAsync(req.userId(), session.id(), input, runConfig)
            .filter(Event::finalResponse)
            .map(Event::stringifyContent)
            .blockingLast("");                           // handler thread blocks; pool sized to flow control

        String resultId = results.publish(req.resultMessage(answer, key)).get();  // wait for the publish
        dedup.complete(key, resultId);
        reply.ack();
    } catch (PoisonMessage e) {
        dedup.fail(key, e.getMessage());
        deadLetter.publish(msg, e);                      // explicit, with the reason attached
        reply.ack();
    } catch (Exception e) {
        dedup.release(key);
        log.warn("transient failure attempt={} key={}", attempt, key, e);
        reply.nack();                                    // retry policy applies the backoff
    }
}

The dedup store needs a conditional write: Firestore, Spanner, Postgres or Redis with an expiring claim all work. A claim with a lease prevents two workers running the same turn concurrently after a redelivery; the completion record makes a later redelivery a cheap ack. Keep completion records at least as long as the subscription's message retention (default 7 days, up to 31), or a late redelivery will find no record.

Be careful with nacks on a held claim: every nack counts as a delivery attempt. A producer duplicate that arrives while the original turn runs would otherwise be nacked repeatedly and dead-lettered for work that succeeded, which is why the handler acks it. For a redelivery of the same message, make sure the retry backoff times the maximum attempts outlasts the claim lease.

Idempotent tools

Dedup at the message level does not make the turn safe on its own. Suppose the worker claims the key, the agent calls the refund tool, and the process dies before completion is recorded. The claim's lease expires, another worker reclaims the message and the model, which is not deterministic, may decide to call the refund tool again. The protection has to sit in the tool: pass the idempotency key through to every side-effecting call, and let the downstream system reject duplicates.

public static Map<String, Object> issueRefund(
        @Schema(name = "orderId", description = "Order to refund") String orderId,
        @Schema(name = "amountCents", description = "Amount in cents") long amountCents,
        @Schema(name = "toolContext") ToolContext toolContext) {  // injected by ADK, not shown to the model
    String key = (String) toolContext.state().get("idempotency_key");   // set when the session was created
    RefundResult r = payments.refund(orderId, amountCents, key + ":refund:" + orderId);
    return Map.of("status", r.status(), "refundId", r.id(), "duplicate", r.wasDuplicate());
}

Most payment and ticketing APIs accept an idempotency key for exactly this reason; for internal services, add one. Read-only tools need nothing. Tools that send email or messages are the awkward middle: keep a sent-log keyed by the same key and check it first. The broader pattern, including compensation when a step cannot be made idempotent, is in resuming ADK Java agents after a crash.

Sessions across workers

Each conversation maps to one ADK session. Look it up by conversation_id; create it on the first message, writing the idempotency key into session state so tools can read it. For later turns, update that state key before running. Two consequences follow. The session service must be persistent and shared across instances, because the next message for a conversation may land on any worker; the in-memory runner is for tests. And without an ordering key, two turns of one conversation can run at once on different workers and interleave their events in the session, so either enable ordering or take the dedup claim per conversation rather than per message.

Retry or dead-letter

Every failure must be sorted into one of two kinds, because Pub/Sub treats them the same way unless you sort them yourself.

FailureKindAction
Unparseable JSON, unknown schema versionPoisonDead-letter explicitly with the reason, ack
Model refused, safety block on inputPoison (for this input)Publish a 'needs human' result, ack
Model 429 or 5xx, tool timeoutTransientNack; backoff from the retry policy
Dedup store unavailableTransientNack; do not run the agent unguarded
Turn exceeded max ack extensionTransient, but suspiciousLog turn length; consider a step limit

For transient failures, configure the subscription's retry policy with exponential backoff; the backoff is capped at 600 seconds, and if you set only a maximum the minimum defaults to 10 seconds. Without a retry policy, a nacked message is redelivered almost immediately and a model outage becomes a tight loop that burns quota.

As a backstop, attach a dead-letter topic. Pub/Sub forwards a message there after a maximum number of delivery attempts, configurable from 5 to 100 with a default of 5, and the count is approximate because forwarding is best-effort. The Pub/Sub service account needs the publisher role on the dead-letter topic and the subscriber role on the source subscription, or forwarding silently does not happen. The handler reads the attempt count with Subscriber.getDeliveryAttempt(msg), which is populated only when dead lettering is configured. Replaying from the dead-letter topic is covered in dead-letter handling in ADK Java.

What exactly-once delivery does and does not give you

Pub/Sub offers exactly-once delivery as a subscription setting. It applies to pull subscriptions only, and the guarantee holds only when subscribers connect in the same region. It means an acknowledged message will not be redelivered, and the client can learn whether an ack actually succeeded. It does not make your handler run once: a crash before the ack still redelivers, and it says nothing about tool side effects. Treat it as a reduction in duplicate volume, useful when turns are expensive, not as a replacement for the dedup store.

Worked example: sizing a triage worker

A support-triage agent receives 30 tickets per second at peak. Measured turns take 12 seconds at the median and 60 seconds at p99, with about 2.5 model calls per turn. By Little's law, concurrency equals arrival rate times time in system: 30 times 12 is 360 turns in flight on average. Provision for the tail as well, so about 500. At 48 outstanding messages per instance that is 11 instances. Model calls run at 30 times 2.5, so 75 per second, which must fit inside the project's model quota with headroom for retries; if it does not, lower flow control and let the backlog absorb the peak, because Pub/Sub retains unacked messages for days.

Set the max ack extension period to 10 minutes, comfortably above the p99. Set dead-lettering at 5 attempts and a retry backoff of 10 to 600 seconds, so a ten-minute model outage costs a few retries per message rather than hundreds. Alert on subscription backlog age, not just count: old messages mean tickets nobody has answered.

What to do next

  1. Define the message contract with a schema version, a business idempotency key and a conversation ID.
  2. Set flow control from your model quota and measured turn length, and the max ack extension above your p99 turn.
  3. Add a dedup store with claim, complete and release, and ack only after the result is published.
  4. Thread the idempotency key into every side-effecting tool call.
  5. Replace the in-memory runner with a persistent, shared session service.
  6. Configure a retry policy with exponential backoff and a dead-letter topic, and grant the Pub/Sub service account both roles.
  7. Decide on ordering keys per conversation, and alert on backlog age and dead-letter volume.
  8. Kill a worker mid-turn in staging and confirm no tool side effect happens twice.
Key takeaway: Pub/Sub will deliver some messages more than once and an agent turn is slow, costly and side-effecting, so the worker must make repetition harmless: a business idempotency key, a dedup claim before the agent runs, idempotency keys passed into tools, and an ack only after the result is published. Size flow control from model quota and turn length, extend leases past the slowest turn, back off on transient errors and dead-letter poison messages with the IAM grants in place.