How to Implement Streaming Responses in Embabel: A Complete Guide
To implement streaming responses in Embabel, wrap your PromptRunner with StreamingPromptRunnerBuilder and invoke generateStream() to receive incremental StreamingEvents as the LLM generates tokens.
The embabel/embabel-agent framework provides a dedicated streaming layer that enables agents to consume incremental LLM output—tokens, chunks, or events—instead of waiting for complete responses. Learning how to implement streaming responses in Embabel allows you to build low-latency applications that display AI-generated content in real time. This guide covers the core architecture, implementation patterns, and observability hooks based on the actual source code.
Architecture of the Embabel Streaming Layer
Embabel’s streaming architecture separates concerns between capability detection, runner construction, and event emission. The design ensures that streaming is only attempted when the underlying LLM provider supports it.
Core Components
The framework relies on several key interfaces and classes located in the embabel-agent-api module:
StreamingPromptRunnerBuilder– A Java-friendly façade located atembabel-agent-api/src/main/java/com/embabel/agent/api/streaming/StreamingPromptRunnerBuilder.javathat validates streaming capability and returns aStreamingPromptRunner.Streaminginstance.StreamingPromptRunner– The core Kotlin interface atembabel-agent-api/src/main/kotlin/com/embabel/agent/api/common/streaming/StreamingPromptRunner.ktexposinggenerateStream(...)andcreateObjectStream(...)methods.StreamingCapabilityDetector– Runtime detector atembabel-agent-api/src/main/kotlin/com/embabel/agent/spi/support/streaming/StreamingCapabilityDetector.ktthat checkssupportsStreamingon thePromptRunner.StreamingLlmOperationsFactory– Factory atembabel-agent-api/src/main/kotlin/com/embabel/agent/core/internal/streaming/StreamingLlmOperationsFactory.ktthat creates provider-specific implementations.StreamingLlmOperationsImpl– Concrete implementation atembabel-agent-api/src/main/kotlin/com/embabel/agent/spi/support/streaming/StreamingLlmOperationsImpl.ktthat forwards requests to Spring AI and converts raw bytes intoStreamingEvents.StreamingToolLoop– Enables tool execution atembabel-agent-api/src/main/kotlin/com/embabel/agent/spi/loop/streaming/StreamingToolLoop.ktallowing tool calls to interleave with token streams.
Data Flow
When you implement streaming responses in Embabel, the framework executes the following sequence:
- Configuration – An LLM provider is configured (e.g., OpenAI
bestmodel). - Capability detection –
StreamingCapabilityDetectorverifiessupportsStreamingon thePromptRunner. - Builder usage –
new StreamingPromptRunnerBuilder(runner).streaming()returns a typed streaming object. - Streaming request – You call
generateStream(messages, options, outputClass, ...)orcreateObjectStream(...). - Factory creation –
StreamingLlmOperationsFactorybuilds aStreamingLlmOperationsinstance that wraps the Spring AIChatClient. - Event emission –
StreamingLlmOperationsImplwraps each token/chunk in aStreamingEventand emits it to your consumer. - Tool loop – If tools are invoked,
StreamingToolLooppauses the stream, executes the tool, then resumes. - Observability –
EmbabelSpanEventListenerandEmbabelMetricsEventListenercapture timestamps, token counts, and errors.
Implementing Streaming Responses in Java
To implement streaming responses in Embabel in a Java application, obtain a PromptRunner instance (typically injected by Spring), then use the builder pattern to create a streaming runner.
import com.embabel.agent.api.common.PromptRunner;
import com.embabel.agent.api.streaming.StreamingPromptRunnerBuilder;
import com.embabel.common.core.streaming.StreamingEvent;
import java.util.List;
public class StreamingDemo {
public static void main(String[] args) {
// 1️⃣ Obtain a PromptRunner (typically injected by Spring)
PromptRunner runner = ...; // e.g., autowired PromptRunner bean
// 2️⃣ Build a streaming runner
var streaming = new StreamingPromptRunnerBuilder(runner).streaming();
// 3️⃣ Prepare the prompt
List<Message> messages = List.of(
Message.ofUser("Explain the concept of quantum entanglement in simple terms.")
);
// 4️⃣ Consume the stream
streaming.generateStream(messages, null, String.class, null, event -> {
// This consumer is called for every chunk/token
if (event instanceof StreamingEvent.Content content) {
System.out.print(content.content()); // incremental text
} else if (event instanceof StreamingEvent.Error err) {
System.err.println("Stream error: " + err.throwable().getMessage());
} else if (event instanceof StreamingEvent.Complete) {
System.out.println("\n--- Stream finished ---");
}
});
}
}
Step-by-step explanation:
- Step 1 –
PromptRunnerabstracts the LLM client and exposes whether streaming is supported. - Step 2 –
StreamingPromptRunnerBuilderchecksrunner.supportsStreaming(); if unsupported, it throwsUnsupportedOperationException. - Step 3 – Construct a list of
Messageobjects to establish conversation context. - Step 4 –
generateStreamreturns a cold stream; your lambda receivesStreamingEventsubtypes (Content,Error,Complete) as implemented inStreamingLlmOperationsImpl.
Handling Tools During Streaming
For agents that require tool use—such as code execution or retrieval-augmented generation—the StreamingToolLoop class enables tool calls to occur mid-stream. When the LLM requests a tool invocation, the loop pauses token emission, executes the tool, appends the result to the context, and resumes streaming. This interleaving happens transparently when you use the streaming runner, provided the tools are configured in your PromptRunner options.
Observability and Monitoring
Embabel automatically instruments streaming operations through listeners in the embabel-agent-observability module:
EmbabelMetricsEventListener– Located atembabel-agent-observability/src/main/java/com/embabel/agent/observability/metrics/EmbabelMetricsEventListener.java, records token counts and stream durations.EmbabelSpanEventListener– Located atembabel-agent-observability/src/main/java/com/embabel/agent/observability/tracing/EmbabelSpanEventListener.java, creates distributed tracing spans for each streaming event.
These hooks activate automatically when you invoke generateStream, providing fine-grained latency metrics without additional configuration.
Streaming vs. Blocking: When to Use Each
Choose the appropriate mode based on your use case:
| Use Case | Recommended Mode |
|---|---|
| Large responses (multi-paragraph explanations, code generation) | Streaming – reduces time-to-first-token and allows incremental processing. |
| Simple short replies (yes/no answers) | Blocking – use runner.run(...) for simpler API surface. |
| Tool-driven workflows (RAG, function calling) | Streaming with StreamingToolLoop – supports mid-stream tool requests. |
| Observability-heavy environments | Streaming – each token is traced for granular latency analysis. |
Summary
- Wrap your
PromptRunnerwithStreamingPromptRunnerBuilderto validate capabilities and obtain a streaming instance. - Invoke
generateStream()orcreateObjectStream()to receive real-timeStreamingEvents. - Handle
StreamingEvent.Content,StreamingEvent.Error, andStreamingEvent.Completein your consumer lambda. - Use
StreamingToolLoopfor complex workflows requiring tool invocation during generation. - Observability is automatic via
EmbabelMetricsEventListenerandEmbabelSpanEventListener.
Frequently Asked Questions
How do I check if my LLM supports streaming in Embabel?
Embabel performs this check automatically via StreamingCapabilityDetector at embabel-agent-api/src/main/kotlin/com/embabel/agent/spi/support/streaming/StreamingCapabilityDetector.kt. When you call new StreamingPromptRunnerBuilder(runner).streaming(), the builder invokes runner.supportsStreaming(); if the LLM configuration lacks streaming support, it throws UnsupportedOperationException immediately.
What is the difference between generateStream and createObjectStream?
Both methods are defined in StreamingPromptRunner.kt at embabel-agent-api/src/main/kotlin/com/embabel/agent/api/common/streaming/StreamingPromptRunner.kt. Use generateStream for plain text or token streaming, where you receive StreamingEvent.Content chunks. Use createObjectStream when you expect structured output (e.g., JSON objects), where the framework attempts to parse and emit partial or complete objects as they arrive.
Can I use tools with streaming responses?
Yes. The StreamingToolLoop class at embabel-agent-api/src/main/kotlin/com/embabel/agent/spi/loop/streaming/StreamingToolLoop.kt enables tool invocation while a stream is active. When the LLM generates a tool call request, the loop pauses the stream, executes the tool, and resumes generation with the tool result appended to the context, all without terminating the connection.
Where is the core streaming implementation located?
The primary implementation that interacts with Spring AI and emits events is StreamingLlmOperationsImpl.kt, located at embabel-agent-api/src/main/kotlin/com/embabel/agent/spi/support/streaming/StreamingLlmOperationsImpl.kt. This class converts raw byte streams from the LLM provider into typed StreamingEvent objects consumed by your application.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →