基于 Socket.IO 实现 WebSocket 长连接
·
基于 Socket.IO 实现 WebSocket 长连接(Java 版)
本文介绍如何使用 Java 技术栈实现 Socket.IO WebSocket 长连接服务,基于 netty-socketio 框架。
目录
什么是 Socket.IO
Socket.IO 是一个基于事件的实时通信库,提供了:
- 实时双向通信:客户端和服务端可以随时互相发送消息
- 自动重连:连接断开后自动尝试重连
- 心跳检测:自动检测连接状态
- 跨浏览器支持:支持所有主流浏览器
- 房间(Room):支持分组广播
- 命名空间(Namespace):支持多路复用
Java 实现方案
netty-socketio 简介
netty-socketio 是 Socket.IO 的 Java 服务端实现,基于 Netty 框架:
┌─────────────────────────────────────────────────────────┐
│ netty-socketio 架构 │
├─────────────────────────────────────────────────────────┤
│ Application Layer (你的业务代码) │
├─────────────────────────────────────────────────────────┤
│ netty-socketio (Socket.IO 协议实现) │
├─────────────────────────────────────────────────────────┤
│ Netty (高性能网络通信框架) │
├─────────────────────────────────────────────────────────┤
│ JDK NIO (Java 原生非阻塞 IO) │
└─────────────────────────────────────────────────────────┘
核心依赖
<!-- Maven 依赖 -->
<dependencies>
<!-- netty-socketio 核心 -->
<dependency>
<groupId>com.corundumstudio.socketio</groupId>
<artifactId>netty-socketio</artifactId>
<version>2.0.9</version>
</dependency>
<!-- Spring Boot Starter(可选) -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
<version>3.2.0</version>
</dependency>
<!-- Lombok(简化代码) -->
<dependency>
<groupId>org.projectlombok</groupId>
<artifactId>lombok</artifactId>
<version>1.18.30</version>
<scope>provided</scope>
</dependency>
<!-- JSON 处理 -->
<dependency>
<groupId>com.fasterxml.jackson.core</groupId>
<artifactId>jackson-databind</artifactId>
<version>2.16.0</version>
</dependency>
<!-- 日志 -->
<dependency>
<groupId>ch.qos.logback</groupId>
<artifactId>logback-classic</artifactId>
<version>1.4.14</version>
</dependency>
</dependencies>
环境准备
项目结构
socketio-java-demo/
├── pom.xml
├── src/
│ ├── main/
│ │ ├── java/
│ │ │ └── com/
│ │ │ └── example/
│ │ │ └── socketio/
│ │ │ ├── SocketIoApplication.java
│ │ │ ├── config/
│ │ │ │ └── SocketIoConfig.java
│ │ │ ├── handler/
│ │ │ │ ├── ConnectHandler.java
│ │ │ │ ├── MessageHandler.java
│ │ │ │ └── DisconnectHandler.java
│ │ │ ├── model/
│ │ │ │ ├── Message.java
│ │ │ │ └── User.java
│ │ │ └── service/
│ │ │ ├── MessageService.java
│ │ │ └── RoomService.java
│ │ └── resources/
│ │ ├── application.yml
│ │ └── logback.xml
│ └── test/
│ └── java/
└── client/
└── index.html
Gradle 配置
// build.gradle
plugins {
id 'java'
id 'org.springframework.boot' version '3.2.0'
}
group = 'com.example'
version = '1.0.0'
sourceCompatibility = '17'
repositories {
mavenCentral()
}
dependencies {
implementation 'org.springframework.boot:spring-boot-starter-web'
implementation 'com.corundumstudio.socketio:netty-socketio:2.0.9'
compileOnly 'org.projectlombok:lombok'
annotationProcessor 'org.projectlombok:lombok'
testImplementation 'org.springframework.boot:spring-boot-starter-test'
}
tasks.named('test') {
useJUnitPlatform()
}
服务端实现
基础服务端(纯 Java)
package com.example.socketio;
import com.corundumstudio.socketio.Configuration;
import com.corundumstudio.socketio.SocketIOClient;
import com.corundumstudio.socketio.SocketIOServer;
import com.corundumstudio.socketio.listener.ConnectListener;
import com.corundumstudio.socketio.listener.DataListener;
import com.corundumstudio.socketio.listener.DisconnectListener;
import lombok.extern.slf4j.Slf4j;
/**
* Socket.IO 服务端主类
*/
@Slf4j
public class SocketIoServer {
private SocketIOServer server;
public void start() throws InterruptedException {
// 配置 Socket.IO 服务器
Configuration config = new Configuration();
config.setHostname("localhost");
config.setPort(9092);
// 设置跨域
config.setOrigin("*");
// 设置最大帧大小
config.setMaxFramePayloadLength(1024 * 1024);
// 设置 HTTP 超时
config.setPingTimeout(60000);
config.setPingInterval(25000);
// 创建服务器
server = new SocketIOServer(config);
// 添加监听器
addListeners(server);
// 启动服务器
server.start();
log.info("Socket.IO 服务器已启动,监听端口: 9092");
// 保持运行
Thread.sleep(Integer.MAX_VALUE);
}
private void addListeners(SocketIOServer server) {
// 连接监听器
server.addConnectListener(new ConnectListener() {
@Override
public void onConnect(SocketIOClient client) {
log.info("客户端连接: sessionId = {}", client.getSessionId());
// 获取客户端握手参数
String token = client.getHandshakeData().getSingleUrlParam("token");
log.info("客户端 token: {}", token);
// 可以在这里进行身份验证
// 验证失败可以断开连接
// client.disconnect();
// 发送欢迎消息
client.sendEvent("welcome", "欢迎连接到 Socket.IO 服务器");
// 存储客户端信息
// ClientManager.addClient(client);
}
});
// 断开连接监听器
server.addDisconnectListener(new DisconnectListener() {
@Override
public void onDisconnect(SocketIOClient client) {
log.info("客户端断开连接: sessionId = {}", client.getSessionId());
// 清理客户端信息
// ClientManager.removeClient(client.getSessionId());
}
});
// 消息监听器
server.addEventListener("message", String.class, new DataListener<String>() {
@Override
public void onData(SocketIOClient client, String message, AckRequest ackRequest) {
log.info("收到消息: sessionId = {}, message = {}", client.getSessionId(), message);
// 回复客户端
client.sendEvent("message", "服务器回复: " + message);
// 广播给所有客户端
server.getBroadcastOperations().sendEvent("broadcast", message);
// ACK 确认
if (ackRequest.isAckRequested()) {
ackRequest.sendAckData("消息已处理");
}
}
});
}
public void stop() {
if (server != null) {
server.stop();
log.info("Socket.IO 服务器已停止");
}
}
public static void main(String[] args) throws InterruptedException {
SocketIoServer socketIoServer = new SocketIoServer();
socketIoServer.start();
// 添加关闭钩子
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
socketIoServer.stop();
}));
}
}
Spring Boot 集成
配置类
package com.example.socketio.config;
import com.corundumstudio.socketio.Configuration;
import com.corundumstudio.socketio.SocketIOServer;
import com.corundumstudio.socketio.annotation.SpringAnnotationScanner;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
/**
* Socket.IO 配置类
*/
@Slf4j
@org.springframework.context.annotation.Configuration
public class SocketIoConfig {
@Value("${socketio.host}")
private String host;
@Value("${socketio.port}")
private int port;
@Value("${socketio.origin:*}")
private String origin;
@Bean
public SocketIOServer socketIOServer() {
Configuration config = new Configuration();
config.setHostname(host);
config.setPort(port);
config.setOrigin(origin);
// 设置认证
config.setAuthorizationListener(handshakeData -> {
// 实现你的认证逻辑
String token = handshakeData.getSingleUrlParam("token");
// return validateToken(token);
return true; // 测试环境直接返回 true
});
// 设置最大帧大小和HTTP超时
config.setMaxFramePayloadLength(1024 * 1024);
config.setPingTimeout(60000);
config.setPingInterval(25000);
SocketIOServer server = new SocketIOServer(config);
log.info("Socket.IO 服务器配置: host={}, port={}", host, port);
return server;
}
@Bean
public SpringAnnotationScanner springAnnotationScanner(SocketIOServer socketIOServer) {
return new SpringAnnotationScanner(socketIOServer);
}
}
启动类
package com.example.socketio;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.context.annotation.Bean;
import com.corundumstudio.socketio.SocketIOServer;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import javax.annotation.PreDestroy;
@Slf4j
@SpringBootApplication
@RequiredArgsConstructor
public class SocketIoApplication {
private final SocketIOServer socketIOServer;
public static void main(String[] args) {
SpringApplication.run(SocketIoApplication.class, args);
}
@Bean
public SocketIOServer socketIOServer() {
socketIOServer.start();
log.info("Socket.IO 服务器已启动");
return socketIOServer;
}
@PreDestroy
public void onShutdown() {
if (socketIOServer != null) {
socketIOServer.stop();
log.info("Socket.IO 服务器已停止");
}
}
}
配置文件
# application.yml
socketio:
host: 0.0.0.0
port: 9092
origin: "*"
spring:
application:
name: socketio-server
server:
port: 8080
logging:
level:
com.corundumstudio.socketio: DEBUG
com.example.socketio: DEBUG
消息处理器
package com.example.socketio.handler;
import com.corundumstudio.socketio.AckRequest;
import com.corundumstudio.socketio.SocketIOClient;
import com.corundumstudio.socketio.annotation.OnConnect;
import com.corundumstudio.socketio.annotation.OnDisconnect;
import com.corundumstudio.socketio.annotation.OnEvent;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
/**
* Socket.IO 消息处理器
*/
@Slf4j
@Component
public class SocketEventHandler {
// 存储客户端连接
private static final ConcurrentHashMap<UUID, SocketIOClient> clients = new ConcurrentHashMap<>();
// 连接事件
@OnConnect
public void onConnect(SocketIOClient client) {
log.info("客户端连接: sessionId = {}", client.getSessionId());
// 获取客户端信息
String token = client.getHandshakeData().getSingleUrlParam("token");
String username = client.getHandshakeData().getSingleUrlParam("username");
log.info("客户端信息: token={}, username={}", token, username);
// 存储客户端
clients.put(client.getSessionId(), client);
// 设置客户端属性
client.set("username", username != null ? username : "匿名用户");
client.set("connectedAt", System.currentTimeMillis());
// 发送连接成功消息
client.sendEvent("connected", buildResponse(200, "连接成功",
Map.of("sessionId", client.getSessionId())));
// 广播用户上线
broadcastUserStatus(client.getSessionId(), "online");
}
// 断开连接事件
@OnDisconnect
public void onDisconnect(SocketIOClient client) {
log.info("客户端断开: sessionId = {}", client.getSessionId());
// 移除客户端
clients.remove(client.getSessionId());
// 广播用户离线
broadcastUserStatus(client.getSessionId(), "offline");
}
// 自定义消息事件
@OnEvent(value = "message")
public void onMessage(SocketIOClient client, String message, AckRequest ackRequest) {
log.info("收到消息: sessionId={}, message={}", client.getSessionId(), message);
String username = client.get("username");
// 构建消息对象
MessageData messageData = MessageData.builder()
.id(UUID.randomUUID().toString())
.sessionId(client.getSessionId())
.username(username)
.content(message)
.timestamp(System.currentTimeMillis())
.build();
// 广播消息给所有客户端
broadcastMessage(messageData);
// ACK 确认
if (ackRequest.isAckRequested()) {
ackRequest.sendAckData(buildResponse(200, "发送成功", null));
}
}
// 加入房间
@OnEvent(value = "join-room")
public void onJoinRoom(SocketIOClient client, String roomId, AckRequest ackRequest) {
log.info("客户端加入房间: sessionId={}, roomId={}", client.getSessionId(), roomId);
client.joinRoom(roomId);
// 获取房间成员
// 注意: netty-socketio 的房间 API 可能有所不同
String username = client.get("username");
// 通知房间内其他用户
client.getRoomOperations(roomId)
.sendEvent("user-joined", Map.of(
"sessionId", client.getSessionId(),
"username", username,
"roomId", roomId
));
// 发送房间列表给当前用户
// TODO: 实现获取房间成员列表的逻辑
if (ackRequest.isAckRequested()) {
ackRequest.sendAckData(buildResponse(200, "加入房间成功", null));
}
}
// 离开房间
@OnEvent(value = "leave-room")
public void onLeaveRoom(SocketIOClient client, String roomId, AckRequest ackRequest) {
log.info("客户端离开房间: sessionId={}, roomId={}", client.getSessionId(), roomId);
client.leaveRoom(roomId);
String username = client.get("username");
// 通知房间内其他用户
client.getRoomOperations(roomId)
.sendEvent("user-left", Map.of(
"sessionId", client.getSessionId(),
"username", username
));
if (ackRequest.isAckRequested()) {
ackRequest.sendAckData(buildResponse(200, "离开房间成功", null));
}
}
// 发送房间消息
@OnEvent(value = "room-message")
public void onRoomMessage(SocketIOClient client, RoomMessageData data, AckRequest ackRequest) {
log.info("收到房间消息: sessionId={}, roomId={}, message={}",
client.getSessionId(), data.getRoomId(), data.getMessage());
String username = client.get("username");
MessageData messageData = MessageData.builder()
.id(UUID.randomUUID().toString())
.sessionId(client.getSessionId())
.username(username)
.content(data.getMessage())
.roomId(data.getRoomId())
.timestamp(System.currentTimeMillis())
.build();
// 发送到房间
server.getRoomOperations(data.getRoomId())
.sendEvent("room-message", messageData);
if (ackRequest.isAckRequested()) {
ackRequest.sendAckData(buildResponse(200, "发送成功", null));
}
}
// ========== 私有方法 ==========
private void broadcastMessage(MessageData message) {
clients.values().forEach(client -> {
client.sendEvent("broadcast-message", message);
});
}
private void broadcastUserStatus(UUID sessionId, String status) {
clients.values().forEach(client -> {
client.sendEvent("user-status", Map.of(
"sessionId", sessionId,
"status", status
));
});
}
private Map<String, Object> buildResponse(int code, String message, Object data) {
return Map.of(
"code", code,
"message", message,
"data", data != null ? data : "",
"timestamp", System.currentTimeMillis()
);
}
// ========== 内部类 ==========
@Data
@Builder
@AllArgsConstructor
@NoArgsConstructor
public static class MessageData {
private String id;
private UUID sessionId;
private String username;
private String content;
private String roomId;
private Long timestamp;
}
@Data
public static class RoomMessageData {
private String roomId;
private String message;
}
}
客户端实现
HTML/JavaScript 客户端
<!DOCTYPE html>
<html lang="zh-CN">
<head>
<meta charset="UTF-8">
<meta name="viewport" content="width=device-width, initial-scale=1.0">
<title>Socket.IO 聊天室</title>
<style>
* { margin: 0; padding: 0; box-sizing: border-box; }
body {
font-family: -apple-system, BlinkMacSystemFont, 'Segoe UI', Roboto, sans-serif;
background: linear-gradient(135deg, #667eea 0%, #764ba2 100%);
height: 100vh;
display: flex;
justify-content: center;
align-items: center;
}
.container {
width: 90%;
max-width: 500px;
background: white;
border-radius: 20px;
box-shadow: 0 20px 60px rgba(0,0,0,0.3);
overflow: hidden;
display: flex;
flex-direction: column;
height: 80vh;
}
.header {
background: linear-gradient(135deg, #667eea 0%, #764ba2 100%);
color: white;
padding: 20px;
display: flex;
justify-content: space-between;
align-items: center;
}
.status {
display: flex;
align-items: center;
gap: 8px;
padding: 6px 12px;
background: rgba(255,255,255,0.2);
border-radius: 20px;
font-size: 14px;
}
.status-dot {
width: 8px;
height: 8px;
border-radius: 50%;
background: #ccc;
}
.status.connected .status-dot {
background: #4ade80;
}
.messages {
flex: 1;
overflow-y: auto;
padding: 20px;
background: #f8fafc;
}
.message {
margin-bottom: 16px;
display: flex;
flex-direction: column;
}
.message.self {
align-items: flex-end;
}
.message-content {
max-width: 70%;
padding: 12px 16px;
border-radius: 12px;
word-wrap: break-word;
}
.message.self .message-content {
background: linear-gradient(135deg, #667eea 0%, #764ba2 100%);
color: white;
}
.message.other .message-content {
background: white;
box-shadow: 0 2px 8px rgba(0,0,0,0.1);
}
.message-username {
font-size: 12px;
color: #64748b;
margin-bottom: 4px;
}
.input-area {
padding: 16px;
background: white;
border-top: 1px solid #e2e8f0;
display: flex;
gap: 10px;
}
.input-area input {
flex: 1;
padding: 12px 16px;
border: 1px solid #e2e8f0;
border-radius: 24px;
outline: none;
transition: all 0.3s;
}
.input-area input:focus {
border-color: #667eea;
box-shadow: 0 0 0 3px rgba(102, 126, 234, 0.1);
}
.input-area button {
padding: 12px 24px;
background: linear-gradient(135deg, #667eea 0%, #764ba2 100%);
color: white;
border: none;
border-radius: 24px;
cursor: pointer;
font-weight: 600;
transition: all 0.3s;
}
.input-area button:hover {
transform: translateY(-2px);
box-shadow: 0 4px 12px rgba(102, 126, 234, 0.4);
}
.input-area button:disabled {
opacity: 0.5;
cursor: not-allowed;
transform: none;
}
</style>
</head>
<body>
<div class="container">
<div class="header">
<h2>💬 聊天室</h2>
<div class="status" id="status">
<span class="status-dot"></span>
<span id="statusText">未连接</span>
</div>
</div>
<div class="messages" id="messages"></div>
<div class="input-area">
<input type="text" id="messageInput" placeholder="输入消息..." />
<button id="sendBtn" disabled>发送</button>
</div>
</div>
<script src="https://cdn.socket.io/4.7.4/socket.io.min.js"></script>
<script>
const SERVER_URL = 'http://localhost:9092';
class ChatClient {
constructor() {
this.socket = null;
this.username = this.generateUsername();
this.initElements();
this.initSocket();
}
initElements() {
this.messagesContainer = document.getElementById('messages');
this.messageInput = document.getElementById('messageInput');
this.sendBtn = document.getElementById('sendBtn');
this.statusIndicator = document.getElementById('status');
this.statusText = document.getElementById('statusText');
}
initSocket() {
this.socket = io(SERVER_URL, {
query: {
token: 'demo-token',
username: this.username
},
reconnection: true,
reconnectionDelay: 1000,
reconnectionAttempts: 10
});
// 连接成功
this.socket.on('connect', () => {
this.updateStatus(true);
this.addSystemMessage('已连接到服务器');
});
// 连接确认
this.socket.on('connected', (data) => {
console.log('连接确认:', data);
});
// 断开连接
this.socket.on('disconnect', () => {
this.updateStatus(false);
this.addSystemMessage('连接已断开');
});
// 广播消息
this.socket.on('broadcast-message', (data) => {
this.addMessage(data);
});
// 用户状态变化
this.socket.on('user-status', (data) => {
const statusText = data.status === 'online' ? '上线了' : '下线了';
this.addSystemMessage(`用户 ${data.sessionId} ${statusText}`);
});
// 用户加入房间
this.socket.on('user-joined', (data) => {
this.addSystemMessage(`${data.username} 加入了房间`);
});
}
sendMessage() {
const message = this.messageInput.value.trim();
if (!message || !this.socket.connected) return;
this.socket.emit('message', message, (response) => {
console.log('服务器响应:', response);
});
this.messageInput.value = '';
}
addMessage(data) {
const isSelf = data.sessionId === this.socket?.id;
const messageDiv = document.createElement('div');
messageDiv.className = `message ${isSelf ? 'self' : 'other'}`;
messageDiv.innerHTML = `
<div class="message-username">${data.username}</div>
<div class="message-content">${this.escapeHtml(data.content)}</div>
`;
this.messagesContainer.appendChild(messageDiv);
this.scrollToBottom();
}
addSystemMessage(text) {
const messageDiv = document.createElement('div');
messageDiv.style.cssText = `
text-align: center;
color: #94a3b8;
font-size: 12px;
margin: 8px 0;
`;
messageDiv.textContent = text;
this.messagesContainer.appendChild(messageDiv);
this.scrollToBottom();
}
updateStatus(connected) {
if (connected) {
this.statusIndicator.classList.add('connected');
this.statusText.textContent = '已连接';
this.sendBtn.disabled = false;
} else {
this.statusIndicator.classList.remove('connected');
this.statusText.textContent = '未连接';
this.sendBtn.disabled = true;
}
}
scrollToBottom() {
this.messagesContainer.scrollTop = this.messagesContainer.scrollHeight;
}
escapeHtml(text) {
const div = document.createElement('div');
div.textContent = text;
return div.innerHTML;
}
generateUsername() {
return `用户_${Math.random().toString(36).substr(2, 6)}`;
}
}
// 初始化客户端
const chatClient = new ChatClient();
// 绑定事件
document.getElementById('sendBtn').addEventListener('click', () => {
chatClient.sendMessage();
});
document.getElementById('messageInput').addEventListener('keypress', (e) => {
if (e.key === 'Enter') {
chatClient.sendMessage();
}
});
</script>
</body>
</html>
Java 客户端
package com.example.socketio.client;
import com.corundumstudio.socketio.client.SocketIOClient;
import com.corundumstudio.socketio.client.SocketIOClient.Builder;
import lombok.extern.slf4j.Slf4j;
import java.util.concurrent.ExecutionException;
/**
* Socket.IO Java 客户端
*/
@Slf4j
public class SocketIoClientExample {
public static void main(String[] args) {
// 创建客户端
SocketIOClient client = Builder.create("http://localhost:9092")
.setQuery("token=client-token&username=JavaClient")
.reconnection(true)
.reconnectionDelay(1000)
.reconnectionAttempts(10)
.build();
// 连接监听器
client.on("connected", objects -> {
log.info("连接成功: {}", objects);
});
// 消息监听器
client.on("broadcast-message", objects -> {
log.info("收到广播消息: {}", objects);
});
// 连接服务器
try {
client.connect().get();
log.info("客户端已启动");
// 发送消息
client.emit("message", "Hello from Java Client");
// 保持运行
Thread.sleep(60000);
} catch (InterruptedException | ExecutionException e) {
log.error("客户端错误", e);
} finally {
client.disconnect();
}
}
}
Android 客户端
// build.gradle
dependencies {
implementation('io.socket:socket.io-client:2.1.0') {
exclude group: 'org.json', module: 'json'
}
}
package com.example.socketio;
import io.socket.client.IO;
import io.socket.client.Socket;
import io.socket.emitter.Emitter;
import android.os.Bundle;
import androidx.appcompat.app.AppCompatActivity;
import org.json.JSONException;
import org.json.JSONObject;
import java.net.URISyntaxException;
public class ChatActivity extends AppCompatActivity {
private Socket mSocket;
@Override
protected void onCreate(Bundle savedInstanceState) {
super.onCreate(savedInstanceState);
try {
// 创建 Socket 连接
IO.Options options = new IO.Options();
options.query = "token=user-token";
options.reconnection = true;
options.reconnectionDelay = 1000;
options.reconnectionAttempts = 10;
mSocket = IO.socket("http://your-server:9092", options);
// 监听连接事件
mSocket.on(Socket.EVENT_CONNECT, onConnect);
mSocket.on(Socket.EVENT_DISCONNECT, onDisconnect);
mSocket.on("broadcast-message", onMessageReceived);
mSocket.connect();
} catch (URISyntaxException e) {
e.printStackTrace();
}
}
// 连接成功
private Emitter.Listener onConnect = new Emitter.Listener() {
@Override
public void call(Object... args) {
runOnUiThread(() -> {
// 更新 UI 显示连接成功
});
}
};
// 断开连接
private Emitter.Listener onDisconnect = new Emitter.Listener() {
@Override
public void call(Object... args) {
runOnUiThread(() -> {
// 更新 UI 显示断开连接
});
}
};
// 接收消息
private Emitter.Listener onMessageReceived = new Emitter.Listener() {
@Override
public void call(Object... args) {
JSONObject data = (JSONObject) args[0];
runOnUiThread(() -> {
try {
String username = data.getString("username");
String content = data.getString("content");
// 显示消息
} catch (JSONException e) {
e.printStackTrace();
}
});
}
};
// 发送消息
private void sendMessage(String message) {
mSocket.emit("message", message);
}
@Override
protected void onDestroy() {
super.onDestroy();
mSocket.disconnect();
mSocket.off(Socket.EVENT_CONNECT, onConnect);
mSocket.off(Socket.EVENT_DISCONNECT, onDisconnect);
mSocket.off("broadcast-message", onMessageReceived);
}
}
Spring Boot 集成
完整的 Spring Boot 应用
package com.example.socketio.config;
import com.corundumstudio.socketio.Configuration;
import com.corundumstudio.socketio.SocketIOServer;
import com.corundumstudio.socketio.Transport;
import com.corundumstudio.socketio.annotation.SpringAnnotationScanner;
import com.corundumstudio.socketio.listener.ConnectListener;
import com.corundumstudio.socketio.listener.DisconnectListener;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import java.util.EnumSet;
/**
* Socket.IO 配置
*/
@Slf4j
@Configuration
@RequiredArgsConstructor
public class SocketIoConfiguration {
@Value("${socketio.host}")
private String host;
@Value("${socketio.port}")
private int port;
@Value("${socketio.origin:*}")
private String origin;
@Bean
public SocketIOServer socketIOServer() {
Configuration config = new Configuration();
config.setHostname(host);
config.setPort(port);
// 跨域配置
config.setOrigin(origin);
// 支持的传输方式
config.setTransports(EnumSet.of(Transport.WEBSOCKET, Transport.POLLING));
// 认证配置
config.setAuthorizationListener(handshakeData -> {
String token = handshakeData.getSingleUrlParam("token");
return validateToken(token);
});
// 超时配置
config.setPingTimeout(60000);
config.setPingInterval(25000);
// 最大帧大小
config.setMaxFramePayloadLength(1024 * 1024);
SocketIOServer server = new SocketIOServer(config);
// 添加连接监听器
addListeners(server);
return server;
}
@Bean
public SpringAnnotationScanner springAnnotationScanner(SocketIOServer socketIOServer) {
return new SpringAnnotationScanner(socketIOServer);
}
private void addListeners(SocketIOServer server) {
server.addConnectListener(client -> {
log.info("客户端连接: sessionId = {}", client.getSessionId());
});
server.addDisconnectListener(client -> {
log.info("客户端断开: sessionId = {}", client.getSessionId());
});
}
private boolean validateToken(String token) {
// 实现你的 token 验证逻辑
return token != null && !token.isEmpty();
}
}
服务启动
package com.example.socketio;
import com.corundumstudio.socketio.SocketIOServer;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.boot.CommandLineRunner;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.context.annotation.Bean;
@Slf4j
@SpringBootApplication
@RequiredArgsConstructor
public class SocketIoApplication {
private final SocketIOServer socketIOServer;
public static void main(String[] args) {
SpringApplication.run(SocketIoApplication.class, args);
}
@Bean
public CommandLineRunner runner() {
return args -> {
socketIOServer.start();
log.info("Socket.IO 服务器已启动");
// 添加关闭钩子
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
socketIOServer.stop();
log.info("Socket.IO 服务器已停止");
}));
};
}
}
服务器管理器
package com.example.socketio.manager;
import com.corundumstudio.socketio.SocketIOClient;
import com.corundumstudio.socketio.SocketIOServer;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
import java.util.*;
import java.util.concurrent.ConcurrentHashMap;
/**
* Socket.IO 服务器管理器
*/
@Slf4j
@Component
@RequiredArgsConstructor
public class SocketServerManager {
private final SocketIOServer server;
// 客户端存储
private final Map<String, SocketIOClient> clients = new ConcurrentHashMap<>();
// 用户 ID 到 Session ID 的映射
private final Map<String, String> userSessionMap = new ConcurrentHashMap<>();
/**
* 发送消息给指定用户
*/
public void sendToUser(String userId, String event, Object data) {
String sessionId = userSessionMap.get(userId);
if (sessionId != null) {
SocketIOClient client = clients.get(sessionId);
if (client != null) {
client.sendEvent(event, data);
log.debug("发送消息给用户 {}: {}", userId, event);
}
}
}
/**
* 广播消息给所有用户
*/
public void broadcast(String event, Object data) {
server.getBroadcastOperations().sendEvent(event, data);
log.debug("广播消息: {}", event);
}
/**
* 发送消息到指定房间
*/
public void sendToRoom(String roomId, String event, Object data) {
server.getRoomOperations(roomId).sendEvent(event, data);
log.debug("发送消息到房间 {}: {}", roomId, event);
}
/**
* 获取在线用户数
*/
public int getOnlineCount() {
return clients.size();
}
/**
* 获取房间内用户数
*/
public int getRoomSize(String roomId) {
return server.getRoomOperations(roomId).getClients().size();
}
/**
* 断开指定用户
*/
public void disconnectUser(String userId) {
String sessionId = userSessionMap.get(userId);
if (sessionId != null) {
SocketIOClient client = clients.get(sessionId);
if (client != null) {
client.disconnect();
}
}
}
/**
* 添加客户端
*/
public void addClient(String sessionId, SocketIOClient client, String userId) {
clients.put(sessionId, client);
if (userId != null) {
userSessionMap.put(userId, sessionId);
}
}
/**
* 移除客户端
*/
public void removeClient(String sessionId) {
SocketIOClient client = clients.remove(sessionId);
if (client != null) {
String userId = (String) client.get("userId");
if (userId != null) {
userSessionMap.remove(userId);
}
}
}
/**
* 获取所有在线用户
*/
public Set<String> getOnlineUsers() {
return new HashSet<>(userSessionMap.keySet());
}
}
核心概念
事件系统
// 服务端监听事件
@OnEvent(value = "custom-event")
public void onCustomEvent(SocketIOClient client, String data, AckRequest ackRequest) {
// 处理事件
}
// 服务端发送事件
client.sendEvent("event-name", data);
// 广播事件
server.getBroadcastOperations().sendEvent("event-name", data);
// 发送到房间
server.getRoomOperations("room-name").sendEvent("event-name", data);
// 发送 ACK 确认
if (ackRequest.isAckRequested()) {
ackRequest.sendAckData("处理结果");
}
房间(Room)
// 加入房间
client.joinRoom("room-name");
// 离开房间
client.leaveRoom("room-name");
// 获取房间所有客户端
Collection<SocketIOClient> roomClients = server.getRoomOperations("room-name").getClients();
// 发送到房间(不包括发送者)
client.getRoomOperations("room-name").sendEvent("event", data);
// 获取客户端所在的所有房间
Collection<String> rooms = server.getAllRooms();
命名空间(Namespace)
// 创建命名空间
SocketIONamespace chatNamespace = server.addNamespace("/chat");
// 命名空间事件监听
chatNamespace.addConnectListener(client -> {
log.info("用户连接到聊天命名空间");
});
chatNamespace.addEventListener("message", String.class, (client, data, ackRequest) -> {
// 处理消息
chatNamespace.getBroadcastOperations().sendEvent("message", data);
});
// 客户端连接到命名空间
// JavaScript: const socket = io('/chat');
客户端存储
// 设置客户端属性
client.set("userId", "user123");
client.set("username", "张三");
client.set("role", "admin");
// 获取客户端属性
String userId = client.get("userId");
String username = client.get("username");
// 删除客户端属性
client.del("userId");
高级特性
JWT 认证
package com.example.socketio.auth;
import com.corundumstudio.socketio.HandshakeData;
import com.corundumstudio.socketio.listener.AuthorizationListener;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import io.jsonwebtoken.Claims;
import io.jsonwebtoken.Jwts;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
import java.nio.charset.StandardCharsets;
/**
* JWT 认证监听器
*/
@Slf4j
@Component
public class JwtAuthorizationListener implements AuthorizationListener {
private static final String SECRET_KEY = "your-secret-key";
private final ObjectMapper objectMapper = new ObjectMapper();
@Override
public boolean isAuthorized(HandshakeData data) {
try {
// 获取 token
String token = data.getSingleUrlParam("token");
if (token == null || token.isEmpty()) {
log.warn("Token 为空");
return false;
}
// 验证 token
Claims claims = Jwts.parser()
.setSigningKey(SECRET_KEY.getBytes(StandardCharsets.UTF_8))
.parseClaimsJws(token)
.getBody();
// 将用户信息存入 HandshakeData
String userId = claims.getSubject();
String username = claims.get("username", String.class);
data.getHttpHeaders().add("X-User-Id", userId);
data.getHttpHeaders().add("X-Username", username);
log.info("用户认证成功: userId={}, username={}", userId, username);
return true;
} catch (Exception e) {
log.error("Token 验证失败", e);
return false;
}
}
}
消息持久化
package com.example.socketio.service;
import com.corundumstudio.socketio.SocketIOClient;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.data.domain.Page;
import org.springframework.data.domain.PageRequest;
import org.springframework.stereotype.Service;
import java.time.LocalDateTime;
import java.util.List;
import java.util.concurrent.*;
/**
* 消息服务
*/
@Slf4j
@Service
@RequiredArgsConstructor
public class MessageService {
private final MessageRepository messageRepository;
// 消息缓冲队列
private final BlockingQueue<Message> messageBuffer = new LinkedBlockingQueue<>(10000);
// 批量保存线程池
private final ScheduledExecutorService executorService = Executors.newSingleThreadScheduledExecutor();
@PostConstruct
public void init() {
// 定时批量保存消息
executorService.scheduleAtFixedRate(() -> {
flushMessages();
}, 5, 5, TimeUnit.SECONDS);
}
/**
* 添加消息到缓冲区
*/
public void addMessage(String roomId, String userId, String username, String content) {
Message message = Message.builder()
.roomId(roomId)
.userId(userId)
.username(username)
.content(content)
.timestamp(LocalDateTime.now())
.build();
if (!messageBuffer.offer(message)) {
log.warn("消息缓冲区已满,消息将被丢弃");
}
}
/**
* 批量保存消息
*/
private void flushMessages() {
List<Message> messages = new ArrayList<>();
messageBuffer.drainTo(messages, 1000);
if (!messages.isEmpty()) {
try {
messageRepository.saveAll(messages);
log.debug("批量保存消息: {} 条", messages.size());
} catch (Exception e) {
log.error("保存消息失败", e);
}
}
}
/**
* 获取历史消息
*/
public Page<Message> getHistoryMessages(String roomId, int page, int size) {
return messageRepository.findByRoomIdOrderByTimestampDesc(
roomId, PageRequest.of(page, size)
);
}
@PreDestroy
public void destroy() {
flushMessages();
executorService.shutdown();
}
}
Redis 集群支持
<!-- 添加 Redis 适配器依赖 -->
<dependency>
<groupId>com.corundumstudio.socketio</groupId>
<artifactId>netty-socketio</artifactId>
<version>2.0.9</version>
</dependency>
<dependency>
<groupId>org.redisson</groupId>
<artifactId>redisson</artifactId>
<version>3.24.3</version>
</dependency>
package com.example.socketio.config;
import com.corundumstudio.socketio.SocketConfig;
import com.corundumstudio.socketio.SocketIOServer;
import com.corundumstudio.socketio.store.RedissonStoreFactory;
import com.corundumstudio.socketio.store.StoreFactory;
import lombok.extern.slf4j.Slf4j;
import org.redisson.Redisson;
import org.redisson.api.RedissonClient;
import org.redisson.config.Config;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
/**
* Redis 集群配置
*/
@Slf4j
@Configuration
public class RedisClusterConfig {
@Value("${spring.redis.host:localhost}")
private String redisHost;
@Value("${spring.redis.port:6379}")
private int redisPort;
@Bean(destroyMethod = "shutdown")
public RedissonClient redissonClient() {
Config config = new Config();
config.useSingleServer()
.setAddress(String.format("redis://%s:%d", redisHost, redisPort))
.setConnectionPoolSize(64)
.setConnectionMinimumIdleSize(10);
return Redisson.create(config);
}
@Bean
public StoreFactory redisStoreFactory(RedissonClient redissonClient) {
return new RedissonStoreFactory(redissonClient);
}
}
实战案例
在线聊天室
package com.example.socketio.handler;
import com.corundumstudio.socketio.AckRequest;
import com.corundumstudio.socketio.SocketIOClient;
import com.corundumstudio.socketio.annotation.OnEvent;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
/**
* 聊天室处理器
*/
@Slf4j
@Component
@RequiredArgsConstructor
public class ChatRoomHandler {
private final SocketServerManager serverManager;
// 房间成员管理
private final Map<String, Map<String, UserInfo>> roomMembers = new ConcurrentHashMap<>();
@OnEvent("join-room")
public void onJoinRoom(SocketIOClient client, String roomId, AckRequest ackRequest) {
String userId = (String) client.get("userId");
String username = (String) client.get("username");
// 加入房间
client.joinRoom(roomId);
// 记录成员
roomMembers.computeIfAbsent(roomId, k -> new ConcurrentHashMap<>())
.put(userId, new UserInfo(userId, username, System.currentTimeMillis()));
// 通知其他成员
serverManager.sendToRoom(roomId, "user-joined", Map.of(
"userId", userId,
"username", username
));
// 发送当前成员列表
Map<String, UserInfo> members = roomMembers.get(roomId);
client.sendEvent("room-members", members.values());
log.info("用户 {} 加入房间 {}", username, roomId);
if (ackRequest.isAckRequested()) {
ackRequest.sendAckData(Map.of("success", true));
}
}
@OnEvent("leave-room")
public void onLeaveRoom(SocketIOClient client, String roomId, AckRequest ackRequest) {
String userId = (String) client.get("userId");
String username = (String) client.get("username");
// 离开房间
client.leaveRoom(roomId);
// 移除成员记录
Map<String, UserInfo> members = roomMembers.get(roomId);
if (members != null) {
UserInfo removed = members.remove(userId);
if (members.isEmpty()) {
roomMembers.remove(roomId);
}
}
// 通知其他成员
serverManager.sendToRoom(roomId, "user-left", Map.of(
"userId", userId,
"username", username
));
log.info("用户 {} 离开房间 {}", username, roomId);
if (ackRequest.isAckRequested()) {
ackRequest.sendAckData(Map.of("success", true));
}
}
@OnEvent("send-message")
public void onSendMessage(SocketIOClient client, MessageData data, AckRequest ackRequest) {
String userId = (String) client.get("userId");
String username = (String) client.get("username");
// 构建完整消息
ChatMessage message = ChatMessage.builder()
.id(UUID.randomUUID().toString())
.roomId(data.getRoomId())
.userId(userId)
.username(username)
.content(data.getContent())
.type(data.getType())
.timestamp(System.currentTimeMillis())
.build();
// 发送到房间
serverManager.sendToRoom(data.getRoomId(), "new-message", message);
// 保存消息
// messageService.saveMessage(message);
log.info("房间 {} 收到消息: {}", data.getRoomId(), data.getContent());
if (ackRequest.isAckRequested()) {
ackRequest.sendAckData(Map.of(
"success", true,
"messageId", message.getId()
));
}
}
@OnEvent("typing")
public void onTyping(SocketIOClient client, TypingData data) {
String username = (String) client.get("username");
// 转发输入状态
client.getRoomOperations(data.getRoomId())
.sendEvent("user-typing", Map.of(
"userId", client.getSessionId(),
"username", username,
"isTyping", data.isTyping()
));
}
// 数据类
@Data
public static class UserInfo {
private String userId;
private String username;
private Long joinTime;
public UserInfo(String userId, String username, Long joinTime) {
this.userId = userId;
this.username = username;
this.joinTime = joinTime;
}
}
@Data
public static class MessageData {
private String roomId;
private String content;
private String type = "text";
}
@Data
@Builder
public static class ChatMessage {
private String id;
private String roomId;
private String userId;
private String username;
private String content;
private String type;
private Long timestamp;
}
@Data
public static class TypingData {
private String roomId;
private boolean typing;
}
}
实时通知系统
package com.example.socketio.service;
import com.corundumstudio.socketio.SocketIOClient;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.scheduling.annotation.Async;
import org.springframework.stereotype.Service;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
/**
* 通知服务
*/
@Slf4j
@Service
@RequiredArgsConstructor
public class NotificationService {
private final SocketServerManager serverManager;
// 用户在线状态缓存
private final Map<String, Boolean> onlineStatus = new ConcurrentHashMap<>();
/**
* 发送通知给指定用户
*/
@Async
public void sendNotification(String userId, Notification notification) {
if (isUserOnline(userId)) {
serverManager.sendToUser(userId, "notification", notification);
log.info("发送通知给用户 {}: {}", userId, notification.getTitle());
} else {
// 用户离线,存储离线通知
saveOfflineNotification(userId, notification);
log.info("用户 {} 离线,保存离线通知", userId);
}
}
/**
* 广播通知给所有在线用户
*/
@Async
public void broadcastNotification(Notification notification) {
serverManager.broadcast("notification", notification);
log.info("广播通知: {}", notification.getTitle());
}
/**
* 发送系统公告
*/
@Async
public void sendAnnouncement(String title, String content) {
Announcement announcement = Announcement.builder()
.id(UUID.randomUUID().toString())
.title(title)
.content(content)
.timestamp(System.currentTimeMillis())
.build();
serverManager.broadcast("announcement", announcement);
log.info("发送系统公告: {}", title);
}
/**
* 检查用户是否在线
*/
public boolean isUserOnline(String userId) {
return onlineStatus.getOrDefault(userId, false);
}
/**
* 更新用户在线状态
*/
public void updateOnlineStatus(String userId, boolean online) {
onlineStatus.put(userId, online);
log.debug("用户 {} 在线状态: {}", userId, online);
}
private void saveOfflineNotification(String userId, Notification notification) {
// 保存到数据库或缓存
}
@Data
@Builder
public static class Notification {
private String id;
private String title;
private String content;
private String type;
private Map<String, Object> data;
private Long timestamp;
}
@Data
@Builder
public static class Announcement {
private String id;
private String title;
private String content;
private Long timestamp;
}
}
生产环境部署
Docker 部署
# Dockerfile
FROM openjdk:17-jdk-slim
WORKDIR /app
# 复制应用文件
COPY target/socketio-server-*.jar app.jar
# 暴露端口
EXPOSE 8080 9092
# JVM 参数
ENV JAVA_OPTS="-Xms512m -Xmx1g -XX:+UseG1GC"
# 启动命令
ENTRYPOINT ["sh", "-c", "java $JAVA_OPTS -jar app.jar"]
# docker-compose.yml
version: '3.8'
services:
socketio-app:
build: .
ports:
- "8080:8080"
- "9092:9092"
environment:
- SPRING_PROFILES_ACTIVE=prod
- SOCKETIO_HOST=0.0.0.0
- SOCKETIO_PORT=9092
- REDIS_HOST=redis
- REDIS_PORT=6379
depends_on:
- redis
restart: unless-stopped
healthcheck:
test: ["CMD", "curl", "-f", "http://localhost:8080/actuator/health"]
interval: 30s
timeout: 10s
retries: 3
redis:
image: redis:7-alpine
ports:
- "6379:6379"
volumes:
- redis_data:/data
restart: unless-stopped
volumes:
redis_data:
Nginx 配置
upstream socketio_backend {
server 127.0.0.1:9092;
}
server {
listen 80;
server_name your-domain.com;
# WebSocket 升级配置
location /socket.io/ {
proxy_pass http://socketio_backend;
proxy_http_version 1.1;
# WebSocket 关键头
proxy_set_header Upgrade $http_upgrade;
proxy_set_header Connection "upgrade";
# 超时设置
proxy_connect_timeout 7d;
proxy_send_timeout 7d;
proxy_read_timeout 7d;
# 其他代理头
proxy_set_header Host $host;
proxy_set_header X-Real-IP $remote_addr;
proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
proxy_set_header X-Forwarded-Proto $scheme;
# 禁用缓冲
proxy_buffering off;
}
# 静态资源
location / {
root /var/www/html;
index index.html;
try_files $uri $uri/ /index.html;
}
}
K8s 部署
# deployment.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
name: socketio-server
spec:
replicas: 3
selector:
matchLabels:
app: socketio-server
template:
metadata:
labels:
app: socketio-server
spec:
containers:
- name: socketio-server
image: your-registry/socketio-server:latest
ports:
- containerPort: 9092
name: socketio
- containerPort: 8080
name: http
env:
- name: SPRING_PROFILES_ACTIVE
value: "prod"
- name: REDIS_HOST
value: "redis-service"
resources:
requests:
memory: "512Mi"
cpu: "500m"
limits:
memory: "1Gi"
cpu: "1000m"
livenessProbe:
httpGet:
path: /actuator/health/liveness
port: 8080
initialDelaySeconds: 60
periodSeconds: 10
readinessProbe:
httpGet:
path: /actuator/health/readiness
port: 8080
initialDelaySeconds: 30
periodSeconds: 5
---
apiVersion: v1
kind: Service
metadata:
name: socketio-service
spec:
selector:
app: socketio-server
ports:
- port: 9092
targetPort: 9092
name: socketio
- port: 8080
targetPort: 8080
name: http
type: LoadBalancer
最佳实践
1. 连接管理
/**
* 客户端连接管理
*/
@Component
@Slf4j
public class ClientConnectionManager {
private final ConcurrentHashMap<String, ClientInfo> connections = new ConcurrentHashMap<>();
public void addClient(SocketIOClient client) {
String sessionId = client.getSessionId().toString();
ClientInfo info = ClientInfo.builder()
.sessionId(sessionId)
.userId((String) client.get("userId"))
.username((String) client.get("username"))
.connectTime(System.currentTimeMillis())
.lastHeartbeat(System.currentTimeMillis())
.build();
connections.put(sessionId, info);
log.info("添加客户端连接: {}", sessionId);
}
public void removeClient(String sessionId) {
ClientInfo removed = connections.remove(sessionId);
if (removed != null) {
log.info("移除客户端连接: {}", sessionId);
}
}
public void updateHeartbeat(String sessionId) {
ClientInfo info = connections.get(sessionId);
if (info != null) {
info.setLastHeartbeat(System.currentTimeMillis());
}
}
@Scheduled(fixedRate = 60000)
public void cleanInactiveClients() {
long now = System.currentTimeMillis();
long timeout = 120000; // 2 分钟超时
connections.entrySet().removeIf(entry -> {
ClientInfo info = entry.getValue();
if (now - info.getLastHeartbeat() > timeout) {
log.warn("清理不活跃客户端: {}", entry.getKey());
return true;
}
return false;
});
}
}
2. 消息限流
/**
* 消息限流拦截器
*/
@Component
@Slf4j
public class RateLimitInterceptor {
private final LoadingCache<String, RateLimiter> limiters = Caffeine.newBuilder()
.expireAfterWrite(1, TimeUnit.MINUTES)
.build(key -> RateLimiter.create(10)); // 每分钟 10 条消息
public boolean checkRateLimit(String sessionId) {
RateLimiter limiter = limiters.get(sessionId);
return limiter.tryAcquire();
}
@OnEvent("message")
public void onMessage(SocketIOClient client, String message, AckRequest ackRequest) {
if (!checkRateLimit(client.getSessionId().toString())) {
client.sendEvent("error", "发送过于频繁,请稍后再试");
return;
}
// 处理消息...
}
}
3. 监控指标
/**
* Socket.IO 监控指标
*/
@Component
@Slf4j
public class SocketIoMetrics {
private final AtomicLong totalConnections = new AtomicLong(0);
private final AtomicLong totalMessages = new AtomicLong(0);
private final AtomicLong totalDisconnections = new AtomicLong(0);
private final ConcurrentHashMap<String, AtomicLong> eventCounts = new ConcurrentHashMap<>();
@EventListener
public void handleConnect(ConnectEvent event) {
totalConnections.incrementAndGet();
}
@EventListener
public void handleDisconnect(DisconnectEvent event) {
totalDisconnections.incrementAndGet();
}
public void recordMessage(String event) {
totalMessages.incrementAndGet();
eventCounts.computeIfAbsent(event, k -> new AtomicLong(0)).incrementAndGet();
}
@Bean
public MeterRegistryCustomizer<MeterRegistry> metricsCommonTags() {
return registry -> {
Gauge.builder("socketio.connections", totalConnections::get)
.description("当前连接数")
.register(registry);
Gauge.builder("socketio.messages", totalMessages::get)
.description("总消息数")
.register(registry);
};
}
}
4. 异常处理
/**
* 全局异常处理器
*/
@ControllerAdvice
@Slf4j
public class SocketExceptionHandler {
@ExceptionHandler(Exception.class)
public void handleException(SocketIOClient client, Exception e) {
log.error("Socket.IO 异常: sessionId={}, error={}",
client.getSessionId(), e.getMessage(), e);
client.sendEvent("error", Map.of(
"code", 500,
"message", "服务器内部错误"
));
}
}
常见问题
Q1: 连接频繁断开
// 增加超时时间
config.setPingTimeout(60000); // 60秒
config.setPingInterval(25000); // 25秒心跳
// 仅使用 WebSocket
config.setTransports(EnumSet.of(Transport.WEBSOCKET));
Q2: 内存泄漏
@Scheduled(fixedRate = 300000)
public void cleanup() {
// 清理空房间
// 清理过期会话
// 清理未使用的缓存
}
Q3: 消息顺序问题
// 使用序列号保证消息顺序
private final AtomicLong messageSequence = new AtomicLong(0);
public void sendMessage(String content) {
long seq = messageSequence.incrementAndGet();
// 发送带序列号的消息
}
参考资料
作者注:本文档持续更新中,如有问题或建议,欢迎讨论交流。
更多推荐

所有评论(0)