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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 20:14:53