主要是讲解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);
}
});
}
}