Spring Boot中SSE事件需调用SseEmitter.complete后才接收的问题
问题描述
我是Server Sent Events(SSE)的新手,正在Spring Boot应用中尝试实现该功能。目前遇到的问题是,SSE事件似乎要等到调用SseEmitter的complete方法后才会被接收。在以下示例中,我需要等待10秒才能收到全部10个事件以及完成消息,而非每秒接收一个事件。
前端代码
<button id="startevents" type="button" class="btn btn-primary"> Start Events </button> <script> document.getElementById('startevents').addEventListener('click', function() { const eventSource = new EventSource('/events'); eventSource.addEventListener('myevent', function(event) { console.log('myevent: ', event.data); }); eventSource.addEventListener('complete', function(event) { console.log('complete: ', event.data); eventSource.close(); }); eventSource.addEventListener('error', function(event) { console.log('error: ', event); }); }); </script>
后端代码
@GetMapping( "/events" ) public SseEmitter republishEvents() { SseEmitter sseEmitter = new SseEmitter(Long.MAX_VALUE); sseEmitter.onCompletion(() -> { log.info("SSE complete"); }); sseEmitter.onTimeout(() -> { log.warn("SSE timeout"); sseEmitter.complete(); }); sseEmitter.onError((e) -> { log.error("SSE error", e); sseEmitter.complete(); }); ExecutorService executor = Executors.newSingleThreadExecutor(); executor.execute(() -> { try { for ( int i = 0; i < 10; i++ ) { log.debug( "Sending event " + i ); SseEmitter.SseEventBuilder event = SseEmitter.event() .name("myevent") .data( "Event " + i ); sseEmitter.send(event); Thread.sleep( 1000 ); } log.debug( "Sending complete event" ); SseEmitter.SseEventBuilder event = SseEmitter.event() .name("complete") .data( "complete" ); sseEmitter.send(event); sseEmitter.complete(); } catch ( Exception e ) { log.error( "Error sending event", e ); sseEmitter.complete(); } finally { executor.shutdown(); } }); return sseEmitter; }
问题原因
这种现象是因为响应缓冲导致的:Spring Boot默认会缓冲响应内容,直到缓冲区被填满或者响应完成(调用complete())时,才会一次性把所有内容发送给客户端,而非实时推送单个事件。
解决方法
修改后端代码,禁用响应缓冲并在每个事件发送后手动刷新响应,确保事件能立即推送到客户端:
import jakarta.servlet.http.HttpServletResponse; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RestController; import org.springframework.web.servlet.mvc.method.annotation.SseEmitter; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; @RestController public class SseController { private static final Logger log = LoggerFactory.getLogger(SseController.class); @GetMapping("/events") public SseEmitter republishEvents(HttpServletResponse response) { // 配置SSE专属响应头,禁用缓冲 response.setContentType("text/event-stream"); response.setCharacterEncoding("UTF-8"); response.setHeader("Cache-Control", "no-cache, no-store, must-revalidate"); response.setHeader("Pragma", "no-cache"); response.setHeader("Connection", "keep-alive"); SseEmitter sseEmitter = new SseEmitter(Long.MAX_VALUE); sseEmitter.onCompletion(() -> log.info("SSE complete")); sseEmitter.onTimeout(() -> { log.warn("SSE timeout"); sseEmitter.complete(); }); sseEmitter.onError((e) -> { log.error("SSE error", e); sseEmitter.complete(); }); ExecutorService executor = Executors.newSingleThreadExecutor(); executor.execute(() -> { try { for (int i = 0; i < 10; i++) { log.debug("Sending event " + i); SseEmitter.SseEventBuilder event = SseEmitter.event() .name("myevent") .data("Event " + i); sseEmitter.send(event); // 刷新响应,强制推送当前事件 response.flushBuffer(); Thread.sleep(1000); } log.debug("Sending complete event"); SseEmitter.SseEventBuilder completeEvent = SseEmitter.event() .name("complete") .data("complete"); sseEmitter.send(completeEvent); response.flushBuffer(); sseEmitter.complete(); } catch (Exception e) { log.error("Error sending event", e); sseEmitter.complete(); } finally { executor.shutdown(); } }); return sseEmitter; } }
关键改动说明
- 设置SSE响应头:明确指定响应类型为
text/event-stream,同时禁用缓存、保持长连接,告诉浏览器这是一个实时流,不要缓冲内容。 - 手动刷新响应:每次调用
sseEmitter.send()后执行response.flushBuffer(),强制将当前事件内容推送给客户端,避免等待缓冲区填满。
前端代码无需修改,调整后即可实现每秒接收一个事件的预期效果。
内容的提问来源于stack exchange,提问作者Steve Maring
相关产品推荐
相关产品推荐

