基于 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();
    // 发送带序列号的消息
}

参考资料


作者注:本文档持续更新中,如有问题或建议,欢迎讨论交流。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐