# 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