WebSocket 实时监控数据推送:5 秒推送 + 在线用户管理完整工程实现
WebSocket实时监控数据推送工程实现
引言
在实时监控场景中,服务端需要将系统运行指标(CPU使用率、内存占用、在线用户数等)以固定频率推送到前端,实现数据的实时可视化展示。传统的HTTP轮询方案存在延迟高、带宽浪费等问题,而WebSocket协议通过建立全双工通信通道,实现了服务端到客户端的低延迟数据推送。
本文将深入分析WebSocket协议原理,探讨Spring Boot中WebSocket的集成方案,以及实时监控数据5秒推送和在线用户管理的工程实现。
一、WebSocket协议原理
1.1 HTTP vs WebSocket
flowchart LR
subgraph HTTP轮询
direction TB
A1[客户端] -->|请求1| A2[服务端]
A2 -->|响应1| A1
A1 -->|请求2| A2
A2 -->|响应2| A1
A1 -->|请求3| A2
A2 -->|响应3| A1
A1 -.->|每次都建立连接| A2
end
subgraph WebSocket
direction TB
B1[客户端] -->|握手升级| B2[服务端]
B2 -->|101 Switching| B1
B1 <-->|双向通信| B2
B2 -.->|服务端主动推送| B1
B1 -.->|客户端发送消息| B2
end
style B2 fill:#4CAF50,color:#fff
1.2 WebSocket连接建立过程
sequenceDiagram
participant Client as 客户端
participant Server as 服务端
Note over Client,Server: 1. HTTP握手升级阶段
Client->>Server: GET /ws HTTP/1.1
Note right of Client: Upgrade: websocket<br/>Connection: Upgrade<br/>Sec-WebSocket-Key: xxx<br/>Sec-WebSocket-Version: 13
Server->>Client: HTTP/1.1 101 Switching Protocols
Note left of Server: Upgrade: websocket<br/>Connection: Upgrade<br/>Sec-WebSocket-Accept: yyy
Note over Client,Server: 2. WebSocket通信阶段
Server->>Client: 推送监控数据 (每5秒)
Client->>Server: 心跳包 (每30秒)
Server->>Client: 心跳响应
Note over Client,Server: 3. 连接关闭
Client->>Server: Close Frame
Server->>Client: Close Frame
1.3 WebSocket数据帧格式
0 1 2 3
0 1 2 3 4 5 6 7 8 9 0 1 2 3 4 5 6 7 8 9 0 1 2 3 4 5 6 7 8 9 0 1
+-+-+-+-+-------+-+-------------+-------------------------------+
|F|R|R|R| opcode|M| Payload len | Extended payload length |
|I|S|S|S| (4) |A| (7) | (16/64) |
|N|V|V|V| |S| | (if payload len==126/127) |
+-+-+-+-+-------+-+-------------+-------------------------------+
| Extended payload length continued, if payload len == 127 |
+-------------------------------+-------------------------------+
| |Masking-key, if MASK set to 1 |
+-------------------------------+-------------------------------+
| Masking-key (continued) | Payload Data |
+-------------------------------- - - - - - - - - - - - - - - - +
: Payload Data continued ... :
+---------------------------------------------------------------+
二、Spring Boot WebSocket集成
2.1 依赖与配置
<!-- pom.xml -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-websocket</artifactId>
</dependency>
/**
* WebSocket配置类
*/
@Configuration
public class WebSocketConfig {
@Bean
public ServerEndpointExporter serverEndpointExporter() {
return new ServerEndpointExporter();
}
}
2.2 WebSocket服务端点
/**
* 监控数据推送WebSocket端点
* 使用JSR-356标准注解方式
*/
@Component
@ServerEndpoint("/ws/monitor")
@Slf4j
public class MonitorWebSocket {
/**
* 线程安全的会话集合
* key: sessionId, value: Session
*/
private static final ConcurrentHashMap<String, Session> SESSION_MAP = new ConcurrentHashMap<>();
/**
* 连接建立成功
*/
@OnOpen
public void onOpen(Session session) {
SESSION_MAP.put(session.getId(), session);
log.info("WebSocket连接建立: sessionId={}, 当前在线数={}", session.getId(), SESSION_MAP.size());
// 连接建立后立即推送一次当前监控数据
try {
String initData = MonitorDataCollector.getCurrentData();
session.getBasicRemote().sendText(initData);
} catch (IOException e) {
log.error("推送初始监控数据失败", e);
}
}
/**
* 接收客户端消息
*/
@OnMessage
public void onMessage(String message, Session session) {
log.debug("收到客户端消息: sessionId={}, message={}", session.getId(), message);
// 处理心跳响应
if ("ping".equals(message)) {
sendMessage(session, "pong");
}
}
/**
* 连接关闭
*/
@OnClose
public void onClose(Session session) {
SESSION_MAP.remove(session.getId());
log.info("WebSocket连接关闭: sessionId={}, 当前在线数={}", session.getId(), SESSION_MAP.size());
}
/**
* 发生错误
*/
@OnError
public void onError(Session session, Throwable error) {
log.error("WebSocket发生错误: sessionId={}", session.getId(), error);
SESSION_MAP.remove(session.getId());
}
/**
* 广播消息给所有在线客户端
*/
public static void broadcast(String message) {
SESSION_MAP.forEach((sessionId, session) -> {
sendMessage(session, message);
});
}
/**
* 获取当前在线连接数
*/
public static int getOnlineCount() {
return SESSION_MAP.size();
}
/**
* 发送消息到指定会话
*/
private static void sendMessage(Session session, String message) {
if (session.isOpen()) {
try {
session.getBasicRemote().sendText(message);
} catch (IOException e) {
log.error("发送WebSocket消息失败: sessionId={}", session.getId(), e);
}
}
}
}
三、实时监控数据5秒推送
3.1 数据采集与推送架构
flowchart TD
subgraph 数据采集层
A1[系统指标采集] --> A2[MonitorDataCollector]
A2 --> A3[CPU/内存/磁盘/网络]
end
subgraph 定时调度层
B1[ScheduledExecutorService] -->|每5秒触发| B2[数据打包]
B2 --> B3[JSON序列化]
end
subgraph 推送层
B3 --> C1[MonitorWebSocket]
C1 --> C2[遍历在线Session]
C2 --> C3[逐个推送]
end
subgraph 前端展示层
C3 --> D1[WebSocket客户端]
D1 --> D2[数据解析]
D2 --> D3[ECharts图表更新]
end
A2 --> B1
style B1 fill:#4CAF50,color:#fff
style C1 fill:#2196F3,color:#fff
style D3 fill:#FF9800,color:#fff
3.2 监控数据采集器
/**
* 监控数据采集器
* 采集服务器运行指标
*/
@Component
@Slf4j
public class MonitorDataCollector {
private static final ObjectMapper objectMapper = new ObjectMapper();
/**
* 采集当前监控数据
*/
public static String getCurrentData() {
try {
MonitorData data = new MonitorData();
// CPU信息
OperatingSystemMXBean osBean = ManagementFactory.getOperatingSystemMXBean();
data.setCpuUsage(osBean.getSystemLoadAverage());
// 内存信息
Runtime runtime = Runtime.getRuntime();
long totalMemory = runtime.totalMemory();
long freeMemory = runtime.freeMemory();
data.setMemoryTotal(totalMemory);
data.setMemoryUsed(totalMemory - freeMemory);
data.setMemoryUsage((double) (totalMemory - freeMemory) / totalMemory * 100);
// JVM信息
data.setJvmMaxMemory(runtime.maxMemory());
data.setJvmTotalMemory(totalMemory);
data.setThreadCount(ManagementFactory.getThreadMXBean().getThreadCount());
// 在线用户数
data.setOnlineUsers(MonitorWebSocket.getOnlineCount());
// 时间戳
data.setTimestamp(System.currentTimeMillis());
return objectMapper.writeValueAsString(data);
} catch (Exception e) {
log.error("采集监控数据失败", e);
return "{}";
}
}
}
/**
* 监控数据DTO
*/
@Data
public class MonitorData {
/** CPU使用率 */
private Double cpuUsage;
/** 内存总量(字节) */
private Long memoryTotal;
/** 内存已用(字节) */
private Long memoryUsed;
/** 内存使用率(%) */
private Double memoryUsage;
/** JVM最大内存 */
private Long jvmMaxMemory;
/** JVM总内存 */
private Long jvmTotalMemory;
/** 线程数 */
private Integer threadCount;
/** 在线用户数 */
private Integer onlineUsers;
/** 采集时间戳 */
private Long timestamp;
}
3.3 定时推送调度器
/**
* 监控数据定时推送调度器
* 每5秒采集并推送一次监控数据
*/
@Component
@Slf4j
public class MonitorPushScheduler {
private final ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor(
r -> {
Thread t = new Thread(r, "monitor-push-thread");
t.setDaemon(true);
return t;
}
);
@PostConstruct
public void init() {
// 每5秒推送一次监控数据
scheduler.scheduleAtFixedRate(() -> {
try {
String data = MonitorDataCollector.getCurrentData();
MonitorWebSocket.broadcast(data);
} catch (Exception e) {
log.error("推送监控数据失败", e);
}
}, 0, 5, TimeUnit.SECONDS);
log.info("监控数据推送调度器启动,推送间隔: 5秒");
}
@PreDestroy
public void destroy() {
scheduler.shutdown();
}
}
四、前端WebSocket客户端
4.1 WebSocket Composable
/**
* WebSocket客户端composable
* 封装连接、断线重连、心跳、消息处理
*/
import { ref, onMounted, onUnmounted } from 'vue'
interface WebSocketOptions {
/** WebSocket服务地址 */
url: string
/** 心跳间隔(毫秒) */
heartbeatInterval?: number
/** 断线重连间隔(毫秒) */
reconnectInterval?: number
/** 最大重连次数 */
maxReconnectCount?: number
/** 消息处理回调 */
onMessage?: (data: any) => void
}
export function useWebSocket(options: WebSocketOptions) {
const {
url,
heartbeatInterval = 30000,
reconnectInterval = 5000,
maxReconnectCount = 5,
onMessage
} = options
const connected = ref(false)
const reconnectCount = ref(0)
let ws: WebSocket | null = null
let heartbeatTimer: ReturnType<typeof setInterval> | null = null
let reconnectTimer: ReturnType<typeof setTimeout> | null = null
/** 建立连接 */
const connect = () => {
if (ws?.readyState === WebSocket.OPEN) return
const protocol = window.location.protocol === 'https:' ? 'wss:' : 'ws:'
const wsUrl = `${protocol}//${window.location.host}${url}`
ws = new WebSocket(wsUrl)
ws.onopen = () => {
connected.value = true
reconnectCount.value = 0
startHeartbeat()
}
ws.onmessage = (event) => {
try {
const data = JSON.parse(event.data)
onMessage?.(data)
} catch {
// 非JSON消息(如心跳响应)
}
}
ws.onclose = () => {
connected.value = false
stopHeartbeat()
tryReconnect()
}
ws.onerror = () => {
connected.value = false
}
}
/** 断线重连 */
const tryReconnect = () => {
if (reconnectCount.value >= maxReconnectCount) return
reconnectCount.value++
reconnectTimer = setTimeout(connect, reconnectInterval)
}
/** 启动心跳 */
const startHeartbeat = () => {
heartbeatTimer = setInterval(() => {
if (ws?.readyState === WebSocket.OPEN) {
ws.send('ping')
}
}, heartbeatInterval)
}
/** 停止心跳 */
const stopHeartbeat = () => {
if (heartbeatTimer) {
clearInterval(heartbeatTimer)
heartbeatTimer = null
}
}
/** 主动断开连接 */
const disconnect = () => {
stopHeartbeat()
if (reconnectTimer) clearTimeout(reconnectTimer)
ws?.close()
ws = null
}
onMounted(connect)
onUnmounted(disconnect)
return { connected, reconnectCount, connect, disconnect }
}
4.2 监控页面集成
<script setup lang="ts">
import { useWebSocket } from '@/composables/useWebSocket'
import * as echarts from 'echarts'
const cpuData = ref<number[]>([])
const memoryData = ref<number[]>([])
const timeLabels = ref<string[]>([])
const { connected } = useWebSocket({
url: '/ws/monitor',
onMessage: (data: MonitorData) => {
// 更新图表数据
const time = new Date(data.timestamp).toLocaleTimeString()
timeLabels.value.push(time)
cpuData.value.push(data.cpuUsage)
memoryData.value.push(data.memoryUsage)
// 保留最近30个数据点
if (timeLabels.value.length > 30) {
timeLabels.value.shift()
cpuData.value.shift()
memoryData.value.shift()
}
// 更新ECharts图表
updateChart()
}
})
</script>
五、在线用户管理
5.1 在线用户追踪
/**
* 在线用户管理器
* 追踪WebSocket连接与用户的对应关系
*/
@Component
@Slf4j
public class OnlineUserManager {
/**
* 在线用户映射
* key: sessionId, value: 用户信息
*/
private static final ConcurrentHashMap<String, OnlineUser> ONLINE_USERS = new ConcurrentHashMap<>();
/**
* 用户上线
*/
public static void userOnline(String sessionId, Long userId, String username) {
OnlineUser user = new OnlineUser();
user.setSessionId(sessionId);
user.setUserId(userId);
user.setUsername(username);
user.setLoginTime(LocalDateTime.now());
user.setLastHeartbeat(LocalDateTime.now());
ONLINE_USERS.put(sessionId, user);
log.info("用户上线: userId={}, username={}, 当前在线数={}",
userId, username, ONLINE_USERS.size());
}
/**
* 用户下线
*/
public static void userOffline(String sessionId) {
OnlineUser removed = ONLINE_USERS.remove(sessionId);
if (removed != null) {
log.info("用户下线: userId={}, 当前在线数={}",
removed.getUserId(), ONLINE_USERS.size());
}
}
/**
* 更新心跳时间
*/
public static void updateHeartbeat(String sessionId) {
OnlineUser user = ONLINE_USERS.get(sessionId);
if (user != null) {
user.setLastHeartbeat(LocalDateTime.now());
}
}
/**
* 获取在线用户列表
*/
public static List<OnlineUser> getOnlineUsers() {
return new ArrayList<>(ONLINE_USERS.values());
}
/**
* 获取在线用户数
*/
public static int getOnlineCount() {
return ONLINE_USERS.size();
}
/**
* 清理超时连接(心跳超过60秒未更新)
*/
@Scheduled(fixedRate = 60000)
public static void cleanTimeoutConnections() {
LocalDateTime threshold = LocalDateTime.now().minusSeconds(60);
ONLINE_USERS.entrySet().removeIf(entry -> {
if (entry.getValue().getLastHeartbeat().isBefore(threshold)) {
log.warn("清理超时连接: userId={}", entry.getValue().getUserId());
return true;
}
return false;
});
}
}
5.2 完整推送时序
sequenceDiagram
participant Browser as 浏览器
participant WS as WebSocket端点
participant Scheduler as 推送调度器
participant Collector as 数据采集器
participant UserManager as 在线用户管理
Browser->>WS: 建立WebSocket连接
WS->>UserManager: userOnline(sessionId, userId)
WS->>Browser: 推送初始监控数据
loop 每5秒
Scheduler->>Collector: 采集监控数据
Collector-->>Scheduler: 返回JSON数据
Scheduler->>UserManager: 获取在线用户数
UserManager-->>Scheduler: 返回在线数
Scheduler->>WS: broadcast(data)
WS->>Browser: 推送监控数据
end
loop 每30秒
Browser->>WS: 发送心跳(ping)
WS->>UserManager: updateHeartbeat(sessionId)
WS->>Browser: 心跳响应(pong)
end
Browser->>WS: 关闭连接
WS->>UserManager: userOffline(sessionId)
六、安全与稳定性保障
6.1 认证鉴权
/**
* WebSocket握手拦截器
* 在握手阶段验证用户身份
*/
@Component
public class WebSocketAuthInterceptor extends HttpSessionHandshakeInterceptor {
@Override
public boolean beforeHandshake(ServerHttpRequest request, ServerHttpResponse response,
WebSocketHandler wsHandler, Map<String, Object> attributes) {
// 从请求中提取Token
if (request instanceof ServletServerHttpRequest servletRequest) {
String token = servletRequest.getServletRequest().getHeader("Authorization");
if (StringUtils.isBlank(token)) {
token = servletRequest.getServletRequest().getParameter("token");
}
// 验证Token有效性
try {
Object loginId = StpUtil.getLoginIdByToken(token);
if (loginId != null) {
attributes.put("userId", loginId);
return true;
}
} catch (Exception e) {
log.warn("WebSocket认证失败: token={}", token);
}
}
return false;
}
}
6.2 连接数限制
/**
* WebSocket连接数限制
*/
@Component
public class WebSocketConnectionLimiter {
/** 最大连接数 */
private static final int MAX_CONNECTIONS = 500;
/**
* 检查是否允许新连接
*/
public static boolean allowConnection() {
return MonitorWebSocket.getOnlineCount() < MAX_CONNECTIONS;
}
}
结论与建议
核心设计要点
| 设计要点 | 实现方式 | 作用 |
|---|---|---|
| 全双工通信 | WebSocket协议 | 服务端主动推送,低延迟 |
| 定时推送 | ScheduledExecutorService | 5秒间隔采集推送 |
| 心跳保活 | 客户端30秒心跳 | 检测连接存活 |
| 断线重连 | 指数退避重连 | 网络抖动自动恢复 |
| 在线管理 | ConcurrentHashMap | 追踪在线用户状态 |
| 认证鉴权 | 握手拦截器 | 防止未授权连接 |
| 连接限流 | 最大连接数限制 | 防止资源耗尽 |
最佳实践建议
- 心跳机制:客户端定时发送心跳,服务端检测超时连接并清理
- 消息压缩:大数据量推送时启用WebSocket permessage-deflate扩展
- 消息确认:关键消息增加确认机制,确保送达
- 优雅关闭:服务端关闭时先发送关闭帧,等待客户端确认
- 监控告警:对连接数异常、推送延迟进行监控告警