Skip to main content

Non-blocking and Reactive

note

Asynchronous and reactive support is experimental. The asynchronous and reactive APIs on this page are annotated @Experimental; APIs and behavior may still change in future releases. The synchronous and TokenStream APIs are unaffected.

By default, calling an AI Service blocks the calling thread for the whole interaction: the LLM call, tool executions, chat memory access and guardrails all happen before the method returns. That is simple and works well for most applications.

If you are building on a reactive stack (Quarkus/Mutiny, Vert.x, Spring WebFlux) or you need one thread to drive many concurrent interactions, LangChain4j can run the same interaction without blocking. You choose by declaring a different return type on your AI Service method — nothing else about the interface changes.

The four modes

Return typeNature
String, a POJO, Result<T>, …synchronous, blocks the calling thread
TokenStreamstreaming via callbacks
CompletableFuture<T>, CompletionStage<T>one response, non-blocking
Flow.Publisher<AiServiceStreamingEvent>, Flow.Publisher<String>a stream of events, non-blocking
interface Assistant {

// synchronous
String chat(String message);

// one response, without blocking the caller
CompletableFuture<String> chatAsync(String message);

// the answer, streamed token by token
Flow.Publisher<String> chatStreaming(String message);

// everything that happens during the interaction, as events
Flow.Publisher<AiServiceStreamingEvent> chatEvents(String message);
}

Getting one response without blocking

CompletableFuture<String> future = assistant.chatAsync("Tell me a joke");

future.thenAccept(System.out::println);

Cancelling the future (future.cancel(true)) releases the caller, stops further LLM rounds and makes a best-effort attempt to abort the in-flight HTTP call.

note

The future is completed on the model's own thread — for HTTP models, the transport's I/O worker that reads the response. A continuation attached without an explicit executor (thenApply, thenAccept, …) therefore runs on that thread, where blocking degrades throughput for every in-flight call. Keep continuations non-blocking, or pass your own Executor: future.thenApplyAsync(fn, executor).

Streaming the answer

A method returning Flow.Publisher<String> streams the text of the answer, chunk by chunk:

Flow.Publisher<String> publisher = assistant.chatStreaming("Tell me a joke");

The publisher is cold: nothing happens until you subscribe, and each subscription starts a new interaction.

Streaming everything that happens

Flow.Publisher<AiServiceStreamingEvent> surfaces the whole interaction, not just the answer text — the same information the callback-based TokenStream gives you:

EventMeaning
PartialResponseEventa chunk of the answer
PartialThinkingEventa chunk of the model's reasoning
PartialToolCallEvent, CompleteToolCallEventa tool call being assembled, then completed
BeforeToolExecutionEvent, AfterToolExecutionEventa tool about to run, and its result
IntermediateResponseEventthe response that closed one tool-calling round
RetrievedContentsEventcontent retrieved by RAG
ToolCompensatedEventa completed tool was compensated after the interaction was cancelled or failed
RawEventa provider-specific event LangChain4j does not map
FinalResponseEventthe final answer
assistant.chatEvents("What is the weather in Munich?").subscribe(new Flow.Subscriber<>() {

@Override
public void onSubscribe(Flow.Subscription subscription) {
subscription.request(Long.MAX_VALUE);
}

@Override
public void onNext(AiServiceStreamingEvent event) {
// the event types are nested in AiServiceStreamingEvent:
// import dev.langchain4j.service.AiServiceStreamingEvent.PartialResponseEvent;
if (event instanceof PartialResponseEvent partial) {
System.out.print(partial.partialResponse().text());
}
}

@Override
public void onError(Throwable error) { }

@Override
public void onComplete() { }
});

New event types may be added over time, so handle unrecognized ones gracefully — do not write an exhaustive type switch without a default branch.

note

Do not block in onNext. Events are delivered on whichever thread produced them: the model's transport I/O worker for the token-level events, or the thread that completed a tool call for the tool-execution events. Offload heavy per-event work to your own Executor. Events are relayed through a bounded buffer, so a subscriber that falls far enough behind terminates with an IllegalStateException instead of buffering without limit; the size defaults to 16384 events and is configurable with AiServices.builder(...).streamingBufferSize(int).

Third-party reactive types

AI Service methods return JDK types — CompletableFuture and Flow.Publisher — so the API does not tie you to any particular reactive programming library.

Reactor types are supported through the langchain4j-reactor module — Mono<T> for the single-response mode and Flux<AiServiceStreamingEvent> for the reactive one. Adding the dependency is enough; the adapters register themselves via ServiceLoader.

interface Assistant {

Mono<String> answer(String userMessage);

Flux<AiServiceStreamingEvent> events(String userMessage);
}
note

Mutiny's Uni/Multi are not supported yet. The seam is there — the CompletableFutureAdapter and PublisherAdapter SPIs, discovered via ServiceLoader — but no adapter ships for them, so declaring such a return type currently fails. Use CompletableFuture or Flow.Publisher and adapt at the call site.

Flux<String> is served by the older TokenStream-based adapter in the same module — see AI Services. It predates this feature, is not part of the non-blocking path described here, and works with every provider.

What has to be non-blocking

Non-blocking has to hold at every layer: a blocking step anywhere makes the whole call blocking again. Each layer has an asynchronous counterpart alongside its existing blocking method:

LayerBlockingNon-blocking counterpart
Chat modelchat(ChatRequest)chatAsync(ChatRequest), and chat(ChatRequest) returning a Publisher on StreamingChatModel
Embedding modelembed(...)embedAsync(...)
Scoring modelscoreAll(...)scoreAsync(...)
Chat memoryadd, messages, setaddAsync, messagesAsync, setAsync
Chat memory storegetMessages, updateMessages, deleteMessagesgetMessagesAsync, updateMessagesAsync, deleteMessagesAsync
Guardrailsvalidate(...)validateAsync(...)
ToolsToolExecutor.execute(...)ToolExecutor.executeAsync(...)
RAGRetrievalAugmentor.augment(...), retrievers, routers, aggregators, query transformersaugmentAsync(...), retrieveAsync(...), routeAsync(...), aggregateAsync(...), transformAsync(...)
Embedding storesearch(...)searchAsync(...)
Web searchsearch(...)searchAsync(...)
MCPMcpClient.executeTool(...)executeToolAsync(...)
Guardrail executionChatExecutor.execute(...)ChatExecutor.executeAsync(...)
HTTP clientexecute(...)executeAsync(...), stream(...)

A component that has not implemented its counterpart fails loudly rather than silently blocking: the returned future or publisher fails with an AsyncNotSupportedException naming the component and the missing method. AsyncNotSupportedException is an UnsupportedFeatureException, so one catch clause covers both flavours of "not supported", and it is never retried.

Blocking code you cannot avoid

Tools are the common case: a tool that calls a database or a blocking HTTP API cannot be made non-blocking. Such tools are offloaded so they never block the model's own thread.

note

On Java 21 and later the offload executor creates virtual threads, so a blocking tool parks a virtual thread, which unmounts from its carrier instead of occupying it. LangChain4j targets Java 17, and on Java 17-20 the same executor is an unbounded pool of platform threads — a very different resource profile under load. Supply your own bounded Executor (see below) if that matters to you.

This applies to @Tool-annotated methods. A hand-written ToolExecutor that does not override executeAsync fails loudly with an AsyncNotSupportedException instead, like any other component that has not opted in.

Tools run concurrently by default in the asynchronous and reactive modes. To run them one at a time, pass a single-threaded executor:

AiServices.builder(Assistant.class)
.chatModel(model)
.tools(new MyTools())
.executeToolsConcurrently(Executors.newSingleThreadExecutor())
.build();

For RAG, a retriever that has not implemented retrieveAsync fails by default rather than silently blocking. Opt in to offloading it instead with offloadBlocking(true) on DefaultRetrievalAugmentor or EmbeddingStoreContentRetriever.

Defaults that differ from the synchronous modes

The asynchronous and reactive modes are not just the same interaction on another thread — a few defaults are deliberately different. If you are migrating an existing method, these are the ones to check.

Synchronous / TokenStreamCompletableFuture / Flow.Publisher
Multiple tool callsexecuted sequentiallyexecuted concurrently
Tool execution errorsent back to the LLMfails the invocation
Tool argument-parse errorfails the invocationsent back to the LLM
@Moderatesupportedrejected at AI Service creation

Tool errors

The two tool error defaults are reversed on purpose. Sending an execution failure to the LLM hides a bug in your tool from you and invites the model to invent an answer around it, so the asynchronous modes fail the invocation instead. A malformed argument string, on the other hand, is something the model produced and can usually fix when told, so it is sent back rather than failing the call.

Both remain configurable, and an explicitly configured handler is used by every mode:

AiServices.builder(Assistant.class)
.chatModel(model)
.tools(new MyTools())
.toolExecutionErrorHandler(myExecutionErrorHandler)
.toolArgumentsErrorHandler(myArgumentsErrorHandler)
.build();

See Error Handling for what those handlers can do.

Tool concurrency

Because tools always run on an Executor in these modes, several tool calls in one LLM response run at the same time. If your tools are not safe to run concurrently, pass a single-threaded executor: executeToolsConcurrently(Executors.newSingleThreadExecutor()).

Event ordering

On the reactive path a tool starts as soon as its CompleteToolCallEvent is emitted, so it overlaps the rest of the round. A round's BeforeToolExecutionEvent / AfterToolExecutionEvent may therefore arrive before that round's IntermediateResponseEvent. Treat IntermediateResponseEvent as a per-round marker, not as a barrier that all of the round's tool events precede. The callback-based TokenStream always reports the intermediate response first.

Controlling the executor and propagating context

Everywhere LangChain4j runs work off the caller thread — concurrent tool calls, offloaded retrieval, retry backoff — it takes the executor from one pluggable seam, the ExecutorProvider SPI:

public interface ExecutorProvider {
Executor executor();
}

Register one via ServiceLoader, or programmatically for tests and non-DI applications:

ExecutorProvider.set(() -> myExecutor);

Without one, LangChain4j uses a virtual-thread-per-task executor on Java 21 or later, and an unbounded platform-thread pool on Java 17-20.

note

A single invocation now crosses several threads, so ambient ThreadLocal state — MDC logging context, tracing spans, security context — is not automatically propagated the way it is in the fully synchronous mode. Return a context-propagating executor from your provider to make it follow the work: a ManagedExecutor on Quarkus/MicroProfile, a TaskDecorator-wrapped executor on Spring, Context.taskWrapping(executor) for OpenTelemetry, or ContextSnapshot.wrap(executor) for Micrometer. On Spring Boot you can point LangChain4j at the application's executor with one property instead — see Spring Boot below.

InvocationContext is unaffected — it is passed explicitly as a parameter, never through a thread-local.

Model listeners must not block

ChatModelListener and EmbeddingModelListener callbacks are invoked synchronously on the model's own threads and are never offloaded. On the asynchronous and reactive APIs that means the transport's I/O worker. A listener that performs blocking I/O there stalls that worker and degrades throughput for every in-flight call — offload such work to your own executor from inside the callback. See Observability.

Provider support

warning

The non-blocking modes only work with providers that have implemented them. Declaring a CompletableFuture or Flow.Publisher return type against any other provider compiles, and then fails at runtime with an AsyncNotSupportedException naming the component and the missing method. There is no silent fallback to a blocking call — that is the point, but it means the return type you can use depends on your provider.

Support is opt-in per provider and is being rolled out gradually. Today:

ProviderCompletableFuture (chatAsync)Flow.Publisher (reactive chat)
OpenAI — Chat Completions
OpenAI — Responses
Anthropic
Bedrock
every other provider

Beyond chat models: OpenAI implements asynchronous embeddings (embedAsync), Cohere asynchronous scoring (scoreAsync) and Tavily asynchronous web search (searchAsync).

A provider that has not opted in is unaffected on the synchronous and TokenStream APIs — those keep working exactly as before.

The same applies to everything else in the call

A chat model that supports the mode is necessary but not sufficient: a single blocking component anywhere in the interaction fails the call the same way. In practice that means:

ComponentWhat it must implementBundled implementations
Chat memoryaddAsync / messagesAsync / setAsyncMessageWindowChatMemory and InMemoryChatMemoryStore already do
GuardrailsvalidateAsyncnone — a guardrail you write must override it, even if it does no I/O
Content retriever, query router, aggregatorretrieveAsync, routeAsync, aggregateAsyncopt into offloading instead with offloadBlocking(true)
Embedding storesearchAsyncas above, via the retriever's offloadBlocking(true)
Custom ToolExecutorexecuteAsync@Tool-annotated methods are offloaded for you

A guardrail that does no blocking work satisfies the contract in one line:

@Override
public CompletableFuture<InputGuardrailResult> validateAsync(InputGuardrailRequest request) {
return CompletableFuture.completedFuture(validate(request));
}

Spring Boot

Flux<String> keeps working with every provider, including those in the ❌ row: it is served by the TokenStream-based adapter in langchain4j-reactor, not by the non-blocking path described on this page.

Mono<T> and Flux<AiServiceStreamingEvent> come from the langchain4j-reactor module and carry the same provider constraint as the JDK types — see Third-party reactive types above.

To make ambient context follow an asynchronous invocation, let LangChain4j offload to the application's own executor rather than its default one:

langchain4j.executor.use-spring-task-executor=true

Spring's task executor propagates tracing spans, MDC and security context when a TaskDecorator is installed (Micrometer context propagation and Spring Security both install one), and its pool follows spring.task.execution.*. It is off by default because the setting is process-wide, not scoped to one application context.