Horizontal scaling means serving more agent traffic by adding replicas instead of bigger machines. For ADK Java the unit of work is a turn: one call to Runner.runAsync that loads a session, makes several model and tool calls, and emits a Flowable<Event> over many seconds. Turns are long, bursty and mostly waiting on the network, which makes some textbook strategies fit badly and others fit well.
This article compares five strategies, from a plain stateless pool to queue-backed workers and connection-bound live sessions, and gives a decision table and a worked sizing example. It assumes the precondition that no session state lives only in a replica, and it does not repeat the capacity model; both are in ADK Java scale-out, in depth.
Strategy 1: a stateless pool
The default: identical replicas behind a load balancer, every replica able to serve any session because sessions, artifacts and memory live in shared services. A request arrives, the replica loads the session, runs the turn and streams events back on the same connection.
It is the right first strategy. Deploys are simple, a lost replica loses only its in-flight turns, and autoscaling on concurrency works out of the box. Its costs appear at scale. Every turn reloads the full session from the store, because no replica has a reason to cache it. Two messages for the same session can land on different replicas and run concurrently, each appending events to a history the other has not seen; you need a per-session lease to serialise them. And the HTTP connection is pinned to the replica for the whole turn, so a rolling deploy must wait out the longest turn or cut it off.
Strategy 2: session affinity
Session affinity routes every request for a session to the same replica while that replica is healthy. Done with a consistent scheme it fixes two pool problems at once: the replica can keep a warm in-memory copy of the session, so loads become cache hits, and same-session turns can be serialised by a per-session lane inside one JVM instead of a distributed lease.
Cookie stickiness on the load balancer is the crude version. A better one is rendezvous (highest random weight) hashing in a routing layer: score every replica against the session ID and pick the highest. When a replica leaves, only its sessions move, and each moves to its second-highest scorer.
static String owner(String sessionId, List<String> replicas) {
String best = null;
long bestScore = Long.MIN_VALUE;
for (String replica : replicas) {
long score = Hashing.murmur3_128()
.hashString(replica + "|" + sessionId, StandardCharsets.UTF_8).asLong();
if (score > bestScore) { bestScore = score; best = replica; }
}
return best;
}
// Inside the owning replica: one lane per session, so turns never overlap.
private final Cache<String, Scheduler> lanes = Caffeine.newBuilder()
.expireAfterAccess(Duration.ofMinutes(10)).build();
Flowable<Event> runTurn(String userId, String sessionId, Content msg) {
Scheduler lane = lanes.get(sessionId, id -> Schedulers.from(Executors.newSingleThreadExecutor()));
return Flowable.defer(() -> runner.runAsync(userId, sessionId, msg)).subscribeOn(lane);
}Two cautions. Affinity is an optimisation, not a source of truth: the cache must be write-through to the session store, and the lane must be backed by the store's own conflict check, because during a membership change two replicas can briefly both believe they own a session. And a lane per session needs eviction that also shuts down its executor; the sketch above leaks threads, so production code should share a bounded pool and queue per-session work on it. The cost of affinity is uneven load: one very chatty session pins one replica, and scaling events move sessions and cool their caches.
Strategy 3: queue-backed turn workers
The biggest structural change is to stop running the turn on the request path. An intake service validates the message, writes it to a queue keyed by session, and answers at once with a request ID. Workers pull from the queue and run runAsync. Events go to the session store and to a per-session stream, and a gateway relays the stream to the client over server-sent events.
Outcome handle(TurnRequest req) { // called by the queue consumer
Claim claim = turns.claim(req.requestId(), workerId, Duration.ofSeconds(30));
switch (claim.status()) {
case DONE: return Outcome.ACK; // redelivery of a finished turn
case LEASED: return Outcome.NACK_LATER; // another worker is alive and running it
case ACQUIRED: break; // new, or the previous lease expired
}
Content msg = Content.fromParts(Part.fromText(req.text()));
runner.runAsync(req.userId(), req.sessionId(), msg)
.doOnNext(ev -> stream.publish(req.sessionId(), req.requestId(), ev.toJson()))
.doOnError(err -> stream.publishError(req.sessionId(), req.requestId(), err))
.doOnSubscribe(s -> claim.startRenewing()) // heartbeat extends the lease
.doFinally(claim::stopRenewing)
.blockingSubscribe();
turns.markDone(req.requestId());
return Outcome.ACK; // ack only after the turn finishes
}What this buys: intake stays fast and cheap under any load; a burst becomes a backlog instead of rejected requests; per-session ordering comes from the queue (Pub/Sub ordering keys or a Kafka partition keyed by session ID); and workers can be drained during deploys simply by stopping consumption. Workers scale on backlog age, which tracks user-visible delay directly, rather than on CPU, which barely moves for network-bound turns.
What it costs: three more moving parts, at-least-once delivery that makes an idempotency claim mandatory (otherwise a worker crash after the model call replays the turn and charges twice), a claim that records status with a renewable lease rather than a bare flag (a bare 'seen it' flag acks the redelivery of a turn whose worker died, and the turn is lost), a poison-message policy for turns that always fail, and a gateway that resumes streams with Last-Event-ID when clients reconnect. Throttling upstream model quota fits naturally here; see the ADK Java rate limiter.
Strategy 4: live, connection-bound sessions
Bidirectional audio and video sessions use runLive(Session, LiveRequestQueue, RunConfig) and hold a connection for the whole conversation. They cannot be balanced per message: the replica that accepts the WebSocket owns the session until it closes. Scale them on concurrent connections, cap connections per replica, and spread new connections by least connections, not round robin. Deploys are the hard part, because a live session cannot migrate; stop accepting new connections on a draining replica and give existing ones a bounded time to finish, as in zero-downtime deploys.
Strategy 5: split heavy tools, and choosing
Sometimes one part of the agent is the bottleneck: a browser tool that needs a gigabyte per instance, a code-execution sandbox, a retrieval service with its own index. Moving it behind a network interface, as an MCP server or an A2A agent, lets it scale on its own signal while the agent replicas stay small. The price is a network hop per call and a new failure mode, so split only where resource profiles really differ.
| Strategy | Fits when | Main cost | Scale signal |
|---|---|---|---|
| Stateless pool | Default; moderate traffic, short sessions | Session reloads, cross-replica races | In-flight turns |
| Session affinity | Long chatty sessions, large histories | Uneven load, rebalancing | In-flight turns per replica |
| Queue-backed workers | Bursty traffic, long turns, strict ordering | More components, idempotency | Backlog age |
| Live connections | Audio and video streaming | No migration, slow drains | Open connections |
| Split heavy tools | One component dominates memory or CPU | Extra hop, extra service | That component's own load |
Scale signals and resumable delivery
Each tier in the queue-backed design needs its own scale signal, and choosing them is most of the operational work. Intake is cheap per request and scales on request rate or CPU, the one place where CPU is a fair proxy. Workers scale on how long the oldest unprocessed turn has waited, because that number is the delay a user feels before the first token. The gateway scales on open connections, since each holds a socket for the life of a turn. On Kubernetes with Pub/Sub, KEDA's gcp-pubsub scaler supports an OldestUnackedMessageAge mode alongside the default SubscriptionSize:
apiVersion: keda.sh/v1alpha1
kind: ScaledObject
metadata:
name: turn-workers
spec:
scaleTargetRef:
name: turn-worker
minReplicaCount: 3
maxReplicaCount: 40
triggers:
- type: gcp-pubsub
metadata:
subscriptionName: agent-turns-sub
mode: OldestUnackedMessageAge
value: "5" # seconds of backlog age to aim for
authenticationRef:
name: pubsub-workload-identityKeep a floor of replicas so a quiet night does not leave the first morning burst waiting for JVM start-up, and cap the maximum at what the model quota can feed; replicas beyond that only add 429 errors.
Delivery needs one more piece. A client that loses its connection mid-turn reconnects to any gateway replica with the ID of the last event it saw. The gateway replays later events from the per-session stream, then continues live. That works only if every published event carries a monotonically increasing ID per session and the stream retains a few minutes of history, so choose a stream technology with replay, such as Redis Streams or a Kafka topic, rather than fire-and-forget pub/sub.
Worked example: sizing for a burst
A support agent peaks at 30 turns per second, and a turn takes 6 seconds on average. By Little's law about 30 times 6, or 180, turns are in flight. Load tests show a worker holds 40 concurrent turns within its memory budget, so 180 over 40 is 4.5, rounded up to 5, plus one spare: 6 workers. Together they complete 6 times 40 over 6, or 40 turns per second.
Now a marketing email doubles arrivals to 60 per second for one minute. In a stateless pool with admission control, a third of those requests are rejected during the spike. With a queue, the backlog grows at 60 minus 40, or 20 per second, reaching 1,200 turns after 60 seconds; the worst queue wait is 1,200 over 40, or 30 seconds. When arrivals fall back to 30, the backlog drains at 10 per second and is gone in 120 seconds, faster if the autoscaler adds workers as backlog age crosses its target. Whether 30 seconds of extra wait beats rejection is a product decision, but the queue makes it a choice rather than an accident.
Failure modes
- Scaling on CPU. Agent replicas wait on the network, so CPU stays low while latency explodes. Scale on concurrency, backlog age or connections; see Kubernetes deployment.
- Scaling past the model quota. More replicas cannot exceed upstream requests-per-minute limits; they only turn queueing into 429 errors.
- Duplicate turns. Redelivery without an idempotency claim runs the model twice and can repeat side-effecting tool calls.
- Half-finished turns. When an expired lease is taken over, the crashed attempt may already have appended events, even tool results, to the session. Make tools idempotent and have the new attempt check the session tail before repeating work.
- Rebalancing storms. Rapid scale up and down with affinity moves sessions repeatedly; add a stabilisation window to the autoscaler.
- Hot sessions and tenants. A few heavy users can saturate one replica or partition; move them deliberately, as described in agent tenant scaling.
Trade-offs: combining strategies
The strategies are not exclusive; mature deployments combine them. A common shape is queue-backed workers with affinity inside the worker pool: with Kafka, partitioning by session ID sends a session's turns to one consumer while partition assignment is stable, which makes a warm session cache worthwhile. Pub/Sub ordering keys serialise delivery but do not pin a session to one subscriber. Live sessions then run on a separate pool with its own connection limits, and the heaviest tools sit behind their own services.
Each step up the ladder buys burst tolerance or efficiency with operational complexity: more components to deploy, more metrics to watch, more failure paths to test. A team of two running a few turns per second should stay on the stateless pool and spend the effort on admission control and quota handling. Move up only when measurements, not expectations, show the pool's costs: session-store load from reloads, same-session conflicts, rejected bursts or deploys that cut turns off.
What to do next
- Confirm sessions, artifacts and memory live in shared services, then run the stateless pool with an admission limit.
- Measure turn duration, memory per turn and same-session overlap rate from production logs.
- If sessions are long and histories large, add rendezvous-hash routing and a write-through session cache.
- If traffic is bursty or turns exceed a few seconds, prototype the queue-backed design with idempotency claims and a resumable SSE gateway.
- Pick one scale signal per tier and load-test a doubling of arrivals for one minute.
- Split out any tool whose memory or CPU profile dominates a replica.