Spring Boot中注册SseEmitter到Hazelcast后SSE连接无法打开
Spring Boot SSE连接无法打开且客户端收不到通知问题
在Spring Boot应用中实现Server-Sent Events(SSE)功能,/sse/getNotification端点可调用发送通知,但/sse/notifications订阅端点的SSE连接无法打开。尽管已成功注册并将SseEmitter存入Hazelcast缓存,客户端仍无法接收通知。
代码实现
1. SSE注册端点(/sse/notifications)
@GetMapping("/sse/notifications") public ResponseEntity<SseEmitter> subscribeToNotifications(@RequestParam("userId") String userId) { log.info("SSE connection request for user: {}", userId); SseEmitter emitter = webPushService.registerSseEmitter(userId); if (emitter != null) { return ResponseEntity.ok() .header("Content-Type", "text/event-stream") .body(emitter); } else { return ResponseEntity.status(HttpStatus.INTERNAL_SERVER_ERROR).build(); } }
2. 注册SseEmitter
public SseEmitter registerSseEmitter(String userId) { log.info("Attempting to register SseEmitter for user: {}", userId); SseEmitter sseEmitter = new SseEmitter(0L); try { SerializableSseEmitter wrapper = new SerializableSseEmitter(userId, sseEmitter); emitterMap.put(userId, wrapper); // Hazelcast Map cache log.info("Successfully registered SseEmitter for user: {}", userId); // Handle events (completion, timeout, errors) sseEmitter.onCompletion(() -> log.info("SSE connection completed for user: {}", userId)); sseEmitter.onTimeout(() -> log.info("SSE connection timed out for user: {}", userId)); sseEmitter.onError((ex) -> log.error("SSE connection error for user: {}", userId)); } catch (Exception e) { log.error("Error registering SseEmitter for user: {} - {}", userId, e.getMessage()); } return sseEmitter; }
3. 发送通知
public void sendNotifications(String userId, String message) { log.info("sendNotifications for user: {}", userId); SseEmitter emitter = cacheService.getSseEmitter(userId); if (emitter != null) { try { emitter.send(SseEmitter.event().name("notification").data(message)); } catch (Exception e) { cacheService.removeSseEmitter(userId); log.error("Error sending notification to user: {}: {}", userId, e.getMessage()); } } else { log.warn("No active SseEmitter found for user: {}", userId); } }
4. 缓存服务
public SseEmitter getSseEmitter(String userId) { SerializableSseEmitter wrapper = emitterMap.get(userId); // Retrieve from Hazelcast cache if (wrapper != null) { return wrapper.getSseEmitter(); } else { log.warn("No SseEmitter found for user: {}", userId); return null; } } public void removeSseEmitter(String userId) { emitterMap.remove(userId); // Remove from Hazelcast cache }
SerializableSseEmitter类
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter; import java.io.*; public class SerializableSseEmitter implements Serializable { private static final long serialVersionUID = 1L; private transient SseEmitter sseEmitter; // SseEmitter is transient as it is not serializable private String userId; private volatile boolean completed = false; private volatile boolean timedOut = false; private volatile boolean errorOccurred = false; public SerializableSseEmitter(String userId, SseEmitter sseEmitter) { this.userId = userId; this.sseEmitter = sseEmitter; // Set the event handlers to track the state sseEmitter.onCompletion(() -> completed = true); sseEmitter.onTimeout(() -> timedOut = true); sseEmitter.onError((ex) -> errorOccurred = true); } public SseEmitter getSseEmitter() { if (this.sseEmitter == null) { this.sseEmitter = new SseEmitter(0L); } return this.sseEmitter; } public boolean isReady() { return !completed && !timedOut && !errorOccurred; } public String getUserId() { return userId; } public void setCompleted(boolean completed) { this.completed = completed; } public boolean isCompleted() { return completed; } private void writeObject(ObjectOutputStream oos) throws IOException { try { oos.defaultWriteObject(); } catch (IOException e) { throw new IOException("Error serializing SerializableSseEmitter for user: " + userId, e); } } private void readObject(ObjectInputStream ois) throws IOException, ClassNotFoundException { try { ois.defaultReadObject(); this.sseEmitter = new SseEmitter(0L); } catch (IOException | ClassNotFoundException e) { throw new IOException("Error deserializing SerializableSseEmitter for user: " + userId, e); } } @Override public String toString() { return "SerializableSseEmitter{" + "userId='" + userId + '\'' + ", completed=" + completed + ", sseEmitter=" + (sseEmitter != null ? "initialized" : "null") + '}'; } }
检查情况
- 订阅端点:日志显示收到订阅请求,但连接未打开
- 缓存存储:SseEmitter已成功存入Hazelcast的emitterMap
- 通知端点:调用
emitter.send()成功,但客户端收不到事件 - 缓存检索:发送通知时可从缓存正确获取SseEmitter
- 事件处理器:已配置onCompletion、onTimeout、onError,但未触发
预期行为
/sse/notifications端点应为用户注册SseEmitter并打开SSE连接,/sse/getNotification端点应向注册的SseEmitter发送通知,客户端需接收该通知事件。
内容的提问来源于stack exchange,提问作者Siva Suriyan
相关产品推荐
相关产品推荐

