ADK Java's model layer has one method that every model implementation must get right: Flowable<LlmResponse> generateContent(LlmRequest llmRequest, boolean stream) on the abstract class BaseLlm. When stream is false it emits one response. When it is true it emits a sequence, and the shape of that sequence is a contract between your model code and the rest of the framework. The UI wants small pieces as early as possible. The session store wants exactly one complete record of what the model said. The tool dispatcher wants each function call exactly once, with all its arguments.
Other pages on this site cover the consumer side: the runner's Flowable of Event and the event shapes a gateway turns into Server-Sent Events. This one is about the producer side: what a streaming generateContent must emit, in what order, and why. It works through the protocol, the two framework checks that make the final response load-bearing, an adapter for an HTTP streaming provider, tool calls, cancellation, errors, and contract tests. API names were checked against the google/adk-java source on the main branch in October 2026. Builder method names have changed between releases, so confirm them against the version you build with.
Where the stream flag comes from
The stream flag is not something an agent author sets on the model. In the source we read, BaseLlmFlow calls llm.generateContent(finalLlmRequest, context.runConfig().streamingMode() == StreamingMode.SSE). RunConfig.StreamingMode has three values: NONE (the default), SSE and BIDI. Only SSE asks your model to stream. BIDI takes a different path entirely, through BaseLlm.connect(LlmRequest), which returns a BaseLlmConnection for live audio and real-time sessions. So a model has three behaviours to get right: one response when the flag is false, a well-formed sequence when it is true, and either a working live connection or an explicit refusal from connect.
Each LlmResponse the model emits becomes an Event. The flow copies content, partial, errorCode, turnComplete and the other fields onto the event and pushes it down the runner's stream, where everything downstream acts on those fields.
| LlmResponse field | Type | What a streaming producer should do with it |
|---|---|---|
content() | Optional<Content> | Partial: the new delta only. Final: every part of the turn, in order. |
partial() | Optional<Boolean> | true on every delta; unset or false on the single final response. |
finishReason() | Optional<FinishReason> | Set on the final response from the provider's stop reason. |
usageMetadata() | usage metadata | Set once, on the final response, so tokens are not counted twice. |
errorCode() / errorMessage() | finish reason, string | Set on a terminal response when the provider reports a failure in-band. |
turnComplete() / interrupted() | Optional<Boolean> | Mainly meaningful on live connections; the flow copies them without acting on them. |
The emission protocol: deltas, then one aggregate
The bundled Gemini model keeps a text buffer while chunks arrive, emits partial(true) responses as they come, and at the end flushes the buffer and any function calls into one aggregated Content with role model, with partial unset. The protocol recommended here for a custom model is that shape made strict. Emit zero or more partials, each carrying only new text. Then emit exactly one non-partial response with the whole turn: all text, any function calls, finish reason and usage.
Two invariants follow. The partials' text concatenates to the final text, and only the final carries function calls and usage. Hold both and the UI can append then replace, history gets one record, and each tool runs once.
Why the final response is load-bearing
The final aggregate is not a courtesy. Two checks in the framework depend on it. First, BaseSessionService.appendEvent starts with if (event.partial().orElse(false)) return Single.just(event);. Partial events are passed through without being stored. Second, BaseLlmFlow executes function calls only when the event is not partial: its condition skips execution when modelResponseEvent.functionCalls().isEmpty() || modelResponseEvent.partial().orElse(false).
Now consider the three ways a hand-written streaming model gets this wrong. A model that marks every chunk partial and never sends a final response streams perfectly to the screen and leaves no trace in the session. The next turn's LlmRequest is built from history, so the model forgets what it just said, and any tool call it made is never executed. A model that marks every chunk as final stores each fragment as a separate event. History fills with slivers like "The order" and " ships Tuesday", and a function call split across chunks may run with half its arguments. A model whose final response holds only the last delta, rather than the whole turn, stores a truncated answer. All three look fine in a demo.
An adapter over an SSE provider
Here is a streaming model over an HTTP provider that sends Server-Sent Events, one data: line per chunk and a [DONE] sentinel at the end. That framing is common to OpenAI-style chat APIs, but treat it as an assumption and check your provider's documentation. The JSON parsing is hidden behind a small ChunkParser interface so the streaming logic stays readable. Request mapping is covered in implementing a custom LLM and is not repeated.
public final class SseChatModel extends BaseLlm {
private final HttpClient http = HttpClient.newHttpClient();
private final RequestMapper mapper; // LlmRequest -> provider JSON (yours)
private final ChunkParser parser; // provider JSON -> Chunk (yours)
public SseChatModel(String model, RequestMapper mapper, ChunkParser parser) {
super(model);
this.mapper = mapper;
this.parser = parser;
}
@Override
public Flowable<LlmResponse> generateContent(LlmRequest req, boolean stream) {
if (!stream) {
return Flowable.fromCallable(() -> callOnce(req)); // one response, partial unset
}
return Flowable.using(
() -> http.send(mapper.toHttpRequest(req, true),
HttpResponse.BodyHandlers.ofLines()).body(), // Stream<String>
lines -> Flowable.defer(() -> {
StreamAggregator agg = new StreamAggregator();
return Flowable.fromStream(lines)
.filter(l -> l.startsWith("data:"))
.map(l -> l.substring(5).trim())
.takeWhile(d -> !d.equals("[DONE]"))
.map(parser::parse)
.concatMapIterable(agg::accept) // zero or one partial per chunk
.concatWith(Flowable.fromCallable(agg::finish)); // exactly one final
}),
Stream::close) // runs on complete, error and cancel
.subscribeOn(Schedulers.io());
}
@Override
public BaseLlmConnection connect(LlmRequest req) {
throw new UnsupportedOperationException("live (BIDI) sessions are not supported by " + model());
}
}Flowable.using ties the HTTP body's lifetime to the subscription: the stream is closed when the Flowable completes, fails or is cancelled. subscribeOn(Schedulers.io()) keeps the blocking send and line reads off the caller's thread. The defer creates a fresh aggregator per subscription, so a retry never inherits a half-filled buffer from the attempt that failed.
The aggregator: text, tool calls and usage
The aggregator is where the protocol lives. It turns each provider chunk into at most one partial response, and keeps everything it will need for the final one.
final class StreamAggregator {
private final StringBuilder text = new StringBuilder();
private final Map<Integer, ToolCallBuffer> calls = new TreeMap<>(); // provider's call index
private String finishReason;
private Usage usage;
List<LlmResponse> accept(Chunk ch) {
ch.toolCallFragments().forEach(f ->
calls.computeIfAbsent(f.index(), i -> new ToolCallBuffer()).add(f));
if (ch.finishReason() != null) finishReason = ch.finishReason();
if (ch.usage() != null) usage = ch.usage();
if (ch.textDelta() == null || ch.textDelta().isEmpty()) return List.of();
text.append(ch.textDelta());
return List.of(LlmResponse.builder()
.content(Content.builder().role("model").parts(Part.fromText(ch.textDelta())).build())
.partial(true)
.build());
}
LlmResponse finish() {
List<Part> parts = new ArrayList<>();
if (text.length() > 0) parts.add(Part.fromText(text.toString()));
for (ToolCallBuffer b : calls.values()) {
parts.add(Part.fromFunctionCall(b.name(), b.parseArgs())); // throws on malformed JSON
}
LlmResponse.Builder out = LlmResponse.builder()
.content(Content.builder().role("model").parts(parts).build());
if (finishReason != null) out.finishReason(Mappers.finishReason(finishReason));
if (usage != null) out.usageMetadata(Mappers.usage(usage));
return out.build(); // partial left unset
}
}Notice what it does not do. It never emits a partial for a tool-call fragment. Providers stream function-call arguments as pieces of a JSON string, keyed by an index, and a fragment like {"order_id": "A1 is not a function call. The buffer concatenates fragments per index and parses only in finish. If parsing fails there, throwing is correct: a turn with a corrupt tool call should fail loudly, not run a tool with guessed arguments. Usage is attached once, to the final response, so cost dashboards that sum usage across events do not count the turn twice.
Backpressure, threads and cancellation
Backpressure comes for free when you build the stream from Flowable operators rather than Flowable.create with a hand-rolled emitter. fromStream pulls one line per request, so a slow consumer slows the reads and the provider's TCP window fills. Chunks are small and the gaps between them are set by the model's generation speed, so buffering is rarely the problem. Holding a thread is. With Schedulers.io(), each open stream holds one thread for its whole duration, and a two-minute generation holds it for two minutes. Size the pool for concurrent streams, not requests per second, or move to an asynchronous body handler if you serve many long generations at once. The runner-level streaming article covers which thread delivers events to the caller.
Cancellation is the user closing the tab. The gateway disposes its subscription, the disposal travels up through the runner and flow to your Flowable, and using closes the line stream, which closes the connection. Whether generation and billing then stop varies by provider, so check yours. A model that copies chunks into its own unbounded queue on a separate thread breaks this chain, and keeps paying for tokens nobody will read.
Errors before and after the first chunk
Streams fail at two different moments, and they need different handling. Before the first chunk, nothing has reached the user. A connection failure, an HTTP 429 or a 5xx can be retried with backoff, exactly like a non-streaming call. Put the retry on the whole using block, which is safe because defer gives each attempt a fresh aggregator:
return streamOnce(req)
.retryWhen(errors -> errors.zipWith(Flowable.range(1, 3), (e, n) -> {
if (!(e instanceof RetryableBeforeFirstChunk) || n == 3) throw e;
return n;
}).flatMap(n -> Flowable.timer(200L << n, TimeUnit.MILLISECONDS)));After the first chunk, the user has seen text. Retrying silently would show the start of the answer twice, or splice two different generations together. Fail instead. If the provider reports the failure in-band, for example a content-filter stop or a token limit, emit a final response with errorCode and errorMessage set and whatever text arrived, so the turn is recorded with its reason. If the transport dies, let onError propagate. The runner's Flowable errors, the gateway can tell the client to discard the partial text, and the session holds nothing for the turn, which is honest. To tell the two cases apart, wrap the transport exception in RetryableBeforeFirstChunk only while the aggregator has emitted nothing.
Worked example: a streamed turn with a tool call
Take one turn where the user asks where order A17 is. The agent has a get_order tool, the model explains briefly and then calls it. The provider sends six data lines. The table shows what the model emits and what each part of the framework does.
| Provider chunk | LlmResponse emitted | UI | Tools | Session |
|---|---|---|---|---|
| text "Let me " | partial(true), text "Let me " | appends | skipped | skipped |
| text "check." | partial(true), text "check." | appends | skipped | skipped |
| tool #0 name get_order, args {"order_ | none | |||
| tool #0 args id":"A17"} | none | |||
| finish tool_calls, usage | none (buffered) | |||
| [DONE] | final: text "Let me check." + functionCall get_order {order_id: A17} | replaces | runs once | one event |
The session ends up with one model event holding both the text and the call. The flow then runs get_order, appends the function response, and calls generateContent again for the model's answer, which streams the same way. Two model calls, two final events, one tool event: that is what history should look like after a streamed tool-using turn.
Testing the contract
Watching text appear proves nothing about the contract. Test the sequence. RxJava's test() subscriber makes the invariants assertions. Feed the model recorded provider transcripts, including one with a split tool call, through a stub HTTP server or a fake line source. Testing custom LLMs covers the scripted-model side.
@Test
void streamFollowsProtocol() {
List<LlmResponse> out = model.generateContent(request, true).test()
.awaitDone(5, TimeUnit.SECONDS).assertComplete().values();
LlmResponse last = out.get(out.size() - 1);
assertFalse(last.partial().orElse(false), "last response must be the aggregate");
out.subList(0, out.size() - 1).forEach(r ->
assertTrue(r.partial().orElse(false), "only the last response may be final"));
String deltas = out.subList(0, out.size() - 1).stream()
.map(r -> textOf(r)).collect(Collectors.joining());
assertEquals(deltas, textOf(last), "partials must concatenate to the final text");
out.subList(0, out.size() - 1).forEach(r -> {
assertTrue(functionCallsOf(r).isEmpty(), "no tool calls on partials");
assertTrue(r.usageMetadata().isEmpty(), "usage only on the final");
});
}Add three more cases: cancel after the second partial and assert the stub server saw the connection close; fail the transport after one chunk and assert onError with no retry; and return 429 before any chunk and assert exactly one retry. Together they catch the protocol bugs listed below.
Failure modes
- No final response. Streams to screen, missing from history, tools never run. Caught by the last-is-final assertion.
- Every chunk final. History full of fragments, tools run with partial arguments. Caught by the only-last-is-final assertion.
- Final holds the last delta, not the whole turn. Truncated history. Caught by the concatenation assertion.
- Usage on every chunk. Token and cost dashboards overcount. Attach usage once.
- Shared aggregator across retries. A retried stream repeats the first attempt's text. Create it inside
defer. - Silent retry after the first chunk. Duplicated or spliced answers. Retry only before any partial is emitted.
What to do next
- Read your model's
generateContentand write down, for the streaming branch, exactly whenpartial(true)is set and where the final aggregate is built. - Add the protocol test above against a recorded transcript that includes a split tool call.
- Check that usage is attached only to the final response, then compare a day of token totals before and after.
- Make sure the HTTP body is closed on cancel: wrap it in
Flowable.usingand assert closure in a test. - Split errors into before-first-chunk (retry with backoff) and after (fail, or emit an
errorCodefinal). - Implement
connector make it throw a clearUnsupportedOperationException, so aBIDIrun config fails fast. - Load-test with long generations and size the I/O pool for concurrent open streams.