# drogon_chat_server_cluster **Repository Path**: ljc8/drogon_chat_server_cluster ## Basic Information - **Project Name**: drogon_chat_server_cluster - **Description**: 基于drogon的WebSocker消息同步 - **Primary Language**: Unknown - **License**: MulanPSL-2.0 - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 0 - **Created**: 2026-04-13 - **Last Updated**: 2026-06-29 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # AOD Server Cluster - WebSocket 集群聊天服务 ## 📋 目录 - [项目概述](#项目概述) - [快速入门](#快速入门) - [架构设计](#架构设计) - [核心组件详解](#核心组件详解) - [数据流与交互](#数据流与交互) - [Redis 数据结构设计](#redis-数据结构设计) - [清理机制详解](#清理机制详解) - [关键算法与实现](#关键算法与实现) - [部署指南](#部署指南) - [故障排查](#故障排查) - [性能优化](#性能优化) - [开发规范](#开发规范) --- ## 项目概述 ### 简介 AOD Server Cluster 是一个基于 **Drogon Framework** 和 **C++20** 构建的高性能、分布式 WebSocket 集群聊天服务。它支持多进程、多实例部署,通过 Redis 实现跨进程的消息广播和状态同步。 ### 核心特性 ✅ **高性能**: 基于 Drogon 异步框架,单实例可支持数万并发连接 ✅ **集群支持**: 多进程/多实例部署,通过 Redis Pub/Sub 实现消息互通 ✅ **自动清理**: 三层清理机制(TTL + 心跳 + 惰性清理)防止脏数据 ✅ **实时统计**: 准确的在线人数统计和用户列表查询 ✅ **容错设计**: 进程崩溃后自动清理残留数据 ✅ **Docker 友好**: 支持容器化部署,自动处理 DNS 解析 ### 技术栈 | 组件 | 版本/说明 | |------|----------| | **语言** | C++20 | | **Web 框架** | Drogon Framework | | **事件循环** | Trantor | | **数据库** | Redis (用于集群通信和状态存储) | | **构建工具** | CMake 3.5+ | | **JSON 库** | JsonCpp (内置于 Drogon) | | **Redis 客户端** | hiredis / Drogon Redis Client | ### 应用场景 - 多人在线聊天室 - 实时协作工具 - 在线客服系统 - 游戏内聊天 - 即时通讯(IM)后端 --- ## 快速入门 ### 环境要求 - **操作系统**: Linux (推荐 Ubuntu 20.04+) / macOS / Windows (WSL2) - **编译器**: GCC 10+ / Clang 12+ (支持 C++20) - **CMake**: 3.5+ - **Redis**: 6.0+ (推荐 7.0+) - **Drogon**: 1.9+ (需安装或作为子模块) ### 编译步骤 ``` bash # 1. 克隆项目 git clone https://gitee.com/ljc8/drogon_chat_server_cluster.git cd drogon_chat_server_cluster # 2. 创建构建目录 mkdir build && cd build # 3. 配置 CMake cmake .. -DCMAKE_BUILD_TYPE=RelWithDebInfo # 4. 编译 make -j$(nproc) # 5. 运行 ./server_cluster_AOD ``` ### Docker 部署 [Docker Compose 部署](deploy-docker-compose/deploy-docker-compose.md) ### 配置文件 #### config.json ```json { "app": { "number_of_threads": 1, "enable_session": false, "session_timeout": 0, "session_same_site" : "Null", "session_cookie_key": "JSESSIONID", "session_max_age": -1, "document_root": "./", "home_page": "index.html", "use_implicit_page": true, "implicit_page": "index.html", "upload_path": "uploads", "file_types": [ "gif", "png", "jpg", "js", "css", "html", "ico", "swf", "xap", "apk", "cur", "xml", "webp", "svg" ], "mime": {}, "locations": [ { "default_content_type": "text/plain", "alias": "", "is_case_sensitive": false, "allow_all": true, "is_recursive": true, "filters": [] } ], "max_connections": 100000, "max_connections_per_ip": 0, "load_dynamic_views": false, "dynamic_views_path": [ "./views" ], "dynamic_views_output_path": "", "json_parser_stack_limit": 1000, "enable_unicode_escaping_in_json": true, "float_precision_in_json": { "precision": 0, "precision_type": "significant" }, "log": { "use_spdlog": false, "logfile_base_name": "", "log_size_limit": 100000000, "max_files": 0, "log_level": "DEBUG", "display_local_time": false }, "run_as_daemon": false, "handle_sig_term": true, "relaunch_on_error": false, "use_sendfile": true, "use_gzip": true, "use_brotli": false, "static_files_cache_time": 5, "idle_connection_timeout": 60, "server_header_field": "", "enable_server_header": true, "enable_date_header": true, "keepalive_requests": 0, "pipelining_requests": 0, "gzip_static": true, "br_static": true, "client_max_body_size": "1M", "client_max_memory_body_size": "64K", "client_max_websocket_message_size": "128K", "reuse_port": false, "enabled_compressed_request": false, "enable_request_stream": false }, "plugins": [ { "name": "drogon::plugin::PromExporter", "dependencies": [], "config": { "path": "/metrics" } }, { "name": "drogon::plugin::AccessLogger", "dependencies": [], "config": { "use_spdlog": false, "log_path": "", "log_format": "", "log_file": "access.log", "log_size_limit": 0, "use_local_time": true, "log_index": 0 } } ], "custom_config": { "server": { "listen_address": "0.0.0.0", "listen_port": 5555, "thread_num": 4 }, "redis": { "host": "host.docker.internal", "port": 6379, "client_name": "RedisClientChat" }, "connection_manager": { "connection_ttl_seconds": 30, "heartbeat_interval_seconds": 10, "cleanup_interval_seconds": 60, "scan_count": 100 }, "json": { "unicode_escaping": false }, "logging": { "level": "DEBUG", "display_local_time": true } } } ``` ### 测试连接 使用 wscat 或浏览器测试: ```bash # 安装 wscat npm install -g wscat # 连接到房间 1025 wscat -c "ws://localhost:5555/chat/1025?userName=Alice" ``` 发送消息: ```json { "type": "message", "content": "Hello, World!" } ``` --- ## 架构设计 ### 整体架构图 ``` ┌─────────────────────────────────────────────────────────────┐ │ Client Layer │ │ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐ │ │ │ Browser │ │ Mobile │ │ Desktop │ │ IoT │ │ │ └────┬─────┘ └────┬─────┘ └────┬─────┘ └────┬─────┘ │ └───────┼─────────────┼─────────────┼─────────────┼──────────┘ │ │ │ │ └─────────────┴──────┬──────┴─────────────┘ │ WebSocket ┌────────────────────────────┼────────────────────────────────┐ │ Server Cluster Layer │ │ │ │ ┌──────────────────────────────────────────────────────┐ │ │ │ Process Instance #1 │ │ │ │ ┌─────────────┐ ┌──────────────┐ ┌────────────┐ │ │ │ │ │ ChatWS │ │ Connection │ │ Health │ │ │ │ │ │ Controller │ │ Manager │ │ Check │ │ │ │ │ └──────┬──────┘ └──────┬───────┘ └────────────┘ │ │ │ │ │ │ │ │ │ │ └────────────────┘ │ │ │ │ Local Connections │ │ │ └────────────────────────┬─────────────────────────────┘ │ │ │ │ │ ┌────────────────────────┼─────────────────────────────┐ │ │ │ Process Instance #2 │ │ │ │ ┌─────────────┐ ┌──────────────┐ │ │ │ │ │ ChatWS │ │ Connection │ │ │ │ │ │ Controller │ │ Manager │ │ │ │ │ └──────┬──────┘ └──────┬───────┘ │ │ │ │ │ │ │ │ │ │ └────────────────┘ │ │ │ └────────────────────────┬─────────────────────────────┘ │ └────────────────────────────┼────────────────────────────────┘ │ ┌────────▼────────┐ │ Redis Server │ │ │ │ • Pub/Sub │ │ • Hash (conn) │ │ • ZSET (room) │ │ • ZSET (process) │ └──────────────────┘ ``` ### 分层架构 #### 1. Controller 层 (`controllers/`) **职责**: 处理 WebSocket 连接、消息路由、协议解析 - **AOD_ChatWS**: WebSocket 控制器 - `handleNewConnection()`: 新连接建立 - `handleNewMessage()`: 消息接收与分发 - `handleConnectionClosed()`: 连接断开处理 - **api_v1_health**: HTTP 健康检查接口 #### 2. Data Manager 层 (`dataManager/AOD/`) **职责**: 连接管理、Redis 交互、状态同步 - **ConnectionManager**: 单例管理器 - 本地连接池管理 - Redis 注册与清理 - 消息广播(本地 + 跨进程) - 在线人数统计 #### 3. Utils 层 (`utils/`) **职责**: 工具函数、辅助类 - **RedisManager**: Redis 连接池管理(预留) - **read_etc_hosts**: Docker DNS 解析辅助 --- ## 核心组件详解 ### 1. ChatWS 控制器 **文件**: `controllers/AOD_ChatWS.cc` #### 1.1 连接建立流程 ```cpp void ChatWS::handleNewConnection(const HttpRequestPtr &req, const WebSocketConnectionPtr &wsConnPtr) { // 1. 解析 roomId (支持正则匹配 /chat/1025) std::string roomId = parseRoomIdFromPath(req->path()); // 2. 获取用户名 (优先查询参数,否则使用 IP) std::string userName = req->getParameter("userName"); if (userName.empty()) { userName = "IP_" + req->getPeerAddr().toIpPort(); } // 3. 生成唯一连接ID std::string connId = generateConnId(wsConnPtr); // 4. 保存到连接上下文 wsConnPtr->setContext(std::make_shared(connId)); // 5. 注册到 ConnectionManager auto& manager = ConnectionManager::getInstance(); manager.addConnection(connId, roomId, userName, wsConnPtr); // 6. 异步获取在线人数并发送欢迎消息 manager.getRoomOnlineCount(roomId, [wsConnPtr, connId, roomId, userName, &manager](int count) { Json::Value welcome; welcome["type"] = "system"; welcome["roomId"] = roomId; welcome["userName"] = userName; welcome["content"] = "Welcome to room " + roomId; welcome["timestamp"] = trantor::Date::now().toFormattedString(false); welcome["senderConnId"] = connId; welcome["senderProcessId"] = manager.getProcessId(); welcome["onlineCount"] = count; wsConnPtr->sendJson(welcome); // 广播加入消息 ChatMessage joinMsg; joinMsg.type = "join"; joinMsg.roomId = roomId; joinMsg.userName = userName; joinMsg.content = userName + " joined"; joinMsg.timestamp = trantor::Date::now().toFormattedString(false); joinMsg.senderConnId = connId; joinMsg.senderProcessId = manager.getProcessId(); joinMsg.onlineCount = count; manager.broadcastMessage(joinMsg); }); } ``` **关键点**: - ✅ 使用正则表达式解析路径参数 (`WS_ADD_PATH_VIA_REGEX`) - ✅ 连接ID = 指针地址 + 时间戳,保证唯一性 - ✅ 异步获取在线人数,避免阻塞连接建立 - ✅ 立即发送欢迎消息 + 广播加入通知 #### 1.2 消息处理流程 ```cpp void ChatWS::handleNewMessage(const WebSocketConnectionPtr &wsConnPtr, std::string &&message, const WebSocketMessageType &type) { // 1. 解析 JSON Json::Reader reader; Json::Value json; if (!reader.parse(message, json)) { sendError(wsConnPtr, "Invalid JSON format"); return; } // 2. 获取连接信息 std::string connId = getConnIdFromContext(wsConnPtr); auto& manager = ConnectionManager::getInstance(); ConnectionInfo* info = manager.getConnectionInfo(connId); if (!info) { sendError(wsConnPtr, "Connection not found"); return; } // 3. 根据消息类型处理 std::string msgType = json.get("type", "message").asString(); if (msgType == "message") { // 普通聊天消息 handleChatMessage(json, info, manager); } else if (msgType == "switch_room") { // 切换房间 handleSwitchRoom(json, info, manager, wsConnPtr); } else if (msgType == "get_users") { // 获取用户列表 handleGetUsers(info, manager, wsConnPtr); } } ``` **消息类型**: | 类型 | 说明 | 请求格式 | 响应 | |------|------|---------|------| | `message` | 普通聊天消息 | `{"type":"message","content":"..."}` | 广播给房间内所有人 | | `switch_room` | 切换房间 | `{"type":"switch_room","roomId":"1026"}` | 离开旧房间,加入新房间 | | `get_users` | 获取用户列表 | `{"type":"get_users"}` | 返回用户名数组 | #### 1.3 连接断开处理 ```cpp void ChatWS::handleConnectionClosed(const WebSocketConnectionPtr &wsConnPtr) { std::string connId = getConnIdFromContext(wsConnPtr); auto& manager = ConnectionManager::getInstance(); ConnectionInfo* info = manager.getConnectionInfo(connId); if (!info) return; std::string roomId = info->roomId; std::string userName = info->userName; // 1. 从管理器移除(自动清理 Redis) manager.removeConnection(connId); // 2. 异步获取剩余人数并广播离开消息 manager.getRoomOnlineCount(roomId, [roomId, userName, connId, &manager](int count) { ChatMessage leaveMsg; leaveMsg.type = "leave"; leaveMsg.roomId = roomId; leaveMsg.userName = userName; leaveMsg.content = userName + " left the room"; leaveMsg.timestamp = trantor::Date::now().toFormattedString(false); leaveMsg.senderConnId = connId; leaveMsg.senderProcessId = manager.getProcessId(); leaveMsg.onlineCount = count; manager.broadcastMessage(leaveMsg); }); } ``` **关键点**: - ✅ 先移除连接,再获取人数(确保人数准确) - ✅ 异步回调中广播,避免阻塞断开流程 - ✅ `onlineCount` 反映的是**离开后**的剩余人数 --- ### 2. ConnectionManager 连接管理器 **文件**: `dataManager/AOD/ConnectionManager.cc` #### 2.1 单例模式 ```cpp ConnectionManager& ConnectionManager::getInstance() { static ConnectionManager instance; return instance; } ``` **特点**: - 线程安全(C++11 静态局部变量保证) - 延迟初始化 - 自动析构 #### 2.2 数据结构 ```cpp // 本地连接池 std::unordered_map m_map_connIdLocalConnections; // 房间订阅计数 std::unordered_map m_map_roomConnectionCount; // Redis 客户端 nosql::RedisClientPtr m_redisClient; std::shared_ptr m_redisSub; // 订阅的房间集合 std::unordered_set m_set_subscribedRooms; // 定时器 trantor::TimerId m_heartbeatTimerId; trantor::TimerId m_cleanupTimerId; ``` #### 2.3 添加连接 ```cpp void ConnectionManager::addConnection(const std::string& connId, const std::string& roomId, const std::string& userName, const WebSocketConnectionPtr& conn) { std::lock_guard lock(m_connectionsMutex); // 1. 创建连接信息 ConnectionInfo info; info.connId = connId; info.roomId = roomId; info.userName = userName; info.conn = conn; info.lastPing = std::chrono::steady_clock::now(); // 2. 存入本地连接池 m_map_connIdLocalConnections[connId] = info; // 3. 增加房间计数 m_map_roomConnectionCount[roomId]++; int firstInRoom = m_map_roomConnectionCount[roomId]; // 4. 如果是房间第一个连接,订阅 Redis 频道 if (firstInRoom == 1) { std::thread([this, roomId]() { this->subscribeRoom(roomId); }).detach(); } // 5. 注册到 Redis registerToRedis(info); } ``` **流程图**: ``` addConnection() ├─ 加锁 (m_connectionsMutex) ├─ 创建 ConnectionInfo ├─ 存入 m_map_connIdLocalConnections ├─ m_map_roomConnectionCount[roomId]++ ├─ if (firstInRoom == 1): │ └─ subscribeRoom(roomId) [异步线程] └─ registerToRedis(info) [异步 Redis 命令] ``` #### 2.4 Redis 注册 ```cpp void ConnectionManager::registerToRedis(const ConnectionInfo& info) const { // 使用秒级时间戳(避免浮点数精度问题) const auto now_seconds = std::chrono::system_clock::now() .time_since_epoch().count() / 1000000000; const double timestamp = static_cast(now_seconds); // 1. HSET: 存储连接详细信息 m_redisClient->execCommandAsync( callback, onRedisError, "HSET ws:conn:%s roomId %s userName %s processId %s lastSeen %lld", info.connId.c_str(), info.roomId.c_str(), info.userName.c_str(), m_processId.c_str(), now_seconds); // 2. EXPIRE: 设置 TTL (30秒) m_redisClient->execCommandAsync( callback, onRedisError, "EXPIRE ws:conn:%s %d", info.connId.c_str(), CONNECTION_TTL_SECONDS); // 3. ZADD: 添加到房间有序集合 m_redisClient->execCommandAsync( callback, onRedisError, "ZADD ws:room:%s %f %s", info.roomId.c_str(), timestamp, info.connId.c_str()); // 4. ZADD: 添加到进程有序集合 m_redisClient->execCommandAsync( callback, onRedisError, "ZADD ws:process:%s %f %s", m_processId.c_str(), timestamp, info.connId.c_str()); } ``` **Redis 操作详解**: | 命令 | Key | Value | 用途 | |------|-----|-------|------| | `HSET` | `ws:conn:{connId}` | `{roomId, userName, processId, lastSeen}` | 存储连接元数据 | | `EXPIRE` | `ws:conn:{connId}` | `30` | 自动过期(防脏数据) | | `ZADD` | `ws:room:{roomId}` | `{score=timestamp, member=connId}` | 房间成员列表(支持按时间排序) | | `ZADD` | `ws:process:{processId}` | `{score=timestamp, member=connId}` | 进程成员列表 | **为什么使用 ZSET 而不是 SET?** ✅ **支持自动过期**: 通过分数(时间戳)可以使用 `ZREMRANGEBYSCORE` 批量删除过期成员 ✅ **有序性**: 可以按加入时间排序 ✅ **高效查询**: `ZCARD` 获取成员数量 O(1),`ZRANGE` 获取成员列表 O(log N + M) --- ### 3. 消息广播机制 #### 3.1 双层广播架构 ```cpp void ConnectionManager::broadcastMessage(const ChatMessage& msg, const std::string& excludeConnId) const { // 1. 本地广播(本进程内的连接) broadcastLocal(msg, excludeConnId); // 2. Redis 广播(其他进程的订阅者) Json::Value json = msg.toJson(); std::string str = Json::writeString(builder, json); redisPublish(msg.roomId, str); } ``` **流程图**: ``` broadcastMessage() ├─ broadcastLocal() │ ├─ 遍历 m_map_connIdLocalConnections │ ├─ 过滤: roomId 匹配 && connId != excludeConnId │ └─ wsConn->sendJson(msg) │ └─ redisPublish() └─ PUBLISH chat:room:{roomId} {json_message} └─ 其他进程的 Redis Subscriber 收到消息 └─ handleRedisMessage() └─ broadcastLocal() [在目标进程中] ``` #### 3.2 Redis Pub/Sub 订阅 ```cpp void ConnectionManager::subscribeRoom(const std::string& roomId) { std::lock_guard lock(m_subscribeMutex); if (m_set_subscribedRooms.count(roomId)) { return; // 已订阅 } std::string channel = "chat:room:" + roomId; m_redisSub->subscribe(channel, [this](const std::string& chan, const std::string& msg) { this->handleRedisMessage(chan, msg); }); m_set_subscribedRooms.insert(roomId); } ``` **订阅策略**: - **懒订阅**: 只有当房间有第一个连接时才订阅 - **动态取消**: 当房间最后一个连接断开时取消订阅 - **去重**: 每个房间只订阅一次(即使有多个连接) #### 3.3 消息去重 ```cpp void ConnectionManager::handleRedisMessage(const std::string& channel, const std::string& message) const { // 解析 roomId std::string roomId = channel.substr(strlen("chat:room:")); ChatMessage msg = ChatMessage::fromJson(Json::parse(message)); // 忽略本进程发送的消息(避免重复广播) if (msg.senderProcessId == m_processId) { LOG_DEBUG << "Ignored local message"; return; } // 广播给本进程的本地连接 broadcastLocal(msg, ""); } ``` **去重原理**: ``` Process A 发送消息 ├─ broadcastLocal() → 发送给 A 的本地连接 └─ redisPublish() → PUBLISH chat:room:1025 Process B 收到消息 ├─ handleRedisMessage() ├─ 检查 senderProcessId != B.m_processId ✅ └─ broadcastLocal() → 发送给 B 的本地连接 Process A 也收到自己的消息 ├─ handleRedisMessage() ├─ 检查 senderProcessId == A.m_processId ❌ └─ 忽略(避免重复发送) ``` --- ## 数据流与交互 ### 1. 完整消息流 ``` Client A (Process 1) Redis Client B (Process 2) │ │ │ │ 1. WebSocket Connect │ │ ├────────────────────────>│ │ │ │ │ │ 2. addConnection() │ │ │ ├─ 本地存储 │ │ │ ├─ HSET ws:conn:A │ │ │ ├─ ZADD ws:room:R │ │ │ └─ subscribe(R) │ │ │ │ │ │ 3. Send Message │ │ ├──── {"type":"message"}─>│ │ │ │ │ │ 4. broadcastMessage() │ │ │ ├─ broadcastLocal() │ │ │ │ └─ 发送给Process1│ │ │ └─ PUBLISH chat:R │ │ │ ├──── Message ──────────>│ │ │ │ 5. handleRedisMessage() │ │ │ ├─ 检查 processId │ │ │ └─ broadcastLocal() │ │ │ └─ 发送给 B │ │ │ │ 6. Client B receives │<───────────────────────┤ │<────────────────────────┼────────────────────────┤ ``` ### 2. 连接断开流 ``` Client A ConnectionManager Redis │ │ │ │ 1. WebSocket Close │ │ ├─────────────────────────────>│ │ │ │ │ │ 2. handleConnectionClosed() │ │ │ ├─ getConnectionInfo(A) │ │ │ └─ removeConnection(A) │ │ │ │ │ │ 3. removeConnection() │ │ │ ├─ 从本地 map 移除 │ │ │ ├─ ZREM ws:room:R A │ │ │ ├─ ZREM ws:process:P A │ │ │ └─ DEL ws:conn:A │ │ │ │ │ │ 4. getRoomOnlineCount(R) │ │ │ └─ ZCARD ws:room:R ─────>│ │ │ │<─── 返回剩余人数 ──────┤ │ │ │ │ 5. broadcastMessage(leave) │ │ │ ├─ broadcastLocal() │ │ │ └─ PUBLISH chat:room:R │ │ │ │ │ │ 6. 其他客户端收到离开通知 │ │ │<─────────────────────────────┼────────────────────────┤ ``` ### 3. 心跳刷新流 ``` Heartbeat Timer (每10秒) │ ├─ 遍历 m_map_connIdLocalConnections │ ├─ For each connection: │ ├─ HSET ws:conn:{id} lastSeen {now} │ ├─ EXPIRE ws:conn:{id} 30 │ ├─ ZADD ws:room:{roomId} {now} {connId} ← 关键!更新分数 │ └─ ZADD ws:process:{pid} {now} {connId} ← 关键!更新分数 │ └─ 日志: "Heartbeat sent for N connections" ``` **为什么需要更新 ZSET 分数?** 如果不更新分数,清理时会误删活跃连接: ``` 错误流程(不更新分数): T0: 连接注册,ZADD ws:room:1025 1000 connA T10: 心跳只更新 TTL,ZSET 分数仍是 1000 T35: 惰性清理,阈值 = 1035 - 30 = 1005 比较: 1000 < 1005 → 被误删!❌ 正确流程(更新分数): T0: 连接注册,ZADD ws:room:1025 1000 connA T10: 心跳更新 ZSET,ZADD ws:room:1025 1010 connA T20: 心跳更新 ZSET,ZADD ws:room:1025 1020 connA T35: 惰性清理,阈值 = 1035 - 30 = 1005 比较: 1020 > 1005 → 保留!✅ ``` --- ## Redis 数据结构设计 ### 1. 键命名规范 ``` ws:conn:{connId} → Hash (连接元数据) ws:room:{roomId} → ZSET (房间成员) ws:process:{processId} → ZSET (进程成员) chat:room:{roomId} → Pub/Sub Channel (消息频道) ``` ### 2. 详细结构 #### 2.1 连接 Hash (`ws:conn:{connId}`) ```redis HSET ws:conn:140119147300512_1776093584952596097 roomId "1025" \ userName "Alice" \ processId "proc_1776093568790014397_5096" \ lastSeen 1776093584 EXPIRE ws:conn:140119147300512_1776093584952596097 30 ``` **字段说明**: | 字段 | 类型 | 说明 | |------|------|------| | `roomId` | String | 所属房间ID | | `userName` | String | 用户名 | | `processId` | String | 所属进程ID | | `lastSeen` | Integer | 最后活跃时间(秒级时间戳) | **TTL**: 30秒(由心跳每10秒刷新) #### 2.2 房间 ZSET (`ws:room:{roomId}`) ```redis ZADD ws:room:1025 1776093584.000000 140119147300512_1776093584952596097 \ 1776093587.000000 140119147369184_1776093587583242581 ``` **成员**: `connId` **分数**: 秒级时间戳(用于判断过期) **常用操作**: ```redis # 获取在线人数 ZCARD ws:room:1025 # 获取所有成员 ZRANGE ws:room:1025 0 -1 # 删除过期成员(分数 < 阈值) ZREMRANGEBYSCORE ws:room:1025 -inf 1776093554.000000 ``` #### 2.3 进程 ZSET (`ws:process:{processId}`) ```redis ZADD ws:process:proc_1776093568790014397_5096 1776093584.000000 140119147300512_1776093584952596097 \ 1776093587.000000 140119147369184_1776093587583242581 ``` **用途**: - 进程崩溃后快速清理该进程的所有连接 - 定期清理时只扫描本进程的数据(性能优化) --- ## 清理机制详解 ### 三层清理架构 ``` ┌─────────────────────────────────────────────────────────┐ │ 第一层: TTL 自动过期 │ │ • ws:conn:* 键设置 30 秒 TTL │ │ • 心跳每 10 秒刷新 TTL │ │ • 进程崩溃后 30 秒自动删除 │ └─────────────────────────────────────────────────────────┘ ↓ ┌─────────────────────────────────────────────────────────┐ │ 第二层: 定期后台清理 (60秒) │ │ • 清理本进程的 ws:process:{pid} │ │ • 清理本进程涉及的 ws:room:{roomId} │ │ • 使用 ZREMRANGEBYSCORE 批量删除 │ └─────────────────────────────────────────────────────────┘ ↓ ┌─────────────────────────────────────────────────────────┐ │ 第三层: 惰性清理 (查询时触发) │ │ • getRoomOnlineCount() 前清理该房间 │ │ • getRoomUsers() 前清理该房间 │ │ • 确保返回数据的准确性 │ └─────────────────────────────────────────────────────────┘ ``` ### 1. TTL 自动过期 **配置**: ```cpp static constexpr int CONNECTION_TTL_SECONDS = 30; static constexpr int HEARTBEAT_INTERVAL_SECONDS = 10; ``` **原理**: ``` T0: 连接注册 → EXPIRE ws:conn:A 30 T10: 心跳 → EXPIRE ws:conn:A 30 (重置) T20: 心跳 → EXPIRE ws:conn:A 30 (重置) T30: 如果心跳失败 → Redis 自动删除 ws:conn:A ``` **优势**: - ✅ 无需手动清理 - ✅ 进程崩溃后自动清理 - ✅ Redis 原生支持,性能高 ### 2. 定期后台清理 **定时器**: ```cpp m_cleanupTimerId = drogon::app().getLoop()->runEvery( CLEANUP_INTERVAL_SECONDS, // 60秒 [this]() { performExpiredMembersCleanup(); } ); ``` **清理逻辑**: ```cpp void ConnectionManager::performExpiredMembersCleanup() const { // 计算过期阈值 const auto now_seconds = std::chrono::system_clock::now() .time_since_epoch().count() / 1000000000; const auto expireThreshold = static_cast( now_seconds - CONNECTION_TTL_SECONDS); // 1. 清理本进程的 ws:process:{pid} cleanExpiredMembersFromProcessSet(m_processId, expireThreshold); // 2. 清理本进程涉及的房间 cleanExpiredMembersFromLocalRooms(expireThreshold); } ``` **为什么只清理本进程数据?** 在多进程场景下: - ❌ **全量清理**: 每个进程都扫描所有房间 → 重复工作,性能差 - ✅ **按进程隔离**: 每个进程只清理自己的数据 → 高效,无竞争 ### 3. 惰性清理 **触发时机**: ```cpp void ConnectionManager::getRoomOnlineCount(const std::string& roomId, std::function &&callback) const { // 先清理过期成员 lazyCleanRoomReferences(roomId); // 再查询人数 m_redisClient->execCommandAsync( callback, onError, "ZCARD ws:room:%s", roomId.c_str()); } ``` **清理实现**: ```cpp void ConnectionManager::lazyCleanRoomReferences(const std::string& roomId) const { const auto now_seconds = std::chrono::system_clock::now() .time_since_epoch().count() / 1000000000; const double expireThreshold = static_cast( now_seconds - CONNECTION_TTL_SECONDS); m_redisClient->execCommandAsync( [](const nosql::RedisResult& removedResult) { long long removed = removedResult.asInteger(); if (removed > 0) { LOG_DEBUG << "Lazy cleaned " << removed << " expired members from ws:room:" << roomId; } }, onRedisError, "ZREMRANGEBYSCORE ws:room:%s -inf %f", roomId.c_str(), expireThreshold); } ``` **优势**: - ✅ 按需清理,减少不必要的操作 - ✅ 确保查询结果准确 - ✅ 分散清理压力,避免集中爆发 --- ## 关键算法与实现 ### 1. 时间戳精度问题 **问题**: 纳秒级时间戳转换为 double 时精度丢失 ```cpp // ❌ 错误:纳秒级时间戳 const auto now_ns = std::chrono::system_clock::now() .time_since_epoch().count(); // 1776093584952596097 const double timestamp = static_cast(now_ns); // 转换为: 1.77609253798e+18 (科学计数法,精度丢失) // ✅ 正确:秒级时间戳 const auto now_seconds = std::chrono::system_clock::now() .time_since_epoch().count() / 1000000000; // 1776093584 const double timestamp = static_cast(now_seconds); // 转换为: 1776093584.000000 (精确) ``` **影响范围**: - `registerToRedis()`: 注册时的初始分数 - `startHeartbeatTimer()`: 心跳更新的分数 - `performExpiredMembersCleanup()`: 清理阈值计算 - `lazyCleanRoomReferences()`: 惰性清理阈值计算 **修复**: 所有地方统一使用秒级时间戳(除以 1000000000) ### 2. SCAN 替代 KEYS **问题**: `KEYS` 命令在大数据量时阻塞 Redis ``` cpp // ❌ 错误:使用 KEYS(阻塞) m_redisClient->execCommandAsync(..., "KEYS ws:conn:*"); // ✅ 正确:使用 SCAN(非阻塞) m_redisClient->execCommandAsync(..., "SCAN %s MATCH ws:conn:* COUNT %d", cursor.c_str(), SCAN_COUNT); ``` **SCAN 迭代示例**: ```cpp void ConnectionManager::scanAndSetTTLForAllConnections() const { std::string cursor = "0"; // 第一次 SCAN m_redisClient->execCommandAsync( [this, cursor](const nosql::RedisResult& result) { processScanResultForTTL(cursor, result); }, onRedisError, "SCAN %s MATCH ws:conn:* COUNT %d", cursor.c_str(), SCAN_COUNT); } void ConnectionManager::processScanResultForTTL( const std::string& cursor, const nosql::RedisResult& result) const { auto array = result.asArray(); std::string newCursor = array[0].asString(); auto keys = array[1].asArray(); // 处理这批键 for (const auto& key : keys) { // 设置 TTL m_redisClient->execCommandAsync(..., "EXPIRE ws:conn:%s %d", connId.c_str(), 30); } // 继续扫描 if (newCursor != "0") { m_redisClient->execCommandAsync( [this, newCursor](...) { ... }, onRedisError, "SCAN %s MATCH ws:conn:* COUNT %d", newCursor.c_str(), SCAN_COUNT); } else { LOG_INFO << "SCAN completed"; } } ``` **优势**: - ✅ 非阻塞,每次只返回少量键(COUNT 100) - ✅ 适合大数据量场景 - ✅ 不会导致 Redis 卡顿 ### 3. 并发获取用户名 **问题**: `HMGET` 不支持跨多个 Hash 读取 ```cpp // ❌ 错误:HMGET 只能从一个 Hash 读取 HMGET ws:conn:id1 userName ws:conn:id2 userName // Redis 报错: ERR wrong number of arguments // ✅ 正确:并发执行多个 HGET for (size_t i = 0; i < connIds.size(); ++i) { m_redisClient->execCommandAsync( [users, remaining, callback, index=i](...) { (*users)[index] = result.asString(); if (--(*remaining) == 0) { callback(*users); } }, onError, "HGET ws:conn:%s userName", connIds[i].c_str()); } ``` **实现细节**: ```cpp auto users = std::make_shared>(); users->resize(connIds.size()); // 预分配空间,保持顺序 auto remaining = std::make_shared>( static_cast(connIds.size())); for (size_t i = 0; i < connIds.size(); ++i) { const size_t index = i; // 捕获索引 m_redisClient->execCommandAsync( [users, remaining, callback, index](const nosql::RedisResult& result) { if (result.type() == nosql::RedisResultType::kString) { (*users)[index] = result.asString(); } else { (*users)[index] = "Unknown"; } // 所有请求完成后返回 if (--(*remaining) == 0) { callback(*users); } }, onError, "HGET ws:conn:%s userName", connIds[index].c_str()); } ``` **优势**: - ✅ 并发执行,性能好 - ✅ 使用原子计数器确保全部完成 - ✅ 预分配空间保持顺序 --- ## 部署指南 ### 1. 单机多进程部署 ```bash # 启动第一个实例 ./server_cluster_AOD --port=8080 & # 启动第二个实例 ./server_cluster_AOD --port=8081 & # 两个实例共享同一个 Redis,自动组成集群 ``` ### 2. Docker Compose 部署 [Docker Compose 部署](deploy-docker-compose/deploy-docker-compose.md) --- #### 2.1 项目部署文件结构 ```bash deploy-docker-compose/ ├── Dockerfile # Docker 镜像构建文件 ├── docker-compose.yml # Docker Compose 配置(连接外部 Redis) ├── docker-compose-独立镜像.yml # Docker Compose 配置(连接外部 Redis) ├── nginx.conf # Nginx 负载均衡配置 ├── start_cluster.sh # 容器启动脚本 └── deploy-docker-compose.md # 详细部署文档 ``` --- #### 2.2 Docker 多阶段构建详解 **Dockerfile 说明** (`deploy-docker-compose/Dockerfile`): 采用多阶段构建策略,减小最终镜像体积: **构建阶段 (builder)**: - 基础镜像: `drogonframework/drogon:latest` - 工作目录: `/server_cluster_AOD` - 复制所有源代码和配置文件 - 执行 CMake 配置和编译 - 生成可执行文件 `server_cluster_AOD` **运行阶段**: - 基础镜像: `drogonframework/drogon:latest`(运行时环境) - 工作目录: `/app` - 从 builder 阶段复制可执行文件和 config.json - 创建日志目录 `/app/logs` - 暴露端口 5555 - 配置健康检查(30秒间隔,调用 `/api/v1/health`) - 复制并执行启动脚本 `start_cluster.sh` **构建命令**: ```bash cd deploy-docker-compose docker build -t chat-server-image -f Dockerfile .. ``` **关键点**: - ✅ 使用多阶段构建,避免将编译工具链打包到最终镜像 - ✅ 健康检查确保容器异常时自动重启 - ✅ 启动脚本支持环境变量配置(REDIS_HOST, REDIS_PORT, SERVER_PORT) #### 2.3 Docker Compose 部署模式 **模式一:连接外部 Redis** (`docker-compose.yml`) 适用场景:已有独立的 Redis 服务(如云数据库 Redis) **配置特点**: - 使用 YAML 锚点 (`&chat-server-config`) 定义模板,避免重复配置 - 启动 3 个聊天服务器实例(chat-server-1/2/3) - 通过 `extra_hosts` 配置 `host.docker.internal` 访问宿主机 Redis - Nginx 作为负载均衡器,统一入口(80端口) - 资源限制:每个容器最多 512MB 内存,1 CPU - 健康检查:30秒间隔,失败3次后重启 **启动步骤**: ```bash cd deploy-docker-compose # 1. 创建 Docker 网络(首次需要) docker network create my-net # 2. 创建日志目录 mkdir -p logs/nginx # 3. 启动所有服务 docker-compose up -d # 4. 查看服务状态 docker-compose ps # 5. 查看实时日志 docker-compose logs -f chat-server-1 docker-compose logs -f nginx # 6. 停止服务 docker-compose down ``` **访问方式**: - WebSocket: `ws://localhost/chat/1025?userName=Alice` - HTTP API: `http://localhost/api/v1/health` - Nginx 健康检查: `http://localhost/nginx-health` **模式二:每一个服务创建独立镜像** (`docker-compose-独立镜像.yml`) 适用场景:每一个服务环境不同,独立部署 **启动命令**: ```bash docker-compose -f docker-compose-独立镜像.yml up -d ``` #### 2.4 Nginx 负载均衡配置详解 **nginx.conf 核心配置** (`deploy-docker-compose/nginx.conf`): **上游服务器配置**: ```nginx upstream chat_backend { ip_hash; # 会话保持(WebSocket 必需) server chat-server-1:5555 max_fails=3 fail_timeout=30s; server chat-server-2:5555 max_fails=3 fail_timeout=30s; server chat-server-3:5555 max_fails=3 fail_timeout=30s; keepalive 32; # 保持长连接池 } ``` **关键配置项说明**: | 配置 | 值 | 说明 | |------|-----|------| | `ip_hash` | - | 同一客户端 IP 始终路由到同一后端,保证 WebSocket 会话不中断 | | `max_fails` | 3 | 连续失败3次标记为不可用 | | `fail_timeout` | 30s | 失败后30秒内不再尝试该后端 | | `keepalive` | 32 | 保持32个空闲连接到后端,减少握手开销 | **WebSocket 特殊配置** (`location /chat`): ```nginx proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection "upgrade"; proxy_read_timeout 86400s; # 24小时超时 proxy_send_timeout 86400s; proxy_buffering off; # 关闭缓冲,实时转发 proxy_request_buffering off; client_max_body_size 5M; # 单条消息最大5MB ``` **路由规则**: - `/nginx-health`: Nginx 自身健康检查(返回 "healthy") - `/api/v1/health/`: 后端健康检查(短超时 5秒) - `/api/*`: REST API 请求转发 - `/chat*`: WebSocket 连接(长超时 24小时) - `/`: 默认路由(转发到后端) **性能优化**: - Gzip 压缩(JSON、文本等) - epoll 事件模型 - worker_connections 4096 - tcp_nopush + tcp_nodelay #### 2.5 启动脚本说明 **start_cluster.sh** (`deploy-docker-compose/start_cluster.sh`): ```bash #!/bin/bash REDIS_HOST=${REDIS_HOST:-127.0.0.1} REDIS_PORT=${REDIS_PORT:-6379} SERVER_PORT=${SERVER_PORT:-5555} echo "Starting chat server..." echo "Redis: $REDIS_HOST:$REDIS_PORT" echo "Server port: $SERVER_PORT" # 单进程模式(Docker 推荐) exec ./server_cluster_AOD $SERVER_PORT $REDIS_HOST $REDIS_PORT ``` **功能**: - 从环境变量读取配置(支持 Docker `-e` 参数覆盖) - 默认值:REDIS_HOST=127.0.0.1, REDIS_PORT=6379, SERVER_PORT=5555 - 使用 `exec` 替换当前进程,确保信号正确传递(支持 `docker stop` 优雅关闭) **历史遗留代码**(已注释): - 原支持单机多进程模式(PROCESS_NUM 参数) - Docker 环境下推荐单容器单进程,通过多个容器实现集群 #### 2.6本地开发(最快上手) ```bash # 1. 编译 mkdir build && cd build cmake .. && make -j$(nproc) # 2. 启动 Redis redis-server --daemonize yes # 3. 运行服务 ./server_cluster_AOD 5555 127.0.0.1 6379 # 4. 测试连接 wscat -c "ws://localhost:5555/chat/1025?userName=Test" ``` #### 2.7 Docker 单机测试 ```bash cd deploy-docker-compose docker-compose -f docker-compose-独立镜像.yml up -d wscat -c "ws://localhost/chat/1025?userName=DockerTest" ``` #### 2.8 生产环境部署 ```bash # 1. 准备服务器(Ubuntu 22.04) sudo apt update sudo apt install -y docker.io docker-compose # 2. 克隆代码 git clone https://gitee.com/ljc8/drogon_chat_server_cluster.git cd drogon_chat_server_cluster/deploy-docker-compose # 3. 修改配置(nginx.conf、docker-compose.yml) vim docker-compose.yml # 4. 启动服务 docker-compose up -d # 5. 验证部署 curl http://localhost/api/v1/health docker-compose ps ``` ### 3. Kubernetes 部署 [待完善](./deploy-k8s/deploy-k8s.md) --- ## 故障排查 ### 1. 在线人数为 0 **症状**: `ZCARD ws:room:1025` 返回 0 **可能原因**: 1. ✅ 时间戳精度问题(已修复) 2. ✅ 心跳未更新 ZSET 分数(已修复) 3. ❌ 清理阈值计算错误 **排查步骤**: ```redis # 1. 检查房间是否有成员 ZRANGE ws:room:1025 0 -1 WITHSCORES # 2. 检查连接是否存在 HGETALL ws:conn:{connId} # 3. 检查 TTL TTL ws:conn:{connId} # 4. 查看日志 grep "ZCARD Result" logs/*.log grep "Lazy cleaned" logs/*.log ``` ### 2. 消息未广播 **症状**: 发送消息后其他客户端收不到 **排查**: ```redis # 1. 检查 Redis Pub/Sub SUBSCRIBE chat:room:1025 # 2. 手动发布测试 PUBLISH chat:room:1025 '{"type":"test"}' # 3. 查看订阅情况 PUBSUB CHANNELS chat:room:* ``` **日志检查**: ```bash grep "PUBLISH Result" logs/*.log grep "handleRedisMessage" logs/*.log grep "Ignored local message" logs/*.log ``` ### 3. 内存泄漏 **监控**: ```bash # 监控 Redis 内存 redis-cli INFO memory # 监控进程内存 ps aux | grep server_cluster_AOD ``` **常见原因**: - WebSocket 连接未正确关闭 - Redis 订阅未取消 - 定时器未停止 **解决方案**: ```cpp void ConnectionManager::shutdown() { // 1. 停止定时器 if (m_heartbeatTimerId != trantor::InvalidTimerId) { drogon::app().getLoop()->invalidateTimer(m_heartbeatTimerId); } if (m_cleanupTimerId != trantor::InvalidTimerId) { drogon::app().getLoop()->invalidateTimer(m_cleanupTimerId); } // 2. 清理所有连接 std::vector connIds = getAllConnectionIds(); for (const auto& connId : connIds) { removeConnection(connId); } } ``` --- ## 性能优化 ### 1. 连接池优化 **当前**: 每个进程一个 Redis 连接 **优化**: 使用连接池(多线程场景) ```cpp // 未来优化方向 std::vector redisPool_; ``` ### 2. 批量操作 **当前**: 逐个执行 Redis 命令 **优化**: 使用 Pipeline 批量执行 ```cpp // 目前不支持pipeline // 伪代码 m_redisClient->pipeline([this](auto& pipe) { for (const auto& conn : connections) { pipe.execCommand("HSET ..."); pipe.execCommand("EXPIRE ..."); } }); ``` ### 3. 缓存优化 **当前**: 每次查询都访问 Redis **优化**: 本地缓存 + 定期同步 ```cpp // 本地缓存房间人数 std::unordered_map roomOnlineCountCache_; // 定期同步 void syncRoomCounts() { for (const auto& [roomId, _] : m_map_roomConnectionCount) { getRoomOnlineCount(roomId, [this, roomId](int count) { roomOnlineCountCache_[roomId] = count; }); } } ``` --- ## 开发规范 ### 1. 代码风格 - 使用 **snake_case** 命名函数和变量 - 使用 **PascalCase** 命名类和结构体 - 使用 **驼峰命名** 命名常量 ### 2. 日志规范 ``` cpp LOG_TRACE << "详细调试信息"; LOG_DEBUG << "调试信息"; LOG_INFO << "重要业务流程"; LOG_WARN << "警告信息"; LOG_ERROR << "错误信息"; ``` ### 3. 错误处理 ```cpp // 始终检查返回值 if (!m_redisClient) { LOG_ERROR << "Redis client is null"; return; } // 异步回调中检查 callback if (nullptr == callback) { LOG_ERROR << "Invalid callback"; return; } ``` ### 4. 线程安全 ```cpp // 访问共享数据时加锁 std::lock_guard lock(m_connectionsMutex); // 原子操作 std::atomic remaining{count}; ``` --- ## 附录 ### A. Git 提交历史亮点 | Commit | 说明 | |--------|------| | `1705516` | 修复时间戳精度问题,统一使用秒级时间戳 | | `b2f3f8b` | 优化用户名获取,改用并发 HGET | | `d929e6f` | 重构清理机制,按进程隔离清理 | | `ae037ad` | 修复 SCAN 格式字符串问题 | | `32bc242` | 添加定期清理和惰性清理机制 | | `5bbc4ef` | 实现 WebSocket 路径正则匹配 | ### B. 参考资料 - [Drogon 官方文档](https://github.com/drogonframework/drogon/wiki) - [Redis 官方文档](https://redis.io/documentation) - [Trantor 事件循环](https://github.com/an-tao/trantor) ### C. 许可证 本项目采用 MIT 许可证。 --- **维护者**: XXX **最后更新**: 2026-04-13```