Spring Boot 3开启SSE流时的优雅关闭方案
解决Spring Boot 3优雅关闭时SSE(及WebSocket)长连接阻塞问题
一、SSE连接的主动关闭方案
默认的graceful shutdown会等待所有请求完成,但SSE的无限Flux流不会主动终止,导致应用等待超时。要解决这个问题,需要主动跟踪并关闭所有活跃的SSE连接:
1. 实现连接跟踪组件
创建一个组件管理所有活跃SSE连接的Disposable(用于终止Flux流):
@Component public class SseConnectionTracker { private final Map<String, Disposable> activeConnections = new ConcurrentHashMap<>(); public void registerConnection(String connectionId, Disposable disposable) { activeConnections.put(connectionId, disposable); } public void unregisterConnection(String connectionId) { activeConnections.remove(connectionId); } public void closeAllConnections() { activeConnections.values().forEach(Disposable::dispose); activeConnections.clear(); } }
2. 修改SSE端点代码,注册并清理连接
在SSE接口中,将Flux的Disposable注册到跟踪组件,并在流结束时自动清理:
@RestController public class NotificationController { private final SseConnectionTracker connectionTracker; private final NotificationService notificationService; public NotificationController(SseConnectionTracker connectionTracker, NotificationService notificationService) { this.connectionTracker = connectionTracker; this.notificationService = notificationService; } @RequestMapping(path = "/notifications", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public Flux<ServerSentEvent<OnlineNotificationEvent>> notificationsSse(HttpServletRequest request) { // 用Session ID作为唯一连接标识,也可以结合客户端IP+UUID生成更唯一的ID String connectionId = request.getSession().getId(); return notificationService.getContinuousNotifications() .map(event -> ServerSentEvent.builder(event).build()) // 订阅时注册连接 .doOnSubscribe(disposable -> connectionTracker.registerConnection(connectionId, disposable)) // 流结束(正常/异常/取消)时移除连接 .doFinally(signal -> connectionTracker.unregisterConnection(connectionId)) .onErrorResume(e -> Flux.empty()); } }
3. 监听上下文关闭事件,主动终止所有连接
通过Spring的ContextClosedEvent触发连接关闭:
@Component public class GracefulShutdownListener implements ApplicationListener<ContextClosedEvent> { private final SseConnectionTracker connectionTracker; public GracefulShutdownListener(SseConnectionTracker connectionTracker) { this.connectionTracker = connectionTracker; } @Override public void onApplicationEvent(ContextClosedEvent event) { // 关闭所有活跃的SSE连接 connectionTracker.closeAllConnections(); } }
二、该方案是否适用于WebSocket?
思路通用,但实现细节不同。WebSocket需要跟踪WebSocketSession而非Flux的Disposable:
1. 实现WebSocket连接跟踪组件
@Component public class WebSocketConnectionTracker { // 用ConcurrentHashMap的keySet保证线程安全 private final Set<WebSocketSession> activeSessions = ConcurrentHashMap.newKeySet(); public void addSession(WebSocketSession session) { activeSessions.add(session); } public void removeSession(WebSocketSession session) { activeSessions.remove(session); } public void closeAllSessions() { activeSessions.forEach(session -> { try { // 发送服务重启的关闭状态,让前端可以自动重连 session.close(CloseStatus.SERVICE_RESTART); } catch (IOException e) { LoggerFactory.getLogger(getClass()).error("关闭WebSocket会话失败", e); } }); activeSessions.clear(); } }
2. 在WebSocket处理器中注册会话
@Configuration @EnableWebSocket public class WebSocketConfig implements WebSocketConfigurer { private final WebSocketConnectionTracker connectionTracker; public WebSocketConfig(WebSocketConnectionTracker connectionTracker) { this.connectionTracker = connectionTracker; } @Override public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) { registry.addHandler(new NotificationWebSocketHandler(), "/ws/notifications") .setAllowedOrigins("*"); } private class NotificationWebSocketHandler extends TextWebSocketHandler { @Override public void afterConnectionEstablished(WebSocketSession session) throws Exception { connectionTracker.addSession(session); super.afterConnectionEstablished(session); } @Override public void afterConnectionClosed(WebSocketSession session, CloseStatus status) throws Exception { connectionTracker.removeSession(session); super.afterConnectionClosed(session, status); } } }
3. 扩展优雅关闭监听器,处理WebSocket连接
修改之前的GracefulShutdownListener,加入WebSocket的关闭逻辑:
@Component public class GracefulShutdownListener implements ApplicationListener<ContextClosedEvent> { private final SseConnectionTracker sseConnectionTracker; private final WebSocketConnectionTracker webSocketConnectionTracker; public GracefulShutdownListener(SseConnectionTracker sseConnectionTracker, WebSocketConnectionTracker webSocketConnectionTracker) { this.sseConnectionTracker = sseConnectionTracker; this.webSocketConnectionTracker = webSocketConnectionTracker; } @Override public void onApplicationEvent(ContextClosedEvent event) { sseConnectionTracker.closeAllConnections(); webSocketConnectionTracker.closeAllSessions(); } }
关键说明
- 默认优雅关闭不生效的原因:SSE/WebSocket属于长连接,请求不会主动结束,Spring的优雅关闭会等待请求完成直到超时(默认30秒)。
- 连接标识要保证唯一:可以用Session ID、客户端IP+UUID、自定义请求头标识,避免不同客户端连接混淆。
- 前端适配:主动关闭连接后,前端可监听SSE的
error/close事件、WebSocket的指定关闭状态,自动触发重连逻辑。
内容的提问来源于stack exchange,提问作者Lubo
相关产品推荐
相关产品推荐

