# websocket-spring-boot-starter **Repository Path**: tangshitai/websocket-spring-boot-starter ## Basic Information - **Project Name**: websocket-spring-boot-starter - **Description**: No description available - **Primary Language**: Java - **License**: Apache-2.0 - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 0 - **Created**: 2026-01-14 - **Last Updated**: 2026-07-03 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # websocket-spring-boot-starter ## 项目简介 **Maven 坐标**: `net.f3322.layman:websocket-spring-boot-starter:1.0.0` 基于 Spring Boot 的 WebSocket 自动配置 Starter(JSR-356),提供连接管理、三阶段过滤器链、按消息类型的业务职责链、下行消息缓存(内存或文件分卷)及运维门面。下文与**当前仓库源码**一致;若实现变更请以代码为准。 --- ## 快速开始 ### 1. 添加依赖 ```xml net.f3322.layman websocket-spring-boot-starter 1.0.0 ``` ### 2. 启用 WebSocket 在启动类上添加 `@EnableWS`: ```java @SpringBootApplication @EnableWS public class MyApplication { public static void main(String[] args) { SpringApplication.run(MyApplication.class, args); } } ``` `@EnableWS` 通过 `@Import` 导入 `WsAutoConfiguration`(核心自动配置)和 `WsBuiltinRegistrar`(内置过滤器/处理器扫描注册),同时标注 `@EnableWebSocket`。 ### 3. 创建端点 ```java @ServerEndpoint(value = "/ws/{uid}", configurator = WsConfigurator.class) @Component public class MyEndpoint extends AbstractWebSocketEndpoint { // 必须的构造函数注入 public MyEndpoint(WsConnectionManager connectionManager, WsMsgHandler handlerManager, WsMsgCacheManager cacheManager, List receivePreProcessFilters, List connectHandlers, WsExceptionHandler exceptionHandler, WsExecutorManager executorManager) { super(connectionManager, handlerManager, cacheManager, receivePreProcessFilters, connectHandlers, exceptionHandler, executorManager); } @Override protected String getServiceName() { return "myService"; } @Override protected boolean enableGroupCache() { return true; } @Override protected boolean enableServiceNameCache() { return true; } } ``` 端点路径参数 `{uid}` 必选,`{group}` 可选。 ### 4. 业务消息处理(MessageHandleFilter) ```java @Component @Order(10) public class MyChatHandler extends AbstractMessageHandleFilter { @Override public String supportType() { return "chat"; } @Override protected void doHandle(MessageHandleContext context, MessageHandleFilterChain chain) throws Exception { WebSocketMessage message = context.getMessage(); // 业务处理;需要继续链时:chain.doHandle(context); } } ``` 默认由 `DefaultWsMsgHandler` 收集全部 `MessageHandleFilter` Bean,按 `supportType()` 分组(转小写),组内以 `@Order` 排序。未匹配到的消息类型回退到 `*` 通配符处理器(`DefaultUnknownMessageHandler`)。 ### 5. 发送消息(WsMsgSender) ```java @Autowired private WsMsgSender wsMsgSender; public void demo(WsConnection conn) { WebSocketMessage msg = new WebSocketMessage("notice", conn.getUid()); wsMsgSender.sendToConnection(conn, msg); wsMsgSender.sendToUid(targetUid, msg); wsMsgSender.sendToService("myService", msg); wsMsgSender.sendToServiceAndGroup("myService", "myGroup", msg); wsMsgSender.sendToServiceAndUser("myService", "userId", msg); } ``` --- ## 自动配置机制 `WsAutoConfiguration`(`@ConditionalOnClass(javax.websocket.server.ServerEndpoint.class)`)注册以下 Bean(均为 `@ConditionalOnMissingBean`): | Bean | 默认实现 | 说明 | |------|---------|------| | `WsConfigurator` | `WsConfigurator` | JSR-356 Configurator(DI + 握手客户端信息) | | `ServerEndpointExporter` | Spring 标准 | 自动发现 `@ServerEndpoint` | | `WsConnectionManager` | `DefaultWsConnectionManager` | 连接管理(含心跳/超时/孤儿桶清理) | | `WsMsgCacheManager` | 由 `storeType` 决定 | `memory`→`DefaultWsMsgCacheManager`;`file`→`FileWsMsgCacheManager` | | `WsMsgHandler` | `DefaultWsMsgHandler` | 消息分发入口 | | `WsMsgSender` | `DefaultWsMsgSender` | 消息发送器(集成发送前过滤器链) | | `WsExceptionHandler` | `DefaultWsExceptionHandler` | 统一异常处理 | | `WsService` | `DefaultWsService` | 运维门面 | | `WsExecutorManager` | `DefaultWsExecutorManager` | 线程池管理 | `WsBuiltinRegistrar` 实现 `ImportBeanDefinitionRegistrar`,扫描 `net.f3322.layman.websocket` 包下实现 `WsConnectFilter`、`WsReceivePreFilter`、`WsSendMsgPreFilter`、`MessageHandleFilter` 的类并注册为 Bean(无需 `@Component`)。 --- ## 包结构与类说明 ``` net.f3322.layman.websocket ├── annotation/ │ └── EnableWS # Starter 启用注解 ├── autoconfigure/ │ ├── WsAutoConfiguration # 核心自动配置 │ ├── WsProperties # 全部配置属性(前缀 websocket) │ ├── WsConfigurator # 端点 Configurator(DI + 握手解析) │ └── WsBuiltinRegistrar # 内置过滤器/处理器扫描注册 ├── endpoint/ │ └── AbstractWebSocketEndpoint # 抽象端点(@OnOpen/@OnClose/@OnMessage/@OnError) ├── session/ │ ├── WsConnection # 逻辑连接接口 │ ├── DefaultWsConnection # 默认连接实现 │ ├── WsConnectionManager # 连接管理器接口 │ └── DefaultWsConnectionManager # 连接管理器(心跳/超时/孤儿桶清理) ├── filter/ │ ├── onopen/ # 建连过滤器链 │ │ ├── WsConnectFilter/WsConnectContext/WsConnectFilterChain │ │ └── impl/(7 个内置过滤器,按 @Order 排序) │ ├── onmessage/ # 入站预处理过滤器链 │ │ ├── WsReceivePreFilter/WsReceivePreContext/WsReceivePreFilterChain │ │ └── impl/(3 个内置过滤器) │ └── send/ # 发送前过滤器链 │ ├── WsSendMsgPreFilter/WsSendMsgPreContext/WsSendMsgPreFilterChain │ └── impl/(4 个内置过滤器) ├── handler/ │ ├── WsMsgHandler/DefaultWsMsgHandler # 消息分发器 │ └── business/ │ ├── MessageHandleFilter/AbstractMessageHandleFilter │ └── impl/(heartbeat/pong/forward/unknown 内置处理器) ├── message/ │ ├── WebSocketMessage # 消息基类 POJO │ ├── ForwardWebSocketMessage # 转发消息子类 │ └── ForwardScope # 转发范围枚举(UID/SERVICE_BROADCAST/SERVICE_GROUP/SERVICE_USER) ├── sender/ │ ├── WsMsgSender # 消息发送器接口 │ └── DefaultWsMsgSender # 默认发送器(集成发送前过滤器链) ├── executor/ │ ├── WsExecutorManager/DefaultWsExecutorManager # 线程池管理 ├── exception/ │ ├── WsException/AbsWsExceptionHandler/DefaultWsExceptionHandler # 异常处理 ├── cache/ │ ├── WsMsgCacheManager/AbstractWsMsgCacheManager # 缓存管理器接口/抽象基类 │ ├── DefaultWsMsgCacheManager # 内存型缓存(ConcurrentHashMap + CopyOnWriteArrayList) │ ├── FileWsMsgCacheManager # 文件型缓存(JSONL 分卷 .ws 文件) │ ├── FileWsVolumeFiles # 文件分卷名解析 │ ├── WsCachedMessageIterator/WsCachedMessageIterators # 迭代器 │ └── WsGroupCompositeStorageKey # 组级复合存储键(TLV + Base64) ├── service/ │ └── impl/DefaultWsService # 运维门面 └── support/ └── HandshakeClientInfoSupport # 握手中客户端信息解析 ``` --- ## 三阶段过滤器链 ### 阶段一:建连过滤器链(@OnOpen 触发) | 顺序 | 过滤器 | 功能 | |------|--------|------| | 1 | `DefaultWsConnectValidationFilter` | 校验 uid/service 非空 | | 2 | `DefaultWsConnectQueryParamsFilter` | 解析 URI 查询参数 | | 3 | `DefaultWsConnectClientInfoFilter` | 记录客户端 IP/端口 | | 4 | `DefaultWsConnectSecurityUserIdFilter` | 从 OAuth2 提取 userId | | 5 | `DefaultWsConnectCreationFilter` | 创建 DefaultWsConnection 并注册 | | 6 | `DefaultWsConnectSuccessFilter` | 下发 connect_success 通知 | | LOWEST | `DefaultWsConnectCachedMsgFilter` | 回放下行缓存 | 链末尾的 `validateConnectChainAllFiltersExecuted()` 确保链上每个配置的过滤器都实际执行,防止提前终止。 ### 阶段二:入站预处理过滤器链(@OnMessage 触发) | 顺序 | 过滤器 | 功能 | |------|--------|------| | HIGHEST | `DefaultWsReceivePreJsonConvertFilter` | JSON→WebSocketMessage 反序列化 | | — | `DefaultWsReceivePreValidationFilter` | 校验 type/uid/data 必填 | | — | `DefaultWsReceivePreCachingFilter` | 上行消息缓存(骨架) | `DefaultWsReceivePreJsonConvertFilter` 当 `type=forward`(大小写不敏感)时反序列化为 `ForwardWebSocketMessage`。 ### 阶段三:发送前过滤器链(逐条连接执行) | 顺序 | 过滤器 | 功能 | |------|--------|------| | 1 | `DefaultWsSendPreValidationFilter` | 校验必填字段 | | 2 | `DefaultWsSendPreSerializeFilter` | JSON 序列化 → context.payload | | 3 | `DefaultWsSendPreCachingFilter` | 按策略写入下行缓存 | | LOWEST | `DefaultWsSendPreTerminalFilter` | 调用 connection.sendMessage() 写出(不再调用 doFilter) | --- ## 连接管理 `DefaultWsConnectionManager` 维护三层索引: - `allSessions`: `ConcurrentHashMap` — sessionId → 连接 - `uidToSessionIds`: `ConcurrentHashMap>` — uid → sessionId 集合 - `userIdToSessionIds`: `ConcurrentHashMap>` — userId → sessionId 集合 **心跳机制**:`ScheduledExecutorService` 周期性(`heartbeatInterval`)执行: - 空闲 > `connectionTimeout`:关闭并移除 - 空闲 > `heartbeatInterval`:下行 `type=ping` 探活(不刷新活跃时间) - 客户端上行 `type=pong` 刷新活跃时间 **孤儿缓存桶清理**:同一心跳周期末尾,枚举缓存桶 key,清除无存活连接的桶。 --- ## 下行消息缓存 ### 写入与回放策略 | 服务级缓存 | 组级缓存 | 组非空 | 写入目标 | 回放来源 | |-----------|---------|--------|----------|---------| | 开 | 关 | 任意 | 服务级桶 | 服务级桶 | | 开 | 开 | 是 | 组级桶 | 组级桶 | | 开 | 开 | 否 | 不写入 | 不重放 | | 关 | 任意 | 任意 | 不写入 | 不重放 | ### 存储介质 | 类型 | 实现类 | 特点 | |------|--------|------| | `memory`(默认) | `DefaultWsMsgCacheManager` | ConcurrentHashMap + CopyOnWriteArrayList,单桶 FIFO 上限 `maxCachedMessages` | | `file` | `FileWsMsgCacheManager` | JSONL 分卷文件(UTF-8),按 `fileVolumeMaxBytes` 滚卷,公平锁+桶锁 | ### 文件型缓存详情 - 根目录默认 `java.io.tmpdir/websocket-layman-msg-cache/` - 服务桶在 `services/`,组桶在 `groups/` - 文件名:`{序号}_{stem}.ws`,stem 为 URL-Safe Base64(serviceType) 或 `WsGroupCompositeStorageKey.encode()` - 超出卷大小时下一条整行写入新卷;单条消息 UTF-8 长度超过卷上限时拒绝写入 - 迭代器须 `close()` 释放读锁 ### 组级复合存储键(WsGroupCompositeStorageKey) 两段 TLV(Tag(1B)+Length(4B int)+Value(UTF-8))拼接 → URL-Safe Base64,组为空时用保留字面量 `NULL-DEFAULT`。 --- ## 消息转发机制 `ForwardWebSocketMessage` + `ForwardMessageHandler`,客户端上行 `type=forward`: | ForwardScope | JSON 键 | 目标参数 | 功能 | |-------------|---------|----------|------| | `UID` | `uid` | forwardTargetUid | 发往指定 uid | | `SERVICE_BROADCAST` | `service_broadcast` | forwardServiceType | 发往指定服务所有连接 | | `SERVICE_GROUP` | `service_group` | forwardServiceType + forwardGroupType | 发往指定服务+组 | | `SERVICE_USER` | `service_user` | forwardServiceType + forwardTargetUserId | 发往指定服务+用户 | --- ## 配置说明 全部属性前缀 `websocket`,可通过 Spring Boot relaxed 绑定。 ```yaml websocket: connection-timeout: 180000 # 连接超时(ms),默认 3 分钟 heartbeat-interval: 30000 # 心跳间隔(ms),默认 30 秒 cache: max-cached-messages: 100 # 单桶 FIFO 上限(仅 memory) read-batch-size: 500 # 迭代器批读最大条数 store-type: memory # memory 或 file file-store-directory: "" # file 根目录 file-volume-max-bytes: 104857600 # 单卷最大字节(默认 100MB) file-volume-max-bytes-jitter-percent: 0 # 0~50,单卷上限随机抖动 error-handler: default-error-code: 500 default-error-message: "发生意外错误" error-message-format: '{"code": ${code}, "message": "${message}"}' include-details: false executor: core-pool-size: 10 max-pool-size: 20 keep-alive-seconds: 60 queue-capacity: 1000 thread-name-prefix: "websocket-msg-" work-queue-class-name: "java.util.concurrent.LinkedBlockingQueue" rejected-execution-handler-class-name: "java.util.concurrent.ThreadPoolExecutor$AbortPolicy" security: principal-account-field: "account" # OAuth2 principal 的账户字段名 ``` --- ## 运维门面(WsService) `WsService`(默认实现 `DefaultWsService`)以 `Map` 返回运维数据,**不包含**内置 REST Controller;宿主应用可自行暴露 HTTP 或定时任务调用: - **缓存**:`getCacheConfig`、`getServiceCacheConfig`、`getCacheStats`、`clearServiceLevelCache`、`clearServiceGroupCache`、`clearAllCache` - **配置与统计**:`getServerConfig`、`getConnectionStats`、`getServiceConnectionStats`、`getThreadPoolStats`、`getComprehensiveStats` - **连接**:`getAllConnections`、`getConnectionsByUid` --- ## 扩展点 | 扩展目标 | 操作 | |----------|------| | 自定义连接管理器 | 注册 `WsConnectionManager` 类型 Bean | | 自定义缓存管理器 | 继承 `AbstractWsMsgCacheManager`,注册 `WsMsgCacheManager` Bean;`getCacheStats()` 需提供 `serviceStats`/`serviceAndGroupStats` 键 | | 自定义消息发送器 | 注册 `WsMsgSender` Bean | | 自定义消息分发器 | 注册 `WsMsgHandler` Bean(高级用法) | | 自定义异常处理器 | 继承 `AbsWsExceptionHandler`,注册 `WsExceptionHandler` Bean | | 自定义线程池 | 注册 `WsExecutorManager` Bean | | 添加过滤器 | 实现对应 Filter 接口 + `@Component` + `@Order` | --- ## 常见问题 - **连接失败**:检查 `@EnableWS`、端口、`ServerEndpoint` 路径及日志 - **缓存不写入**:确认端点 `enableServiceNameCache` / `enableGroupCache` 与路径组是否满足"组为空不写入"等条件 - **文件缓存占锁过久**:避免长时间持有 `WsCachedMessageIterator` 不关闭 - **依赖注入不生效**:端点构造函数中必须注入所有参数并传给父类 ## 许可证 MIT License ## 作者 - Email: layman@f3322.net - GitHub: https://github.com/layman-f3322/websocket-spring-boot-starter