Streaming SSE with Spring WebFlux Permalink to this section
Part of Java & Spring SSE Implementation, under Backend Stream Generation & Connection Management.
Returning Flux<ServerSentEvent<T>> from a WebFlux controller produces a working event stream in five lines. A production stream needs more: replay from Last-Event-ID, filtering so users see only their own events, heartbeats, a deliberate policy for subscribers that fall behind, and cleanup when they leave — all without a single blocking call on the event loop. This guide assembles those pieces into one endpoint.
Symptom & Developer Intent Permalink to this section
Teams usually reach for this guide after a first version shows one of these:
- Under load, all streams on an instance freeze for seconds at a time, then resume together.
- Memory grows steadily with connected clients even though event volume is flat.
- Clients that reconnect miss everything that happened while they were away.
- Idle streams are closed by the load balancer after 60 seconds.
Sinks.many()throwsEmitFailureHandlererrors, or events silently vanish, once more than one subscriber is attached.
The intent is a reactive endpoint that never blocks, has bounded memory per subscriber, resumes correctly, and cleans up after itself.
Root Cause Analysis Permalink to this section
Most WebFlux stream failures trace back to one of three misunderstandings.
- Blocking on the event loop. WebFlux runs a handful of event-loop threads, each serving thousands of connections. A JDBC query, a blocking HTTP client or
Thread.sleepinsidemaporflatMapstalls every stream on that thread. - Unbounded buffering. A subscriber that reads slower than events are produced accumulates items in Reactor’s queues unless an operator says otherwise. The default for many sinks and operators is to buffer, which becomes a per-subscriber memory leak.
- The wrong sink.
Sinks.many().unicast()accepts one subscriber; areplay()sink holds its entire history unless bounded; amulticast().onBackpressureBuffer()sink applies one shared buffer to all subscribers, so the slowest one determines what everyone else can receive.
Step-by-Step Resolution Permalink to this section
Step 1 — Choose a sink that isolates subscribers Permalink to this section
@Component
class EventHub {
// directBestEffort: each subscriber receives what it has demand for; others are unaffected.
private final Sinks.Many<AppEvent> sink = Sinks.many().multicast().directBestEffort();
private final AtomicLong seq = new AtomicLong();
void publish(String userId, String type, Object payload) {
AppEvent e = new AppEvent(seq.incrementAndGet(), userId, type, payload);
replayStore.append(e); // durable, bounded (step 3)
sink.emitNext(e, Sinks.EmitFailureHandler.busyLooping(Duration.ofMillis(50)));
}
Flux<AppEvent> live() { return sink.asFlux(); }
}
busyLooping retries briefly if two threads emit at once, which Sinks does not allow concurrently. Per-subscriber buffering is added downstream in step 4, where each subscriber gets its own policy.
Step 2 — Keep blocking work off the loop Permalink to this section
If an event needs enrichment from a blocking source, move that call to a scheduler built for blocking:
Flux<AppEvent> enriched = hub.live()
.flatMap(e -> Mono.fromCallable(() -> jdbcLookup(e)) // blocking call…
.subscribeOn(Schedulers.boundedElastic()), 16); // …off the event loop, max 16 in flight
Better still, enrich once in publish() before emitting, so the work happens once per event rather than once per subscriber.
Step 3 — Replay from Last-Event-ID, then follow live Permalink to this section
@GetMapping(path = "/api/events", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
Flux<ServerSentEvent<Object>> events(@AuthenticationPrincipal Jwt jwt,
@RequestHeader(value = "Last-Event-ID", required = false) Long lastId) {
String user = jwt.getSubject();
Flux<AppEvent> live = hub.live().filter(e -> e.userId().equals(user));
Flux<AppEvent> replay = lastId == null
? Flux.empty()
: replayStore.after(user, lastId); // reactive query, bounded window
// Subscribe to live first (cached), then emit replay, then live — dropping any overlap.
Flux<AppEvent> liveCached = live.replay(512).autoConnect(1);
AtomicLong highest = new AtomicLong(lastId == null ? 0 : lastId);
Flux<AppEvent> ordered = liveCached.take(0) // triggers the live subscription now
.thenMany(Flux.concat(replay, liveCached))
.filter(e -> e.seq() > highest.get())
.doOnNext(e -> highest.set(e.seq()));
return ordered.map(this::toSse);
}
private ServerSentEvent<Object> toSse(AppEvent e) {
return ServerSentEvent.builder().id(Long.toString(e.seq())).event(e.type()).data(e.payload()).build();
}
Connecting the live subscription before the replay query closes the window in which an event could be published after the query ran and before the live subscription existed. The monotonic filter removes duplicates where the two overlap.
Step 4 — Add backpressure, heartbeats and cleanup Permalink to this section
return Flux.merge(
ordered.onBackpressureBuffer(256,
dropped -> log.debug("dropping for slow subscriber"),
BufferOverflowStrategy.ERROR) // overflow ends this stream only
.map(this::toSse),
Flux.interval(Duration.ofSeconds(15))
.map(i -> ServerSentEvent.<Object>builder().comment("hb").build()))
.doOnSubscribe(s -> metrics.streamOpened(user))
.doFinally(signal -> metrics.streamClosed(user, signal)); // cancel = client left
For notifications, ERROR on overflow is right: the stream ends, the client reconnects and replays. For state-shaped data such as prices, replace the buffer with onBackpressureLatest(), which keeps only the newest item for a slow subscriber.
doFinally receives SignalType.CANCEL when the client disconnects — Netty notices the closed channel and cancels the subscription — which is the reactive equivalent of a close handler. Release per-connection resources there.
Step 5 — Set the response headers proxies need Permalink to this section
@Bean
WebFilter sseHeaders() {
return (exchange, chain) -> {
if (exchange.getRequest().getPath().value().startsWith("/api/events")) {
exchange.getResponse().getHeaders().set("X-Accel-Buffering", "no");
exchange.getResponse().getHeaders().setCacheControl("no-cache");
}
return chain.filter(exchange);
};
}
If you prefer functional endpoints to annotated controllers, the same pipeline returns from a HandlerFunction with ServerResponse.ok().contentType(MediaType.TEXT_EVENT_STREAM).body(flux, ServerSentEvent.class). The header filter, backpressure and cleanup are unchanged; only the routing differs. Whichever style you use, keep the stream endpoint’s route separate from the rest of the API so that proxy, compression and tracing exclusions can target it by path rather than by guesswork.
Finally, set a server-side maximum lifetime. A stream that has been open for many hours holds a subscription, an authentication decision made at connect time and a place in whatever replay window the store keeps. Completing the Flux after, say, an hour with .take(Duration.ofHours(1)) forces a clean reconnect that re-authenticates and resumes from Last-Event-ID at negligible cost.
Validation & Monitoring Permalink to this section
# Replay: reconnect with an old id and confirm ordered, gap-free ids.
curl -sN -H "Authorization: Bearer $T" -H 'Last-Event-ID: 51' http://localhost:8080/api/events \
| grep --line-buffered '^id:' | head
# Blocking detector: BlockHound fails the test if anything blocks on the event loop.
// In tests: install BlockHound so a blocking call inside the pipeline throws immediately.
@BeforeAll static void blockHound() { BlockHound.install(); }
Monitor reactor.netty connection gauges, per-stream close reasons from doFinally (cancel, error, complete) and overflow errors. Overflow errors that correlate with a particular client network indicate genuinely slow subscribers; overflow across all subscribers means the publisher is outrunning the event loop.
Production Checklist Permalink to this section
Frequently Asked Questions Permalink to this section
Should the controller return Flux<T> or Flux<ServerSentEvent<T>>?
Return ServerSentEvent when you need ids, event names, retry or comments — which a production stream always does. A plain Flux of objects produces data-only events with no id, so reconnects cannot resume.
How do I detect a client disconnect in WebFlux?
The subscription is cancelled when Netty sees the connection close. Handle it in doFinally or doOnCancel. Detection happens when the channel closes or the next write fails, so heartbeats keep it prompt.
Can one Sinks.Many serve thousands of subscribers?
Yes; multicast sinks are designed for it. Keep per-subscriber work light, filter early, and avoid operators that buffer unboundedly after the sink.
Why use onBackpressureBuffer with ERROR instead of DROP_OLDEST?
For event types that must not be lost, silently dropping items corrupts the client's view. Ending the stream forces a reconnect and a replay from the last delivered id, which restores correctness.