A ParallelAgent turns one user request into several model calls at once. That is the point: the user waits for the slowest branch instead of the sum of all of them. It is also how a working agent turns into a stream of 429 errors. Model providers enforce requests-per-minute and tokens-per-minute quotas per project, and a fan-out multiplies your request rate by its width at exactly the moment traffic is highest. Worse, the default ParallelAgent merges branch streams so that one failed branch cancels its siblings, so a single throttled call can throw away the work of every other branch.

This article shows how to keep a fan-out under its quota in the Agent Development Kit for Java: the traffic arithmetic, why the limit belongs in a model wrapper, a BaseLlm wrapper with permits and a token budget, retries that do not duplicate output, and sharing across requests and replicas. API details were read from the adk-java source on 2026-10-03; anything not confirmed there is flagged in the text. Fan-out design and failure containment have their own articles: the fan-out pattern and handling failures in parallel agents.

Where the limiter sits

User requests40 per minuteParallelAgent4 branches + synthesizerbranch callsRateLimitedLlm wrapperpermits, token budget, retryadmittedModel providerRPM and TPM quota429 or usageReconcileactual tokens, backoffShared limit storeper replica or centralbudget shareOne limiter per model and quota, shared by every branch of every request in the processPermits bound concurrency; the token budget bounds tokens per minute; both release in doFinally
Every branch of every request calls the same rate-limited model wrapper, which admits calls within the concurrency and token budget and reconciles actual usage.

From user requests to branch calls

Plan capacity in branch calls, not user requests. Suppose an agent fans out to four reviewers and then runs one synthesizer, and each LLM agent makes one model call per run. One user request costs five model calls. At 40 user requests per minute that is 200 calls per minute. If each call averages 5,000 input tokens and 800 output tokens, the agent uses about 1.16 million tokens per minute. A quota of 1,000 requests and 1 million tokens per minute is comfortable on the first axis and already exceeded on the second.

Concurrency follows from Little's law: calls in flight equal arrival rate times latency. Two hundred calls per minute is 3.3 per second; at 8 seconds per call about 27 calls are in flight on average. Averages hide bursts. A fan-out issues its four calls in the same millisecond, so 10 users who click at once produce 40 simultaneous calls, then a burst of synthesizer calls a few seconds later. Branches that call tools and loop add more calls per run, and RunConfig's maximum LLM calls per invocation is the built-in ceiling on that.

So you need three limits: concurrent calls, tokens per minute, and how much of both one request may take, shared by every branch of every request that uses the model.

Where to enforce the limit

There are three places you might enforce a limit, and only one fits. The scheduler is tempting: ParallelAgent.builder().scheduler(...) accepts any RxJava scheduler, and the default is Schedulers.io(). But a bounded scheduler bounds threads in one agent, not calls: a branch holds its thread while tools run, and separately built parallel agents do not share the bound. Use it as a bulkhead, not a rate limiter.

Model callbacks are the second candidate. But the after-model callback does not run when the call throws, so a 429 leaks a permit unless an error callback also releases it, and a cancelled branch may run neither. The signatures also move: on adk-java main BeforeModelCallback receives an LlmRequest.Builder, where earlier releases passed an LlmRequest. Code that depends on the exact shape breaks on upgrade.

The model itself is the right place. LlmAgent.builder().model(...) accepts a BaseLlm as well as a model name, and BaseLlm has two abstract methods: Flowable<LlmResponse> generateContent(LlmRequest, boolean stream) and BaseLlmConnection connect(LlmRequest). A wrapper that delegates to the real model sees every call, from every agent configured with it, and can tie the permit to the lifetime of the response stream. The general technique of writing your own model class is covered in implementing a custom LLM.

A rate-limited model wrapper

The wrapper holds a fair semaphore for concurrency and a token budget for throughput. It acquires both before subscribing to the real model, and releases the permit and reconciles tokens once in doFinally, which runs on completion, error or cancellation, so a branch cancelled by a failing sibling still returns its permit.

import com.google.adk.models.BaseLlm;
import com.google.adk.models.BaseLlmConnection;
import com.google.adk.models.LlmRegistry;
import com.google.adk.models.LlmRequest;
import com.google.adk.models.LlmResponse;
import io.reactivex.rxjava3.core.Flowable;
import java.time.Duration;
import java.util.concurrent.Semaphore;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicLong;

public final class RateLimitedLlm extends BaseLlm {
  private final BaseLlm delegate;
  private final Semaphore permits;        // concurrent calls
  private final TokenBudget budget;       // tokens per minute
  private final int estimatedTokens;      // reserved per call
  private final Duration maxWait;

  public RateLimitedLlm(String modelName, int maxConcurrent, TokenBudget budget,
                        int estimatedTokens, Duration maxWait) {
    super(modelName);
    this.delegate = LlmRegistry.getLlm(modelName);
    this.permits = new Semaphore(maxConcurrent, true);   // fair: FIFO waiters
    this.budget = budget;
    this.estimatedTokens = estimatedTokens;
    this.maxWait = maxWait;
  }

  @Override
  public Flowable<LlmResponse> generateContent(LlmRequest request, boolean stream) {
    return Flowable.defer(() -> {
      if (!permits.tryAcquire(maxWait.toMillis(), TimeUnit.MILLISECONDS)) {
        return Flowable.error(new RateLimitTimeout("no permit within " + maxWait));
      }
      if (!budget.reserve(estimatedTokens, maxWait)) {
        permits.release();
        return Flowable.error(new RateLimitTimeout("token budget exhausted"));
      }
      AtomicBoolean released = new AtomicBoolean();
      AtomicLong usedTokens = new AtomicLong(-1);      // latest reported total
      // defer makes the call cold, so a retry re-issues the request
      return withRetry(Flowable.defer(() -> delegate.generateContent(request, stream)), 3)
          .doOnNext(r -> Usage.totalTokens(r).ifPresent(usedTokens::set))
          .doFinally(() -> {                 // complete, error or cancel
            if (released.compareAndSet(false, true)) {
              long used = usedTokens.get();
              if (used >= 0) budget.reconcile(estimatedTokens, used);  // once per call
              permits.release();
            }
          });
    });
  }

  @Override
  public BaseLlmConnection connect(LlmRequest request) {
    throw new UnsupportedOperationException("live connections are not rate limited");
  }
}

Flowable.defer moves acquisition to subscription time, so a stream that is built but never run acquires nothing. Waiting is bounded by maxWait, after which the branch fails with your own RateLimitTimeout type, which a branch guard can record as a timeout instead of cancelling the fan-out. Blocking in tryAcquire is cheap on virtual threads and acceptable on the io scheduler, never on a computation scheduler. Live connections are refused explicitly rather than passed through unlimited.

Tokens are reserved as an estimate before the call and reconciled when usage arrives. LlmResponse.usageMetadata() returns an optional usage object from the google-genai library; your small Usage.totalTokens helper reads its total-token field (confirm the accessor name against your genai version) and returns empty when no usage is reported.

A token budget that reconciles

The token budget is a token bucket in model tokens. It refills continuously, bursts up to one minute's worth, and may go negative after reconciliation so an underestimate is repaid before new calls start. Set it to 85 to 90 percent of the provider's quota. The generic bucket, per-tenant keys and their trade-offs are covered in rate limiting in ADK Java.

public final class TokenBudget {
  private final long capacity;          // tokens per minute, e.g. 900_000 of a 1M quota
  private final double refillPerMs;     // exact: 30_000 TPM refills 0.5 per ms
  private double available;
  private long lastRefill = System.currentTimeMillis();

  public TokenBudget(long tokensPerMinute) {
    this.capacity = tokensPerMinute;
    this.refillPerMs = tokensPerMinute / 60_000.0;
    this.available = tokensPerMinute;
  }

  public synchronized boolean reserve(long tokens, java.time.Duration maxWait)
      throws InterruptedException {
    long deadline = System.currentTimeMillis() + maxWait.toMillis();
    while (true) {
      refill();
      if (available >= tokens) { available -= tokens; return true; }
      long waitMs = Math.min(deadline - System.currentTimeMillis(),
                             (long) Math.ceil((tokens - available) / refillPerMs));
      if (waitMs <= 0) return false;
      wait(waitMs);                       // releases the monitor while waiting
    }
  }

  public synchronized void reconcile(long reserved, long actual) {
    available = Math.min(capacity, available + reserved - actual);
    notifyAll();
  }

  private void refill() {
    long now = System.currentTimeMillis();
    available = Math.min(capacity, available + (now - lastRefill) * refillPerMs);
    lastRefill = now;
  }
}

Wire one instance per model and quota and pass it to every agent that uses that model, including the synthesizer. Two wrappers around the same model with separate budgets each think they own the whole quota.

Retries that do not duplicate output

Even a well-sized limiter meets throttling, because other services share the project and providers may throttle below the published number. Retry rate-limit errors, not most others, and follow one streaming rule: a streamed call that has already emitted partial responses would emit them again, duplicating text in the branch's events and output key. Retry only before the first response arrives.

// Retry a throttled call only if it has not emitted anything yet.
static Flowable<LlmResponse> withRetry(Flowable<LlmResponse> call, int maxAttempts) {
  return Flowable.defer(() -> {
    AtomicBoolean emitted = new AtomicBoolean();
    return call
        .doOnNext(r -> emitted.set(true))
        .retryWhen(errors -> errors.zipWith(
            Flowable.range(1, maxAttempts),
            (err, attempt) -> {
              if (emitted.get() || !isRateLimited(err) || attempt == maxAttempts) throw err;
              return attempt;
            })
            .flatMap(attempt -> Flowable.timer(backoffMillis(attempt), TimeUnit.MILLISECONDS)));
  });
}

static long backoffMillis(int attempt) {      // full jitter: random in [0, base * 2^attempt)
  long cap = Math.min(30_000L, 500L << attempt);
  return java.util.concurrent.ThreadLocalRandom.current().nextLong(cap);
}

The wrapper applies withRetry to a deferred, and therefore cold, delegate call after the permit is acquired, so a retry re-issues the request and keeps its place. isRateLimited is yours to write, because the exception type depends on the model client. Honour a Retry-After value if the error exposes one, keep attempts to two or three, and let the branch guard record the failure after that.

Fairness inside and across requests

A global limiter is fair to calls, not to users: one request fanning out to 20 branches can take every permit. Bound width per request by splitting data-driven branches into groups of at most k, run as a SequentialAgent of ParallelAgents. Give interactive and batch traffic separate wrappers with budgets carved from the same quota, so a backfill cannot starve users. The semaphore's fair mode serves waiters in arrival order.

Because a permit is scoped to one model call, a branch does not hold it while tools run. That prevents deadlock when a tool is an AgentTool whose inner agent uses the same model and would otherwise wait for a permit its parent holds.

Sharing a quota across replicas

Quotas are per project, and most deployments run several replicas. The simplest scheme divides the quota statically: with three replicas and a 1 million token quota, each gets a 300,000 token budget and a third of the concurrency. It needs no coordination but wastes headroom when load is uneven, and every scale-up must change the share. A central token bucket, for example in Redis updated by an atomic script, shares the quota exactly at the cost of a round trip per call and a dependency whose outage must fail open or closed by decision. Adaptive concurrency is the middle path: each replica cuts its limit multiplicatively on 429s and grows it additively on success, converging on a share without coordination.

Worked example: sizing the review agent

Return to the five-call review agent: 40 user requests per minute, 5,800 tokens per call, a project quota of 1,000 requests and 1 million tokens per minute, and three replicas. Tokens are the binding constraint. Budget 90 percent, 900,000 tokens per minute, split three ways: 300,000 per replica. At 5,800 tokens per call each replica can sustain about 51 calls per minute, about 10 user requests, so 30 user requests per minute across the fleet. Demand is 40. The limiter will hold the excess in queues until maxWait expires, which is the correct behaviour: it turns quota exhaustion into slow responses and recorded timeouts rather than cascades of 429s and cancelled fan-outs.

The real fixes are a higher quota, shorter prompts, three reviewers instead of four, or a smaller model with its own quota for the reviewers. For concurrency, Little's law gives about 7 calls in flight per replica at 51 calls per minute and 8 seconds each; set 10 permits per replica so a two-request burst of fan-outs is admitted without queueing. Set maxWait to what the user will tolerate, perhaps 20 seconds, and the per-call estimate from a week of measured usage.

Failure modes

FailureCausePrevention
Permits leak until everything blocksRelease in an after-model callback that does not run on error or cancelRelease in doFinally around the response stream
Duplicated text after a retryStreamed call retried after it emittedRetry only before the first response
Retry storm after an outageEvery branch retries at once with fixed delaysJittered exponential backoff, few attempts
Fan-out cancelled by one 429Error reaches the mergeTyped limiter errors, recorded by a branch guard
Quota exceeded despite a limiterSeveral wrappers or replicas each assume the full quotaOne budget per model per process; divide or centralise across replicas
Users starved by one large requestGlobal limiter, unbounded fan-out widthCap width per request; separate traffic classes
Deadlock with nested agentsPermit held across tool executionScope permits to one model call

Observability

Export permits in use, time to acquire, reserved versus actual tokens, and rate-limit errors by agent. Time to acquire rises before any 429 appears, so alert on it. Load test the full agent, not the model, because the fan-out's burst shape is what breaks quotas; running branches on virtual threads keeps blocked waiters cheap.

Trade-offs

Waiting protects the quota and adds latency; failing fast protects latency and drops work. Blocking permits are simple but hold a thread per waiter. Narrower fan-outs cost less and give up parallel speed-up. Make each choice explicitly and write it down next to the agent.

What to do next

  1. Count model calls per user request for every agent tree, then multiply by peak traffic and average tokens per call.
  2. Wrap each model in one shared RateLimitedLlm and pass that instance to every agent that uses it, including the synthesizer.
  3. Release permits in doFinally and test it: fail and cancel a branch, then confirm permits return to the starting count.
  4. Add retries that stop after the first emitted response, with jittered backoff and a provider-specific isRateLimited.
  5. Bound fan-out width per request and split interactive from batch traffic.
  6. Decide how replicas share the quota: static shares, a central bucket or adaptive limits.
  7. Dashboard time to acquire and reserved versus actual tokens; load test with the full agent.
Key takeaway: A ParallelAgent multiplies model calls by its width, so plan quota in branch calls and tokens, not user requests. Enforce the limit in one shared BaseLlm wrapper per model, not in the scheduler or in callbacks: acquire a fair permit and a token reservation at subscription, release in doFinally so errors and cancellation return the permit, and reconcile reserved tokens against reported usage. Retry throttled calls only before they emit, with jittered backoff, cap fan-out width per request, and decide explicitly how replicas share the project quota.