Skip WebSocket for AI Streaming: SSE Delivers Typewriter Effects with Half the Code in Spring Boot
This comprehensive guide demonstrates why Server-Sent Events (SSE) is the superior choice over WebSocket for unidirectional server push scenarios like AI token streaming, covering Spring Boot integration patterns, production hardening (async thread pools, Nginx buffering, heartbeats, reconnection), and cluster deployment with Redis Pub/Sub.
01 Server Push: Only Four Real Options
When building AI applications last year, a common requirement emerged: the LLM generates tokens one by one in the backend, and the frontend must display them in real time like ChatGPT. The first instinct is WebSocket, but hands-on work reveals steep costs: extra dependencies, handshake upgrades, heartbeats, reconnection logic, and Nginx Upgrade header configuration — all for a simple "server speaks, client listens" unidirectional need.
Switching to SSE cut backend code by more than half, and the browser handles reconnection automatically. Notably, ChatGPT, Claude, and Ernie Bot all use SSE for their streaming outputs.
02 What Is SSE?
SSE (Server-Sent Events) is part of the HTML5 standard. No magic: it's a plain HTTP GET where the server sets Content-Type: text/event-stream, keeps the connection open, and writes text chunks followed by flush.
HTTP/1.1 200 OK
Content-Type: text/event-stream;charset=UTF-8
Cache-Control: no-cache
Connection: keep-alive
X-Accel-Buffering: no
retry: 10000
event: message
id: 1
data: {"content":"你"}
event: message
id: 2
data: {"content":"好"}
: this is a comment, used as heartbeat keep-alive
event: done
id: 3
data: [DONE]Four fields to remember: data: message body; multiple data: lines are concatenated with newlines. event: event name, defaults to message; frontend can use addEventListener for custom events. id: event ID; on reconnect the browser sends Last-Event-ID, enabling native resume. retry: tells the browser how many milliseconds to wait before reconnecting.
Critical: two messages must be separated by a blank line. Missing that blank line is the most common beginner pitfall — the server sends data but the frontend receives nothing.
03 Frontend: 20 Lines to Run
<script>
// Auto GET with Accept: text/event-stream
const es = new EventSource('/sse/subscribe/1001');
// Default event
es.onmessage = (e) => {
console.log('Received:', e.data, 'id =', e.lastEventId);
};
// Custom event (matches backend .name("done"))
es.addEventListener('done', (e) => {
console.log('Stream ended:', e.data);
es.close(); // Active close, otherwise browser keeps reconnecting
});
// Reconnect is automatic, interval driven by retry:
es.onerror = () => {
console.warn('Connection error, readyState =', es.readyState);
// 0=CONNECTING 1=OPEN 2=CLOSED; after CLOSED no auto-reconnect, must rebuild manually
};
</script>Convenient, but EventSource has a hard limitation : it cannot set custom request headers. Authentication must use Cookies or URL tokens. Tokens in query strings end up in browser history, Nginx access logs, and proxy logs — if unavoidable, use one-time, short-lived, single-purpose tickets.
04 When SSE Truly Replaces WebSocket
SSE fits scenarios that are unidirectional, text-based, server-initiated :
AI LLM token streaming (most frequent production use case)
In-app messages, system notifications
Real-time dashboards / large-screen refresh
Long-running task progress bars (export, batch, transcoding)
Real-time log streams (ops console, CI build output)
Config/state broadcast (notify all online clients to refresh after config change)
Conversely, avoid SSE for:
Client sends frequent messages : chat rooms, games, collaborative editing
Large binary payloads : file chunks, audio/video streams
Complex sub-protocols : high-frequency bidirectional with custom protocols like collaborative editing
05 Three Spring Boot SSE Implementations
Method 1: SseEmitter (Most Common, Production First Choice)
Built into Spring MVC, zero extra dependencies, excellent lifecycle management.
@GetMapping(
value = "/subscribe/{userId}",
produces = MediaType.TEXT_EVENT_STREAM_VALUE
)
public SseEmitter subscribe(@PathVariable String userId) {
// 0L = no timeout; timeout control via heartbeat and active cleanup
SseEmitter emitter = new SseEmitter(0L);
emitterPool.put(userId, emitter);
emitter.onCompletion(() -> emitterPool.remove(userId));
emitter.onTimeout(() -> emitterPool.remove(userId));
emitter.onError(e -> emitterPool.remove(userId));
return emitter;
}Method 2: WebFlux Reactive (Most Elegant Code)
@GetMapping(
value = "/flux",
produces = MediaType.TEXT_EVENT_STREAM_VALUE
)
public Flux<ServerSentEvent<String>> stream() {
return Flux.interval(Duration.ofSeconds(1))
.map(seq -> ServerSentEvent.<String>builder()
.id(String.valueOf(seq))
.event("tick")
.data("第 " + seq + " 秒")
.build());
}Method 3: Raw HttpServletResponse (Not Recommended, but Good to Understand)
@GetMapping("/raw")
public void raw(HttpServletResponse response)
throws IOException, InterruptedException {
response.setContentType("text/event-stream;charset=UTF-8");
response.setHeader("Cache-Control", "no-cache");
response.setHeader("X-Accel-Buffering", "no"); // Disable Nginx buffering
PrintWriter writer = response.getWriter();
for (int i = 0; i < 10; i++) {
writer.write("data: 第 " + i + " 条
"); // Blank line mandatory
writer.flush(); // Must flush, otherwise not streaming
Thread.sleep(1000);
}
writer.close();
}Manual HttpServletResponse is only for learning the protocol; you must handle exceptions, disconnects, timeouts, serialization yourself — SseEmitter already does all that.
06 Building a Production-Ready Push Service
The key is the connection manager . It must solve four problems: multiple tabs per user, heartbeat keep-alive, dead connection cleanup, and event ID generation.
@Slf4j
@Component
public class SseEmitterManager {
// userId -> multiple connections (same account may open many tabs)
private final Map<String, CopyOnWriteArrayList<SseEmitter>> emitters = new ConcurrentHashMap<>();
private final AtomicLong eventId = new AtomicLong(0);
private final ScheduledExecutorService heartbeat = Executors.newSingleThreadScheduledExecutor(r -> {
Thread t = new Thread(r, "sse-heartbeat");
t.setDaemon(true);
return t;
});
public SseEmitterManager() {
// Heartbeat every 25s to prevent proxies from killing idle connections
heartbeat.scheduleAtFixedRate(this::ping, 25, 25, TimeUnit.SECONDS);
}
public SseEmitter connect(String userId) {
SseEmitter emitter = new SseEmitter(0L);
emitters.computeIfAbsent(userId, k -> new CopyOnWriteArrayList<>()).add(emitter);
emitter.onCompletion(() -> remove(userId, emitter));
emitter.onTimeout(() -> remove(userId, emitter));
emitter.onError(e -> remove(userId, emitter));
try {
emitter.send(SseEmitter.event()
.id(String.valueOf(eventId.incrementAndGet()))
.name("connected")
.data("连接已建立")
.reconnectTime(3000L)); // Tell browser to reconnect after 3s
} catch (IOException e) {
remove(userId, emitter);
}
return emitter;
}
// Push: failure = dead connection, clean immediately to avoid zombies
public void send(String userId, String eventName, Object data) {
List<SseEmitter> list = emitters.get(userId);
if (list == null || list.isEmpty()) {
log.debug("[SSE] User {} no online connections, message dropped", userId);
return;
}
long id = eventId.incrementAndGet();
List<SseEmitter> dead = new ArrayList<>();
for (SseEmitter emitter : list) {
try {
emitter.send(SseEmitter.event()
.id(String.valueOf(id))
.name(eventName)
.data(data)
.reconnectTime(3000L));
} catch (IOException | IllegalStateException e) {
dead.add(emitter); // Connection dead, collect for batch cleanup
}
}
dead.forEach(e -> remove(userId, e));
}
// Heartbeat: comment line not received by any frontend listener, purely for keep-alive
private void ping() {
for (Map.Entry<String, CopyOnWriteArrayList<SseEmitter>> entry : emitters.entrySet()) {
List<SseEmitter> dead = new ArrayList<>();
for (SseEmitter emitter : entry.getValue()) {
try {
emitter.send(SseEmitter.event().comment("ping"));
} catch (Exception ex) {
dead.add(emitter);
}
}
dead.forEach(e -> remove(entry.getKey(), e));
}
}
}Controller stays thin, plus an ops stats endpoint:
@RestController
@RequestMapping("/sse")
@RequiredArgsConstructor
public class SseController {
private final SseEmitterManager manager;
@GetMapping(
value = "/subscribe/{userId}",
produces = MediaType.TEXT_EVENT_STREAM_VALUE
)
public SseEmitter subscribe(
@PathVariable String userId,
@RequestHeader(value = "Last-Event-ID", required = false) String lastEventId
) {
// On reconnect browser sends Last-Event-ID; use it to replay missed messages
if (lastEventId != null) {
// TODO: replay from DB/Redis by id
}
return manager.connect(userId);
}
@PostMapping("/push/{userId}")
public Map<String, Object> push(
@PathVariable String userId,
@RequestBody Map<String, Object> body
) {
manager.send(userId, "message", body);
return Map.of(
"ok", true,
"online", manager.onlineUsers().contains(userId)
);
}
@GetMapping("/stats")
public Map<String, Object> stats() {
return Map.of(
"connections", manager.count(),
"users", manager.onlineUsers()
);
}
}07 No Frontend Needed: curl Is Enough
# Terminal 1: subscribe (must use -N to disable curl's own buffering, otherwise stream invisible)
curl -N -H "Accept: text/event-stream" http://localhost:8080/sse/subscribe/1001
# Terminal 2: push
curl -X POST http://localhost:8080/sse/push/1001 \
-H "Content-Type: application/json" \
-d '{"msg":"你好, SSE"}'The -N flag is critical when debugging; without it curl buffers the response and you'll think the server isn't pushing.
08 Four Production Pitfalls You Must Handle
Pitfall 1: Default Async Thread Pool Has No Limit
SseEmitteruses async Servlet — it doesn't occupy business threads but consumes async threads. Spring Boot defaults to SimpleAsyncTaskExecutor — creates a new thread per request with no upper bound — long connections explode the thread count. Production must replace it:
@Configuration
public class AsyncConfig implements WebMvcConfigurer {
@Override
public void configureAsyncSupport(AsyncSupportConfigurer configurer) {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(50);
executor.setMaxPoolSize(200);
executor.setQueueCapacity(1000);
executor.setThreadNamePrefix("sse-");
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
executor.initialize();
configurer.setTaskExecutor(executor);
configurer.setDefaultTimeout(0L); // No timeout, rely on heartbeat & active cleanup
}
}On Spring Boot 3.2+ / JDK 21, simply enable virtual threads:
spring:
threads:
virtual:
enabled: truePitfall 2: Nginx Buffers the Stream
By default Nginx buffers upstream responses, turning "streaming" into "wait for all then return once". Server adds X-Accel-Buffering: no, Nginx side disables proxy_buffering:
location /sse/ {
proxy_pass http://backend;
proxy_buffering off; # Critical
proxy_cache off;
proxy_read_timeout 3600s; # Long connections need long timeout
proxy_http_version 1.1;
proxy_set_header Connection "";
}Same issue appears in Spring MVC's ResponseBodyEmitter buffering and various gateways (APISIX, cloud SLB) — symptom: works locally, breaks behind gateway.
Pitfall 3: Missing Heartbeat, Connection Silently Killed
ISPs, NATs, load balancers all drop "idle" connections. Send a comment line every 25 seconds : ping — cheapest keep-alive; comment lines are invisible to all frontend listeners.
Pitfall 4: Messages Lost After Reconnect
Browser auto-reconnects but only sends Last-Event-ID , does not replay data . If message loss is unacceptable, persist events by ID (DB or Redis), and on subscribe replay everything after the received Last-Event-ID.
09 AI Streaming: The Most Common Production Pattern
Pipe LLM token callbacks into SseEmitter for the ChatGPT typewriter effect:
@GetMapping(
value = "/ai/chat",
produces = MediaType.TEXT_EVENT_STREAM_VALUE
)
public SseEmitter chat(@RequestParam String prompt) {
SseEmitter emitter = new SseEmitter(5 * 60 * 1000L); // 5 min for streaming chat
CompletableFuture.runAsync(() -> {
try {
aiClient.streamChat(prompt, new AiClient.TokenListener() {
@Override
public void onToken(String token) {
try {
emitter.send(SseEmitter.event()
.name("delta")
.data(Map.of("content", token)));
} catch (IOException e) {
// Client disconnected, stop pushing
}
}
@Override
public void onFinish(String fullText) {
try {
emitter.send(SseEmitter.event()
.name("done")
.data("[DONE]"));
} catch (IOException ignored) {}
emitter.complete();
}
@Override
public void onError(Throwable t) {
emitter.completeWithError(t);
}
});
} catch (Exception e) {
emitter.completeWithError(e);
}
});
emitter.onTimeout(emitter::complete);
emitter.onError(e -> emitter.complete());
return emitter;
}Frontend:
const es = new EventSource(`/ai/chat?prompt=${encodeURIComponent(prompt)}`);
let answer = '';
es.addEventListener('delta', (e) => {
answer += JSON.parse(e.data).content;
renderMarkdown(answer); // Typewriter effect
});
es.addEventListener('done', () => es.close());Real-world detail: EventSource only does GET; long prompts exceed URL length limits. Correct pattern: POST context + large prompt first to obtain a sessionId, then GET /ai/chat?sessionId=xxx to open the SSE stream.
10 Cluster Deployment: SseEmitter Lives in JVM Memory
SseEmitter resides in a single JVM's heap. Multi-instance deployment: user connects to node A, event originates on node B — calling send on B won't reach the connection on A.
Standard solution: middleware for cross-node broadcast, each node pushes only its own connections :
// Publisher: any node publishes events to Redis channel
public void broadcastCluster(String eventName, Object data) {
redisTemplate.convertAndSend("sse:broadcast",
Map.of("event", eventName, "data", data));
}
// Subscriber: every node listens, pushes only to users connected locally
@Component
public class SseBroadcastListener implements MessageListener {
private final SseEmitterManager manager;
public SseBroadcastListener(SseEmitterManager manager) {
this.manager = manager;
}
@Override
public void onMessage(Message message, byte[] pattern) {
// Deserialize and broadcast by event name to this node's online users
manager.broadcast(eventName, data);
}
}If "offline must not lose" is required, Pub/Sub alone isn't enough — it's fire-and-forget, nodes offline lose events. Must also persist events by ID to DB/Stream , combined with Last-Event-ID for replay.
11 Summary
SSE isn't new tech; it's just a "GET that never ends". Choose it not because it's trendy, but because:
Direction match : unidirectional push is its home turf; don't bring a bidirectional sledgehammer.
Ultra-low cost : one HTTP request, text/event-stream, native browser support, auto-reconnect.
Simple rollout : Spring Boot + SseEmitter does it all, zero extra deps.
Streaming-friendly : de facto standard for AI typewriter effect, saves effort on both ends.
But for production, four things are non-negotiable: Replace default async thread pool (or enable virtual threads), disable Nginx buffering, send heartbeat every 25s, use Last-Event-ID for replay on reconnect; in clusters, wire Redis Pub/Sub or MQ for broadcast.
Next time you face "server must push messages", don't rush to WebSocket — ask: does the client need to speak? If no, SSE is enough.
Signed-in readers can open the original source through BestHub's protected redirect.
This article has been distilled and summarized from source material, then republished for learning and reference. If you believe it infringes your rights, please contactand we will review it promptly.
Java Tech Enthusiast
Sharing computer programming language knowledge, focusing on Java fundamentals, data structures, related tools, Spring Cloud, IntelliJ IDEA... Book giveaways, red‑packet rewards and other perks await!
How this landed with the community
Was this worth your time?
0 Comments
Thoughtful readers leave field notes, pushback, and hard-won operational detail here.
