当前位置:首页 > 随笔 > 正文内容

在 Spring Boot 中构建 WebSocket 实时通信服务

访客 随笔 2026年9月20日 11

依赖与环境

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 实时通信链路已正常运作。

相关文章

可以按小时收费的VPS

很多 VPS 提供商都支持 按小时计费(hourly billing),想短期试用 / 临时搭建节点、测试网络、短期项目等场景非常合适。下面是当前最主流且靠谱的按小时 VPS 选项,分别按不同需求场景整理: 1. Vultr(全球节点,包括日本) 按小时计费 可选机房:东京 / 大阪 / 洛杉矶 / 法兰克福 / 伦敦 … 支持 PayPal(部分情况),但更常用信用卡/PayPal+卡价格参考$...

在 iPhone 上下载国外App

地区/国家限制App Store 会根据 Apple ID 的国家或地区限制应用下载。如果你的 Apple ID 绑定的是中国大陆,就可能无法下载 OpenAI 官方的 ChatGPT 应用,因为它在大陆 App Store 不上架。解决办法:换成美国、加拿大、香港等地区的 Apple ID。或者在现有 Apple ID 上更改地区。注册一个国外 Apple ID(推荐)比如注册 美国区 Appl...

Node.js 中的异步编程:回调与 Promise

Node.js 是一个基于 JavaScript 构建的单线程、非阻塞运行环境,它通过异步编程机制来高效处理多个操作。在执行如文件读取、API 请求或数据库查询等任务时,Node.js 不会等待这些操作完成,而是使用回调函数和 Promise 来避免阻塞主线程。 回调方式实现异步 那么当异步操作完成后,Node.js 如何知道接下来要做什么呢?这就要用到 回调函数(callback)。 回调本质上...

Selenium自动化测试入门指南

Selenium自动化测试入门指南

什么是自动化测试? 自动化测试是指利用软件工具自动执行测试用例,模拟用户操作,如打开网页、点击链接、输入文本等,并验证结果是否符合预期。 其主要优点包括: 大幅减少人工成本 测试速度快 可以在非工作时间运行 支持持续集成和交付 然而,它也存在一些局限性,例如开发成本较高、不适合快速变化的项目、依赖稳定的UI界面等。 自动化测试的应用条件 适合引入自动化测试的情况包括: 手动测试耗时且需要大量...

MariaDB Galera集群故障快速恢复指南

OpenStack控制节点采用三节点MariaDB Galera集群架构。当数据库集群因故障重启时,有时会出现Galera集群无法正常启动的问题。虽然有多种方法可以恢复数据库服务,但如何实现快速启动同时确保数据完整性呢? 通过分析日志发现,MariaDB Galera集群节点宕机时会在日志中输出以下信息: [Note] WSREP: 新集群视图:全局状态: 874d8e7e-5980-11e8-8...

Android 中 EventBus 的通信机制与实现原理深度解析

EventBus 核心设计思想 EventBus 是一个基于观察者模式的事件总线框架,广泛应用于 Android 平台以实现组件解耦。它通过中心化的消息分发机制,使不同层级、不同线程的对象能够以"发布-订阅"方式通信,避免了传统接口回调或广播带来的强依赖问题。 核心角色说明 事件(Event):任意 Java 对象,作为数据载体,如网络状态变更通知、用户登录信息等。 发布者(Publi...

发表评论

访客

◎欢迎参与讨论,请在这里发表您的看法和观点。