Most tools answer once. The model asks for the weather, the method returns a map, and the turn moves on. Some work does not fit that shape: watching a price, counting people in a camera feed, tailing a build log, reporting a long export as it progresses. For those, ADK has streaming tools, which keep running after the call and push intermediate results back to the agent so it can narrate or react without the user asking again. ADK Java supports them in live (bidirectional) sessions, and the docs mark the feature as experimental.

This article explains exactly what the Java runtime does with a streaming tool, because the behaviour is not what most people guess: results are not delivered as function responses, a second call to the same tool does not do what you expect, and errors vanish into a log. Everything below was checked against the google/adk-java main branch at the time of writing (the classes Functions, FunctionTool, Runner and BaseLlmFlow). It is a young API; re-read those classes after every upgrade. For the event stream itself see streaming in ADK Java and streaming events.

Two meanings of streaming a tool result

People say streaming tool results and mean one of two different things. The first is a tool whose output is a stream: it starts, then produces many values over seconds or hours while the conversation keeps going. The second is a tool that produces one result slowly, where you want to show progress to a human while it runs. ADK Java covers the first case directly with Flowable-returning tools in runLive. The second case, especially outside live mode, needs a different pattern, covered later.

The distinction matters because a live streaming tool feeds the model, while a progress indicator usually belongs in the user interface. Streaming twenty progress percentages into the model's context buys twenty model turns that say nothing.

How the runtime recognises and runs a streaming tool

A streaming tool is an ordinary FunctionTool whose Java method returns a parameterised Flowable. FunctionTool.isStreaming() reads the method's generic return type and returns true when the raw type is assignable to Flowable. There is no annotation and no separate class. Return Flowable<Map<String, Object>> (the test suite uses Flowable<ImmutableMap<String, Object>>, which also passes the check), because the runtime casts the result to a flowable of maps.

When the live flow handles a function call, Functions.processFunctionLive checks three cases in order. A call named stopStreaming with a functionName argument stops a running tool. A FunctionTool that is streaming is started in the background. Anything else is called normally and answers once. The diagram shows the second path.

One streaming tool call inside runLiveLive modelemits functionCallprocessFunctionLiveisStreaming() is truecallLiveinvokes your methodYour Flowableemits MapsImmediate responsestatus: results pendingturn continuessubscribe(onNext, onError, onComplete)Disposable kept in activeStreamingToolsEach Map becomes user Contenttext: Function X returned: {...}onNextLiveRequestQueue (session input)same queue the user audio and text usemodel reactsstopStreaming(functionName)dispose task, close inputStream, removecancelsResults do not arrive as a function response: they re-enter the conversation as user-role text.
The call returns a pending status at once; each later emission is injected into the live request queue as user text, where the model sees it like any other input.

Two details in that path shape everything else. First, the model's function call is answered immediately with the map {"status": "The function is running asynchronously and the results are pending."}, so the model can keep talking. Second, each emitted map is turned into a string, Function <name> returned: <map.toString()>, wrapped in a user-role Content and sent through the session's LiveRequestQueue with content(...). The model therefore receives results as if the user had typed them. Write the map so that its toString() reads well, because that string is literally the prompt.

Writing a streaming tool

Here is a complete price watcher. The method polls a quote source every few seconds, emits only when the price moves by more than a threshold, and finishes on its own after an hour so a forgotten stream cannot run forever. stopStreaming is declared as a normal tool with an empty body; the runtime intercepts it by name before your method would run.

public final class MarketTools {
  private static final QuoteClient quotes = QuoteClient.create();

  @Schema(description = "Watch a stock and report significant price moves until stopped.")
  public static Flowable<Map<String, Object>> watchPrice(
      @Schema(name = "symbol", description = "Ticker, e.g. GOOG") String symbol,
      @Schema(name = "minMovePct", description = "Report moves at least this large") double minMovePct) {
    AtomicReference<Double> last = new AtomicReference<>();
    return Flowable.interval(0, 5, TimeUnit.SECONDS, Schedulers.io())
        .take(Duration.ofHours(1).toSeconds() / 5)          // hard upper bound on lifetime
        .map(tick -> quotes.lastPrice(symbol))              // blocking call, on the io scheduler
        .filter(price -> {
          Double prev = last.get();
          boolean moved = prev == null || Math.abs(price - prev) / prev * 100 >= minMovePct;
          if (moved) last.set(price);
          return moved;
        })
        .map(price -> Map.<String, Object>of("symbol", symbol, "price", price))
        .onErrorReturn(e -> Map.<String, Object>of("symbol", symbol, "error", "quote feed failed: " + e.getMessage()));
  }

  @Schema(description = "Stop a running streaming tool.")
  public static void stopStreaming(
      @Schema(name = "functionName", description = "Name of the streaming tool to stop") String functionName) {}
}

LlmAgent agent = LlmAgent.builder()
    .name("market_watch")
    .model(LIVE_MODEL)   // a Gemini model that supports the Live API; check current model names
    .instruction("When the user asks to watch a stock, call watchPrice. Mention a move in one short "
        + "sentence. Call stopStreaming with functionName watchPrice when asked to stop.")
    .tools(FunctionTool.create(MarketTools.class, "watchPrice"),
           FunctionTool.create(MarketTools.class, "stopStreaming"))
    .build();

Driving it uses the live runner. Runner.runLive(session, liveRequestQueue, runConfig) returns a Flowable<Event>; the queue carries user text and audio in, and the events carry model output back. The streaming mode lives on RunConfig.

RunConfig runConfig = RunConfig.builder()
    .setStreamingMode(RunConfig.StreamingMode.BIDI)
    .build();
LiveRequestQueue input = new LiveRequestQueue();
Disposable events = runner.runLive(session, input, runConfig)
    .subscribe(event -> ui.render(event), err -> log.error("live session failed", err));

input.content(Content.fromParts(Part.fromText("Watch GOOG and tell me about 1% moves.")));
// ... later
input.content(Content.fromParts(Part.fromText("Stop watching.")));
input.close();

Note the onErrorReturn in the tool. Without it, an error in the flowable reaches the subscriber's error handler, which in the current source only logs it. The model is never told, so it keeps promising updates that will not come. Converting failures into a final emitted map is the only way the model learns the stream died.

Tools that consume the input stream

The second kind of streaming tool consumes the session's own input, for example video frames the client sends while the user talks. Declare a parameter of type LiveRequestQueue and name it inputStream. Three pieces of code cooperate here. When runLive starts, the runner scans the agent's tools and, for every FunctionTool whose method has a parameter of type LiveRequestQueue, registers a fresh queue under the tool's name. The live flow then copies every incoming LiveRequest into each registered tool queue as well as forwarding it to the model. Finally, when the tool is called, FunctionTool fills the parameter whose name is inputStream with that queue.

@Schema(description = "Report how many people are visible whenever the number changes.")
public static Flowable<Map<String, Object>> countPeople(
    @Schema(name = "inputStream") LiveRequestQueue inputStream) {
  return inputStream.get()
      .filter(req -> req.blob().isPresent()
          && "image/jpeg".equals(req.blob().get().mimeType().orElse("")))
      .sample(1, TimeUnit.SECONDS)                    // at most one frame per second
      .observeOn(Schedulers.io())
      .map(req -> detector.countPeople(req.blob().get().data().orElseThrow()))
      .distinctUntilChanged()                         // speak only when the count changes
      .map(n -> Map.<String, Object>of("people", n));
}

The @Schema(name = "inputStream") annotation is not decoration. Parameter names are matched by string, and Java only keeps real parameter names in bytecode when compiled with -parameters. Without the annotation or the flag the parameter is called arg0, so the runtime neither hides it from the model's function declaration nor fills it with the queue. Check the getter shapes on Blob against your google-genai version; the ones above follow the Optional-returning style of current releases.

Worked example: one watch, start to stop

A worked timeline makes the moving parts concrete. The user says: watch GOOG and tell me about one percent moves.

  1. The model emits functionCall watchPrice {symbol: GOOG, minMovePct: 1}.
  2. The runtime calls your method, subscribes on the spot, stores the Disposable under watchPrice and answers the call with the pending status. The model says it is watching.
  3. Five seconds later the first price passes the filter (there is no previous price). The map {symbol=GOOG, price=171.2} becomes the user text Function watchPrice returned: {symbol=GOOG, price=171.2}. The model hears it and may say the opening price.
  4. Forty minutes of small moves emit nothing, so no tokens are spent. The filter belongs in the tool.
  5. A 1.3 percent drop is emitted. The model, primed by the instruction, says one sentence.
  6. The user says stop. The model calls stopStreaming {functionName: watchPrice}. The runtime disposes the subscription, which cancels Flowable.interval upstream, removes the entry and replies Successfully stopped streaming function watchPrice.

If the user never says stop, the take bound ends the stream after an hour; on completion the runtime removes the entry from the active map itself.

Designing what to emit

Every emission costs a model turn's worth of input tokens and may provoke speech. Design emissions the way you would design alerts for an on-call engineer.

  • Emit changes, not samples. distinctUntilChanged, thresholds and debouncing turn a firehose into a handful of meaningful events.
  • Rate-limit at the source. sample or throttleLatest bound the rate no matter how fast the upstream runs. Two emissions a second in a voice session means the model interrupts itself.
  • Keep maps small and self-describing. The model sees toString() output; three named fields beat a nested object.
  • Bound the lifetime. Use take, takeUntil or a timeout so a stream ends even if nobody calls stop.
  • Emit a terminal message. A final map such as {status=finished} or an error map tells the model the stream is over; plain completion is invisible to it.
  • Keep blocking work off the session threads. Put blocking calls behind Schedulers.io() or your own executor; see async tool execution.

Outside live mode

Outside runLive, a Flowable return is not special. The regular call path handles Maybe, Single and plain values; a Flowable falls through to the generic conversion, which cannot turn it into a map and wraps the object itself as {result: ...}. Nothing subscribes, so no work runs and the values never reach the model. Do not register the same streaming method for runAsync agents.

For a slow single result in request-response mode, use one of these patterns instead.

  • Long-running tool plus client updates. Register the method with LongRunningFunctionTool.create(Jobs.class, "startExport"). It returns a ticket and status at once, the agent tells the user the job started, and your client later sends a message carrying a function response with the same call id and name when the job finishes or reaches a milestone the model should know about.
  • Progress out of band. Write progress to your job store and push it to the UI over your own channel. The model only gets the final result. This is usually right for percentages.
  • Split the work. If the model genuinely needs partial results, give it a paged tool such as fetchResults(jobId, cursor) and let it call again. Each page is a normal one-shot result.

Time limits for the one-shot path are covered in tool timeout handling.

Failure modes

These are the failures to plan for, all visible in the current source.

  • Silent stream errors. The background subscriber logs errors and does nothing else. Use onErrorReturn or onErrorResumeNext to emit an error map, and alert on the log line so the failure is visible to you as well.
  • One slot per tool name. Active streams are keyed by tool name. A second watchPrice call for another symbol replaces the stored Disposable without disposing the first, so the first stream keeps emitting and stopStreaming can no longer reach it. Either accept a list of symbols in one call and restart the stream when the list changes, or register separate tools.
  • Input queue gone after a stop. Stopping an input-consuming tool closes its queue and removes the entry, and normal completion removes the entry too. A later call in the same session then gets null for inputStream. Null-check it and emit an explanatory map, or restart the live session.
  • Stop never fires. stopStreaming must be registered as a tool with a functionName parameter, and the instruction must tell the model the exact name to pass. A typo returns a polite no active streaming function message and the stream keeps running.
  • Token burn. Count emissions per session and cap them in the tool.
  • Lost results on reconnect. Emissions travel through the live queue, not the function-response history. If the client reconnects into a new live session, the streams of the old one are gone; persist anything that matters to your own store.

Testing streaming tools

Test the flowable on its own first, without a model. Inject the quote client and the scheduler instead of using statics, drive it with RxJava's TestScheduler, advance virtual time and assert the exact maps emitted with TestSubscriber.assertValues. Filters, thresholds and lifetime bounds become cheap, deterministic checks. For the runtime wiring, the repository's StreamingToolTest drives runLive with a fake live model and checks that emissions arrive as user content and that stopStreaming cancels the stream; copy its structure for an end-to-end test of your own agent.

Trade-offs

ApproachModel sees partial resultsWorks in runAsyncMain cost
Flowable streaming toolYes, as user textNoTokens per emission; experimental API
Long-running tool plus updatesWhen you send themYesClient must track call ids
Progress out of bandNoYesSeparate channel to build
Paged tool callsYes, when it asksYesExtra model turns, latency

What to do next

  1. Decide whether your partial results are for the model or for the user. Only the first needs a streaming tool.
  2. Return Flowable<Map<String, Object>>, bound its lifetime, filter to meaningful changes and add onErrorReturn.
  3. Register stopStreaming and name the exact function in the agent instruction.
  4. For input-consuming tools, annotate the parameter with @Schema(name = "inputStream") and null-check it.
  5. Do not call the same streaming tool twice concurrently; take a list or use separate tools.
  6. Unit-test the flowable with TestScheduler, then run one end-to-end live test modelled on StreamingToolTest.
  7. After each ADK upgrade, re-read Functions.processFunctionLive; this is where the behaviour described here lives.
Key takeaway: In ADK Java a streaming tool is a FunctionTool returning a Flowable, and it streams only under runLive. The call is answered at once with a pending status and every emission re-enters the conversation as user text, so filter, rate-limit and bound the stream in the tool, convert errors into emitted maps, register stopStreaming, and never run two copies of the same tool. For progress outside live mode, use long-running tools or an out-of-band channel.