在 Spring Boot 中构建 WebSocket 实时通信服务
依赖与环境
Spring Boot 内置的 Tomcat 和 Jetty 容器已经原生支持 JSR-356 WebSocket 规范。因此,在大多数情况下,只需引入 spring-boot-starter-web 即可直接使用 WebSocket 功能,无需额外添加复杂的依赖。
注册端点导出器
在使用 Spring Boot 内置容器时,必须注册 ServerEndpointExporter Bean。该组件负责扫描并注册所有被 @ServerEndpoint 注解标记的类。如果项目打包为 WAR 并部署到外部 Servlet 容器,则此配置可以省略。
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.web.socket.server.standard.ServerEndpointExporter;
@Configuration
public class WebSocketInfrastructureConfig {
@Bean
public ServerEndpointExporter endpointExporter() {
return new ServerEndpointExporter();
}
}
会话管理与端点实现
在标准的 JSR-356 规范中,每次客户端建立连接时,容器都会为 @ServerEndpoint 类创建一个全新的实例。这会导致 Spring 的依赖注入(如 @Autowired)在端点实例中失效。为了解决这一生命周期冲突并保持架构的整洁,我们将会话(Session)的管理逻辑剥离到一个独立的 Spring 单例组件中。
会话管理器
通过静态变量持有 Spring Bean 实例,使得 WebSocket 端点能够以静态方式调用会话管理方法,同时享受 Spring 容器的生命周期管理。
import org.springframework.stereotype.Component;
import javax.annotation.PostConstruct;
import javax.websocket.Session;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
@Component
public class WebSocketSessionManager {
private static WebSocketSessionManager instance;
private final Map<String, Session> activeConnections = new ConcurrentHashMap<>();
@PostConstruct
public void init() {
instance = this;
}
public static void register(Session session) {
instance.activeConnections.put(session.getId(), session);
}
public static void unregister(Session session) {
instance.activeConnections.remove(session.getId());
}
public static void broadcast(String payload) {
instance.activeConnections.values().forEach(session -> {
try {
session.getAsyncRemote().sendText(payload);
} catch (Exception e) {
// 记录发送失败的日志
}
});
}
public static int getOnlineCount() {
return instance.activeConnections.size();
}
}
WebSocket 端点
定义具体的 WebSocket 路径和事件处理逻辑。此处将路径设定为 /ws/events,并委托 WebSocketSessionManager 处理连接的注册与注销。
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
import javax.websocket.OnClose;
import javax.websocket.OnError;
import javax.websocket.OnMessage;
import javax.websocket.OnOpen;
import javax.websocket.Session;
import javax.websocket.server.ServerEndpoint;
@Slf4j
@Component
@ServerEndpoint("/ws/events")
public class EventNotificationEndpoint {
private String sessionId;
@OnOpen
public void handleConnectionOpen(Session session) {
this.sessionId = session.getId();
WebSocketSessionManager.register(session);
log.info("新连接建立: {}, 当前在线数: {}", sessionId, WebSocketSessionManager.getOnlineCount());
// 发送欢迎消息
session.getAsyncRemote().sendText("连接成功,您的会话ID为: " + sessionId);
}
@OnMessage
public void handleIncomingMessage(String message, Session session) {
log.info("接收到客户端 [{}] 消息: {}", sessionId, message);
// 针对特定客户端回复确认信息
session.getAsyncRemote().sendText("服务端已接收: " + message);
}
@OnClose
public void handleConnectionClose(Session session) {
WebSocketSessionManager.unregister(session);
log.info("连接已断开: {}, 当前在线数: {}", sessionId, WebSocketSessionManager.getOnlineCount());
}
@OnError
public void handleTransportError(Session session, Throwable throwable) {
log.error("连接 [{}] 发生异常", sessionId, throwable);
WebSocketSessionManager.unregister(session);
}
}
定时推送任务测试
为了验证广播功能,可以创建一个 Spring 定时任务,周期性地向所有在线客户端推送系统状态或模拟数据。
import org.springframework.scheduling.annotation.EnableScheduling;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
@Component
@EnableScheduling
public class MetricsPushTask {
private static final DateTimeFormatter FORMATTER = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
@Scheduled(fixedRate = 5000)
public void pushSystemMetrics() {
String timestamp = LocalDateTime.now().format(FORMATTER);
String payload = String.format("{\"type\":\"heartbeat\", \"time\":\"%s\"}", timestamp);
WebSocketSessionManager.broadcast(payload);
}
}
客户端连接验证
启动 Spring Boot 应用后,可以使用 ApiFox、Postman 或浏览器开发者工具建立 WebSocket 连接。连接地址格式为:
ws://localhost:8080/ws/events
连接成功后,客户端将立即收到欢迎消息,并且每隔 5 秒会接收到来自 MetricsPushTask 推送的 JSON 格式心跳数据。通过观察控制台日志和客户端接收到的数据流,即可确认 WebSocket 实时通信链路已正常运作。
