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() throws EmitFailureHandler errors, 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.

What runs on the Netty event loop Stack of the Netty event loop, the subscriber pipeline, a mapping step, and a blocking database call, showing how a blocking call in the pipeline stalls every stream on that loop. What runs on the Netty event loop Netty event loop one thread, many streams shared by thousands of connections Flux pipeline map, filter, merge runs on the loop by default Blocking JDBC call 200 ms stalls every stream on this loop Socket writes queued behind it delivered in a burst afterwards
One blocking call inside the pipeline parks a shared event-loop thread. Every stream assigned to that thread freezes until it returns.
  1. 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.sleep inside map or flatMap stalls every stream on that thread.
  2. 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.
  3. The wrong sink. Sinks.many().unicast() accepts one subscriber; a replay() sink holds its entire history unless bounded; a multicast().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.

Replay and live without a gap Sequence diagram of a reconnecting client whose request subscribes to the live hub first, then runs the replay query, then emits replayed events followed by buffered live events, filtering duplicates. Replay and live without a gap Client Controller Hub Replay store GET Last-Event-ID: 51 subscribe (buffer live) events after 51 live 58 (buffered) 52 … 58 52 … 58, then live from 59
Event 58 was published while the replay query ran. It sits in the live buffer, is emitted after the replay, and the monotonic filter drops it if the query already returned it.

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(); }
Longest stall seen by any stream during a 5-minute load test Bar chart comparing the maximum stream stall with a blocking lookup on the event loop, with the lookup moved to boundedElastic, and with enrichment done once at publish time. Longest stall seen by any stream during a 5-minute load test Blocking call in map() 7.4 s boundedElastic per subscriber 310 ms Enrich once at publish 40 ms maximum gap between frames on any stream, 5,000 subscribers
A 200 ms blocking call on the loop becomes multi-second stalls under load. Moving it off the loop fixes the stalls; moving it to publish time also removes the per-subscriber cost.

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.