i007.cc

i007.cc

优先队列-降维打击

05.价值资料

使用asio访问网络,数据库,使用定时器开发游戏服务器的最佳实践

这是一个完整的实战指南,把前面所有讨论落地成可用的代码结构。


一、项目结构与线程模型

game_server/
├── main.cpp
├── net/
│   ├── server.hpp          # TCP 接入,accept 循环
│   └── session.hpp         # 单个连接,读写协程
├── game/
│   ├── room.hpp            # Actor 模型,strand 保护
│   ├── player.hpp          # 玩家内存状态
│   └── game_loop.hpp       # Tick 主循环
├── db/
│   ├── redis_client.hpp    # aedis 封装
│   └── mysql_client.hpp    # MySQL 异步封装
├── proto/
│   └── messages.proto      # Protobuf 消息定义
└── util/
    ├── timer_wheel.hpp     # 时间轮(大量定时器时用)
    └── object_pool.hpp     # 减少内存分配

 

cpp
// main.cpp —— 线程模型的核心决策

int main() {
    // io_context 数量 = CPU 核心数
    // 每个 io_context 跑一个线程,减少跨线程竞争
    const int THREAD_COUNT = std::thread::hardware_concurrency();
    
    std::vector<std::unique_ptr<asio::io_context>> io_ctxs;
    std::vector<std::thread> threads;
    
    for (int i = 0; i < THREAD_COUNT; i++) {
        io_ctxs.push_back(std::make_unique<asio::io_context>());
    }
    
    // 每个 Room 绑定到一个 io_context(通过 strand)
    // 同一 Room 的所有操作天然串行,不需要锁
    auto& room_ctx = *io_ctxs[0];  // Room 专用
    auto& net_ctx  = *io_ctxs[1];  // 网络 I/O 专用
    auto& db_ctx   = *io_ctxs[2];  // 数据库 I/O 专用
    
    for (auto& ctx : io_ctxs) {
        // work_guard 防止 io_context 在没有任务时退出
        auto guard = asio::make_work_guard(*ctx);
        threads.emplace_back([&ctx, g = std::move(guard)] {
            ctx.run();
        });
    }
    
    // 启动服务
    GameServer server(net_ctx, room_ctx, db_ctx);
    server.start(8888);
    
    for (auto& t : threads) t.join();
    return 0;
}

 


二、网络层:Session 生命周期管理

cpp
// net/session.hpp

class GameSession : public std::enable_shared_from_this<GameSession> {
public:
    // 消息包头:4字节长度 + 2字节命令ID + 2字节序列号
    struct Header {
        uint32_t length;   // body 长度(网络字节序)
        uint16_t cmd_id;
        uint16_t seq;
    };
    static constexpr size_t HEADER_SIZE = sizeof(Header);
    static constexpr size_t MAX_BODY_SIZE = 64 * 1024;  // 64KB

    explicit GameSession(tcp::socket socket)
        : socket_(std::move(socket))
        , heartbeat_timer_(socket_.get_executor())
        , session_id_(generate_session_id())
        , send_strand_(socket_.get_executor()) {}

    void start() {
        // 同时启动读协程和心跳检测
        co_spawn(socket_.get_executor(),
            [self = shared_from_this()]() { return self->read_loop(); },
            asio::detached);
        
        schedule_heartbeat();
    }

    // 线程安全的发送:通过 strand 序列化写操作
    void send(std::vector<uint8_t> data) {
        asio::post(send_strand_, [self = shared_from_this(),
                                  data = std::move(data)]() mutable {
            self->send_queue_.push(std::move(data));
            if (!self->writing_) {
                self->do_write();
            }
        });
    }

    // 注册消息处理器(由 Room 设置)
    void set_message_handler(
        std::function<void(uint16_t, uint16_t, std::span<const uint8_t>)> handler) {
        message_handler_ = std::move(handler);
    }

    void disconnect(std::string_view reason) {
        if (!connected_.exchange(false)) return;  // 防止重复断开
        LOG_INFO("session {} disconnect: {}", session_id_, reason);
        heartbeat_timer_.cancel();
        socket_.close();
        if (disconnect_handler_) disconnect_handler_(session_id_);
    }

    uint64_t session_id() const { return session_id_; }

private:
    awaitable<void> read_loop() {
        try {
            while (connected_) {
                // 读包头
                Header header;
                co_await async_read(socket_,
                    asio::buffer(&header, HEADER_SIZE),
                    use_awaitable);
                
                // 字节序转换
                uint32_t body_len = ntohl(header.length);
                uint16_t cmd_id   = ntohs(header.cmd_id);
                uint16_t seq      = ntohs(header.seq);
                
                if (body_len > MAX_BODY_SIZE) {
                    disconnect("oversized packet");
                    co_return;
                }
                
                // 读消息体
                std::vector<uint8_t> body(body_len);
                if (body_len > 0) {
                    co_await async_read(socket_,
                        asio::buffer(body),
                        use_awaitable);
                }
                
                on_message(cmd_id, seq, body);
            }
        } catch (const std::exception& e) {
            disconnect(e.what());
        }
    }

    void on_message(uint16_t cmd, uint16_t seq,
                    const std::vector<uint8_t>& body) {
        reset_heartbeat();
        if (message_handler_) {
            message_handler_(cmd, seq, std::span{body});
        }
    }

    void do_write() {
        if (send_queue_.empty()) {
            writing_ = false;
            return;
        }

        // 批量发送:把队列里所有消息合并成一次 writev
        std::vector<asio::const_buffer> buffers;
        while (!send_queue_.empty()) {
            write_batch_.push_back(std::move(send_queue_.front()));
            send_queue_.pop();
            buffers.push_back(asio::buffer(write_batch_.back()));
            
            if (buffers.size() > 64) break;  // 单批次上限
        }

        writing_ = true;
        async_write(socket_, buffers,
            asio::bind_executor(send_strand_,
                [self = shared_from_this()](auto ec, auto) {
                    self->write_batch_.clear();
                    if (ec) {
                        self->disconnect(ec.message());
                        return;
                    }
                    self->do_write();  // 继续发剩余
                }));
    }

    void schedule_heartbeat() {
        heartbeat_timer_.expires_after(std::chrono::seconds(30));
        heartbeat_timer_.async_wait([self = shared_from_this()](auto ec) {
            if (!ec) self->disconnect("heartbeat timeout");
        });
    }

    void reset_heartbeat() {
        heartbeat_timer_.cancel();
        schedule_heartbeat();
    }

    tcp::socket socket_;
    steady_timer heartbeat_timer_;
    uint64_t session_id_;
    
    // 写操作必须通过 strand 串行化,否则并发写 socket 会崩
    asio::strand<asio::any_io_executor> send_strand_;
    std::queue<std::vector<uint8_t>> send_queue_;
    std::vector<std::vector<uint8_t>> write_batch_;
    bool writing_ = false;
    
    std::atomic<bool> connected_{true};
    std::function<void(uint16_t, uint16_t, std::span<const uint8_t>)> message_handler_;
    std::function<void(uint64_t)> disconnect_handler_;
};

 


三、Room:Actor 模型 + Strand

cpp
// game/room.hpp

class Room : public std::enable_shared_from_this<Room> {
public:
    explicit Room(asio::io_context& ctx, uint64_t room_id)
        : strand_(asio::make_strand(ctx))  // 所有 Room 操作在同一个 strand 上
        , tick_timer_(strand_)
        , save_timer_(strand_)
        , room_id_(room_id) {}

    // 玩家加入:在 strand 上执行,保证线程安全
    awaitable<void> add_player(std::shared_ptr<GameSession> session,
                               uint64_t player_id) {
        // co_spawn 到 strand,保证后续操作串行
        co_await asio::post(strand_, use_awaitable);
        
        // 加载玩家数据(这是冷路径,可以等 DB)
        auto player_data = co_await load_player_data(player_id);
        
        auto player = std::make_shared<Player>(player_id, player_data);
        player->set_session(session);
        players_[player_id] = player;
        
        // 设置消息路由:网络消息投递到 strand 上处理
        session->set_message_handler(
            [self = shared_from_this(), player_id]
            (uint16_t cmd, uint16_t seq, std::span<const uint8_t> body) {
                // 把消息 post 到 strand,保证游戏逻辑串行
                asio::post(self->strand_,
                    [self, player_id, cmd, seq,
                     body_vec = std::vector<uint8_t>(body.begin(), body.end())]() {
                        self->on_player_message(player_id, cmd, seq, body_vec);
                    });
            });
        
        // 通知房间内其他玩家有人加入
        broadcast_player_joined(player_id);
        
        LOG_INFO("player {} joined room {}", player_id, room_id_);
    }

    void start() {
        start_tick();
        start_periodic_save();
    }

private:
    // ─── 游戏主循环 ───────────────────────────────────────────

    void start_tick() {
        tick_timer_.expires_after(TICK_INTERVAL);
        tick_timer_.async_wait(
            asio::bind_executor(strand_, [self = shared_from_this()](auto ec) {
                if (!ec) self->do_tick();
            }));
    }

    void do_tick() {
        auto now = Clock::now();
        auto dt = std::chrono::duration<float>(now - last_tick_).count();
        last_tick_ = now;

        // 监控 Tick 耗时,超出预算立即告警
        auto tick_start = Clock::now();

        // ── 热路径:纯内存操作,绝不访问 DB ──
        process_pending_inputs();   // 消费本帧积累的输入
        update_physics(dt);
        update_skills(dt);
        check_collisions();
        broadcast_delta_state();    // 只广播变化量
        // ── 热路径结束 ──

        auto tick_cost = Clock::now() - tick_start;
        if (tick_cost > TICK_INTERVAL * 0.8f) {
            LOG_WARN("room {} tick overtime: {}ms", room_id_,
                std::chrono::duration_cast<std::chrono::milliseconds>(tick_cost).count());
        }

        // 调度下一帧(用 expires_at 而非 expires_after,防止时间漂移)
        tick_timer_.expires_at(tick_timer_.expiry() + TICK_INTERVAL);
        tick_timer_.async_wait(
            asio::bind_executor(strand_, [self = shared_from_this()](auto ec) {
                if (!ec) self->do_tick();
            }));
    }

    void process_pending_inputs() {
        // Tick 开始时一次性消费所有待处理输入
        // 输入已经在 strand 上投递过来了,这里直接处理
        for (auto& [pid, player] : players_) {
            while (auto input = player->pop_input()) {
                apply_player_input(pid, *input);
            }
        }
    }

    // ─── 数据库访问:冷路径 ──────────────────────────────────

    awaitable<PlayerData> load_player_data(uint64_t player_id) {
        // 并发加载不同数据源
        using namespace asio::experimental::awaitable_operators;
        
        auto [base, inventory, skills] = co_await (
            redis_client_->hgetall("player:" + std::to_string(player_id)) &&
            redis_client_->lrange("inventory:" + std::to_string(player_id), 0, -1) &&
            mysql_client_->query_async<SkillRow>(
                "SELECT * FROM player_skills WHERE player_id = ?", player_id)
        );
        
        co_return PlayerData{
            .base      = parse_base_data(base),
            .inventory = parse_inventory(inventory),
            .skills    = skills
        };
    }

    // ─── 定时存档:延迟写回 ──────────────────────────────────

    void start_periodic_save() {
        save_timer_.expires_after(std::chrono::seconds(30));
        save_timer_.async_wait(
            asio::bind_executor(strand_, [self = shared_from_this()](auto ec) {
                if (!ec) {
                    co_spawn(self->strand_,
                        [self]() { return self->flush_dirty_players(); },
                        asio::detached);
                    self->start_periodic_save();  // 重新调度
                }
            }));
    }

    awaitable<void> flush_dirty_players() {
        // 找出所有脏数据玩家,并发写回
        std::vector<awaitable<void>> tasks;
        for (auto& [pid, player] : players_) {
            if (player->is_dirty()) {
                tasks.push_back(save_player(pid, player));
            }
        }
        
        // 并发执行所有存档任务
        for (auto& task : tasks) {
            co_await std::move(task);
        }
    }

    awaitable<void> save_player(uint64_t pid, std::shared_ptr<Player> player) {
        try {
            auto snapshot = player->take_snapshot();  // 快照,不阻塞游戏逻辑
            
            co_await redis_client_->hmset(
                "player:" + std::to_string(pid),
                snapshot.to_redis_fields());
            
            // 关键数据同步写 MySQL
            if (snapshot.has_currency_change()) {
                co_await mysql_client_->execute_async(
                    "UPDATE player_currency SET gold=?, diamond=? WHERE id=?",
                    snapshot.gold, snapshot.diamond, pid);
            }
            
            player->clear_dirty();
        } catch (const std::exception& e) {
            LOG_ERROR("save player {} failed: {}", pid, e.what());
            // 存档失败不影响游戏,下次 Tick 继续标记 dirty
        }
    }

    // ─── 消息分发 ────────────────────────────────────────────

    void on_player_message(uint64_t player_id, uint16_t cmd,
                           uint16_t seq, const std::vector<uint8_t>& body) {
        // 此时已经在 strand 上,可以安全访问 players_
        auto it = players_.find(player_id);
        if (it == players_.end()) return;
        
        // 把消息放入玩家输入队列,等下一个 Tick 处理
        // 这样保证输入处理在 Tick 的固定时间点发生
        it->second->push_input(cmd, seq, body);
    }

    void broadcast_delta_state() {
        auto delta = world_state_.compute_delta(last_broadcast_);
        if (delta.empty()) return;
        
        auto packet = serialize_delta(delta);
        for (auto& [pid, player] : players_) {
            if (auto session = player->session()) {
                session->send(packet);
            }
        }
        last_broadcast_ = world_state_.snapshot();
    }

    static constexpr auto TICK_INTERVAL = std::chrono::milliseconds(66);
    
    asio::strand<asio::io_context::executor_type> strand_;
    steady_timer tick_timer_;
    steady_timer save_timer_;
    
    uint64_t room_id_;
    std::unordered_map<uint64_t, std::shared_ptr<Player>> players_;
    
    std::shared_ptr<RedisClient>  redis_client_;
    std::shared_ptr<MysqlClient>  mysql_client_;
    
    WorldState world_state_;
    WorldState last_broadcast_;
    Clock::time_point last_tick_ = Clock::now();
};

 


四、Redis 客户端封装

cpp
// db/redis_client.hpp

class RedisClient {
public:
    explicit RedisClient(asio::io_context& ctx, const RedisConfig& cfg)
        : ctx_(ctx), config_(cfg) {}

    awaitable<void> connect() {
        conn_ = std::make_shared<aedis::connection>(ctx_);
        
        aedis::connection::config aedis_cfg;
        aedis_cfg.addr.host = config_.host;
        aedis_cfg.addr.port = config_.port;
        
        co_await conn_->async_run(aedis_cfg, {}, use_awaitable);
    }

    // HGETALL:返回 hash 的所有字段
    awaitable<std::unordered_map<std::string, std::string>>
    hgetall(const std::string& key) {
        aedis::request req;
        req.push("HGETALL", key);
        
        aedis::response<std::vector<std::string>> resp;
        co_await conn_->async_exec(req, resp, use_awaitable);
        
        auto& vec = std::get<0>(resp).value();
        std::unordered_map<std::string, std::string> result;
        for (size_t i = 0; i + 1 < vec.size(); i += 2) {
            result[vec[i]] = vec[i + 1];
        }
        co_return result;
    }

    // HMSET:批量设置 hash 字段
    awaitable<void> hmset(const std::string& key,
                          const std::vector<std::pair<std::string, std::string>>& fields) {
        aedis::request req;
        req.push_range("HMSET", key, fields);
        
        aedis::ignore_response resp;
        co_await conn_->async_exec(req, resp, use_awaitable);
    }

    // 原子扣减(防止超卖/刷金币)
    awaitable<bool> deduct_if_enough(const std::string& key,
                                     const std::string& field,
                                     int64_t amount) {
        // Lua 脚本保证原子性
        static const std::string SCRIPT = R"(
            local val = tonumber(redis.call('HGET', KEYS[1], ARGV[1]))
            if val == nil or val < tonumber(ARGV[2]) then
                return 0
            end
            redis.call('HDECRBY', KEYS[1], ARGV[1], ARGV[2])
            return 1
        )";
        
        aedis::request req;
        req.push("EVAL", SCRIPT, 1, key, field, std::to_string(amount));
        
        aedis::response<int64_t> resp;
        co_await conn_->async_exec(req, resp, use_awaitable);
        co_return std::get<0>(resp).value() == 1;
    }

    // 带重试的执行
    template<typename Request, typename Response>
    awaitable<void> exec_with_retry(Request& req, Response& resp,
                                    int max_retries = 3) {
        for (int i = 0; i < max_retries; i++) {
            try {
                co_await conn_->async_exec(req, resp, use_awaitable);
                co_return;
            } catch (const std::exception& e) {
                if (i == max_retries - 1) throw;
                
                // 指数退避
                steady_timer backoff(ctx_);
                backoff.expires_after(std::chrono::milliseconds(100 * (1 << i)));
                co_await backoff.async_wait(use_awaitable);
                
                // 重连
                co_await connect();
            }
        }
    }

private:
    asio::io_context& ctx_;
    RedisConfig config_;
    std::shared_ptr<aedis::connection> conn_;
};

 


五、MySQL 客户端封装

MySQL 官方异步 API 不完善,最实用的方案是用线程池桥接:

cpp
// db/mysql_client.hpp

class MysqlClient {
public:
    explicit MysqlClient(asio::io_context& ctx, const MysqlConfig& cfg)
        : ctx_(ctx)
        , pool_(std::make_unique<ConnectionPool>(cfg, POOL_SIZE)) {}

    // 把同步 MySQL 操作包装成 awaitable
    // 在线程池里执行,不阻塞 io_context 线程
    template<typename Row>
    awaitable<std::vector<Row>> query_async(std::string_view sql, auto&&... args) {
        co_return co_await asio::async_compose
            decltype(use_awaitable),
            void(std::exception_ptr, std::vector<Row>)>(
            
            [this, sql = std::string(sql),
             args_tuple = std::make_tuple(std::forward<decltype(args)>(args)...)]
            (auto& self) mutable {
                // 提交到线程池执行
                asio::post(thread_pool_,
                    [this, sql, args_tuple, &self]() mutable {
                        try {
                            auto conn = pool_->acquire();
                            auto result = std::apply(
                                [&](auto&&... a) {
                                    return execute_query<Row>(conn, sql, a...);
                                }, args_tuple);
                            
                            // 结果投递回 io_context
                            asio::post(ctx_,
                                [result = std::move(result), &self]() mutable {
                                    self.complete(nullptr, std::move(result));
                                });
                        } catch (...) {
                            auto eptr = std::current_exception();
                            asio::post(ctx_,
                                [eptr, &self]() mutable {
                                    self.complete(eptr, {});
                                });
                        }
                    });
            },
            use_awaitable, ctx_);
    }

    awaitable<void> execute_async(std::string_view sql, auto&&... args) {
        // 类似上面,省略
    }

    // 事务:必须在同一个连接上执行所有语句
    awaitable<void> transaction(
        std::function<void(MysqlConnection&)> operations) {
        
        co_return co_await asio::async_compose
            decltype(use_awaitable), void(std::exception_ptr)>(
            
            [this, operations = std::move(operations)](auto& self) mutable {
                asio::post(thread_pool_,
                    [this, operations, &self]() mutable {
                        try {
                            auto conn = pool_->acquire();
                            conn->execute("BEGIN");
                            try {
                                operations(*conn);
                                conn->execute("COMMIT");
                            } catch (...) {
                                conn->execute("ROLLBACK");
                                throw;
                            }
                            asio::post(ctx_, [&self] { self.complete(nullptr); });
                        } catch (...) {
                            auto eptr = std::current_exception();
                            asio::post(ctx_, [eptr, &self] { self.complete(eptr); });
                        }
                    });
            },
            use_awaitable, ctx_);
    }

private:
    static constexpr int POOL_SIZE = 8;
    
    asio::io_context& ctx_;
    asio::thread_pool thread_pool_{4};  // MySQL 操作专用线程池
    std::unique_ptr<ConnectionPool> pool_;
};

 


六、定时器的正确使用

cpp
// 三种常见定时器模式

// ─── 模式1:固定间隔(Tick)─── 用 expires_at 防止漂移
void schedule_next_tick() {
    tick_timer_.expires_at(tick_timer_.expiry() + TICK_INTERVAL);
    //          ^^^^^^^^^^^^^^^^^^^^^^^^^^^
    // 关键:基于上次过期时间而非当前时间
    // 用 expires_after 会导致每次延迟叠加,时间漂移
    tick_timer_.async_wait(...);
}

// ─── 模式2:一次性延迟 ───
awaitable<void> skill_cooldown(uint64_t player_id, uint32_t skill_id,
                               std::chrono::milliseconds duration) {
    steady_timer timer(co_await this_coro::executor);
    timer.expires_after(duration);
    co_await timer.async_wait(use_awaitable);
    
    // 冷却结束,通知客户端
    notify_skill_ready(player_id, skill_id);
}

// ─── 模式3:超时竞争 ─── 最重要的模式
awaitable<std::optional<PlayerData>>
load_with_timeout(uint64_t player_id,
                  std::chrono::milliseconds timeout) {
    using namespace asio::experimental::awaitable_operators;
    
    steady_timer timer(co_await this_coro::executor);
    timer.expires_after(timeout);
    
    // 两个 awaitable 竞争,谁先完成用谁
    auto result = co_await (
        load_player_data(player_id) ||
        timer.async_wait(use_awaitable)
    );
    
    if (result.index() == 0) {
        co_return std::get<0>(result);  // 正常加载完成
    } else {
        LOG_WARN("load player {} timeout", player_id);
        co_return std::nullopt;         // 超时
    }
}

// ─── 大量定时器:时间轮 ───
// 上万个技能 CD 用 asio::steady_timer 会产生大量 timer 对象
// 时间轮更高效:O(1) 添加/删除,O(1) 触发检查
class TimerWheel {
    static constexpr int SLOTS = 256;          // 256 格
    static constexpr auto SLOT_DURATION = 10ms; // 每格 10ms,覆盖 2.56 秒
    
    std::array<std::vector<TimerCallback>, SLOTS> slots_;
    int current_slot_ = 0;
    steady_timer wheel_timer_;

public:
    void add_timer(std::chrono::milliseconds delay, TimerCallback cb) {
        int slot_count = delay / SLOT_DURATION;
        int target_slot = (current_slot_ + slot_count) % SLOTS;
        slots_[target_slot].push_back(std::move(cb));
    }

    awaitable<void> run() {
        while (true) {
            wheel_timer_.expires_after(SLOT_DURATION);
            co_await wheel_timer_.async_wait(use_awaitable);
            
            current_slot_ = (current_slot_ + 1) % SLOTS;
            for (auto& cb : slots_[current_slot_]) cb();
            slots_[current_slot_].clear();
        }
    }
};

 


七、错误处理与优雅关闭

cpp
// 统一的错误处理策略

// 协程内部:try/catch,局部恢复
awaitable<void> safe_save_player(uint64_t pid) {
    try {
        co_await save_player(pid);
    } catch (const RedisException& e) {
        // Redis 失败:标记 dirty,下次重试
        LOG_ERROR("redis save failed: {}", e.what());
        players_[pid]->mark_dirty();
    } catch (const MysqlException& e) {
        // MySQL 失败:报警,人工介入
        LOG_CRITICAL("mysql save failed for player {}: {}", pid, e.what());
        alert_oncall("mysql_save_failure", pid, e.what());
    }
}

// 顶层协程:用 co_spawn 的第三个参数捕获未处理异常
co_spawn(ctx,
    room->game_loop(),
    [room_id](std::exception_ptr eptr) {
        if (eptr) {
            try {
                std::rethrow_exception(eptr);
            } catch (const std::exception& e) {
                LOG_CRITICAL("room {} crashed: {}", room_id, e.what());
                // 重启 Room 或通知所有玩家回到大厅
                recover_room(room_id);
            }
        }
    });

// 优雅关闭
awaitable<void> graceful_shutdown() {
    LOG_INFO("shutting down...");
    
    // 1. 停止接受新连接
    acceptor_.close();
    
    // 2. 给所有玩家发关服通知,给 5 秒时间
    broadcast_server_shutdown(5);
    steady_timer wait(ctx_);
    wait.expires_after(5s);
    co_await wait.async_wait(use_awaitable);
    
    // 3. 强制保存所有玩家数据
    std::vector<awaitable<void>> save_tasks;
    for (auto& [rid, room] : rooms_) {
        save_tasks.push_back(room->force_save_all());
    }
    for (auto& task : save_tasks) {
        co_await std::move(task);
    }
    
    // 4. 断开所有连接
    for (auto& [sid, session] : sessions_) {
        session->disconnect("server shutdown");
    }
    
    LOG_INFO("shutdown complete");
    ctx_.stop();
}

 


八、一张架构总览

io_context(net)          io_context(room)       io_context(db)
                    │                        │                       │
              ┌─────▼─────┐           ┌──────▼──────┐        ┌──────▼──────┐
              │  Server   │           │    Room     │        │  DB Clients │
              │  accept   │           │  (strand)   │        │Redis/MySQL  │
              └─────┬─────┘           └──────┬──────┘        └─────────────┘
                    │                        │
              ┌─────▼─────┐    消息          │  co_await(冷路径)
              │  Session  │  post到strand ──►│
              │ read_loop │                  │  Tick 纯内存(热路径)
              │ send()    │◄── 广播 ─────────┤
              └───────────┘                  │  定时存档(co_await DB)
                                            └──────────────────────

 

核心原则只有一条:冷路径用协程等 DB,热路径(Tick)只摸内存,所有线程安全靠 Strand 而不是锁。这套结构能稳定支撑单节点数万并发,横向扩展只需要加 Room 节点。

发表回复