Revision note (2026-09-15). This article replaces an earlier version that had four errors: the server named its events business-event while the browser example listened with onmessage, which never fires for named events; the onerror handler called EventSource.close() and then promised automatic reconnection; the prerequisites said Spring Boot 2.3+ although Sinks needs Reactor 3.4, which arrived with Spring Boot 2.4; and limitRate() was described as a way to cap connections. Each point is corrected below, and every behavior claimed here was reproduced with the code and versions listed.
What this article covers
You want a Spring Boot endpoint that pushes events to browsers over one long-lived HTTP response, and you want to know what happens when a client is slow, disconnects, or reconnects. This is a lab implementation. It is small enough to read in full, it is tested, and its limits are listed at the end. It is not a drop-in production component.
Tested versions:
| Component | Version |
|---|---|
| Spring Boot | 3.5.16 |
| Spring WebFlux | 6.2.19 |
| Reactor Core | 3.7.19 |
| Reactor Netty HTTP | 1.2.18 |
| Java | 17 (Amazon Corretto 17.0.14) |
| Browser check | Google Chrome via Playwright 1.63.0 |
The version floor matters. Sinks.many() was introduced in Reactor Core 3.4.0. Spring Boot 2.3.12 manages Reactor Dysprosium-SR20, which is reactor-core 3.3.x and has no Sinks class. Spring Boot 2.4.0 manages 2020.0.1, which is reactor-core 3.4.x. So the minimum for this code is Spring Boot 2.4, and the article was tested on 3.5.16.
When SSE fits
Server-Sent Events is a one-direction stream of UTF-8 text from server to browser over a normal HTTP response with Content-Type: text/event-stream. The browser's EventSource object handles parsing, reconnection, and the Last-Event-ID header for you. Pick it when the server talks and the client mostly listens: notifications, progress, dashboards. Pick WebSocket when the client also needs to send a stream, or when payloads are binary. Pick polling when the intermediaries between browser and server cannot hold connections open.
The server
Three classes do the work: an event record, a publisher that owns the shared sink, and a controller.
public record BusinessEvent(long id, String payload, Instant createdAt) {
}
The publisher keeps one hot Sinks.Many for every connected client, plus a small in-memory ring of recent events so a reconnecting browser can catch up.
@Component
public class BusinessEventPublisher {
public record Outcome(BusinessEvent event, Sinks.EmitResult result) {
}
private final Sinks.Many<BusinessEvent> sink = Sinks.many().multicast().directBestEffort();
private final ArrayDeque<BusinessEvent> replay;
private final int replaySize;
private final AtomicLong nextId = new AtomicLong();
public BusinessEventPublisher(SseProperties properties) {
this.replaySize = properties.replaySize();
this.replay = new ArrayDeque<>(properties.replaySize());
}
public synchronized Outcome publish(String payload) {
BusinessEvent event = new BusinessEvent(nextId.incrementAndGet(), payload, Instant.now());
Sinks.EmitResult result = sink.tryEmitNext(event);
remember(event); // kept even when nobody is connected, so reconnects can catch up
return new Outcome(event, result);
}
public Flux<BusinessEvent> stream(Long lastEventId) {
Flux<BusinessEvent> live = sink.asFlux();
if (lastEventId == null) {
return live;
}
return Flux.fromIterable(missedSince(lastEventId)).concatWith(live);
}
private void remember(BusinessEvent event) {
if (replay.size() == replaySize) {
replay.pollFirst();
}
replay.addLast(event);
}
private synchronized List<BusinessEvent> missedSince(long lastEventId) {
List<BusinessEvent> missed = new ArrayList<>();
for (BusinessEvent event : replay) {
if (event.id() > lastEventId) {
missed.add(event);
}
}
return missed;
}
}
Two choices here come straight from the tests in the next sections: the sink is directBestEffort() rather than onBackpressureBuffer(), and publish() is synchronized.
The controller sets the SSE id and event fields and merges a heartbeat comment so idle connections are not cut by proxies.
@RestController
public class EventStreamController {
public static final String EVENT_NAME = "business-event";
@GetMapping(path = "/events", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<ServerSentEvent<BusinessEvent>> events(
@RequestHeader(name = "Last-Event-ID", required = false) Long lastEventId) {
Flux<ServerSentEvent<BusinessEvent>> connected = Flux.just(
ServerSentEvent.<BusinessEvent>builder().comment("connected").build());
Flux<ServerSentEvent<BusinessEvent>> business = publisher.stream(lastEventId)
.map(event -> ServerSentEvent.builder(event)
.id(Long.toString(event.id()))
.event(EVENT_NAME)
.build());
Flux<ServerSentEvent<BusinessEvent>> heartbeat = Flux.interval(properties.heartbeat())
.map(tick -> ServerSentEvent.<BusinessEvent>builder().comment("keep-alive").build());
return connected.concatWith(Flux.merge(business, heartbeat));
}
@PostMapping(path = "/publish", consumes = MediaType.TEXT_PLAIN_VALUE)
public ResponseEntity<Map<String, Object>> publish(@RequestBody String payload) {
BusinessEventPublisher.Outcome outcome = publisher.publish(payload);
return ResponseEntity.accepted().body(Map.of(
"id", outcome.event().id(),
"result", outcome.result().name()));
}
}
A WebFlux test confirms the wire format: the first element is the connected comment, then each published event arrives with event: business-event and an id equal to the event id (streamStartsWithAConnectedCommentAndSendsNamedEventsWithIds). A second test publishes three events with no client connected, then connects with Last-Event-ID set to the first id and receives exactly the second and third (reconnectingWithLastEventIdReplaysTheMissedEvents).
The browser client, and the bug that hides every event
The server sends event: business-event. The SSE specification routes an event with an event field to listeners registered for that name. onmessage only receives events that have no event field. So this client, which is what the earlier version of the article showed, logs nothing:
// Receives nothing from this server: the events are named.
eventSource.onmessage = (event) => console.log(event.data);
The working client registers for the name and leaves reconnection to the browser:
const es = new EventSource('/events');
es.addEventListener('business-event', (event) => {
console.log('id=' + event.lastEventId + ' data=' + event.data);
});
es.onopen = () => console.log('open');
// Do not call es.close() here. After a transient failure the EventSource
// moves to CONNECTING (readyState 0) and the browser retries by itself.
// close() is for "this tab is done listening", and it is permanent.
es.onerror = () => console.log('error, readyState=' + es.readyState);
The earlier version closed the connection inside onerror. That is exactly the wrong place: the first hiccup ends the stream forever, and no reconnection can happen afterwards.
What Chrome actually did
A Playwright script drove a real Google Chrome against the running jar, killed the server process, restarted it, and finally called close(). Counters were read from the page after each step. Here is the recorded sequence:
| Step | named events | onmessage | opens | errors | readyState | Note |
|---|---|---|---|---|---|---|
| Page loaded | 0 | 0 | 1 | 0 | OPEN (1) | |
| 3 events published | 3 | 0 | 1 | 0 | OPEN (1) | ids 1, 2, 3 delivered to the named listener only |
| Server process killed | 3 | 0 | 1 | 1 | CONNECTING (0) | browser is retrying on its own |
| Server restarted | 3 | 0 | 2 | 1 | OPEN (1) | server log: client reconnected with Last-Event-ID=3 -> replaying 0 event(s) |
| 1 event published | 4 | 0 | 2 | 1 | OPEN (1) | |
close() called | 4 | 0 | 2 | 1 | CLOSED (2) | |
| Event published, then server killed | 4 | 0 | 2 | 1 | CLOSED (2) | no error, no reconnect: closed means closed |
Two details are worth keeping. First, onmessage stayed at zero through the whole run, which is the concrete proof that the old client would have shown nothing. Second, Chrome sent Last-Event-ID: 3 on its own when it reconnected. The replay found nothing because the restarted process started numbering from 1 again, which is one of the limitations listed at the end.
Reactor Sinks: four behaviors that decide your design
These are unit tests against reactor-core 3.7.19 with no Spring involved. The buffer size is 16 because Reactor rounds the requested size up to at least 8 and to a power of two, so a size of 4 does not mean 4.
1. The buffered multicast sink warms up, then rejects. With Sinks.many().multicast().onBackpressureBuffer(16, false) and no subscriber yet, the first 16 tryEmitNext calls return OK and the 17th returns FAIL_ZERO_SUBSCRIBER. The first subscriber then receives those 16 buffered elements.
2. One subscriber without demand starves every other subscriber. Two subscribers on the same buffered sink: one requests unbounded, one requests exactly one element and stops. After the shared element, 16 more emissions return OK, the 17th returns FAIL_OVERFLOW, and the unbounded subscriber has received only the first element. The buffered multicast sink delivers in lockstep with its slowest subscriber. For SSE that means one stalled browser tab freezes the stream for everyone.
3. emitNext on overflow terminates the sink for everyone, permanently. Calling sink.emitNext(value, Sinks.EmitFailureHandler.FAIL_FAST) in the overflow state above delivers an overflow error to the fast subscriber, later tryEmitNext calls return FAIL_TERMINATED, and a brand-new subscriber gets the error immediately instead of a stream. The earlier version of this article used emitNext with a handler that only retried on FAIL_NON_SERIALIZED; on overflow it would have killed the stream for all clients.
4. directBestEffort() drops per subscriber and keeps the sink alive. Same two subscribers, five emissions: the fast one gets all five, the one without demand gets the first only, no error anywhere, and a subscriber that joins afterwards receives the next emission. This is why the publisher above uses directBestEffort(). A slow browser misses live events instead of stopping the others, and it recovers through Last-Event-ID replay when it reconnects.
There is a fifth fact about threads. Four threads calling tryEmitNext 20,000 times each on one sink produced 59,207 FAIL_NON_SERIALIZED results out of 80,000 in one run. Sinks do not serialize concurrent emitters. Either serialize yourself (the synchronized publish method above) or use emitNext(value, Sinks.EmitFailureHandler.busyLooping(Duration)), which retried until all 80,000 elements were delivered in the same test.
Backpressure and the limitRate() correction
limitRate(n) shapes how many elements an operator requests from its upstream at a time. It says nothing about how many HTTP clients may connect. In the earlier version it was offered as an answer to "how do I limit concurrent SSE connections", and that is wrong. Connection limits belong in the reverse proxy or in a counter you keep yourself, for example by incrementing on doOnSubscribe and decrementing on doOnCancel and refusing new subscriptions above a threshold.
With directBestEffort(), backpressure inside the JVM is resolved by dropping for the subscriber without demand. Whether the HTTP layer signals demand depends on Reactor Netty's write buffer and the client's TCP receive rate, which this lab did not measure.
Spring MVC or WebFlux
The earlier version said SseEmitter "uses servlet threads and can lead to scalability issues". That is too coarse. SseEmitter relies on Servlet asynchronous request processing, so an open connection does not pin a request thread while idle; the write happens on whichever thread calls send(). WebFlux gives you the same stream as a Flux, which composes naturally with Reactor sources, and Reactor Netty handles many idle connections with a small event loop. If your application is already Spring MVC, SseEmitter is a reasonable choice. This lab uses WebFlux because the publisher is a Reactor sink.
Running it across several instances
The sink in this lab lives in one JVM. A browser connected to instance A never sees an event published on instance B. The usual fix is a broker fan-out: each instance subscribes to a channel (Redis Pub/Sub, a Kafka topic with one consumer group per instance, RabbitMQ fan-out exchange) and pushes what it receives into its local sink. The earlier version included a Redis listener snippet for this. It was not built or tested for this revision, so it is not shown here; treat it as a design direction, not a verified component.
Proxies and timeouts
Long-lived responses need the reverse proxy to stop buffering and to allow long read timeouts. For NGINX that is proxy_buffering off and a large proxy_read_timeout on the location, per the ngx_http_proxy_module documentation. The heartbeat comment every 15 seconds in the controller exists so that idle connections keep traffic flowing through such proxies. No proxy was part of this lab.
Limits of this implementation
- Event ids restart at 1 when the process restarts, so
Last-Event-IDreplay across a restart returns nothing. Ids would need to be monotonic across processes, for example from a database sequence or a broker offset. - The replay ring is in memory (100 events) and there is a small window between reading the ring and subscribing to the live sink where an event can be delivered twice or missed. Clients should de-duplicate on id.
directBestEffort()drops events for a client that has no demand. That client learns nothing until it reconnects.- No authentication, no per-client connection limit, no metrics. Each is a few lines but none is here.
- Fan-out across instances is described, not implemented.
Reproduce it
gradle test # 7 tests: 3 WebFlux endpoint tests, 4 Reactor Sinks behavior tests
gradle bootJar
cd browser-test && npm install && npm run check # drives Google Chrome, kills and restarts the server
The browser check writes results/latest.json with the counter table above.
Sources
- Spring Framework reference: Server-Sent Events in WebFlux
- MDN: Using server-sent events
- WHATWG HTML Standard: server-sent events
- Project Reactor reference: Sinks
- Reactor Core API: Sinks.MulticastSpec
- Spring Boot 2.3.12 dependency versions and Spring Boot 2.4.0 dependency versions
- NGINX ngx_http_proxy_module: proxy_buffering
