A2A streaming lets a remote agent report progress while it works: the caller sends one message and receives a stream of task status changes and artifact chunks instead of waiting for a single late reply. Turning streaming on is a configuration matter. Consuming the stream correctly is not: the events arrive on another thread, artifacts can arrive in pieces, the stream can end in a state that is not finished, and connections drop.

This recipe covers the client side in Java with the a2a-java SDK: what is on the wire, building the client, a consumer that turns events into state, reassembling artifacts, deciding when a call is over, resubscribing after a drop, cancelling, and running it inside an ADK application. Wiring RemoteA2AAgent and the server side are covered in ADK Java and A2A integration. Names here were checked against a2a-java 0.3.2.Final (groupId io.github.a2asdk, protocol 0.3), the line the google-adk-a2a module pins. The newer org.a2aproject.sdk line may rename types and methods, so re-check against whatever you pin.

What arrives on the wire

In protocol 0.3 over JSON-RPC, a streaming call is a POST of a message/stream request. The server answers with a Server-Sent Events response: each data: line carries one JSON-RPC response whose result is one of four kinds, distinguished by a kind field. A task is a full snapshot with an ID, a context ID, a status and any artifacts. A status-update carries a new TaskStatus (state, optional message, timestamp) and a final flag. An artifact-update carries an artifact plus append and lastChunk flags. A message is a direct reply, used by agents that answer without creating a task.

A typical long task streams a task snapshot in SUBMITTED, then status updates in WORKING with progress text, then artifact updates, then a final status update in COMPLETED with final set. Reconnecting to a task already in flight is a separate method, tasks/resubscribe, which streams the same event kinds from that point forward. Neither method replays events you missed, so the client must be able to rebuild state from a snapshot.

One streaming call: SSE events in, task state and artifacts outYour Java serviceor an ADK orchestratora2a-java ClientJSON-RPC transportRemote A2A agentmessage/streamsendMessagePOSTSSE eventsTask storeserver sideConsumer callback on the transport thread: keep it fast, hand off workTaskEventfirst snapshot: taskIdStatus updateWORKING, isFinal()Artifact updateappend, lastChunkMessageEventdirect reply, no taskArtifact bufferkeyed by artifactIdCompletableFuturedone on final stateStream drops (error handler, or silence)watchdog fires: resubscribe(taskId)Deadline passes or user abortscancelTask(taskId), then close()
The client sends once and receives four event kinds on a transport thread. The consumer buffers artifact chunks, completes a future on a final state, and recovers from drops with resubscribe or cancel.

Building the client

Add a2a-java-sdk-client from io.github.a2asdk at the version your google-adk-a2a release pins, so the protocol types match. Then resolve the agent card and build a client. ClientConfig streams by default (setStreaming defaults to true), but the client only streams if the card also advertises capabilities.streaming, so log what the card says at startup.

AgentCard card = new A2ACardResolver(baseUrl).getAgentCard();
log.info("agent {} streaming={}", card.name(), card.capabilities().streaming());

ClientConfig config = new ClientConfig.Builder()
        .setStreaming(true)
        .setAcceptedOutputModes(List.of("text/plain", "application/json"))
        .build();

Client client = Client.builder(card)
        .clientConfig(config)
        .withTransport(JSONRPCTransport.class, new JSONRPCTransportConfig())
        .build();

Build one Client per remote agent and reuse it; it holds the card and the transport. Consumers and an error handler can be registered on the builder for every call, or passed per call. Per call is the better default for streaming, because each call needs its own state: its own buffer, its own future, its own task ID.

A consumer that tracks state

A consumer is a BiConsumer<ClientEvent, AgentCard>. The SDK wraps wire events in three client event types: MessageEvent with getMessage(), TaskEvent with getTask(), and TaskUpdateEvent with getTask() and getUpdateEvent(), where the update is a TaskStatusUpdateEvent or a TaskArtifactUpdateEvent. The recipe keeps all per-call state in one object:

final class StreamCall {
    final CompletableFuture<Result> done = new CompletableFuture<>();
    final Map<String, StringBuilder> artifacts = new ConcurrentHashMap<>();
    volatile String taskId;
    volatile String contextId;
    volatile TaskState lastState = TaskState.SUBMITTED;
    volatile long lastEventNanos = System.nanoTime();
    volatile Throwable lastError;
    volatile boolean incomplete;      // chunks lost in a reconnect gap

    void onEvent(ClientEvent event, AgentCard card) {
        lastEventNanos = System.nanoTime();
        if (event instanceof MessageEvent m) {
            done.complete(Result.message(text(m.getMessage().getParts())));
        } else if (event instanceof TaskEvent t) {
            remember(t.getTask());
            onState(t.getTask().getStatus());
        } else if (event instanceof TaskUpdateEvent u) {
            remember(u.getTask());
            if (u.getUpdateEvent() instanceof TaskArtifactUpdateEvent a) {
                onArtifact(a);
            } else if (u.getUpdateEvent() instanceof TaskStatusUpdateEvent s) {
                onState(s.getStatus());
                if (s.isFinal()) finish();
            }
        }
    }

    void remember(Task task) {
        taskId = task.getId();
        contextId = task.getContextId();
    }

    void onState(TaskStatus status) {
        lastState = status.state();
        progress.publish(taskId, status.state(), status.message());  // UI, logs
        if (status.state().isFinal()
                || status.state() == TaskState.INPUT_REQUIRED
                || status.state() == TaskState.AUTH_REQUIRED) {
            finish();
        }
    }

    void finish() {
        done.complete(Result.task(taskId, contextId, lastState, assembled(), incomplete));
    }

    void onError(Throwable t) {          // loud drop: hand off, never block here
        lastError = t;
        recoveryQueue.offer(this);       // drained by the worker that calls recover()
    }
}

Three habits keep this safe. Capture the task ID from the first event that carries one, since every recovery path needs it. Make the callback fast: it runs on the transport's thread, so push progress to a queue or executor rather than calling a database or a model inline. And make completion idempotent: CompletableFuture.complete ignores every call after the first, so a final status update followed by a terminal state cannot finish the call twice.

Reassembling streamed artifacts

Artifacts can stream in pieces. Each TaskArtifactUpdateEvent carries an Artifact record (artifactId(), name(), parts()) and two flags. When isAppend() is true, the parts extend the artifact with that ID; when false, they replace it. isLastChunk() marks the final piece. Both flags are nullable Booleans in 0.3, so compare with Boolean.TRUE.equals(...) rather than unboxing.

void onArtifact(TaskArtifactUpdateEvent a) {
    Artifact art = a.getArtifact();
    String chunk = text(art.parts());
    if (Boolean.TRUE.equals(a.isAppend())) {
        artifacts.computeIfAbsent(art.artifactId(), k -> new StringBuilder()).append(chunk);
    } else {
        artifacts.put(art.artifactId(), new StringBuilder(chunk));
    }
    if (Boolean.TRUE.equals(a.isLastChunk())) {
        progress.artifactReady(taskId, art.artifactId(), art.name());
    }
}

static String text(List<Part<?>> parts) {
    StringBuilder sb = new StringBuilder();
    for (Part<?> p : parts) {
        if (p instanceof TextPart t) sb.append(t.getText());
    }
    return sb.toString();
}

This recipe builds artifacts only from the update events it receives, keyed by artifactId. The task object on a TaskUpdateEvent may also carry artifacts the SDK has aggregated; whether and how it does was not confirmed for this version, so do not also append from it, or you risk counting a chunk twice. Pick one source. If you need structured output, buffer the raw parts and parse once the last chunk arrives, never per chunk.

Knowing when the call is over

The JSON-RPC transport opens the SSE request asynchronously, so sendMessage returns before the stream ends. The caller waits on the future, with a deadline:

Result call(Client client, String prompt, Duration deadline) throws Exception {
    StreamCall call = new StreamCall();
    Message msg = A2A.toUserMessage(prompt);
    client.sendMessage(msg, List.of(call::onEvent), call::onError, null);
    try {
        return call.done.get(deadline.toMillis(), TimeUnit.MILLISECONDS);
    } catch (TimeoutException e) {
        cancelQuietly(client, call.taskId);
        throw e;
    }
}

The future completes on any of four conditions: a direct message reply; a status update with isFinal() set; a terminal state (COMPLETED, FAILED, CANCELED, REJECTED, and UNKNOWN, which TaskState.isFinal() also treats as terminal); or INPUT_REQUIRED or AUTH_REQUIRED. The last two end the stream without finishing the task. The agent is waiting for you, so the result must say so, and the caller must reply with a message carrying the same task and context IDs:

Message followUp = new Message.Builder()
        .role(Message.Role.USER)
        .parts(List.of(new TextPart(answer)))
        .taskId(result.taskId())
        .contextId(result.contextId())
        .build();

Treat REJECTED and UNKNOWN as failures in your result type. Callers that only check for COMPLETED will otherwise report an empty success.

Dropped streams and resubscribe

Streams fail in two ways. Loud failures reach the streaming error handler: connection reset, an HTTP error, a malformed event. Quiet failures do not: a proxy or load balancer closes an idle connection, or the stream simply stops, and no callback fires. Handle both with the same recovery, triggered either by the error handler or by a watchdog that notices no event for longer than the agent's expected heartbeat.

void recover(Client client, StreamCall call) {
    if (call.taskId == null) {             // dropped before the first event
        call.done.completeExceptionally(new IllegalStateException("no task id"));
        return;                            // resending risks a duplicate task
    }
    try {
        Task snap = client.getTask(new TaskQueryParams(call.taskId), null);
        List<Artifact> arts = snap.getArtifacts();
        if (arts != null && !arts.isEmpty()) {          // snapshot is authoritative after a gap
            call.artifacts.clear();
            for (Artifact a : arts) {
                call.artifacts.put(a.artifactId(), new StringBuilder(text(a.parts())));
            }
        } else if (!call.artifacts.isEmpty()) {
            call.incomplete = true;                     // gap chunks are not replayed
        }
        call.onState(snap.getStatus());   // may already be final
        if (!call.done.isDone()) {
            client.resubscribe(new TaskIdParams(call.taskId),
                    List.of(call::onEvent), call::onError, null);
        }
    } catch (Exception e) {
        retryLater(call, e);               // bounded backoff, then fail
    }
}

Fetch the snapshot before resubscribing. The task may have finished while you were disconnected, and resubscribing to a finished task gains nothing. Because resubscribe does not replay, any artifact chunks streamed during the gap are lost from your buffer; if the snapshot contains the complete artifact, replace the buffer from it, and if not, mark the result incomplete rather than returning a truncated report. Run the watchdog on a scheduler, check it every few seconds, and set the idle limit above the server's longest quiet period between events.

Do not blindly resend the original message after a drop. Without a task ID the server sees a new request and starts a second task, which for an agent with side effects means doing the work twice.

Cancellation and shutdown

When the deadline passes or the user abandons the request, call cancelTask(new TaskIdParams(taskId), null) so the remote agent stops spending tokens and tool calls. Cancellation is a request: the server may refuse a task that is already finishing, and the returned Task shows the state it reached. Log it either way. Call client.close() when the service shuts down, not after each call, since the client is meant to be reused.

Inside an ADK application

Inside an ADK orchestrator you usually do not hand-write this consumer: RemoteA2AAgent drives the client and converts events to ADK events, and the integration article explains its streaming switches. The recipe still applies in two places. First, when a plain Java service, a batch job or a test harness calls an A2A agent directly, which is common for evaluation and for non-agent backends. Second, as the mental model for debugging the orchestrator: when remote replies arrive all at once, are cut short or hang, the cause is one of the cases above. For how ADK surfaces streamed events to your own code, see ADK Java streaming events, and for mapping remote failures into ADK errors, A2A error propagation.

Worked example: a summary that survives a proxy timeout

A reporting service asks a research agent for a three-section market summary with a 90-second deadline and a 20-second idle watchdog. The first event is a task snapshot in SUBMITTED; the consumer records task t-81 and context c-12. Status updates in WORKING arrive with messages such as gathering sources and drafting section 1, which are pushed to the UI. Artifact updates for artifact summary arrive with append true, one per paragraph.

Forty seconds in, a corporate proxy cuts the idle connection while the agent is waiting on a slow search tool. No error fires. At 60 seconds the watchdog sees 20 seconds without events and calls recover: the snapshot shows WORKING with a partial artifact, so the consumer replaces its buffer from the snapshot and resubscribes. Remaining chunks arrive, then a status update in COMPLETED with final set. The future completes at 78 seconds with the full artifact. Without the watchdog the call would have sat until the 90-second deadline, then been cancelled, wasting the agent's work.

Failure modes

SymptomLikely causeFix
Whole reply arrives at the endCard does not advertise streaming, or client config disables itLog card capabilities at startup; enable on both sides
Caller hangs until deadlineStream died quietly; no event, no errorIdle watchdog plus snapshot and resubscribe
Call returns success with no outputINPUT_REQUIRED, REJECTED or UNKNOWN treated as doneModel every end state in the result type
Artifact text duplicatedChunks appended from both update events and the task objectChoose one source of truth
Artifact truncated after reconnectResubscribe does not replay missed chunksRebuild from snapshot or mark incomplete
Remote work runs twiceOriginal message resent after a dropRecover by task ID; never resend blindly
Throughput collapses under loadSlow I/O inside the consumer callbackHand events to a queue or executor
Unboxing NullPointerExceptionappend or lastChunk absent on the wireCompare with Boolean.TRUE.equals

Trade-offs

Streaming costs an open connection per in-flight call and more client code; it buys progress, earlier partial results and early cancellation. For calls that finish in a second or two, a blocking call is simpler and fine. For work that runs for minutes or hours, an open stream is fragile; push notifications, configured with setPushNotificationConfig, let the agent call you back instead, at the cost of exposing an authenticated webhook. Many teams stream for interactive use and fall back to polling getTask for background jobs. Protocol-level behaviour is covered in A2A streaming.

What to do next

  1. Pin a2a-java to the version your google-adk-a2a release uses and log card capabilities at startup.
  2. Write a per-call state object with a future, an artifact buffer and the task and context IDs.
  3. Complete the future on final updates, terminal states and INPUT_REQUIRED or AUTH_REQUIRED, and model each in the result.
  4. Add an idle watchdog that snapshots with getTask and resubscribes by task ID.
  5. Cancel on deadline or user abort, and never resend a message to recover.
  6. Test by killing the connection mid-stream with a proxy and confirm the result is complete or flagged incomplete.
Key takeaway: Consuming an A2A stream in Java means treating each call as a small state machine: capture the task ID from the first event, buffer artifact chunks by artifactId from update events only, complete a future on final updates, terminal states or input-required, recover quiet drops with a watchdog that snapshots and resubscribes by task ID, and cancel rather than abandon. Never resend a message to recover a stream.