主要是讲解SSE使用和websocket使用

WebSocket

独立协议ws,wss,可以发送文本和二进制。

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-websocket</artifactId>
</dependency>

STOMP协议

类似订阅发布模式,通过主题可以广播,点对点。

  • 核心配置类
@Configuration
@EnableWebSocketMessageBroker // 启用WebSocket消息代理
public class WebSocketStompConfig implements WebSocketMessageBrokerConfigurer {

    // 注册STOMP端点
    @Override
    public void registerStompEndpoints(StompEndpointRegistry registry) {
        // 客户端通过 "/ws" 路径连接WebSocket服务器
        registry.addEndpoint("/ws") // 端点
                .setAllowedOriginPatterns("*") // 允许跨域,生产环境请指定具体域名
                .withSockJS(); // 启用SockJS作为备用方案,以支持不支持WebSocket的旧浏览器
    }

    // 配置消息代理
    @Override
    public void configureMessageBroker(MessageBrokerRegistry registry) {
        // 设置消息代理的前缀,客户端订阅的地址需以此开头,如 /topic/public
        registry.enableSimpleBroker("/topic");
        // 设置应用程序前缀,客户端发送消息的地址需以此开头,如 /app/send
        registry.setApplicationDestinationPrefixes("/app");
    }
}
  • 业务控制器
@Controller
public class WebSocketController {

    // 前端发送消息至 /app/sendMessage,将自动路由至此方法
    @MessageMapping("/sendMessage")
    // 方法返回值将被发送给所有订阅了 /topic/public 的客户端
    @SendTo("/topic/public")
    public String processMessage(String message) throws Exception {
        // 你可以在这里进行消息处理,例如存储到数据库
        return "服务端回显: " + message;
    }
}
WebSocketHandler (原始处理器)

比较灵活,需要自定义双方数据交换模式

  • 配置类
@ServerEndpoint("/websocket")
public class WebSocketServer {
 
    @OnOpen
    public void onOpen(Session session) {
        System.out.println("Connection opened: " + session.getId());
        sessions.add(session);
    }
 
    @OnMessage
    public void onMessage(Session session, String message) throws IOException {
        System.out.println("Received message: " + message);
        session.getBasicRemote().sendText("Server received: " + message);
    }
 
    @OnClose
    public void onClose(Session session) {
        System.out.println("Connection closed: " + session.getId());
        sessions.remove(session);
    }
 
    private static final Set<Session> sessions = Collections.synchronizedSet(new HashSet<Session>());
}
  • 启用配置
@Configuration
@EnableWebSocket
public class WebSocketConfig implements WebSocketConfigurer {
    @Override
    public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) {
        registry.addHandler(new WebSocketServer(), "/websocket").setAllowedOrigins("*");
    }
}

SSE

比较轻量级的-只有服务端往客户端发送数据的一种协议,基于http,发送文本。

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-data-redis</artifactId>
</dependency>
  • SseEmitter单体
@RestController
public class SseController {
    private final Map<String, SseEmitter> emitterMap = new ConcurrentHashMap<>();

    @GetMapping("/sse/connect/{userId}")
    public SseEmitter connect(@PathVariable String userId) {
        SseEmitter emitter = new SseEmitter(60_000L); // 超时时间60秒
        emitterMap.put(userId, emitter);

        // 连接关闭或超时后的清理
        emitter.onCompletion(() -> emitterMap.remove(userId));
        emitter.onTimeout(() -> emitterMap.remove(userId));
        emitter.onError(e -> emitterMap.remove(userId));

        // 发送初始连接成功消息(可选)
        try {
            emitter.send(SseEmitter.event().name("init").data("connected"));
        } catch (IOException e) {
            emitter.completeWithError(e);
        }
        return emitter;
    }

    // 主动推送数据
    public void pushToUser(String userId, Object data) {
        SseEmitter emitter = emitterMap.get(userId);
        if (emitter != null) {
            try {
                emitter.send(SseEmitter.event().name("message").data(data));
            } catch (IOException e) {
                // 发送失败表示连接已断开,移除
                emitterMap.remove(userId);
            }
        }
    }
}
  • 分布式

保证分布式sse稳定运行。

连接:当用户携带唯一id连接sse时,负载均衡到一个后端业务服务器进行连接,连接完成之后,将用户信息和服务器ID存入redis

消息推送到前端:定时推送,定时服务实例拿到需要推送的sse连接。通过openfein调用指定的后端业务服务器,实现消息精确发送。

@Component
public class SseManager {
    private final Map<String, SseEmitter> localEmitters = new ConcurrentHashMap<>();
    private final RedisTemplate<String, String> redisTemplate;
    private final String instanceId;

    public SseManager(RedisTemplate<String, String> redisTemplate, 
                      @Value("${instance.id}") String instanceId) {
        this.redisTemplate = redisTemplate;
        this.instanceId = instanceId;
    }

    public SseEmitter connect(String userId) {
        SseEmitter emitter = new SseEmitter(120_000L);
        localEmitters.put(userId, emitter);
        redisTemplate.opsForHash().put("sse:user:node", userId, instanceId);
        redisTemplate.expire("sse:user:node", Duration.ofMinutes(2));

        emitter.onCompletion(() -> {
            localEmitters.remove(userId);
            redisTemplate.opsForHash().delete("sse:user:node", userId);
        });
        emitter.onTimeout(() -> {
            localEmitters.remove(userId);
            redisTemplate.opsForHash().delete("sse:user:node", userId);
        });
        return emitter;
    }

    public boolean pushToUser(String userId, Object data) {
        String targetInstance = (String) redisTemplate.opsForHash().get("sse:user:node", userId);
        if (targetInstance == null) return false;
        if (targetInstance.equals(instanceId)) {
            SseEmitter emitter = localEmitters.get(userId);
            if (emitter != null) {
                try {
                    emitter.send(data);
                    return true;
                } catch (IOException e) {
                    localEmitters.remove(userId);
                    redisTemplate.opsForHash().delete("sse:user:node", userId);
                }
            }
        } else {
            // 远程调用目标实例的推送接口(如 RestTemplate)
            callRemoteInstance(targetInstance, userId, data);
        }
        return false;
    }

    @Scheduled(fixedRate = 30000)
    public void heartbeat() {
        localEmitters.forEach((userId, emitter) -> {
            try {
                emitter.send(":heartbeat\n\n");
            } catch (IOException e) {
                localEmitters.remove(userId);
                redisTemplate.opsForHash().delete("sse:user:node", userId);
            }
        });
    }
}