Using Spring SseEmitter without Thread Exhaustion Permalink to this section
Part of Java & Spring SSE Implementation, under Backend Stream Generation & Connection Management.
SseEmitter releases the request thread as soon as the controller returns, which makes it look free. The cost moves to whichever thread calls send(). A broadcaster that loops over a few thousand emitters on one thread, or a design that parks a thread per emitter waiting for data, runs out of threads long before it runs out of connections. This guide diagnoses the three ways Spring MVC services exhaust threads with SSE and fixes each one.
Symptom & Developer Intent Permalink to this section
- Under load, ordinary REST endpoints on the same service start timing out, while SSE clients look connected.
- Thread dumps show hundreds of threads in
SocketOutputStream.socketWriteorNioEndpoint$NioSocketWrapper.doWrite, all called fromSseEmitter.send. - Event latency climbs with the number of subscribers: the thousandth emitter receives each event seconds after the first.
server.tomcat.threads.maxis hit, yet CPU is mostly idle.- A single client on a bad network visibly delays every other client.
The intent is a Spring MVC service that holds many thousands of streams, delivers each event to every subscriber within a bounded time regardless of slow clients, and leaves request threads free for the rest of the API.
Root Cause Analysis Permalink to this section
SseEmitter.send() is synchronous. It serialises the event, writes it to the response and flushes. If the client’s TCP receive window is full, the write blocks until it drains or the socket times out. Three designs turn that blocking into exhaustion:
- Sequential broadcast. One thread loops over every emitter. The loop is only as fast as the slowest socket, and a stalled client can hold it for the full socket write timeout.
- Thread per emitter. Each connection gets a worker that blocks on a queue and then sends. Idle connections each hold a parked platform thread; 5,000 clients need 5,000 threads.
- Sends on the wrong pool.
@Asyncmethods orCompletableFuture.runAsyncwithout an explicit executor may run on a shared pool, so blocked SSE sends compete with request handling.
Step-by-Step Resolution Permalink to this section
Step 1 — Confirm the diagnosis with a thread dump Permalink to this section
# Group threads by the top frame to see where they are waiting.
jcmd $(pgrep -f my-service.jar) Thread.print > threads.txt
grep -A3 '"' threads.txt | grep -E 'at (java|org|sun)' | sort | uniq -c | sort -rn | head
# 412 at sun.nio.ch.SocketDispatcher.write0 … ← blocked on client sockets
Many threads in socket write with a caller in SseEmitter.send confirms that client writes, not your business logic, are holding threads.
Step 2 — Send on virtual threads, one task per emitter per event Permalink to this section
On Java 21, a virtual thread parked in a socket write costs a few hundred bytes, not a megabyte of stack. Moving each send onto its own virtual thread removes the head-of-line blocking and the thread ceiling together.
@Component
class Broadcaster {
private final ExecutorService sendPool = Executors.newVirtualThreadPerTaskExecutor();
private final Map<SseEmitter, Client> clients = new ConcurrentHashMap<>();
void broadcast(SseEventBuilder event) {
for (Client c : clients.values()) c.enqueue(event); // never blocks the caller
}
}
Step 3 — Preserve ordering with a per-emitter queue Permalink to this section
Firing one independent task per event would let two events for the same client race and arrive out of order. Give each client a small queue and a single drainer at a time:
final class Client {
private final SseEmitter emitter;
private final Queue<SseEventBuilder> queue = new ConcurrentLinkedQueue<>();
private final AtomicBoolean draining = new AtomicBoolean();
private final AtomicInteger depth = new AtomicInteger();
private final ExecutorService pool;
private final Runnable onDead;
Client(SseEmitter emitter, ExecutorService pool, Runnable onDead) {
this.emitter = emitter; this.pool = pool; this.onDead = onDead;
}
void enqueue(SseEventBuilder e) {
if (depth.incrementAndGet() > 256) { // bounded: a slow client is evicted
emitter.completeWithError(new IllegalStateException("client too slow"));
onDead.run();
return;
}
queue.add(e);
if (draining.compareAndSet(false, true)) pool.execute(this::drain);
}
private void drain() {
try {
SseEventBuilder e;
while ((e = queue.poll()) != null) {
depth.decrementAndGet();
emitter.send(e); // may park this virtual thread; that's fine
}
} catch (IOException | IllegalStateException ex) {
onDead.run();
} finally {
draining.set(false);
if (!queue.isEmpty() && draining.compareAndSet(false, true)) pool.execute(this::drain);
}
}
}
At most one virtual thread per client is ever sending, events stay in order, and a client whose queue exceeds 256 pending events is disconnected rather than allowed to hold memory indefinitely. It will reconnect and replay from Last-Event-ID.
Step 4 — Do not park threads waiting for data Permalink to this section
Remove any per-emitter loop that waits for the next event:
// Anti-pattern: one thread per client, parked most of the time.
executor.submit(() -> {
while (true) {
Price p = queue.take(); // blocks a thread per connection
emitter.send(p);
}
});
Events should be pushed to clients when they happen, as in step 2, so an idle connection consumes no thread at all — platform or virtual.
Step 5 — Raise the connection ceiling to match Permalink to this section
With threads no longer the constraint, Tomcat’s connection cap becomes one:
server:
tomcat:
max-connections: 50000 # default 8192
accept-count: 1000
threads:
max: 200 # request threads; streams no longer need more
spring:
threads:
virtual:
enabled: true # request handling on virtual threads too (Boot 3.2+)
And raise the process’s file descriptor limit above max-connections, or the kernel will refuse sockets first.
Validation & Monitoring Permalink to this section
Load test with one deliberately slow client to prove isolation:
# 3,000 normal subscribers (see the k6 guide), plus one client that reads 100 bytes per second.
curl -sN --limit-rate 100 http://localhost:8080/api/prices/stream > /dev/null &
# Meanwhile, ordinary endpoints must stay fast.
hey -z 60s -c 20 http://localhost:8080/api/health
The k6 load-testing guide shows how to hold thousands of SSE connections and measure per-event latency.
Export three metrics: active clients, total queued events across clients, and evictions for slowness. Queued events near zero with occasional spikes is healthy; a steadily rising total means the send pool or the network cannot keep up with the publish rate.
Production Checklist Permalink to this section
Frequently Asked Questions Permalink to this section
Is SseEmitter.send thread-safe?
Concurrent sends to the same emitter are serialised internally, but interleaving them from several threads makes event order unpredictable. Use one drainer per emitter so events leave in the order they were published.
What if I cannot use Java 21?
Use a dedicated bounded thread pool for sends, never the request pool, keep the per-client queue and single-drainer design, and evict clients whose queues fill. Size the pool to cover the slow-client fraction you expect, typically a few dozen threads.
Why evict slow clients instead of buffering more?
An unbounded buffer turns a slow client into a memory leak and delivers ever-staler events. Disconnecting it lets EventSource reconnect and replay from its last id, which is a faster route to a current view.
Does WebFlux avoid this problem entirely?
WebFlux never blocks a thread on a socket write, so exhaustion cannot happen the same way. It still needs an explicit per-subscriber backpressure policy, which is the reactive equivalent of the bounded queue.