A ParallelAgent runs several sub-agents at the same time and interleaves their events into one stream. It is the natural shape for research fan-outs, multi-source retrieval and multiple-perspective review. It is also the place where one flaky dependency does the most damage, because by default the first branch to throw takes every other branch down with it.

This article reads the failure behavior directly from the adk-java source, then builds the pieces that make a fan-out degrade gracefully: a wrapper agent that records each branch's outcome in session state, a quorum check before the synthesis step, retries that do not duplicate work, and scheduler isolation. The code targets the public adk-java API; check it against the version you depend on.

Advertisement

What a ParallelAgent actually does

The implementation is short enough to reason about completely. In runAsyncImpl, the agent takes its sub-agents, builds a child InvocationContext whose branch includes its own name, and for each sub-agent calls subAgent.runAsync(ctx).subscribeOn(scheduler). The scheduler defaults to Schedulers.io() and can be replaced through the builder's scheduler(...) method. The per-branch streams are combined with Flowable.merge, and the merged stream ends early via takeUntil when an event from a direct sub-agent carries escalate = true. runLiveImpl returns an error: live, bidirectional streaming is not supported for this agent.

When each sub-agent starts, BaseAgent appends its own name to the branch, so a branch looks like research_fanout.docs_agent. The branch is what isolates conversation history: an LLM agent in one branch does not see the model turns of its siblings. Session state is not isolated. All branches read and write the same state map, which matters below.

ParallelAgent: branches merged into one event stream, so one error ends them allParallelAgentmerge + takeUntil(escalate)branch: docssubscribeOn(io)branch: ticketsthrows 429branch: websubscribeOn(io)without guards: onError reaches merge, the other two subscriptions are cancelledGuardedBrancherror to status eventGuardedBranchstatus:tickets = failedGuardedBranchtimeout, retry predicateMerge step (SequentialAgent next)quorum check in beforeAgentCallback, then synthesizeSession state is shared by all branches; conversation history is filtered per branch.
Top: the default behavior, where an error in one branch reaches merge and cancels the siblings. Bottom: each branch wrapped so failures become state, and a merge step that decides whether enough succeeded.

The failure semantics that follow

Flowable.merge, unlike mergeDelayError, propagates the first onError immediately and cancels the other subscriptions. Everything else follows from that and from shared state.

What happens in one branchEffect on the fan-out
Model or tool throws (a 429, a network error, a JSON parse failure)The whole ParallelAgent errors; siblings are cancelled mid-flight; the error reaches the caller of the runner
Branch emits an event with escalate settakeUntil completes the merged stream; siblings are cancelled, but without an error
Invocation exceeds RunConfig maxLlmCallsincrementLlmCallsCount throws LlmCallsLimitExceededException in whichever branch crossed the limit, and it propagates like any other error
Branch hangs on a slow callNothing fails; the fan-out waits for the slowest branch, forever if nothing sets a deadline
Branch returns plausible nonsenseNothing fails; its outputKey holds bad data that the next step trusts
Two branches use the same outputKeyBoth succeed; the key holds whichever wrote last, and the order changes from run to run

Two consequences are easy to miss. First, events a branch emitted before a sibling failed have typically already been appended to the session by the runner, so a failed fan-out leaves partial state behind. Design the next step to treat state from a failed invocation as untrusted. Second, cancellation is cooperative. Cancelling an RxJava subscription stops downstream delivery, but a blocking HTTP call already running on an io thread may keep running until it returns, and its side effects still happen.

Advertisement

The contract: every branch reports an outcome

Graceful degradation needs a record of what each branch did, written in a place the next step can read. Use session state with a fixed naming convention: branch docs writes its result to result_docs through its outputKey, and the guard writes status_docs as ok, failed or timeout, plus error_docs with a short, non-sensitive reason. Distinct keys per branch also remove the last-writer-wins race. The conventions for moving data through state are covered in passing context between sequential steps.

GuardedBranch: errors and timeouts become events

A small custom agent wraps each real branch. It runs the delegate, appends an ok status event if the delegate completes, and on error or timeout replaces the rest of the stream with a single status event. Because the guard never lets an error escape, merge never sees one.

import com.google.adk.agents.BaseAgent;
import com.google.adk.agents.InvocationContext;
import com.google.adk.events.Event;
import com.google.adk.events.EventActions;
import io.reactivex.rxjava3.core.Flowable;
import java.time.Duration;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;

public final class GuardedBranch extends BaseAgent {
  private final BaseAgent delegate;
  private final String key;          // e.g. "docs"
  private final Duration deadline;

  public GuardedBranch(BaseAgent delegate, String key, Duration deadline) {
    super("guard_" + key, "Runs " + delegate.name() + " and records its outcome.",
          List.of(delegate), List.of(), List.of());
    this.delegate = delegate;
    this.key = key;
    this.deadline = deadline;
  }

  @Override
  protected Flowable<Event> runAsyncImpl(InvocationContext ctx) {
    return delegate.runAsync(ctx)
        .map(this::dropEscalation)
        .concatWith(Flowable.defer(() -> Flowable.just(status(ctx, "ok", null))))
        .timeout(deadline.toMillis(), TimeUnit.MILLISECONDS)
        .onErrorResumeNext(err -> Flowable.just(status(ctx,
            err instanceof TimeoutException ? "timeout" : "failed",
            err.getClass().getSimpleName())));
  }

  @Override
  protected Flowable<Event> runLiveImpl(InvocationContext ctx) {
    return Flowable.error(new UnsupportedOperationException("live mode not supported"));
  }

  private Event status(InvocationContext ctx, String outcome, String reason) {
    Map<String, Object> delta = new ConcurrentHashMap<>();
    delta.put("status_" + key, outcome);
    if (reason != null) delta.put("error_" + key, reason);
    return Event.builder()
        .id(Event.generateEventId())
        .invocationId(ctx.invocationId())
        .author(name())
        .branch(ctx.branch().orElse(null))
        .actions(EventActions.builder().stateDelta(new ConcurrentHashMap<>(delta)).build())
        .timestamp(System.currentTimeMillis())
        .build();
  }

  private Event dropEscalation(Event e) {
    return e; // policy hook: see "Escalation" below
  }
}

Notes on the design. The delegate is registered as the guard's only sub-agent, and an agent can have only one parent, so construct a fresh delegate for each guard. Flowable.timeout(long, TimeUnit) as used here is an inactivity timeout: it fires when no event arrives within the deadline, which catches a hung model or tool call but not a branch that keeps emitting slowly for minutes. If you need a hard wall-clock budget for the whole branch, add a separate overall limit, for example by merging with a Flowable.timer that emits a timeout status and ends the stream. The recorded reason is the exception's class name, not its message, because messages can carry prompts, URLs or tokens into session state and from there into logs.

Escalation is a policy decision

Because ParallelAgent stops everything when a direct sub-agent escalates, and the guard is now the direct sub-agent, an escalation event authored by the delegate no longer ends the fan-out; only an escalation authored by the guard would. Decide explicitly which you want. If one branch finding a blocking condition, such as a policy violation in the input, should stop all work, have the guard re-emit an escalating status event when it sees one. If escalation inside a branch should only end that branch, leave it as is and record status_x = escalated. Either way, write the choice down; it is the most common source of fan-outs that stop for no visible reason.

The quorum check before synthesis

Put the parallel agent inside a SequentialAgent whose next step synthesizes the results. Before that step calls a model, a beforeAgentCallback counts successful branches. Returning content skips the model call and uses the content as the step's output; returning empty lets it run. Callback mechanics are covered in the callback architecture guide.

// Registered on the synthesis LlmAgent with .beforeAgentCallback(Quorum::require2of3)
static Maybe<Content> require2of3(CallbackContext ctx) {
  List<String> keys = List.of("docs", "tickets", "web");
  long ok = keys.stream()
      .filter(k -> "ok".equals(ctx.state().get("status_" + k)))
      .count();
  ctx.state().put("sources_ok", ok);
  if (ok >= 2) return Maybe.empty();                 // enough evidence: run the model
  return Maybe.just(Content.fromParts(Part.fromText(
      "Not enough sources answered (" + ok + " of 3). Please retry shortly.")));
}

The synthesis instruction should also name which sources succeeded, so the answer can say it is based on two of three sources rather than silently pretending all three agreed. A key such as sources_ok in state, read through a placeholder in the instruction, does that.

Retries without duplicated work

It is tempting to add .retry(2, this::isTransient) to the guard. RxJava will resubscribe and run the delegate again from the start, which means a second copy of every event the first attempt already emitted, a second round of tool calls, and a second set of side effects. Retries belong as close to the failing call as possible: in the tool, with backoff and an idempotency key, or in the model client for 429 and 503 responses. Retry at branch level only when the branch is read-only and its first attempt failed before emitting anything, and never retry LlmCallsLimitExceededException or validation errors, which will fail again.

When several branches share one model quota, a 429 in one branch is a signal about all of them. Independent retries with the same backoff turn one throttle into a synchronized storm. Add jitter, and share a circuit breaker per upstream so that once it opens, every branch fails fast into a failed status instead of waiting.

Isolating slow branches

All branches of all parallel agents in the process default to the shared io scheduler, which grows threads without bound. A dependency that hangs can therefore consume threads until the JVM struggles. Give each fan-out, or each class of downstream, its own bounded scheduler through the builder, for example ParallelAgent.builder().name("research_fanout").subAgents(...).scheduler(Schedulers.from(executor)).build() with a fixed-size executor. That is the bulkhead pattern applied to agent branches: one slow source can exhaust its own pool but not its neighbors'.

Worked example: a three-source research fan-out

A support assistant answers engineering questions by searching internal docs, the ticket history and the public web in parallel, each through an LLM agent with a search tool and an outputKey of result_docs, result_tickets and result_web. Each is wrapped in a GuardedBranch with a 20-second deadline and runs on a six-thread executor.

On one request the ticket search tool receives a 429 from its backend. Without guards, the exception cancels the docs and web branches at about 4 seconds and the user sees an error. With guards, the tickets branch writes status_tickets = failed and error_tickets = ClientException at 4 seconds, the docs branch finishes at 7 seconds, the web branch at 11. The quorum callback counts two successes, the synthesis step runs with an instruction noting that ticket history was unavailable, and the trace shows exactly one failed span. The metric that matters, answers produced with fewer than all sources, increments by one.

Testing the failure paths

Failures are rare in development, so inject them. A test agent that returns Flowable.error(new RuntimeException("boom")), another that returns Flowable.never(), and a third that emits an escalating event give you the three shapes that matter. Run each through the in-memory runner inside the fan-out and assert on session state: the right status_ keys, no stray result_ key from a failed branch, and a synthesis step that either ran or returned the fallback message. Add one test with two branches writing the same outputKey to prove your naming convention prevents it.

Failure modes

  • Unguarded branches. One transient error fails the whole request. Guard every branch, including the ones that never fail today.
  • No deadline. The fan-out is as slow as the slowest dependency on its worst day. Set per-branch deadlines and, for tail latency, consider request hedging on the slowest source.
  • Retrying the whole branch. Duplicate events and side effects. Retry inside tools.
  • Shared output keys. Nondeterministic results that pass every test that runs branches in a fixed order.
  • Trusting partial state. A failed invocation leaves some keys written. Clear or version per-invocation keys.
  • Leaking exception messages into state. Store classes or error codes, not messages.

What to do next

  1. List every ParallelAgent you run and mark which branches can fail independently.
  2. Adopt the result_x and status_x naming convention and remove any shared outputKey.
  3. Wrap each branch in a guard that records ok, failed or timeout, with a deadline per branch.
  4. Add a quorum callback on the synthesis step and make the answer name its sources.
  5. Move retries into tools and model clients, with jitter and a shared circuit breaker per upstream.
  6. Give each fan-out a bounded scheduler, and add failure-injection tests for error, hang and escalation.
Key takeaway: ADK Java's ParallelAgent merges its branches with Flowable.merge, so the first error in any branch cancels all the others, an escalation from a direct sub-agent stops the fan-out, and nothing bounds a slow branch. History is isolated per branch but session state is shared. Make each branch report an outcome in its own state keys through a guard that converts errors and timeouts into status events, check a quorum before synthesis, keep retries inside tools, isolate branches on bounded schedulers, and test the error, hang and escalation paths deliberately.