Java & Spring SSE Implementation Permalink to this section
Part of Backend Stream Generation & Connection Management.
Spring gives you two ways to serve Server-Sent Events, and they sit on different concurrency models. Spring MVC’s SseEmitter runs on the Servlet stack: the request is switched into async mode, the request thread is released, and you push events from whatever thread you like. Spring WebFlux returns a Flux<ServerSentEvent<T>> on a non-blocking event loop, and the framework subscribes, serialises and writes as items arrive. Both work in production; both have characteristic failure modes — thread exhaustion and forgotten timeouts on MVC, accidental blocking and unbounded buffers on WebFlux. This guide covers the protocol mapping, complete implementations on each stack, the Java client side, the proxy and container settings that break streams, and how to size a Spring service for tens of thousands of open connections. It assumes Spring Boot 3.x on Java 21, where virtual threads change the MVC trade-offs considerably.
How It Works Permalink to this section
On the wire, Spring produces exactly the frames the event stream format defines. Both stacks give you a builder that maps onto the four fields plus comments:
The difference is ownership. With SseEmitter, your code holds an object and calls send() on it from some thread; if that thread blocks on a slow client, it stays blocked. With WebFlux, the framework pulls items from your Flux as the connection can accept them; backpressure propagates upstream through Reactor operators. On Servlet containers the MVC async request has a timeout; on WebFlux the stream lives until the Flux completes or the client disconnects.
HTTP/1.1 200
Content-Type: text/event-stream
Transfer-Encoding: chunked
retry:3000
id:42
event:price
data:{"sym":"ACME","bid":101.2}
:hb
Spring writes field:value without a space after the colon, which is valid: the specification strips one optional leading space, so both forms parse identically in every compliant client.
Server-Side Implementation Permalink to this section
Spring MVC with SseEmitter Permalink to this section
The controller creates an emitter with an explicit timeout, registers it, and returns immediately. Events are sent from a separate component — here a broadcaster fed by an application event or a message listener.
@RestController
@RequestMapping("/api")
class PriceStreamController {
private final PriceBroadcaster broadcaster;
PriceStreamController(PriceBroadcaster broadcaster) { this.broadcaster = broadcaster; }
@GetMapping(path = "/prices/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
SseEmitter stream(@RequestHeader(value = "Last-Event-ID", required = false) String lastId) {
// Explicit timeout: without one, the container default (often 30 s) ends the stream.
SseEmitter emitter = new SseEmitter(Duration.ofMinutes(30).toMillis());
broadcaster.register(emitter, lastId);
return emitter; // the request thread is released here
}
}
@Component
class PriceBroadcaster {
private final Set<SseEmitter> emitters = ConcurrentHashMap.newKeySet();
private final AtomicLong seq = new AtomicLong();
void register(SseEmitter emitter, String lastId) {
emitters.add(emitter);
Runnable remove = () -> emitters.remove(emitter);
emitter.onCompletion(remove); // normal end, including after complete()
emitter.onTimeout(emitter::complete); // async timeout: end cleanly, client reconnects
emitter.onError(e -> remove.run()); // I/O failure: the client is gone
try {
emitter.send(SseEmitter.event().reconnectTime(3000).comment("connected"));
} catch (IOException e) {
emitters.remove(emitter);
}
}
@EventListener
void onPrice(PriceChanged p) {
SseEventBuilder evt = SseEmitter.event()
.id(Long.toString(seq.incrementAndGet()))
.name("price")
.data(p, MediaType.APPLICATION_JSON);
for (SseEmitter e : emitters) {
try {
e.send(evt);
} catch (IOException | IllegalStateException ex) {
emitters.remove(e); // dead or completed emitter: drop it
}
}
}
@Scheduled(fixedRate = 15_000)
void heartbeat() {
for (SseEmitter e : emitters) {
try { e.send(SseEmitter.event().comment("hb")); }
catch (IOException | IllegalStateException ex) { emitters.remove(e); }
}
}
}
Two details in the broadcaster are easy to miss. send() is synchronous: it writes and flushes on the calling thread, so the loop above is only as fast as the slowest client. For more than a few hundred subscribers, hand each send to an executor — on Java 21 a virtual-thread-per-task executor is the natural choice — so one stalled socket cannot delay everyone else. And the heartbeat is not optional: it is the only way a broadcaster learns that an idle client has gone, because the failure surfaces as an IOException on the next write. Using Spring SseEmitter without thread exhaustion develops both.
Replaying missed events on reconnect Permalink to this section
The controller already receives Last-Event-ID. To honour it, keep a bounded replay buffer in the broadcaster and send everything newer than the client’s id before adding the emitter to the live set. The ordering matters: send the replay and register under the same lock that the publisher takes, or an event published between the two steps is either lost or delivered twice.
private final Deque<Sent> recent = new ArrayDeque<>(); // bounded replay window
private final Object lock = new Object();
record Sent(long id, SseEventBuilder event) {}
void register(SseEmitter emitter, String lastId) {
long after = lastId == null ? Long.MAX_VALUE : Long.parseLong(lastId);
synchronized (lock) { // publisher takes the same lock
try {
for (Sent s : recent) if (s.id() > after) emitter.send(s.event());
} catch (IOException e) {
return; // client vanished mid-replay
}
emitters.add(emitter);
}
// …callbacks as before…
}
void publish(PriceChanged p) {
synchronized (lock) {
long id = seq.incrementAndGet();
SseEventBuilder evt = SseEmitter.event().id(Long.toString(id)).name("price").data(p);
recent.addLast(new Sent(id, evt));
while (recent.size() > 1_000) recent.removeFirst();
emitters.forEach(e -> executor.execute(() -> safeSend(e, evt)));
}
}
A client whose id is older than the oldest buffered event has missed more than the window holds; send it a resync event so it refetches state, as described in implementing a replay buffer for Last-Event-ID. Across several instances the buffer belongs in shared storage — a Redis Stream or Kafka topic — not in each JVM.
Spring WebFlux with Flux<ServerSentEvent> Permalink to this section
On WebFlux the stream is a declarative pipeline. A Sinks.Many bridges imperative producers into it, and each subscriber gets its own view with its own backpressure policy.
@RestController
class ReactivePriceController {
private final Sinks.Many<Price> sink = Sinks.many().multicast().directBestEffort();
@EventListener
void onPrice(PriceChanged p) {
sink.tryEmitNext(p.price()); // never blocks; slow subscribers miss items
}
@GetMapping(path = "/api/prices/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
Flux<ServerSentEvent<Object>> stream() {
Flux<ServerSentEvent<Object>> prices = sink.asFlux()
.onBackpressureLatest() // per subscriber: keep only the newest
.map(p -> ServerSentEvent.builder()
.id(Long.toString(p.seq()))
.event("price")
.data(p)
.build());
Flux<ServerSentEvent<Object>> heartbeats = Flux.interval(Duration.ofSeconds(15))
.map(i -> ServerSentEvent.builder().comment("hb").build());
return Flux.merge(prices, heartbeats)
.startWith(ServerSentEvent.builder().retry(Duration.ofSeconds(3)).comment("connected").build());
}
}
directBestEffort() drops an item for a subscriber that has no outstanding demand instead of failing the emission, and onBackpressureLatest() gives each subscriber a one-element buffer holding the newest value — the right policy for prices and dashboards, the wrong one for notifications, where you want onBackpressureBuffer(n) and a disconnect on overflow. Streaming SSE with Spring WebFlux covers replay with Last-Event-ID, per-user filtering and disconnect hooks.
Client-Side Consumption Permalink to this section
Browsers consume a Spring stream with a plain EventSource; nothing Spring-specific is needed. Java services consuming another service’s stream have three options: WebFlux’s WebClient, which decodes ServerSentEvent natively; the JDK HttpClient with a line-based body handler; or OkHttp’s okhttp-sse module.
// WebClient: typed, reactive, handles the event-stream framing for you.
WebClient client = WebClient.create("https://prices.example.com");
Flux<ServerSentEvent<Price>> events = client.get()
.uri("/api/prices/stream")
.accept(MediaType.TEXT_EVENT_STREAM)
.retrieve()
.bodyToFlux(new ParameterizedTypeReference<ServerSentEvent<Price>>() {});
events
.retryWhen(Retry.backoff(Long.MAX_VALUE, Duration.ofSeconds(1)).maxBackoff(Duration.ofSeconds(30)))
.subscribe(e -> handle(e.id(), e.event(), e.data()));
Without WebFlux on the classpath, the JDK’s own HttpClient is enough. Read the body as lines, accumulate fields until a blank line, and dispatch — the same parsing algorithm the browser runs:
HttpClient http = HttpClient.newHttpClient();
HttpRequest req = HttpRequest.newBuilder(URI.create("https://prices.example.com/api/prices/stream"))
.header("Accept", "text/event-stream")
.header("Last-Event-ID", lastId) // resume where the previous connection stopped
.build();
HttpResponse<Stream<String>> res = http.send(req, HttpResponse.BodyHandlers.ofLines());
StringBuilder data = new StringBuilder();
String event = "message";
for (Iterator<String> it = res.body().iterator(); it.hasNext(); ) {
String line = it.next();
if (line.isEmpty()) { // blank line: dispatch the event
if (data.length() > 0) handle(event, data.substring(0, data.length() - 1));
data.setLength(0);
event = "message";
} else if (line.startsWith("data:")) data.append(stripSpace(line.substring(5))).append('\n');
else if (line.startsWith("event:")) event = stripSpace(line.substring(6));
else if (line.startsWith("id:")) lastId = stripSpace(line.substring(3));
}
Run it on a virtual thread and wrap it in a reconnect loop with backoff; the iterator blocks while the stream is idle and ends when the server closes the response.
WebClient does not track Last-Event-ID for you. Keep the last id in a variable and send it as a header when you resubscribe, or the retry starts from whatever the server considers “now”. Consumption patterns for other runtimes are in non-browser SSE clients.
Edge Cases & Network Interference Permalink to this section
Spring services usually sit behind several layers that each have an opinion about long responses.
- Async request timeout.
spring.mvc.async.request-timeoutsets the default forSseEmitterinstances created without one. Unset, the container default applies — 30 seconds on Tomcat — and streams end and reconnect every 30 seconds. Always pass a timeout to the constructor. - Response compression.
server.compression.enabled=truecompresses matching MIME types;text/event-streamis not in Spring Boot’s default list, but teams often addtext/*wholesale. Keep event streams out of it, or buffered compression will batch events. - Security filters and sessions. Spring Security’s filters run once at connect; an expired session does not end an open stream. Re-check authorisation periodically for long-lived streams or bound the emitter timeout to the token lifetime.
- Reverse proxies. nginx buffers by default. Set
X-Accel-Buffering: noon the response — a@ControllerAdviceorResponseEntityheader does it — and see proxy and CDN configuration for SSE. - Actuator and tracing. Observation and tracing filters may hold the request span open for the life of the stream, producing thirty-minute spans. Exclude stream endpoints, or end the span at first byte.
Mitigation checklist:
Performance & Scale Considerations Permalink to this section
The capacity question on MVC is “how many threads does an idle stream cost?” — ideally zero. After the controller returns, it holds none: the connection is parked by the container’s NIO connector. Threads are consumed only while send() runs. The container’s connection cap is then the limit: Tomcat’s server.tomcat.max-connections defaults to 8,192, so a Spring MVC service stops accepting new streams at that point unless it is raised.
The settings that matter:
# application.yml
spring:
threads:
virtual:
enabled: true # Boot 3.2+: request handling and @Async on virtual threads
mvc:
async:
request-timeout: 30m # default for emitters created without a timeout
server:
tomcat:
max-connections: 50000 # each open stream holds one connection
keep-alive-timeout: 75s
Beyond the connector, the limits are the ones every SSE server shares: file descriptors per process, kernel socket buffers, and heap per connection (a few kilobytes on either stack when idle). Raise the file descriptor limit for the JVM as described in tuning file descriptor limits. For multi-instance deployments, feed every instance from a broker — Redis, Kafka or NATS — rather than from in-process events, exactly as in Redis pub/sub fan-out.
Validation & Debugging Permalink to this section
# Frames arrive promptly and the stream outlives the default 30-second async timeout.
curl -sN -H 'Accept: text/event-stream' http://localhost:8080/api/prices/stream \
| while IFS= read -r l; do echo "$(date +%T) $l"; done
# Count emitters and threads during a load test.
curl -s localhost:8080/actuator/metrics/jvm.threads.live | jq '.measurements[0].value'
curl -s localhost:8080/actuator/metrics/sse.emitters.active | jq '.measurements[0].value'
Automated tests are straightforward on both stacks. WebTestClient can bind to a WebFlux controller or a running MVC application and decode the stream as typed events, so a test asserts on the first few events and then cancels:
@SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT)
class PriceStreamTest {
@Autowired WebTestClient client;
@Autowired ApplicationEventPublisher events;
@Test
void streamsPricesWithIds() {
Flux<ServerSentEvent<String>> body = client.get().uri("/api/prices/stream")
.accept(MediaType.TEXT_EVENT_STREAM)
.exchange().expectStatus().isOk()
.returnResult(new ParameterizedTypeReference<ServerSentEvent<String>>() {})
.getResponseBody();
StepVerifier.create(body.filter(e -> "price".equals(e.event())))
.then(() -> events.publishEvent(new PriceChanged(new Price("ACME", 101.2))))
.assertNext(e -> assertThat(e.id()).isNotBlank())
.thenCancel()
.verify(Duration.ofSeconds(5));
}
}
Register a gauge for active emitters (Gauge.builder("sse.emitters.active", emitters, Set::size).register(registry)), and a counter for removals by reason — completion, timeout, error. A timeout count equal to connection count every thirty seconds means the emitter timeout is not being applied. In a thread dump, look for many threads parked inside SseEmitter.send or OutputStream.write: that is one slow client per thread, and the cue to move sends to an executor.
Production Checklist Permalink to this section
Frequently Asked Questions Permalink to this section
Should I use SseEmitter or WebFlux for a new service?
If the service is already Spring MVC and its data sources are blocking (JDBC, blocking clients), use SseEmitter with virtual threads. If the pipeline is reactive end to end — R2DBC, reactive Kafka, WebClient — WebFlux handles more streams per instance with less code.
Why does my SseEmitter stream stop after 30 seconds?
The async request timed out. Emitters without a constructor timeout use spring.mvc.async.request-timeout, or the container default when that is unset. Pass an explicit timeout and complete the emitter in onTimeout so the client reconnects cleanly.
How does Spring tell me a client disconnected?
Usually on the next write: send() throws an IOException, or onError fires. An idle stream with no writes never finds out, which is one reason heartbeats are required.
Can Kotlin coroutines return an SSE stream?
Yes. On WebFlux, a controller can return Flow<ServerSentEvent<T>> and Spring adapts it to a Flux, so suspend functions and flow operators replace Reactor operators while the wire format and backpressure behaviour stay the same.
How do I send an event to one specific user?
Key the emitter registry by user id — a map from user to a set of emitters, since users open several tabs — and look up the set when publishing. Across instances, route by user through the broker so the event reaches whichever instance holds that user's streams.
Do virtual threads make SseEmitter scale like WebFlux?
They remove thread exhaustion as the limiting factor, because a blocked send parks a cheap virtual thread. WebFlux still uses less memory per stream and gives operator-level backpressure, but the gap is small enough that the existing codebase should usually decide.