You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.20 17:47:11