文章

高性能C++服务架构分析

本文拆解一个C++游戏任务服务的高并发架构,涵盖epoll多Reactor、四层异步流水线、反向索引过滤、Redis分布式锁与冷热分离等核心设计,适合后端工程师学习高并发实践。

高性能C++服务架构分析

摘要:本文档从高性能架构视角深度拆解 xxxStateMachineService(游戏任务状态机服务)的设计与实现,涵盖服务定位与上下游拓扑、反向索引驱动的请求过滤机制、epoll 多 Reactor 网络模型、四层异步流水线架构、Redis 分布式锁与乐观并发控制、数据存储与序列化策略、冷热分离方案、Pipeline 使用与限制,以及性能监控体系。适合想从成熟项目中学习高并发服务设计的后端工程师阅读。

引言

本文基于一个真实生产项目的游戏任务状态机服务,对其核心架构设计进行复盘分析。所有代码示例和架构描述均源自实际源码,服务名称和业务数据已做脱敏处理。

本文由作者提供素材和源码,经 AI 辅助整理、扩展和润色后成文。作者负责提供架构知识、源码路径并纠正事实性错误,AI 负责组织结构和行文表达。

由于脱敏需要,文中不包含具体的性能基线数据、线上事故案例和团队成员信息。阅读时请关注各层的设计决策和取舍逻辑,而非具体数字的精确性。

特别提醒:任何架构复盘都存在 “后见之明” 偏差——写出来的因果链比实际演化过程更整洁。本文第七章的层间依赖分析是作者在复盘时梳理的逻辑关系,真实系统在多年迭代中形成的依赖更接近网状,建议读者将其理解为脉络梳理而非绝对推导。

一、业务概述:这个服务在做什么

1.1 服务定位

这个服务是 任务系统的核心计算引擎 ,它接收用户行为数据变化(如 “玩家充值了100元”),判断是否触发了某个任务的某个条件,如果满足就推进任务状态、发放奖励、推送进度通知。

1.2 业务模型

1
2
3
4
一个活动任务 (HFSM)
  ├─ 子任务1 (FSM): "登录1天"      → 状态: 待领取→进行中→待领奖→完成
  ├─ 子任务2 (FSM): "对战3局"      → 状态: 待领取→进行中→待领奖→完成
  └─ 子任务3 (FSM): "充值100元"    → 状态: 待领取→进行中→待领奖→完成

每个子任务是一个独立的状态机,包含 6 种状态:

1
默认(0) → 待领取(1) → 进行中(2) → 待领奖(3) → 完成(4)/放弃(5)/超时(6)

每个状态下可配置 规则动作组 (RuleActGroup)

GroupType 含义 配置示例
1 领取规则 “玩家登录” → 创建子任务
2 完成规则 “积分≥100” → 标记为待领奖
3 领奖规则 “玩家点击领奖” → 发奖
4 超时规则 “超过7天” → 标记为超时
5 放弃规则 “玩家点击放弃” → 标记为放弃
6 循环规则 “完成后重置” → 回到待领取

子任务的 6 个状态各自实现为独立的状态类,新增状态只需新增子类。规则和动作各自抽象为接口,运营可任意组合”什么条件触发什么动作”。

以上是业务层面的抽象模型。接下来从网络拓扑角度,看本服务在整体系统中的位置。

二、服务拓扑:上游与下游

2.1 上游调用方

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
                    ┌─────────────────────────┐
                    │   TKxxxStateMachine     │
                    │       Service           │
    ────────────────│  (本服务 = Logic 层)    │───────────────►
                    │                         │
                    └─────────────────────────┘
         ▲                      ▲                      ▲
         │                      │                      │
    ┌────┴────┐           ┌─────┴──────┐         ┌─────┴──────┐
    │xxxBroker│           │xxxDisplay  │         │ xxxConfig  │
    │(对接层) │           │(展示层)     │         │ (配置层)   │
    └─────────┘           └────────────┘         └────────────┘
    推送DAO数据变化         查询任务进度列表         任务配置推送
    人工操作请求            用户注册信息同步         轨迹配置同步
    采集单笔充值                                   迁移指令
  • xxxBroker(对接层):通用数据管道,将 DAT(Data Access Tier,数据访问层,内部代号 MGW)的 DAO 数据变化全量推送过来
  • xxxDisplay(展示层):游戏大厅/活动中心查询用户当前任务状态
  • xxxConfig(配置层):运营配置任务后推送配置更新

2.2 下游依赖方

1
2
3
4
5
6
7
8
9
10
                    ┌─────────────────────────┐
                    │   TKxxxStateMachine     │
                    │       Service           │
                    └───────────┬─────────────┘
            ┌─────────┬─────────┼─────────┬─────────┐
            ▼         ▼         ▼         ▼         ▼
       ┌────────┐ ┌────────┐ ┌──────┐ ┌──────┐ ┌─────────┐
       │Broker  │ │Display │ │ ECA  │ │Detail│ │xxxConfig│
       │进度推送 │ │进度推送│  │ 发奖 │ │ 日志 │ │ 轨迹追踪 │
       └────────┘ └────────┘ └──────┘ └──────┘ └─────────┘
  • xxxBroker/Display:推送任务进度变化,让客户端能实时刷新
  • ECABroker:执行 ECA(Event-Condition-Action)发奖方案
  • DetailSrv:写入 MySQL 操作详情日志
  • xxxConfig:发送轨迹追踪数据供运维排查

2.3 连接池管理

1
2
3
4
5
6
7
8
// TKxxxStateMachineService.h 中每个下游维护独立连接池
CDisSockConnPool m_csConfigSrvCP;   // 到 xxxConfig
CDisSockConnPool m_csBrokerSrvCP;   // 到 xxxBroker
CDisSockConnPool m_csDisplaySrvCP;  // 到 xxxDisplay
CDisSockConnPool m_csECAConfigCP;   // 到 ECA配置服务
CDisSockConnPool m_csECABrokerCP;   // 到 ECA发放服务
CDisSockConnPool m_csDetailSrvCP;   // 到 详情服务
CDisSockConnPool m_csMSGBrokerCP;   // 到 消息服务

每个下游独立连接池,支持连接复用和故障隔离。

2.4 分层职责设计

本节说明 DAT、Broker、StateMachine 三层之间的职责划分。StateMachine 之所以收到全量推送,根本原因在于上游两层都不感知业务。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
游戏服务器 (GS)
    │
    │  "玩家123充值了100金币"、"玩家456登录了"、"玩家789对战了一局" …
    │  每秒可能有海量数据变化事件
    ▼
┌──────────────┐
│  DAT 数据层  │  ← 负责存储所有用户的所有属性值(积分、金币、道具…)
│   (MGW)      │     这是一个通用数据存储服务,不知道"任务"的存在
└──────┬───────┘
       │
       │  "某个属性的值变化了" ,广播推送出去
       │  每条推送 = (用户ID, DtID=数据类型, AtID=账户类型, OriID=资源ID, 新值, 变化值)
       │  这些推送是全量的、无差别的,DAT 不关心谁在消费
       ▼
┌──────────────┐
│  xxxBroker   │  ← 对接层,负责接收 DAT 的推送,再分发给下游消费者
│   (对接层)    │     Broker 同样不关心任务配置——它的职责仅仅是"广播"
└──────┬───────┘
       │
       │  全量推送:所有属性变化都发给 StateMachine
       ▼
┌──────────────────────┐
│  TKxxxStateMachine   │  ← 本服务。在这里才真正知道"哪个任务关心哪个属性"
│     (任务引擎)        │
└──────────────────────┘

核心设计理念:DAT 和 Broker 都是通用的、无业务感知的数据管道。它们不关心 “哪个属性变化对应哪个任务条件”。这种分层职责划分的好处是:

  1. DAT/MGW 只做存储和推送,逻辑简单,性能极高
  2. Broker 只做数据路由和分发,同样不关心业务
  3. 所有的 “业务智能” 集中在 StateMachine 这一个服务

代价就是 StateMachine 会收到全量推送,必须自己在内部做高效过滤。

具体做法是反向索引(倒排表),将 (DtID, AtID, OriID) → [HFSMID] 的映射预加载到内存,推送到达时 O(1) 查找命中哪些任务。详细机制见 3.4 节

从全链路看,这本质上是发布-订阅模式——DAT 作为发布者广播属性变化,Broker 充当消息代理,StateMachine 是订阅者,通过反向索引筛选自己关心的消息。


三、四层异步流水线

在剥离外部依赖(不连接 Redis、MongoDB、下游服务)的简化条件下,裸框架吞吐可达数十万 QPS 量级(具体数值受硬件规格和消息大小影响,此处仅作量级参考)。以下分析其架构设计如何支撑这一并发能力。

3.1 整体架构

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
 ┌─────────────────────────────────────────────────────────────┐
 │                    四层异步流水线架构                         │
 │                                                             │
 │  第1层: 网络I/O层 (epoll 多 Reactor)                         │
 │  ┌─────────────┐   ┌─────────────┐   ┌─────────────┐       │
 │  │ Connecter-1 │   │ Connecter-2 │...│ Connecter-N │       │
 │  │ epoll_wait  │   │ epoll_wait  │   │ epoll_wait  │       │
 │  │ 收包 + 解析  │   │ 收包 + 解析  │   │ 收包 + 解析  │       │
 │  └──────┬──────┘   └──────┬──────┘   └──────┬──────┘       │
 │         │                 │                 │               │
 │         └──────────┬──────┴─────────────────┘               │
 │                    │ TKMSGBIND_ASYNC 投递                    │
 │                    ▼                                        │
 │  第2层: Handler 消息队列 (CTKAsyncMsgFunctor)                │
 │  ┌──────────────┐ ┌──────────────┐ ┌──────────────┐       │
 │  │asynUserOpe   │ │asynBrokerGet │ │asynGetHFSMCA │        │
 │  │ N线程        │ │ N线程         │ │ 少量线程     │       │
 │  └──────┬───────┘ └──────┬───────┘ └──────┬───────┘       │
 │         │                │                │                │
 │         └────────┬───────┴────────────────┘                │
 │                  │ AddWorkData                              │
 │                  ▼                                          │
 │  第3层: 业务处理层 (CTKAsyncWorkObject)                       │
 │  ┌──────────────┐ ┌──────────────┐ ┌──────────────┐       │
 │  │asynDoProcess │ │asynPushData  │ │asynPushProc  │        │
 │  │ N线程        │ │ M线程         │ │ M线程        │       │
 │  └──────┬───────┘ └──────┬───────┘ └──────┬───────┘       │
 │         │                │                │                │
 │         └────────┬───────┴────────────────┘                │
 │                  ▼                                         │
 │  第4层: 存储访问层(连接池 — 为上层提供数据读写能力)          │
 │  ┌──────────────┐ ┌──────────────┐ ┌──────────────┐       │
 │  │Redis 连接池  │ │MongoDB连接池  │ │MySQL 连接池   │       │
 │  │按需伸缩      │ │              │ │              │       │
 │  └──────────────┘ └──────────────┘ └──────────────┘       │
 └─────────────────────────────────────────────────────────────┘

四层的职责划分是严格的

  • 第 1 层只做网络 I/O,不做业务 — 收包、解析、ACK、投递
  • 第 2 层做请求分发和泳道隔离 — 按 UID 哈希分流,同用户串行
  • 第 3 层做业务计算 — 状态机、规则引擎、动作执行
  • 第 4 层是基础设施层 — 连接池管理、数据读写

每一层只依赖于直接相邻的下一层,不跨层调用。

3.2 第1层:网络I/O层

关于 epoll/IOCP 的原理、触发模式、Reactor 模式定义等基础知识,详见 IO多路复用 epoll & IOCP。本节聚焦本服务如何使用这套框架。

3.2.1 框架如何对接 epoll

框架层通过条件编译屏蔽平台差异,Linux 下使用 epoll、Windows 下使用 IOCP。业务代码不感知底层实现,通过 MessageHandler::Handle() 接收处理好的消息即可。

框架的层次关系:

1
2
3
4
5
6
7
8
9
10
业务代码
  Handler::Handle(pMsg)                ← 拿到解好的消息,只关心业务
────────────────────────────────
框架层(C++)
  SockServer / ConnecterImpl           ← 封装 epoll_wait 循环 + 事件分发
  MessageHandler                       ← 按消息类型路由到具体 Handler
  ObjectPool<SockConn>                 ← 连接对象复用
────────────────────────────────
系统调用
  epoll_create / epoll_ctl / epoll_wait ← Linux 内核接口

业务开发者只需 service.SockServer::Start() 一行,框架内部自动完成 epoll_create、epoll_ctl 注册、epoll_wait 循环等所有细节。

3.2.2 多 Reactor 架构

本服务使用多 Reactor 模型,即多个 epoll 事件循环线程并行运行:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
┌─────────────────────────────────────────────────────────────┐
│                    多 Reactor 架构                           │
│                                                             │
│  Accepter 线程 (专门 accept 的 Reactor)                      │
│  ┌─────────────────────────────┐                             │
│  │ while:                      │                             │
│  │   conn = accept(listen_fd)  │  只做一件事:接受新连接       │
│  │   ConnPool.assign(conn)     │  然后分配给空闲的 Connecter  │
│  └─────────────────────────────┘                             │
│         │                                                   │
│         └────── 分配到最空闲的 Connecter ──────┐              │
│                                               ▼              │
│  Connecter-1 (Reactor)          Connecter-2 (Reactor)       │
│  ┌──────────────────────┐      ┌──────────────────────┐     │
│  │ while:               │      │ while:               │     │
│  │  events = epoll_wait │      │  events = epoll_wait │     │
│  │  for (each event):   │      │  for (each event):   │     │
│  │   if (可读)          │      │   if (可读)           │     │
│  │    RecvNetData()     │      │    RecvNetData()     │     │
│  │    Handle(pMsg)      │      │    Handle(pMsg)      │     │
│  │   if (可写)           │      │   if (可写)          │     │
│  │    SendNetData()     │      │    SendNetData()     │     │
│  └──────────────────────┘      └──────────────────────┘     │
│   管理 5000 个连接               管理 5000 个连接             │
└─────────────────────────────────────────────────────────────┘

为什么用多个 Reactor 而不是一个?

  • 单 Reactor:1 个线程处理 10000 个连接,单核 CPU 是瓶颈
  • 多 Reactor:N 个线程各处理 10000/N 个连接,利用多核 CPU
  • 这里的 N 通常等于或略大于 CPU 核数

3.2.3 消息在 epoll 线程内的路径

1
2
3
4
5
// service/TKxxxStateMachineMain.cpp
CTKxxxStateMachineHandler handler;
CTKxxxStateMachineService service(&handler);  // 继承自 SockServer
service.Init();
service.SockServer::Start();  // 启动 Accepter + N 个 Connecter

SockServer 内部持有 Accepter、Connecter 和 ObjectPool(各组件的详细实现见 IO多路复用 epoll & IOCP 第八章),业务开发者只需 SockServer::Start() 一行启动全部线程。

消息在 epoll 线程内的路径

1
2
3
4
5
6
epoll_wait 返回就绪事件
  → RecvNetData()        // 从 socket 读数据到 m_pRecvBuf
  → 解析 TKHEADER        // 获取 dwType (消息类型)、dwSerial (追踪号)
  → Handle(pMsg, conn)   // 分发到 Handler
     → 快速路径:OnReqPushData → AddWorkData 投递 + 立即 SendMsg(ACK)
     → 慢速路径:OnReqUserOperationSyn → 同步执行后 SendMsg 回复

核心原则:epoll 线程绝不做耗时操作。 收到请求后要么立即回 ACK,要么投递到异步队列。这保证了每个 epoll 循环的延迟可控(通常 < 1ms)。

3.3 第2层:Handler 队列层

消息到达后按来源→消息类型逐级分发,每一级只处理自己关心的类型。以下展开泳道路由和两种处理路径。

3.3.1 泳道隔离

通过对 UID 做哈希,将同一个用户的请求永远路由到同一个队列:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
                    epoll 线程投递消息
                          │
                    ┌─────▼─────┐
                    │ Hash(UID) │  → 同 UID 消息固定落同一队列
                    └─────┬─────┘
          ┌────────┬──────┼──────┬────────┐
          ▼        ▼      ▼      ▼        ▼
        Queue-0  Queue-1  ...  Queue-N-1  Queue-N
        ┌──────┐┌──────┐     ┌──────┐  ┌──────┐
        │UID=3 ││UID=7 │     │UID=5 │  │UID=1 │
        │UID=9 ││UID=13│     │UID=11│  │UID=4 │
        └──┬───┘└──┬───┘     └──┬───┘  └──┬───┘
           │       │            │         │
           ▼       ▼            ▼         ▼
        Thread0  Thread1     ThreadN-1  ThreadN

价值:同 UID 的请求串行处理,天然避免了同一个用户的任务状态被并发修改。不同 UID 之间完全并行。这是一种用空间换正确性的设计:不需要在处理逻辑中加细粒度锁,因为同一用户永远不会被两个线程同时处理。

3.3.2 两种消息处理路径

不同消息对响应时机的需求不同,因此设计了两种路径。但不管哪种路径,epoll 线程 都不执行耗时业务。

异步路径(fire-and-forget):适用于 DAO 推送、人工操作等场景。epoll 线程收到消息后立即回 ACK,然后投递到队列异步处理,处理完不再回复上游。

1
2
3
4
5
6
7
8
9
10
11
12
13
// 快速投递 + 立即 ACK — OnReqPushData
int CTKxxxStateMachineHandler::OnReqPushData(PTKHEADER pMsg, SockConn *pSockConn)
{
    // 数据拷贝到队列,立即返回
    g_tkxxxService->m_asynPushData.AddWorkData(pMsg, pMsg->dwLength + sizeof(TKHEADER));

    // 在 epoll 线程内回复 ACK——上游不用等任务计算
    TKHEADER header = {0};
    header.dwType = TK_ACK | pMsg->dwType;
    header.dwParam = TK_ACKRESULT_SUCCESS;
    SendMsg(pSockConn, &header);
    return 0;
}

上游只需确认消息已送达,快速 ACK 使得整条链路不被最慢的业务操作拖累。

同步路径(request-response):适用于查询任务进度、获取统计数据等场景。epoll 线程不回复,投递到队列后立即释放;队列线程完成业务计算后,通过同一个连接回包。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
// 查询任务进度 — OnReqProcessByHFSMID
int CxxxStateMachineHandler::OnReqProcessByHFSMID(PTKHEADER pMsg, SockConn *pSockConn)
{
    // ① 消息被路由到泳道队列,由队列线程执行,当前位于泳道线程中

    // ② 业务逻辑(查询Redis → 规则计算 → 序列化)
    CxxBuffer bufHFSMProc, bufSessionCache;
    g_xxxService->GetSingleHFSMProc(bufHFSMProc, bufSessionCache, ...);

    // ③ 组装响应,由队列线程通过原连接 SendMsg 返回
    xxACK_GET_USR_PROGRESS stAck = {0};
    stAck.header.dwType = TK_ACK | pMsg->dwType;
    bufAck.Append(&stAck, sizeof(stAck));
    bufAck.Append(bufHFSMProc);
    SendMsg(pSockConn, (PTKHEADER)bufAck.GetBufPtr());
}

3.3.3 SendMsg 的跨线程实现

两种路径都涉及在不同线程中调用 SendMsg:异步 ACK 在 epoll 线程,同步回包在队列线程。它们能安全共存的原因是:SendMsg 并不直接操作 socket。

SendMsg 内部通过 SockConn::AddMsg() 把待发送数据放入该连接的发送队列。真正的 send() 系统调用由 epoll 线程在下次 epoll_wait 返回可写事件时统一执行。业务线程调用 SendMsg 只是入队,不会并发写 socket。

这也是为什么同步路径不需要让 epoll 线程等待:查询操作涉及 Redis 读取和规则计算(5-20ms),如果放在 epoll 线程同步执行,该 Reactor 下的所有其他连接都会被阻塞。投递到队列线程后,epoll 线程立即释放继续服务其他连接,业务完成后再通过 SendMsg 入队回包,epoll 线程负责最终发出。

3.3.4 反压保护

每个队列通过 InitWithMaxWaitCnt(MAX_QUEUE_SIZE) 设置最大等待数。当队列积压超过上限时,新的请求会被丢弃,防止内存无限增长导致 OOM。被丢弃的消息由上游 Broker 的定期重推机制或定时任务从 DAT 回源补偿,不会永久丢失。

3.4 第3层:业务处理层

第三层是服务的核心——状态机计算引擎。

前面两层(网络 I/O、泳道队列)把所有消息送到了这里,接下来要回答的问题是:一条推送怎么从 “某个属性变了” 变成 “某个任务进度推进了、奖励发放了、通知推送了”。

3.4.1 一条推送的完整旅程

数据来源详见 2.4 节:每条推送包含 (DtID=数据类型, AtID=账户类型, OriID=资源ID, 新值)。接下来分三步处理。

第一步:反向索引 —— 哪些任务关心这个属性?

“反向索引” 是相对 “正向” 说的。正向是 “一个任务配置了哪些属性作为条件”(任务→属性),反向是 “一个属性的变化要通知哪些任务”(属性→任务)。本服务在启动时从配置库加载全部任务配置,预构建 (DtID, AtID, OriID) → [HFSMID...] 的哈希表常驻内存。推送到达时直接查表:

1
2
3
CTKFixList<DWORD> *pfxlstHFSMID = nullptr;
if (!DicMgr().GetHFSMIDByDAO(dwDtID, dwAtID, dwOriID, pfxlstHFSMID))
    continue;  // 没有任何任务关心这个属性 —— 大多数推送在此终止

每个属性通常只被 0~3 个任务配置为条件,大多数推送直接返回空列表。如果不用反向索引而是遍历全部任务(可能有上万条配置),每次推送都是 O(任务总数)。DicMgr 之外还有一个 TKDataRequestRouter,是多数据源(RTC 远程调用 + LocalStore 本地存储)的统一代理,内部根据 DtID 路由到实际数据源。

第二步:过滤 —— 命中的任务中,当前用户是否匹配?

反向索引告诉你 “哪些任务关心这个属性”,但未必每个任务都适用于当前用户。以下过滤在创建 HFSM 对象之前完成——不需要查 Redis,只靠内存中的静态配置做判断:

  • 特殊类型:导名单推送、单笔充值等不是通用 DAO 属性变化,走自己的处理路径,不进后续计算
  • AB 实验:当前用户是否在该任务的实验组?
  • 静态预校验:推送的新值是否在任务配置的条件范围内?例如任务要求 “积分≥100”,但当前推送的是 “积分+10” 且已知当前值=50——只比对静态配置和推送值就能判断不满足,跳过
  • 注册信息:用户的平台/渠道是否匹配任务的目标人群?例如任务只对微信区开放,手Q区用户直接跳过

过滤掉一个候选任务,省下后续的对象创建、Redis 读取和规则计算全部开销。

第三步:执行计算。

通过所有过滤的候选任务进入实际处理——创建 HFSM 对象,执行 OnProcess

1
2
CTKxxxHFSMBase *pHFSM = CreateHFSM(dwUID, dwHFSMID);  // 按模型类型分发给 5 种子类
pHFSM->OnProcess(true);  // LoadData → ProcessRule → Lock → ReLoadData → Process

OnProcess 内部的核心流程详见 4.3 节。计算完成后对象即销毁,不缓存、不跨请求复用。

第三条路是 “查询”:用户登录大厅时客户端拉取所有任务进度,不走推送链路,而是直接走 GetHFSMProc 接口(见 3.3.2 的同步路径),内部同样调 OnProcess 做规则校验,但增加了会话缓存优化(见 3.4.2)。

3.4.2 会话缓存:减少对 DAT 的重复查询

会话缓存只在批量查询场景下生效:用户登录大厅时,客户端调 GetHFSMProc 按 Class/Area 拉取该用户所有任务的进度,假设 300 个任务中有 100 个的规则条件依赖 “当前积分” 这个 DAO 属性值。

没有缓存时,100 个 HFSM 各自调 GetMetaData(dwDtID, 3积分, 0, 10001, &value),产生 100 次对 DAT 的远程调用。有会话缓存时:

  • 服务入口处创建一个栈上的 sessionCache 对象
  • 遍历 300 个任务时,每个 HFSM 通过 SetSessionCache(&cache) 共享同一块缓存
  • HFSM 内部查属性时先查缓存:命中直接返回;未命中走 DAT 查询,拿到值后回写缓存
  • 效果:同一属性一个请求内只查 DAT 一次,后续 HFSM 全部内存命中

缓存的生命周期是一次请求内——入口创建,遍历结束后函数返回时自动析构。不跨请求复用,因为两次请求之间用户属性值可能已经变化(另一个活动刚扣了积分),跨请求缓存会造成脏读。

注意:会话缓存替代的是对 DAT 的远程调用,不是对本服务 Redis 的调用。每个 HFSM 的 LoadData(加载任务状态)仍然要走本服务 Redis——缓存的是用户属性值,不是任务进度数据。

优化天花板:DAT 调用次数从 N 压到 1 之后,再往下优化需要要么 DAT 侧提供批量查询接口,要么引入跨实例共享的属性缓存层(如 Redis),都是跨系统的大改动。在当前架构约束下,会话缓存已经是针对这个场景的最优解。

3.4.3 进度通知合并

同一个 HFSM 可能在短时间内经历多次状态变化。典型场景是充值活动——用户一次充 500 元,上游连续推送三条数据变化通知,依次触发子任务 1、2、3 完成。

每条推送处理时,OnProcess 独立完成计算并将结果写入 Redis。三次 OnProcess 各自写入一次:第一条处理完 Redis 里子任务 1 标记为完成,第二条处理完子任务 2 标记为完成,第三条处理完子任务 3 标记为完成。计算和持久化不受合并影响。

OnProcess 完成后调用 PushProcInfo,向下游 Broker/Display 发送进度更新通知。PushProcInfo 内部通过 IsUpdateStatePushProc() 判断是否实际发出:

1
bool IsUpdateStatePushProc();  // true: 发出;false: 跳过

返回 false 时跳过本次通知。下一次 OnProcessLoadData 从 Redis 读到的是累积后的全量状态,完成计算后再次调用 PushProcInfo,此时推送的就是合并后的最新进度。下游不需要知道中间经历了几个步骤。

具体到上例:子任务 1 完成时跳过,子任务 2 完成时再次跳过,子任务 3 完成后发出一次通知,推送内容为三个子任务均已完成。下游 Broker/Display 的消息量从 3 条降为 1 条,在大型活动并发高峰期收益显著。

3.5 第4层:存储访问层

第 4 层不是独立的处理线程,而是 被第 3 层共享调用的基础设施。服务通过三类连接池访问外部存储:

1
2
3
CTKRedisConnPool m_HFSMRedisCP;        // Redis 连接池(含驱动)
CTKMongoConnectionPool m_MongoPoolHFSM; // MongoDB 连接池(含驱动)
CTKMySqlConnectionPool m_mysqlCPDetail; // MySQL 连接池(含驱动)

这些连接池属于对象池模式——高频创建/销毁的连接对象被预分配和复用,配合 SockServer 中的 ObjectPool<SockConn>,服务中所有关键对象均由池化管理。连接池同时承担数据库驱动和连接管理的双重角色。

业务代码通过它们直接执行数据库命令,例如分布式锁通过 HFSMRedisCP().RedisCommandWithReply("set key val ex 2 nx") 一行完成,不需要手动管理 TCP 连接。连接池内部提供:

  • 连接复用:多个业务线程共享同一批连接,借出/归还,不重复建连
  • 心跳保活:空闲超过阈值的连接自动 ping,防止被中间网络设备断开
  • 慢查询监控:超过阈值的操作自动记录日志,便于排查性能问题
  • 连接池伸缩:按 min/max 配置按需创建和释放

各级并发度配置:

层级 队列名 线程数 功能
第2层 asynUserOpe ~N 用户手动操作(领奖/放弃)
第2层 asynBrokerGetHFSM ~N Broker 拉取任务进度
第3层 asynDoProcess ~N 异步状态机计算
第3层 asynPushData ~M DAO 推送驱动状态机
第3层 asynPushProc ~M 异步推送进度到下游
第3层 asynWriteDetail 少量 写详情日志到 MySQL
第3层 asynReissueAward ~N 补发奖励

以上仅列举核心队列。实际服务中还有存钱罐(SaveBox)异步处理、群体协作任务状态推送、轨迹发送、模拟 DAO 推送等 30+ 个异步队列,各自负责独立的业务子域,按需配置线程数。

四、分布式锁与乐观并发控制

第三章的四层流水线解决了消息的接收、分发和计算调度问题。但进入实际计算后,还需要解决 并发安全:泳道保证同一用户的消息串行处理,但 Redis 中的数据可能被多实例同时修改,一旦状态推进触发了奖励发放,重复执行将导致严重事故。本章讨论这个问题的架构解决方案。

4.1 并发冲突场景

1
2
3
4
5
同一用户、同一任务,两条推送几乎同时到达:
  线程A: 推送"积分+100" → 读状态(积分=0) → 判断满足条件(100≥100) → 推进状态 → 发奖
  线程B: 推送"积分+50"  → 读状态(积分=0) → 判断满足条件(50≥50)  → 推进状态 → 再次发奖!
    ↓
  可能重复发奖

泳道隔离只能保证同一实例内同一用户串行。在多实例部署下,同一个用户的请求可能落在不同实例上——需要跨实例的互斥机制。

4.2 分布式锁:Redis SET NX EX

锁的粒度设计为 {HFSMID} + {UID}——不同用户、不同任务之间完全无锁争用,只有”同一用户 + 同一任务”才需要排队:

1
2
3
4
5
6
bool CTKxxxHFSMBase::Lock(int nSecond = 2)
{
    redisReply *pReply = g_tkxxxService->HFSMRedisCP().RedisCommandWithReply(
        "set xxx_HFSM_L_%u_%u 1 ex %d nx", GetID(), m_dwUID, nSecond);
    return (pReply != NULL && pReply->type != REDIS_REPLY_NIL);
}

基于 Redis 的 SET key value NX EX seconds 原子指令——NX 保证只有第一个调用者成功,EX 2 秒自动过期防止进程崩溃导致死锁。服务还有一个通用的 BizLock(const string &strBizKey, int nSecond = 3),用于存钱罐等非 HFSM 场景。

4.3 乐观并发控制:先校验、后加锁

OnProcess 的设计遵循乐观并发控制——大多数推送并不满足任何规则条件,无需进入锁保护区:

1
2
3
4
5
6
7
8
9
10
11
12
13
bool CTKxxxHFSMBase::OnProcess(bool bIsClearCache)
{
    LoadData();           // ① 无锁加载任务状态
    ProcessRule();        // ② 无锁预计算 —— 大多数推送在此直接 return

    if (!Lock()) {        // ③ 只有规则确实满足,才抢锁
        AddReProcHFSM(this);  // 抢锁失败加入重试队列
        return false;
    }
    ReLoadData();         // ④ 锁内重读,防止等待期间数据被其他线程修改
    Process(stProcIn);    // ⑤ 执行状态切换和动作
    Unlock();
}

场景 A:收到 “积分+1” 推送,任务要求 “积分≥100” ——LoadData → ProcessRule → 不满足 → 直接 return,全程不抢锁。

场景 B:收到 “积分+100” 推送,任务要求 “积分≥100” ——LoadData → ProcessRule → 满足 → Lock → ReLoadData → Process。多一次 LoadData 的代价远小于每次推送都抢一次分布式锁。

如果反过来设计(先锁后读),绝大部分推送白白抢了一次锁——日常场景中满足规则条件的推送比例极低——既增加了 Redis 调用量,又延长了其他线程的等待时间。

4.4 重试队列

抢锁失败说明同一任务正在被其他线程处理,当前推送的上下文被序列化后加入超时队列,等待定时重试:

1
2
3
4
5
6
7
8
void CTKxxxStateMachineService::AddReProcHFSM(CTKxxxHFSMBase *pHFSM)
{
    xxxReProc *pReProc = new xxxReProc;
    pReProc->dwUID = pHFSM->GetUID();
    pReProc->dwHFSMID = pHFSM->GetID();
    m_pTCReProcHFSM->Regist(
        new TimeoutEvent(m_dwReProcTime, TIMEOUT_EVENT_ONCE, ReProcHFSM, pReProc));
}

队列底层是一个最小堆(堆顶是最近到期的事件),O(log n) 插入、O(1) 取到期事件。

三种重试队列:

队列 用途
m_pTCReProcHFSM 抢锁失败后重试触发任务计算
m_pTCReReissueHFSM 抢锁失败后重试补发奖励
m_pTCReModSaveBox 抢锁失败后重试存钱罐操作

存钱罐(SaveBox)是本服务中一个独立的子系统,负责管理可积累型奖励的存储和发放(如”活动期间累计登录 7 天领取奖励”)。它与任务状态机共享 Redis 锁机制,但有自己独立的 BizLock、异步队列和重试队列,代码中的 m_asynModDAOSaveBoxm_asynSaveBoxExcuteBill 等均属于该子系统。

五、数据存储策略

第四章程讨论了计算过程中的并发安全,本章讨论数据层面——状态数据怎么存、怎么查、怎么省。

5.1 Key 体系

Redis 被选择的核心原因:原子操作(分布式锁必需)、纯内存读写延迟可忽略、Pipeline 支持批量命令合并。本服务在 Redis 中的五类 Key——每一类对应一个明确的数据需求:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
xxx_HFSM_Pro_{HFSMID}_{UID}     → KV     → JSON 任务状态数据(核心)
   每个用户每个任务的当前进度——状态、循环次数、自定义元数据。每次 OnProcess 都要读写。
   
xxx_HFSM_L_{HFSMID}_{UID}       → KV     → "1"(分布式锁,2 秒过期)
   第四章的 HFSM 锁,{HFSMID}+{UID} 粒度,基于 SET NX EX 实现。
   
xxx_HFSM_Stat_{HFSMID}          → Hash   → 各状态完成人数统计
   运营需要知道"多少人完成了这个任务",Hash field 为状态编号,value 为计数值。
   
xxx_ReissueAward_{HFSMID}_{tag} → Set    → 待补发用户集合
   活动结束后自动补发奖励时使用,Set 存储待处理的 UID 列表。
   
xxx_ColdSink_Tag_Cluster{ID}    → Bitmap → 冷数据下沉标记
   0=热数据在 Redis,1=已下沉到 MongoDB。1 bit/用户,1 亿用户仅需 12.5MB。

5.2 Pipeline:批量查询优化

Pipeline 的使用场景是用户进入大厅时的任务列表查询——Display 服务调用本服务的 GetHFSMProc 接口,按 Class/Area 一次性拉取该用户所有任务的状态(可能 300 个 HFSM)。300 次独立 GET 就是 300 次 Redis 网络往返,Pipeline 将 N 条命令打包一次发送、一次接收,从 N 次 RTT 降为 1 次:

1
2
3
4
5
CTKxxxRedisReplySet replySet(&HFSMRedisCP());
for (int i = 0; i < fxlstHFSMID.GetCount(); i++)
    replySet.PushRedisPipeCommaindF("get xxx_HFSM_Pro_%u_%u", dwHFSMID, dwUID);
for (int i = 0; i < fxlstHFSMID.GetCount(); i++)
    redisReply *pReply = replySet.PopRedisPipeReply();

300 个 GET 之间互不依赖且目标为同一 Redis 实例——这是 Pipeline 的理想场景。不使用 Pipeline 时,仅加载数据就需要约 150ms(300 × 0.5ms RTT),加上后续 OnProcess 的计算耗时,用户登录体验不可接受。

Pipeline 与 Redis 事务(MULTI/EXEC)不同——Pipeline 只合并网络往返,不保证原子性。使用限制:命令数过多时内存占用和连接独占时间线性增长,单条 bigkey 操作会拖慢整条管道。高频推送(每秒几十万条)不走 Pipeline,原因是 Pipeline 期间连接被独占,而高频推送需要用完立刻归还连接,不阻塞其他请求。本服务仅在按 Class/Area 批量查询时使用,命令数受控(5-300 个),不跨用户合并。

5.3 JSON 全量覆写

状态数据存储在 5.1 的核心 Key(xxx_HFSM_Pro_{HFSMID}_{UID})中,格式为 JSON。每个用户在每个任务下有一份,子任务数量少(3-10 个),体量小(几百字节)。选择全量覆写(SET)而非增量更新(HINCRBY)。全量覆写的安全基础来自两层保护——实例内由第三章的泳道隔离保证同一个用户只有一个线程写入,跨实例由第四章的分布式锁(Lock → ReLoadData → Process → SET)保证同一时刻只有一个实例在写。两层叠加后,不存在两个线程同时 SET 同一份数据的情况,因此不需要先 GET 再 SET 的 Read-Modify-Write 模式。采用缩写模式(字段名压缩为单字母)节省约 30% 存储空间,大型 JSON 启用 zlib 压缩后写入。

5.4 冷热分离:问题背景

Redis 虽然快,但内存成本高。实际需要的集群内存 = 数据量 × 副本因子(主从 2 倍)× 预留因子(AOF 重写、fork 等,1.5~2.0 倍)。任务系统的数据有明确的时间衰减特征——大部分用户超过 7 天未登录后,其任务数据几乎不再被访问。继续将这些数据放在昂贵的内存中是不划算的。

5.5 冷热分离流程

1
2
3
4
5
6
7
8
9
10
11
12
用户请求到达
      │
      ▼
Redis GET: xxx_HFSM_Pro_{UID}_{HFSMID}
      │
      ├── 命中(热数据)──▶ 直接处理
      │
      └── 未命中 ──▶ 检查 BITMAP ColdSink_Tag
                        │
                        ├── 已下沉 ──▶ MongoDB 回捞(上浮)
                        │
                        └── 未下沉 ──▶ 空数据(用户从未触发该任务)

下沉由独立的定时任务驱动:单线程按活跃时间分批扫描(每次处理固定数量用户,避免全量扫描造成 CPU 尖峰),投递到多线程下沉队列并行执行。处理每个用户时先 Lock() 检测是否在线(在线则跳过),写入 MongoDB 后回读校验确认一致,再删除 Redis 数据并标记 BITMAP。

5.6 冷热分离设计评价

优点

  1. 成本收益显著:Redis(内存)→ MongoDB(磁盘),存储成本降低 5-10 倍
  2. 对在线用户透明:下沉前 Lock(),在线用户持有锁时自动跳过
  3. 回读校验:写入 MongoDB 后立即回读比对,确保数据正确后才删除 Redis 中的数据
  4. BITMAP 高效标记:1 bit/用户,海量用户仅需极少内存

缺点

  1. 回捞延迟高:从 MongoDB 回读冷数据的延迟比 Redis 高 1-2 个数量级(ms → 10-100ms)
  2. 数据一致性窗口:Redis 删除和 BITMAP 设置不是原子的(影响可控,最多导致重复下沉)
  3. 额外的运维复杂度:需要管理 Redis 和 MongoDB 两套存储,监控下沉/上浮的准确性
  4. 与业务强耦合:下沉逻辑依赖”用户是否在线”的判断(通过 Lock 和注册信息)

与业界常见冷热方案对比

维度 本服务方案 Redis + SSD 纯 Redis 集群
存储成本
回捞延迟 高(10-100ms) 中(1-5ms) 低(<1ms)
实现复杂度 高(自研) 低(Redis Enterprise) 极低
适用场景 冷数据极少被访问 冷数据偶尔访问 所有数据频繁访问

对本服务的适用场景而言,任务数据有明确的冷热边界(7 天未登录 = 冷),冷数据极少回捞(不活跃用户几乎不回归),回捞延迟对回归用户体验影响很小。方案的代价(复杂度、双存储运维)与收益(内存成本大幅降低)匹配。

六、性能监控体系

6.1 三层监控

1
2
3
层1: CTKMsgCapture  — 性能埋点(50+ 个埋点,微秒级耗时):覆盖 Redis 查询、MongoDB 操作、状态机计算等所有关键路径
层2: CTKMsgStat     — 消息级统计(按 ClassID/HFSMID 分维度,毫秒分布直方图):按任务和主题维度聚合,定位热点
层3: 队列积压监控    — 每分钟打印所有异步队列的待处理数:提前发现流量堆积,触发扩容或限流

6.2 CAutoCapture:RAII 零侵入埋点

1
2
3
4
5
6
// 构造函数记录开始时间,析构函数自动计算耗时并上报
bool SomeFunction() {
    CAutoCapture cap(CAPTYPE_REDIS_GET, *g_pCapture);
    // ... 业务逻辑 ...
    return true;  // 析构时自动 Capture 耗时(微秒)
}

6.3 关键埋点

埋点 含义
Redis 单任务查询 查询单个任务状态的耗时
Redis Pipeline 批量查询 批量查询任务列表的耗时
MongoDB 冷数据回捞 从 MongoDB 回读冷数据的耗时
数据下沉 将 Redis 数据写入 MongoDB 的耗时
下沉抢锁失败 下沉时 Lock 失败次数(反映在线冲突)
下沉校验失败 下沉后回读校验失败次数(数据一致性告警)
异步状态机计算 状态机规则判断和状态切换耗时
异步推送进度 向下游推送任务进度变化的耗时
重试次数 抢锁失败后重试的次数(反映锁竞争程度)
冷数据回捞 冷数据从 MongoDB 回到 Redis 的次数

6.4 全功能开关 + 旁路验证

每个重要特性都配有独立开关(如冷热分离 m_nColdDataSinkSwitch、数据缩写 m_nSaveDataAbbrSwitch、空数据存储 m_nSaveEmptyFlagSwitch 等)。新功能上线时先以旁路模式运行(如 m_nColdDataSinkType=0 仅验证不删数据),确认无误后再切为正式模式。这是一种低风险的生产级灰度发布策略。

七、总结

回顾全文,这个服务的高并发能力不是靠某一项技术单独达成的,而是六个维度逐层叠加、相互配合的结果。以下从三个角度做收束:各层的核心取舍、层与层之间的依赖关系,以及从中可提取的通用设计原则。

7.1 各层的核心取舍

架构设计的核心不是”选最好的技术”,而是在约束条件下做正确的取舍。下表汇总全文六个维度的关键决策:

维度 决策 收益 代价
分层职责(第二章) DAT/Broker 不感知业务,全量推送 上游极简、逻辑集中 本服务承受全量推送,必须自建过滤
反向索引(第三章) 预建 (DtID,AtID,OriID) → [HFSMID] 哈希表 O(1) 过滤,大多数推送直接丢弃 内存占用,配置变更需重建索引
泳道隔离(第三章) UID 哈希 → 固定队列 → 同用户串行 业务代码无需细粒度锁 队列积压时同用户所有操作延迟增加
乐观并发控制(第四章) LoadData → ProcessRule → Lock → ReLoadData 日常场景中绝大多数推送不进锁保护区 满足条件的推送多做一次 LoadData + ProcessRule
JSON 全量覆写(第五章) SET 替代 HINCRBY 逻辑简单,无需 Read-Modify-Write 无法增量更新单字段,小改动也需全量序列化
Pipeline 批量查询(第五章) N 条 GET 合并为一次网络往返 N×RTT → 1×RTT 连接独占期间其他请求等待;不适用于高频推送
冷热分离(第五章) Redis → MongoDB + BITMAP 标记 内存成本降低 5-10 倍 回捞延迟从 <1ms 升至 10-100ms;双存储运维复杂度
全功能开关(第六章) 每个特性独立开关 + 旁路验证 变更风险可控,秒级回滚 每个特性增加开关代码和验证逻辑

这些取舍不是孤立的——一个维度的代价,往往被另一个维度的设计所弥补。下面分析它们之间的协同关系。

7.2 层与层之间的协同

以下是几个重要的跨层依赖关系(部分关系在系统演化中自然形成而非刻意设计):

反向索引依赖分层职责。 第二章确立了”上游不感知业务、全量推送”的约束,第三章的反向索引是对这个约束的直接回应。二者互为前提——没有全量推送就不需要反向索引,没有反向索引则全量推送就是灾难。合在一起,它们实现了一个高效的发布-订阅模型:发布者简单,订阅者智能。

泳道隔离是乐观并发控制的基础。 泳道保证同一实例内同一用户串行处理,将并发冲突的范围从”同一实例内 + 跨实例”缩小为”仅跨实例”。如果没有泳道隔离,每个 OnProcess 都需要分布式锁——乐观并发控制”先校验后加锁”的收益将大打折扣,因为同一实例内两个线程同时处理同一用户的消息时会直接数据竞争。

泳道隔离 + 分布式锁共同保障 JSON 全量覆写安全。 全量覆写(SET)的隐含前提是不存在并发写入者——如果有两个线程同时 SET 同一份 JSON,后者覆盖前者,数据丢失。泳道隔离消除了实例内的并发写入,分布式锁(Lock → ReLoadData → Process → SET)消除了跨实例的并发写入。两层机制叠加后,可以用最简单的 SET 语义,无需 Read-Modify-Write、无需乐观锁 CAS、无需 Redis Watch。

分布式锁的语义被冷热分离复用。 冷数据下沉任务在执行前通过 Lock() 检查用户是否在线——在线用户的 HFSM 锁被持有,下沉自动跳过。分布式锁在这里的语义从”并发互斥”扩展到了”存活检测”,是对已有机制的巧妙复用,而非引入新的检测通道。

epoll 线程不阻塞是整条流水线的第一推动力。 回溯全文,四层异步流水线存在的根本原因是:网络线程如果被业务计算阻塞,该 Reactor 下的所有其他连接都会卡住。这个约束驱动了后续全部设计——因为网络层不能阻塞,所以需要异步投递;因为异步投递,所以需要泳道隔离;因为泳道隔离,所以 JSON 覆写安全;因为 JSON 覆写安全,所以存储逻辑简单。这是一条贯穿全文的因果链。

全功能开关是上述所有机制的保险丝。 反向索引、冷热分离、数据缩写、会话缓存——任何一个优化出问题都可能引发连锁反应。全功能开关的价值不在于”功能本身”,而在于让所有其他维度的变更都有一条安全的退路。

7.3 可迁移的设计原则

抛开本服务的具体业务上下文,以下几条原则对后端服务架构设计具有通用参考价值:

原则一:让上游简单,让下游智能。 DAT 只管存储和推送,Broker 只管路由和分发,所有业务判断收敛到 StateMachine。在分布式系统中,与其让每一层都做一点业务逻辑,不如让数据层和代理层保持纯粹,把复杂性集中到一个点。集中的代价由反向索引等机制消化。

原则二:先过滤、后加锁、最后计算。 全文最核心的处理模式——反向索引过滤(O(1) 查表)→ 静态预校验(内存比对)→ ProcessRule(规则计算)→ Lock(分布式锁)→ ReLoadData + Process(最终执行)。这一链条的设计哲学是:在真正开始昂贵操作之前,用尽量廉价的检查排除尽可能多的候选。从反向索引到静态预校验到 ProcessRule,每一步的过滤成本递增、候选集递减。

原则三:用空间换正确性。 泳道隔离多占了队列和线程资源,BITMAP 冷热标记多占了少量内存,但它们换来了确定性的正确性保证——同用户串行处理不需要分布式锁的细粒度协调、冷热状态判断不需要扫描 MongoDB。在并发系统中,确定性比省资源更重要——省下的几 MB 内存远不如一次数据竞争造成的线上事故代价大。

原则四:每个可能出问题的优化都配有开关。 冷热分离、数据缩写、空数据标记、会话缓存——每个特性都独立可控,新功能先以旁路模式运行验证。这不是过度设计,是生产环境的必要条件:没有开关的优化叫赌博,有开关的优化叫实验。

原则五:框架封装平台差异,业务只关心逻辑。 #ifndef WIN32 之下是 epoll 还是 IOCP,业务代码不知道也不应该知道。框架层承担了所有平台适配的复杂度,业务层拿到的是统一的 Handle(pMsg)SendMsg 接口。这种分层在跨平台 C++ 项目中是可复用的标准做法。


全文六个维度,每一个单拿出来都不是新技术——epoll 存在了二十年,Redis 分布式锁是常见模式,泳道隔离是消息队列的常规操作。这套架构的价值在于把这些成熟技术用在了正确的位置,并且每一层的设计都严格遵循了它上下层的约束。这种”不越界”的纪律,以及对每一处取舍的清醒认知,比单个技术选型更难做到,也更有长期价值。

本文由作者按照 CC BY 4.0 进行授权