前端无法接收Server-Sent Events(SSE)问题排查求助
问题背景
前端通过React TypeScript的EventSource订阅后端Spring Boot(Java 17)的SSE事件,订阅接口调用成功、订阅者注册完成,但触发事件时前端收不到消息。后端用RedisCache保存订阅者,用于Pod故障/部署时恢复连接。
前端代码
订阅事件函数(组件useEffect中调用)
/** * This function registers and receives BE events to update data * similar to a webSocket, but only unidirectional from BE to FE * For now, just UNANSWERED_BUCKET is used, but in the future * can be used for other event types * @returns */ function registerUnansweredSSeEvent() { const sourceUnanswered = SseEventService.subscribeToEventsEventSource( context?.myExternalId, SseSubscriptionEvents.UNANSWERED_BUCKET, ); sourceUnanswered.addEventListener( SseSubscriptionEvents.UNANSWERED_BUCKET, (event) => { fillUnansweredAfterCategorization(event); }, ); sourceUnanswered.onopen = (event: any) => logger.info("Connection opened"); sourceUnanswered.onerror = (event: any) => logger.info("Connection error"); return sourceUnanswered; }
EventSource创建函数
subscribeToEventsEventSource( externalId: string | undefined, eventToSubscribe: string, ) { return new EventSource( `${this.baseUrl}/subscribe?personalExternalId=${externalId}&eventToSubscribe=${eventToSubscribe}`, ); }
后端代码
Controller层
@GetMapping(value = "/subscribe", produces = "text/event-stream") public SseEmitter subscribeSseEmitter( @RequestParam("personalExternalId") String personalExternalId, @RequestParam("eventToSubscribe") String eventToSubscribe){ return emitterService.registerClient(personalExternalId, eventToSubscribe); }
SSEService实现
@Service public class SSEService { private final NewRelicLogger logger = NewRelicLogger.getLogger(SSEService.class); private static final AtomicInteger ID_COUNTER = new AtomicInteger(1); //2hs maximum time Timeout. Should we increase this? public static final long DEFAULT_TIMEOUT = 7200000L; private final SlackNotificationService slackNotificationService; private final CacheService cacheService; public SSEService(SlackNotificationService slackNotificationService, CacheService cacheService) { this.slackNotificationService = slackNotificationService; this.cacheService = cacheService; } /** * @param subscriberExternalId: Email of the registered user to subscribe to notifications * @return This method subscribes the user to the events sent from the BE */ @AutomaticLogging public SseEmitter registerClient(String subscriberExternalId, String eventToSubscribe) { var emitter = new SseEmitter(DEFAULT_TIMEOUT); var sseClient = new SseClient(subscriberExternalId != null ? subscriberExternalId : UUID.randomUUID().toString(), emitter); addOrUpdateSseClient(sseClient, SseEventType.valueOf(eventToSubscribe)); emitter.onCompletion(() -> removeFromCache(subscriberExternalId)); emitter.onError(err -> removeAndLogError(sseClient, err.getMessage())); emitter.onTimeout(() -> removeAndLogError(sseClient, "TIMEOUT")); logger.info("New client registered {}", sseClient.getExternalId()); slackNotificationService.logInfo("New client registered " + sseClient.getExternalId()); return emitter; } private void removeFromCache(String subscriberExternalId){ try { cacheService.removeCacheKey(SSE_SUBSCRIBER, subscriberExternalId); } catch (Exception e) { logger.warn("Couldn't remove cache for subscriber: " + subscriberExternalId); } } @AutomaticLogging public void unregisterClient(String subscriberExternalId) { removeFromCache(subscriberExternalId); } /** * @param newClient Here we add a new client removing the previous one * to not keep open connections if the user reloads the page */ public void addOrUpdateSseClient(SseClient newClient, SseEventType newEventType) { String key = newClient.getExternalId(); SseClient registeredClientsObject = cacheService.getSSEClientFromCache(SSE_SUBSCRIBER, key); newClient.getSubscribedEvents().add(newEventType); if (registeredClientsObject == null ){ cacheService.putValueInCache(SSE_SUBSCRIBER, key, newClient); } else { if (!registeredClientsObject.getSubscribedEvents().contains(newEventType)){ List<SseEventType> events = registeredClientsObject.getSubscribedEvents(); events.add(newEventType); newClient.setSubscribedEvents(events); cacheService.putValueInCache(SSE_SUBSCRIBER, key, newClient); } } } /** * @param eventType The type of the event sent to the subscribers * So only subscribed to this eventType receives the notifications * Method used to Broadcast the messages to all subscribers * Use this for generic updates */ public void broadcastSseEmitterMessages(SseEventType eventType) { List<SseClient> clients = cacheService.getAllSSEClientFromCache() .stream() .filter(client -> client.getSubscribedEvents().contains(eventType)) .toList(); for (SseClient client : clients) { sendBroadcasterEmitterMessage(client, eventType); } } /** * @param client Subscriber who will receive the event * @param event event to notify */ @AutomaticLogging private void sendBroadcasterEmitterMessage(SseClient client, SseEventType event) { var sseEmitter = client.getSseEmitter(); try { logger.info("Notify client {}", client.getExternalId()); var eventId = ID_COUNTER.incrementAndGet(); SseEmitter.SseEventBuilder eventBuilder = SseEmitter.event().name(event.name()) .id(String.valueOf(eventId)) .data(event, MediaType.APPLICATION_JSON); sseEmitter.send(eventBuilder); } catch (IOException e) { sseEmitter.completeWithError(e); } } private void removeAndLogError(SseClient client, String error) { logger.error("Error during communication. Unregister client {}, error {}", client.getExternalId(), error); slackNotificationService.logError("Error during communication. Unregister client " + client.getExternalId() + " with error: " + error); removeFromCache(client.getExternalId()); } public void broadcastSseEmitterMessagesChatBuckets(List<SmallChannelDTO> channels, SseEventType eventType) { List<SseClient> clients = cacheService.getAllSSEClientFromCache(); for (SseClient client: clients) { sendBroadcastSseEmitterMessagesChatBuckets(client, channels, eventType); } } /** * @param client Subscriber who will receive the event for unanswered bucket on chat * @param channels the channels to update * @param eventType event to notify */ @AutomaticLogging private void sendBroadcastSseEmitterMessagesChatBuckets(SseClient client, List<SmallChannelDTO> channels, SseEventType eventType) { var sseEmitter = client.getSseEmitter(); try { logger.info("Notify client {} - sseEmitter {} - the event {}", client.getExternalId(), sseEmitter, eventType.name()); var eventId = ID_COUNTER.incrementAndGet(); SseEmitter.SseEventBuilder eventBuilder = SseEmitter.event().name(eventType.name()) .id(String.valueOf(eventId)) .data(channels, MediaType.APPLICATION_JSON); sseEmitter.send(eventBuilder); } catch (IOException e) { sseEmitter.completeWithError(e); } } }
Redis缓存操作代码
public SseClient getSSEClientFromCache(String cacheName, String key) { SpringRedisCache cache = (SpringRedisCache) cacheManager.getCache(cacheName); if (cache != null) { Cache.ValueWrapper valueWrapper = cache.get(key); if (valueWrapper != null) { return (SseClient) valueWrapper.get(); } } return null; } public List<SseClient> getAllSSEClientFromCache() { SpringRedisCache cache = (SpringRedisCache) cacheManager.getCache(SSE_SUBSCRIBER); if (cache != null) { List<String> allKeys = cache.getAllKeys().stream() .filter(c -> c.startsWith(SSE_SUBSCRIBER)) .map(c -> c.substring(SSE_SUBSCRIBER.length() + 1)) //we cut the key so only get the externalID .toList(); List<SseClient> clients = new ArrayList<>(); for (String key : allKeys) { SseClient client = this.getSSEClientFromCache(SSE_SUBSCRIBER, key); if (client != null) { clients.add(client); } } return clients; } return Collections.emptyList(); } public void putValueInCache(String cacheName, Object key, Object value) { SpringRedisCache cache = (SpringRedisCache) cacheManager.getCache(cacheName); if (cache != null) { cache.put(key, value); } }
排查方向
1. SseClient序列化失效问题
Redis缓存中存储的SseClient包含SseEmitter对象,而SseEmitter是Spring请求绑定对象,无法被序列化/反序列化。从Redis取出SseClient时,SseEmitter已失效,发送事件时实际是往无效连接发送,前端收不到消息。
- 验证:在
sendBroadcasterEmitterMessage中打印sseEmitter.isCompleted(),若返回true则说明Emitter已失效。
2. 事件名称不匹配
检查前端监听的SseSubscriptionEvents.UNANSWERED_BUCKET与后端发送的event.name()是否完全一致(大小写、拼写),SSE事件名大小写敏感。
3. 跨域配置缺失
若前后端域名不同,需在Spring Boot中配置SSE跨域支持:
@Configuration public class WebConfig implements WebMvcConfigurer { @Override public void addCorsMappings(CorsRegistry registry) { registry.addMapping("/subscribe") .allowedOrigins("前端域名") .allowedMethods("GET") .allowedHeaders("*") .allowCredentials(true) .maxAge(3600); } }
- 验证:查看浏览器控制台是否有跨域报错。
4. 连接超时或中断
前端EventSource长时间无消息可能自动重连,后端DEFAULT_TIMEOUT设为2小时,但网络波动可能导致连接提前断开。
- 验证:查看前端
onerror是否触发,后端是否有超时/错误日志。
5. 订阅者事件列表错误
检查addOrUpdateSseClient逻辑:已有订阅者时是否正确合并事件列表,确保触发事件时该订阅者的subscribedEvents包含目标事件类型。
- 验证:广播事件前打印目标订阅者的
subscribedEvents,确认包含当前触发的事件类型。
6. 负载均衡环境问题
若后端是多Pod部署,订阅请求落在Pod A,事件触发请求落在Pod B,Pod B从Redis取出的SseClient对应的SseEmitter属于Pod A,无法跨Pod发送SSE消息。需用消息队列(如Redis Pub/Sub、RabbitMQ)实现跨Pod事件广播:
- 方案:每个Pod订阅消息队列,事件触发时发送消息到队列,所有Pod接收后向本地
SseEmitter发送事件。
内容的提问来源于stack exchange,提问作者MatiCepe

