使用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 节点。
