做检索,第一反应通常是 Elasticsearch。但如果把场景收窄:数据量在十万级,字段类型基本是整数与枚举,筛选条件组合极其灵活,要求毫秒级返回,又不想维护一个独立集群,那么"自研一个内存检索引擎"就从一个疯狂的想法变成了一个可以在可控成本内落地的工程问题。
筛选服务(filter)就是在这个前提下自研的一套检索引擎。它提供对标轻量级 Elasticsearch 的检索能力,业务上支撑教育场景的选课搜索,覆盖课程、课程包、老师三类对象的筛选、排序、分页、分组与聚合,并配套完整的数据生产与统计回流链路。
本文按"问题 → 架构 → 存储 → 查询 → 持久化 → 高可用 → 数据链路 → 业务适配 → 改进方向"的顺序,把它的设计拆开讲一遍。
一、问题与边界#
选课场景的检索请求有这么几个特征:
- QPS 高、请求密集。用户在前端不断勾选筛选条件(年级、学科、学期、难度、版本、价格区间、省份、设备、老师……),每一次勾选变化几乎都会触发一次检索。
- 结果集不大,但条件组合极多。一次检索最终返回的往往只有几十条,但参与筛选的字段有十几个,且组合方式由用户自由决定。
- 数据规模可控。课程、老师这类对象是万级到十万级,远没有到必须依赖分布式检索集群的量级。
- 业务字段变化频繁。中台体系(文中称"乐高"体系)会不断调整字段语义、映射规则和排序权重,检索层需要快速适配。
用 Elasticsearch 当然能做,但集群运维、JVM 内存、网络往返、mapping 变更这几项成本跑不掉。而这里的数据规模小到可以完整放进单机内存,一旦接受这个前提,倒排索引 + 位图运算就能把这个场景做到极致:网络、磁盘 IO、序列化开销,全都没有。
这也是整套引擎的设计基调:先承认场景边界,再在边界内把性能与简洁性做到极致,而不是为想象中的规模提前引入复杂度。
二、整体架构#
仓库包含三套独立运行的程序:
| 程序 | 定位 | 说明 |
|---|---|---|
filter-server | 筛选引擎服务 | 常驻 HTTP 服务,提供检索、文档/索引管理、快照导入导出 |
course | 课程数据同步任务 | 离线/定时 ETL,把业务库数据加工后批量写入引擎 |
course-count | 购课数统计任务 | 基于引擎快照并行统计购课数并写回业务库 |
系统由「四层引擎 + 基础组件」构成,整体架构如下:
flowchart TB
subgraph L1["接入层"]
GW["HTTP 服务(Gin)
search / doc / index / db"]
end
subgraph L2["查询处理层"]
PA["Parser
DSL 解析"] --> PL["Planner
查询计划"] --> EX["Executor
执行树"]
EX --> AG["聚合器
桶 + 指标"]
end
subgraph L3["存储引擎层"]
DB[("内存 DB
多索引")]
IDX["倒排索引
哈希 / 跳表"]
BM["文档位图"]
DB --- IDX
DB --- BM
end
subgraph L4["持久化与基础组件"]
ST[("Store
快照多后端")]
ET["Etcd
主备选举"]
NC["Nacos
配置下发"]
end
GW --> PA
EX --> DB
DB --> ST
ET -. 选举 .-> DB
NC -. 配置 .-> GW
各层职责:
| 层次 | 主要组成 | 职责 |
|---|---|---|
| 接入层 | HTTP 服务、路由与控制器 | 对外暴露检索、文档、索引、快照等 HTTP 接口 |
| 查询处理层 | Parser、Planner、Executor、聚合器 | DSL 解析 → 查询计划 → 执行树 → 排序分页分组聚合 |
| 存储引擎层 | 内存 DB、字段倒排索引、文档位图 | 多索引内存库;字段级倒排索引与位图运算 |
| 持久化与基础组件 | Store 抽象、Etcd、Nacos、运行时观测 | 快照的多后端存储;主备选举、配置下发、指标采集 |
| 数据生产链路 | course 同步任务、course-count 统计任务 | 业务数据加工写入引擎;基于快照统计购课数并回流 |
从数据流动的视角看,同一套系统串起了写入、查询、统计三条链路:
flowchart TB
D[("业务 MySQL")] --> E["course 同步任务
数据加工 → 建索引 → 批量写入"]
E --> F[("内存索引
倒排索引 + 位图")]
A["选课请求"] --> B["检索接口
位图筛选 → 排序分页聚合"]
B --> C["结果返回"]
F --> B
F --> G["定时导出快照"]
G --> H["course-count
并行计数"]
H --> I[("统计 MySQL")]
这套结构里有两个设计点想单独说一下:
- 查询处理与存储解耦,读取路径无锁。查询全程只读内存索引,只有写入(文档变更、快照导入)才加锁,所以一次耗时的快照导入不会长时间阻塞查询。
- 数据流入与流出分离。业务数据经 ETL 任务单向写入引擎;引擎的状态则通过快照在实例之间流动。两条路径互不依赖,任一环节出问题都不会同时影响读写。
三、存储引擎:倒排索引与位图#
3.1 文档模型:只存"原始 JSON"#
每个文档保存两样东西:
| 组成 | 内容 | 作用 |
|---|---|---|
JSONValues | 原始 JSON 字符串 | 结果返回时原样拼进响应体,不参与检索 |
fields | 字段名 → []int 的映射 | 检索阶段唯一操作的对象 |
原始 JSON 不解码成结构体,需要检索的字段值另外解析成整型数组存一份。这个分离贯穿整个查询链路:检索阶段只碰 int 数组和位图,绝不反序列化 JSON,只有最终要把结果吐给业务方时,才把原始 JSON 直接作为 json.RawMessage 拼进响应。查询链路上反复解析 JSON 的开销,就这么省掉了。
3.2 字段索引:按查询模式分配成本#
每个字段(FieldIndex)按它被查询的方式,选用不同的索引结构:
flowchart TB
IDX["Index(一个业务索引)"] --> DOC["Docs
文档集合(ID → Doc)"]
IDX --> FI["FieldIndex × N
每个可检索字段一个"]
FI --> H["Hash:值 → 位图
(所有字段)"]
FI --> SK["SkipList:值有序 + 分层跨越位图
(仅 range 字段)"]
H --> U1["承担精确匹配
term / terms"]
SK --> U2["承担范围匹配
range"]
| 结构 | 建立条件 | 承担的查询 |
|---|---|---|
Hash(值 → 位图) | 所有字段 | 精确匹配、多值匹配(term / terms) |
SkipList(有序 + 位图) | 仅 IndexType = "range" 的字段 | 范围匹配(range) |
并非所有字段都付出有序结构的代价,只有声明为 range 的字段才建立跳表,其余字段只保留 值 → 位图 的哈希映射。这就是"按字段的实际查询模式分配存储成本":一个只做枚举筛选的字段,维护有序性纯属浪费。
3.3 范围查询:跳表的分层跨越位图#
范围查询要回答的是"取值落在 [1000, 5000] 区间内的所有文档"。跳表提供了值的有序性,但还需要把区间内每个取值对应的文档集合合并起来,逐个取值求并的话,范围查询就退化成了遍历。
筛选服务的解法是:让跳表的每个节点在每一层维护一个位图,表示该层指针所跨越区间内所有取值的文档并集。
flowchart TB
L2["第 2 层(跨度大)
节点 1000、6000"] --> L1["第 1 层(跨度中)
节点 1000、4000"]
L1 --> L0["第 0 层(原始链表)
节点 1000、2000、…、5000"]
| 层级 | 节点位图的含义 |
|---|---|
| 第 2 层 | 该层指针跨越区间(如 1000 |
| 第 1 层 | 该层指针跨越区间(如 1000 |
| 第 0 层 | 该取值自身的文档集合 |
于是范围查询变成了一条"沿层跳跃 + 合并位图"的流水线:
flowchart TD
A["按区间下界定位起点"] --> B["沿跳表层级向前跳跃
合并沿途的跨越位图"]
B --> C["补上终点节点自身的位图"]
C --> D["得到区间内的全部文档"]
这样一次范围查询只需沿跳表层级向前跳跃,把沿途 O(log n) 个"跨越位图"求并,再补上终点节点自身的位图,用少量位图运算替代了对区间内每个取值的逐个遍历。写入时同步维护这些跨越位图,新文档 ID 会被加到沿途各层节点的位图上;删除时再逐个移除。
代价是空间与写放大:每个节点在每一层都要维护一个位图,写入一个文档需要更新 O(log n) 个位图。对于"批量导入为主、单文档更新为辅"的同步场景,这个取舍是划算的。
3.4 字符串编码为整数:位图化的前提#
位图只能装整数,而业务字段里有大量字符串枚举(学科名、版本名、难度名……)。引擎在写入时把字符串编码成自增整数,一次编码同时解决两个问题:
flowchart TD
V["字符串取值
如「数学」"] --> Q{"编码表中
是否已存在?"}
Q -->|是| HIT["直接返回已有编码"]
Q -->|否| GEN["分配新自增编码"]
GEN --> FWD["写入正向映射
string → int"]
GEN --> REV["写入反向字典
int → string"]
FWD --> USE["该字段可参与位图运算"]
REV --> TR["结果返回前翻译回可读值"]
- 正向映射:让字符串字段也能参与位图运算;
- 反向字典:结果返回前把编码翻译回可读值(业务方看到的是"数学"而不是
7)。
两件事用一次编码一起做完,不用维护两套结构。代价是字典只增不减,删除文档不会回收已分配的编码。对低基数的枚举字段(学科、年级、难度)这是完全合理的取舍;但若用在高基数、高频变更的字段上(比如用户 ID),字典会持续膨胀,这是使用时的边界。
3.5 位图集合运算#
所有筛选结果统一封装为 DocIDSet,底层是 RoaringBitmap(压缩位图),对外只暴露集合语义:
| 运算 | 语义 | 对应查询 |
|---|---|---|
And | 交集 | must、filter |
Or | 并集 | should、多值匹配 |
AndNot | 差集 | must_not |
Count | 基数 | 命中总数(不需物化文档) |
flowchart LR
A["年级 = 三年级
位图 A"]
B["价格 ∈ 1000~5000
位图 B"]
C["老师 ∈ 排除名单
位图 C"]
A --> AND1["A ∩ B"]
B --> AND1
AND1 --> AND2["(A ∩ B) 再排除 C"]
C --> AND2
AND2 --> R["最终文档 ID 集合"]
十万级文档规模下,多个条件的交、并、差就是几次微秒级的位运算,“毫秒级筛选"的物理基础就在这里。
3.6 写入路径:幂等的重建式更新#
flowchart TB
START["写入文档 doc"] --> EXISTS{"该 ID
已存在?"}
EXISTS -->|是| DELOLD["按旧字段值逐一撤位
从各字段位图中移除该 ID"]
EXISTS -->|否| SKIP["跳过"]
DELOLD --> PUT["存入文档集合"]
SKIP --> PUT
PUT --> RENEW["按新字段值重新占位
写入各字段位图 / 跳表"]
RENEW --> DONE["完成(同 ID 重复写入不产生重复)"]
先拆旧、再建新,保证同一个 ID 反复写入不会产生重复条目。代价是每次更新都要遍历所有字段索引(而不是只处理实际变化的字段),更新成本与字段总数成正比。这是"用简单的全量处理换正确性"的取舍,在批量导入为主的场景下非常划算。
四、查询引擎#
4.1 类 Elasticsearch 的 JSON DSL#
对外接口是一套类 ES 的查询 DSL,用法习惯与 ES 保持一致,业务方几乎不用重新学:
{
"index": "index-course",
"query": {
"bool": {
"must": [{ "term": { "field": "lego_grade_id", "value_int": 3 } }],
"filter": [{ "range": { "field": "price", "gte": 1000, "lte": 5000 } }],
"must_not": [{ "terms": { "field": "teacher_id", "values": [1, 2] } }]
}
},
"sort": [{ "price": { "order": "desc" } }],
"from": 0,
"size": 20,
"aggs": { "by_grade": { "terms": { "field": "lego_grade_id" } } }
}请求解析按路径取值,只解析需要的部分,不会为了一次查询就把整个请求体反序列化成结构体。条件的反序列化用注册表模式组织,新增一种条件类型就是多写一个解析器,已有代码不用动。
4.2 三个阶段#
flowchart TD
DSL["JSON DSL"] --> P["Parse
语法校验 + 路径解析"]
P --> PLAN["Planner
模型化为查询计划
查询 / 排序 / 分组 / 聚合"]
PLAN --> BUILD["Builder
组装算子树"]
BUILD --> RUN["Execute
Open → Next → Close"]
RUN --> RES["Chunk 结果
文档 / 分组 / 聚合桶"]
计划层把 DSL 变成一个纯数据模型(查询、排序、分组、聚合),与执行完全解耦。业务侧传入的参数(品牌、设备、省份、年级等)也由这一层承载,作为执行阶段的上下文。业务差异止步于此,执行引擎保持通用。
4.3 执行树#
执行树按"业务语义从外到内"的顺序组装,筛选条件作为叶子节点:
flowchart TD
SEARCH["SearchExec
(返回结果)"] --> LIMIT["LimitExec
切片分页"]
LIMIT --> SORT["SortExec
排序"]
SORT --> AGG["AggregateExec
聚合"]
AGG --> GROUP["GroupExec
分组"]
GROUP --> QUERY["QueryExec
位图筛选"]
QUERY -. 数据从叶子逐层向上流动 .-> SEARCH
执行时数据从叶子向根流动:查询条件在位图上筛出候选集,逐层经过分组、聚合、排序,最后在根节点切出当前页。
这种组装顺序与业务书写顺序一致(“用户最终拿到什么"最先写,最底层的筛选最后写),Builder 的代码顺序和阅读顺序是同一个方向,不需要在大脑里反向建树。
4.4 火山模型#
所有算子实现同一套接口,各自只关心"如何产生自己的数据”:
| 方法 | 职责 | 默认实现 |
|---|---|---|
Open | 准备资源 | 递归向下打开所有子算子 |
Next | 产出一批数据到 Chunk | 各算子自行实现 |
Close | 释放资源 | 递归向上收尾,返回首个错误 |
Open / Close 的递归遍历由基类统一处理,每个算子只需要实现自己的 Next。新增一种查询算子 = 新增一个文件 + 在 Builder 里挂一行,已有算子一行不动。引擎的可扩展性就是从这来的。
4.5 Chunk:算子间的数据单元与惰性物化#
Chunk 是算子之间传递的数据单元,也是整个引擎里最关键的抽象。它同时承载两类数据:
| 字段 | 含义 | 何时使用 |
|---|---|---|
DocIDSet | 命中文档的位图 | 筛选阶段(主要载体) |
Docs | 命中文档对象列表 | 需要排序、分组、输出时 |
Groups | 分组结果 | 分组场景 |
Aggs | 聚合结果 | 聚合场景 |
Total / Sorted / Limited | 结果状态 | 决定各字段是否可信、是否需要物化 |
核心设计是位图与文档分离:筛选阶段全程只操作位图,完全不碰文档对象;只有最后一个需要输出结果的算子才把 ID 批量换成文档。
sequenceDiagram
participant Q as 筛选算子
participant C as Chunk
participant D as 文档集合
Q->>C: 写入命中位图(10 万 → 60 条)
Note over C: 仅位图,零文档解引用
C->>C: 排序算子读取位图基数
C->>D: GetDocs:位图 → ID 数组 → 批量取文档
D-->>C: 文档对象(仅命中的少量文档)
C-->>Q: 结果集
一次典型的选课查询可能从十万文档筛到几十条,这十万次"排除"的代价只是几次位图与运算,全程没有文档解引用,也没有 JSON 解析。毫秒级响应靠的就是这个,GC 压力也顺带小了。
4.6 布尔条件即集合运算#
must / filter / must_not 的语义与集合运算天然对应,实现上就是位图的交与差:
flowchart TB
subgraph MUST["must / filter:逐个条件求交"]
M1["条件 1 位图"] --> MA["∩"]
M2["条件 2 位图"] --> MA
MA --> MR["收敛结果集"]
end
subgraph MUSTNOT["must_not:先并后排"]
N1["排除条件 1"] --> NA["∩"]
N2["排除条件 2"] --> NA
NA --> NR["排除集合"]
end
MR --> FINAL["结果集再排除干扰集合"]
NR --> FINAL
有两个设计细节。
must_not针对"上游已收敛的结果集"做差,不是先在全量数据上排除:布尔算子的子节点顺序是must → filter → must_not,执行到差集时结果集已经很小,参与运算的位图规模跟着缩小。(这个顺序对正确性也有影响,见 11.3。)must与filter的实现完全一致。这套引擎里没有相关性打分环节,两者的区别只剩语义表达,保留两个名字纯粹为了和 ES 习惯对齐。这算是自研引擎的一个隐性好处:不需要为不存在的功能付出复杂度。
五、排序与分组#
5.1 排序:为高频路径开一条捷径#
排序算子对最常见的"按 id 排序"做了特殊处理:升序不做任何操作,降序只把结果反转一次,完全跳过比较排序。
flowchart TB
IN["待排序结果集"] --> Q{"首个排序字段
是 id ?"}
Q -->|是,升序| NOOP["无需操作
(已按 ID 升序)"]
Q -->|是,降序| REV["反转一次
O(n)"]
Q -->|否| CMP["比较排序
O(n log n)"]
捷径为什么能成立?命中结果是从位图取出的,而压缩位图导出的 ID 天然按升序排列,所以"按 id 升序"是顺带送的,存储层不用为它维护任何有序结构。选课场景里"按 id 排序"常被用作默认排序,这条捷径的命中频率不低。
有一个边界要交代:捷径依赖"结果来自位图"这个前提。若某条路径直接把文档列表灌进结果集(比如全量匹配),前提就不成立了,两条路径对有序性的假设并不一致(见 11.5)。
5.2 多字段排序与结果稳定性#
多字段排序逐键比较,遇到字段缺失就跳过该键、交给下一个,全部键都相等时按文档 ID 兜底。
| 规则 | 行为 | 原因 |
|---|---|---|
| 字段缺失 | 跳过该排序键 | 比"缺失一律排最后"更灵活 |
| 多键相等 | 按 ID 兜底排序 | 非稳定排序下若无兜底,相等元素顺序随机,分页会出现重复或漏项 |
第二条尤其要紧:分页是对有序结果切片,全序关系必须唯一,不然同一文档可能在两页里重复出现,也可能有文档一直没被展示。
5.3 虚拟排序字段:把业务权重从索引里挪出来#
“老师亲密值"是一个不在任何索引里的排序字段,它的值由 (老师 ID, 学科, 年级) 三元组从业务配置中查出:
flowchart TD
S["排序字段
teacher_intimacy"] --> K["拼接查询键
老师ID, 学科, 年级"]
K --> CFG["业务配置字典
(随配置热更新)"]
CFG --> V["得到排序权重"]
V --> CMP["参与多字段比较"]
这个设计把**“随业务策略变化的排序权重"从索引结构中剥离**了出来:调整老师排序策略(运营侧常做的事)不用重建索引,也不用重新推数据,改配置即可生效。
代价是排序比较次数为 O(n log n),每次比较都要拼接键并查字典,常数不小。这算是拿"语义灵活性"换来的热路径开销,也是一个明确的优化点(见 11.5)。
5.4 分组:用多叉前缀树表达层级#
多字段分组用一棵前缀树表达:第一层按第一个字段的取值建节点,第二层在各分支下按第二个字段的取值建节点,最深层成为叶子并持有一个文档集合。
flowchart TB
ROOT["根节点
分组字段: 年级 → 学科 → 难度"] --> G1["年级 = 三年级"]
ROOT --> G2["年级 = 四年级"]
G1 --> S1["学科 = 数学"]
G1 --> S2["学科 = 英语"]
S1 --> D1["难度 = 基础
文档集合"]
S1 --> D2["难度 = 进阶
文档集合"]
G2 --> S3["学科 = 数学"]
S3 --> D3["难度 = 基础
文档集合"]
好处很直接:
- 不存在空的组合桶:节点因文档实际出现才被创建,不用预先枚举所有字段取值的笛卡尔积;
- 分组与筛选共享同一次文档遍历,文档流经分组算子时就地分发,省掉一次专门为分组的扫描;
- 层级语义天然表达,“年级 → 学科 → 难度"三层分组正好是树的三层。
这里有个语义细节:文档在分组字段上多值时会被分发到多个分组(比如一节课面向多个年级),所以各组文档数之和会大于实际文档总数,业务侧做统计时必须心里有数。
5.5 分组场景有独立的分页与排序#
分组场景使用的是另一组算子,排序分两级进行:先按分组字段对组排序,再对组内文档排序;分页切片的对象是"组"而不是"文档”。
这与业务语义完全对得上:按年级分组、每组取前 N 门课,用户要的是"每个组各看前 N 条”,不是"全局前 N 条再分组”。语义不同,执行树结构也就不同。所以分页、排序做成了独立算子(GroupLimitExec / GroupSortExec),而不是在同一个算子里加 if。
六、快照持久化#
6.1 关键取舍:存原始文档,不存倒排索引#
flowchart TB
subgraph MEM["内存中的索引"]
HASH["字段哈希索引"]
SKIP["跳表"]
BM["文档位图"]
end
subgraph SNAP["快照内容(不含索引结构)"]
META["索引元数据
字段定义 / 字典 / 排序规则"]
DOCS["原始 JSON 文档列表
(按 ID 有序)"]
TS["时间戳 + 校验信息"]
end
MEM -->|导出| SNAP
SNAP -->|导入:逐文档重建索引| MEM
快照里放的是字段元数据 + 原始 JSON 文档列表,倒排索引和位图一个都不存。导入时逐文档重建索引。
这个取舍有得有失:
收益
- 快照格式与索引实现解耦:换位图库、改跳表结构、调整编码方式,都不影响快照兼容性;
- 快照是"数据的真值”,可以独立校验、独立分析,统计任务直接拿快照当数据源,完全不需要起一个引擎实例;
- 体积小:存索引结构反而会因为每个字段的哈希与位图产生大量冗余。
代价
- 导入需要重建全部索引,是 CPU 密集操作,恢复时间与文档数成正比;
- 导入后的索引质量完全取决于元数据是否被完整回填,这一点在 11.1 会看到实际问题。
另外,只有标记了"需要持久化"的索引才会进入快照,临时索引不参与,避免把中间态数据写进存储。
6.2 序列化链路#
flowchart TD
MEM["内存 DB"] --> ED["ExportedDB 结构体"]
ED --> PB["Protobuf 序列化
时间戳 + 各索引"]
PB --> GZ["gzip 最高压缩级别"]
GZ --> MD5["计算 MD5
与上次比对"]
MD5 --> UP["上传:OSS / Redis / 文件"]
| 环节 | 选择 | 理由 |
|---|---|---|
| 序列化 | Protobuf | 体积小、编解码快 |
| 压缩 | gzip 最高压缩级别 | 导出是低频后台任务,用 CPU 换传输与存储成本 |
| 校验 | 内容 MD5 | 支撑去重与完整性校验 |
6.3 内容去重:没变就不上传#
导出任务可能被频繁触发,但数据往往几分钟甚至几十分钟才变一次。导出流程在压缩完成后计算一次 MD5,与上次的结果比对,一致就直接跳过上传。
flowchart TB
E["导出触发"] --> SER["序列化 + 压缩"]
SER --> C["计算 MD5"]
C --> Q{"与上次 MD5
相同?"}
Q -->|相同| SKIP["跳过上传
(数据未变化)"]
Q -->|不同| UPLOAD["上传并记录新 MD5"]
这里有个不显眼但必要的配合:导出文档列表必须用按 ID 有序的遍历方式。要是无序遍历,“数据完全没变"的两次导出也可能产生字节不同的快照,MD5 比对失效,去重形同虚设。一个性能优化的正确性,依赖另一个看起来无关的细节,这种事在工程里比想象中常见。
上次的 MD5 存在内存里,进程重启后第一次导出必然上传一次——“多传一次"的代价远小于"漏传一次”,这个取舍方向是对的。
6.4 多后端可插拔#
存储层抽象为一个统一接口,提供"导出快照"与"导入快照"两个能力,并自带配置:
| 实现 | 用途 |
|---|---|
| Redis | 作为低延迟的共享存储,实例间快速同步 |
| 对象存储(OSS) | 作为持久化归档,配合元数据库记录版本 |
| 本地文件 | 单机部署与调试 |
| Mock 文件 | 单元测试,为可测试性留的口子 |
多个后端可以同时启用、互为备份,主后端不可用时切到另一个继续同步。按名字分发,调用方不需要知道自己连的是哪种存储。
七、一致性与同步#
7.1 同步机制:通知与数据分离#
实例之间的数据同步不直接推送数据,而是靠"存储配置"这个中间层来驱动:
sequenceDiagram
participant P as 主实例
participant S as 存储后端
participant C as 存储配置
participant R as 副本实例
P->>S: 导出并上传快照
P->>C: 更新存储配置(含新时间戳)
loop 每 12 秒
R->>C: 拉取存储配置
end
C-->>R: 返回配置(时间戳更新)
R->>S: 拉取快照
R->>R: 重建索引,置为就绪
把"通知"和"数据"分开之后,实例之间就不需要知道彼此的存在,新增实例只要能读到存储配置就行。这也是这套机制能支撑多集群双活的原因。
7.2 两条准入规则#
副本在导入前会做两项校验:
| 规则 | 判断 | 解决的问题 |
|---|---|---|
| 版本一致 | 配置中的协议版本与本地常量相等 | 快照格式演进时宁可跳过,也不误读结构未知的数据 |
| 时间戳递增 | 新时间戳 严格大于 当前内存库的时间戳 | 防止旧快照覆盖新数据(重复投递、乱序到达) |
时间戳用的是严格大于而不是大于等于,比较基准是当前内存库里实际的数据版本,不是本地上次看到的时间戳。等于把"回退"从根上堵死了。
7.3 配置获取方式#
| 配置 | 内容 | 获取方式 | 生效延迟 |
|---|---|---|---|
| 存储配置 | 快照时间戳、存储后端 | 轮询 | 最长 12 秒 |
| 业务配置 | 品牌维度的搜索类型、连接信息、密钥 | 轮询 + 配置中心 | 最长 76 秒 |
用轮询而不是长连接监听,取舍很清楚:实现简单,省掉断线重连和事件乱序的维护成本,代价是配置生效有延迟。对"分钟级同步一次"的数据链路来说,这点延迟完全能接受;但如果业务要求"改了配置立刻生效”,轮询就是错的,必须换成事件监听。
另外,数据就绪状态用显式标志表达,导入成功才置位。接入层可以据此拒绝服务或者降级,不至于在数据还没就绪时就把错误结果返回出去。
八、高可用:Etcd 主备选举#
多个实例同时导出快照会产生写冲突,也浪费资源。引擎用 Etcd 的租约选举,保证同一时刻只有一个"主"承担导出职责:
stateDiagram-v2
state "观察中" as Observing
state "竞选中" as Campaigning
state "主实例" as Primary
state "副本" as Replica
[*] --> Observing: 启动
Observing --> Campaigning: 无主 / 自己不是主
Campaigning --> Primary: 抢到租约
Primary --> Replica: 租约失效或被抢占
Replica --> Campaigning: 观察到主已变化
Primary --> Observing: 持续观察主的变化
| 设计点 | 取值 / 做法 | 说明 |
|---|---|---|
| 租约 TTL | 30 秒 | 主实例宕机后,租约到期,其他实例在 30 秒内接管 |
| 竞选入口 | 本地互斥锁防重入 | 避免每次收到观察事件都重复发起竞选 |
| 角色迁移 | 四种状态分别记录日志 | 首次成为主、从主降为副本、副本仍是副本、成为新主 |
| 写职责 | 仅主实例 | 副本只做观察与导入,天然避免快照写冲突 |
表格里第三条想单独说说。这类分布式状态机排障时,难点往往不是逻辑错了,而是不知道自己现在处于什么角色。把角色迁移的四种情况分开打日志,“当前状态"就一目了然了。
九、数据生产链路#
引擎的性能建立在"数据已经在内存里"这个前提上,所以数据怎么写进去,同样重要。
9.1 生产者-消费者流水线#
数据生产被抽象成一个通用的流水线骨架:多个生产者并行产出,汇进一个通道,再串成一条消费者加工链。
flowchart TD
P1["producer1"] --> OUT["producerOut
(汇聚)"]
P2["producer2"] --> OUT
P3["producer3"] --> OUT
OUT --> C1["consumer1"] --> C2["consumer2"] --> CN["consumerN"] --> D["(排空)"]
| 工程细节 | 做法 | 收益 |
|---|---|---|
| 生产者收尾 | 等所有生产者结束后才关闭汇聚通道 | 不会关得太早导致 panic,也不会漏数据 |
| 环节缓冲 | 每级之间使用小容量缓冲通道 | 提供流水线并行度的同时,让背压尽快向上游传导,避免数据积压 |
| 尾部排空 | 专门的协程排空最后一个通道 | 防止消费链末端因无人读取而阻塞 |
| 可观测性 | 每 5 秒输出各环节自定义指标 | 长跑任务能看到进度,而非"干等 + 猜” |
| 优雅退出 | 监听中断信号,取消全链路,超时强退 | 避免同步任务中途被粗暴终止 |
抽象出"生产者 / 消费者"两个接口之后,整条链路的编排就变成了纯配置:
flowchart TD
A["课程生产者"] --> B["基础组装"]
B --> C["外部接口补全
促销 / 销量 / 班级 / 满意度"]
C --> D["字段映射
展示转换 / 排序标志"]
D --> E["建索引 + 批量写入"]
E --> F[("筛选引擎")]
“慢且可能失败的外部依赖"放在中间环节,任一环节出问题,都能从进度输出里马上看出是哪一级卡住、卡在哪条数据上。这就是流水线设计比"一个大函数顺序执行"实在的地方。
9.2 攒批:兼顾批量效率与低流量时延#
单条消息逐条写入引擎,锁竞争和索引更新的开销都很大,所以写入前得先把单条流攒成批。
flowchart TB
IN["单条数据流"] --> TRY{"通道中
有数据?"}
TRY -->|有| BUF["攒入缓冲区"]
BUF --> FULL{"攒满一批?"}
FULL -->|是| FLUSH["发送一批"]
FULL -->|否| TRY
TRY -->|没有| EMPTY{"缓冲区
非空?"}
EMPTY -->|是| FLUSH2["先发送当前缓冲
(避免低流量死等)"]
EMPTY -->|否| BLOCK["阻塞等待一条数据
(避免空转烧 CPU)"]
这个策略是为了对付一个很实际的矛盾:既要批量效率,又不能因为流量低而长时间不发送。纯攒批(攒满才发)在夜间低流量时,数据会迟迟写不进引擎;纯逐条发又把批量收益丢了。先非阻塞尝试,攒不到就先发,真的没数据才阻塞等待,两头都照顾到了。
用在写入引擎时,批量写入把 N 次"加锁 + 更新索引 + 更新位图"合并成一次,收益是线性的。
9.3 并行计数:worker pool + 全或无#
购课数统计任务基于快照做并行计数,做法里有一个明确的正确性取舍:
flowchart TB
START["读取最新快照"] --> CHK{"时间戳
为最新?"}
CHK -->|否| SKIP["跳过本轮"]
CHK -->|是| IMP["导入快照到内存"]
IMP --> DISP["分发任务
(传课程下标)"]
DISP --> W["N 个 worker 并行计数"]
W --> RES{"全部成功?"}
RES -->|是| TX["单事务写回统计库"]
RES -->|否| ABORT["取消本轮
不写任何结果"]
| 设计 | 做法 | 目的 |
|---|---|---|
| 任务分发 | 通道只传下标,worker 自行按索引取数据 | 避免在通道中搬运大对象 |
| 结果收集 | 主协程逐个接收结果,失败立即取消全局 | 快速失败,不浪费算力 |
| 写库 | 单事务提交 | 要么全写进去,要么一条都不写 |
| 失败处理 | 任何一条数据异常则本轮整体作废 | 宁可不出结果,也不写半截脏数据 |
这就是"统计数据的正确性优先于时效性”。对下游要依赖统计结果的业务来说,一份缺失的数据远好过一份错误的汇总。
十、业务适配:把变化关进配置里#
这套引擎被多个品牌和业务线复用,差异主要靠配置来吸收:
| 层次 | 内容 | 变更方式 | 生效延迟 |
|---|---|---|---|
| 基础配置 | 端口、存储后端、日志、Etcd/Nacos 地址 | 本地 TOML + 环境变量 | 需重启 |
| 业务配置 | 品牌维度的搜索类型、MySQL 连接、接口密钥 | TOML + 配置中心 | 最长 76 秒 |
| 业务规则配置 | 字段映射、课程包规则、学科/年级/难度排序、字典、老师亲密值 | JSON 文件 + DB | 随数据同步任务加载 |
| 存储配置 | 快照时间戳、存储后端 | 运行时更新 | 最长 12 秒 |
检索入口按类型分流,走不同的解析与修正逻辑:
flowchart TB
REQ["选课请求
(含 type 与筛选条件)"] --> T{"type"}
T -->|课程 / 加加购| C1["解析课程条件"]
T -->|课程包| C2["解析课程包条件"]
T -->|老师| C3["解析老师条件"]
C1 --> FIX["按品牌修正
排序规则 / 字段映射 / 字典"]
C2 --> FIX
C3 --> FIX
FIX --> PLAN["组装查询计划"]
PLAN --> EXEC["执行引擎检索"]
EXEC --> OUT["结果组装
字典翻译 / 字段裁剪"]
“按品牌修正"这一步是关键,它把品牌差异集中到一个地方:按当前品牌加载排序规则、字段映射、字典翻译表,再去修正查询计划。同一份引擎代码,不同品牌走不同规则。
查询计划里保留了业务参数(品牌、设备、省份、年级等),但执行引擎完全不知道自己服务的是哪个品牌,业务差异被挡在解析层与计划层之内。这个边界划得很干净,也是这套代码能被复用、而不是被条件分支淹没的根本原因。
十一、设计要点与改进方向#
梳理这套架构的过程中,有几个设计细节值得单独拎出来讲。它们不影响"整体架构是否成立"的判断,但都对应明确的改进空间。
11.1 快照元数据的回填不完整#
快照结构里带着字段的索引类型信息,但导出、导入两侧的处理并不对称:
| 环节 | 对字段索引类型的处理 |
|---|---|
| 导出 | 写入字段元数据(含索引类型) |
| 导入 | 重建字段索引时固定传空值,未使用元数据中的索引类型 |
因为只有声明为 range 的字段才会创建跳表,“索引类型为空"就意味着从快照恢复出来的索引丢了全部范围索引,范围查询能力跟着失效。
影响面取决于使用方式。实例启动后如果总是全量重建索引,问题就被掩盖了;如果走"启动导入快照 + 后续增量更新"的路径,价格区间、时间区间这类条件在恢复之后就不再生效——而恢复能力恰恰是快照最主要的用途。
flowchart LR
EXP["导出的元数据
包含索引类型"] -.->|未被使用| IMP["导入重建
索引类型为空"]
IMP --> NORANGE["range 字段无跳表"]
NORANGE --> BROKEN["范围查询不可用"]
11.2 条件构建失败的降级策略偏宽#
构建查询条件时,某个字段不存在、类型不支持,或者缺所需的索引,当前的处理是记一条警告、把该条件丢弃,接着执行。
容错的初衷能理解,某个字段暂时不可用时,不至于让整个查询跟着失败。但筛选条件是用户明确表达的约束,把它丢掉,代价是结果集变大、语义变了:业务侧看到的是"这个筛选没生效”,而不是"接口报错了”。这种问题在测试环境往往发现不了(字段齐全、索引完整),只会在生产的边缘场景里冒出来,而且只留下一行警告日志。
更稳妥的做法是分级:结构性错误(字段不存在、类型不支持)直接让请求失败;只有降级不会改变结果语义时,才走警告继续。
flowchart TB
B["构建筛选条件"] --> E{"构建成功?"}
E -->|成功| OK["加入执行树"]
E -->|失败| W["记录警告"]
W --> DROP["丢弃该条件,继续执行"]
DROP --> SEM["结果集变大
筛选语义静默改变"]
SEM --> OBS["业务侧观测:筛选未生效"]
把 11.1 和 11.2 连起来看,是一条完整的链路:快照恢复缺元数据 → 条件构建失败 → 错误降级为日志 → 查询结果静默改变。整条路径上没有任何一处抛异常或返回错误,而这恰恰是最需要防范的一类缺陷。
11.3 排除条件的执行依赖#
must_not 的语义是"对上游已收敛的结果集做差集",实现上对空集做了短路:当前结果集或排除集合为空时直接返回。
flowchart TB
S["进入 must_not 算子"] --> Q1{"当前结果集
为空?"}
Q1 -->|是| R1["直接返回
(排除条件不生效)"]
Q1 -->|否| Q2{"排除集合
为空?"}
Q2 -->|是| R2["直接返回
(结果集不变)"]
Q2 -->|否| DIFF["执行差集"]
也就是说,如果查询里只有 must_not、没有 must / filter 先产出结果集,排除条件实际不会生效。而全量匹配算子只设置文档列表、不设位图,也没法当兜底的起点。
目前这个问题被 DSL 结构掩盖住了,因为布尔查询与全量匹配是互斥分支,而布尔查询内部通常必然存在 must。但将来要是支持"排除若干 ID 后返回全部数据"这类查询,它会立刻变成一个非常安静的缺陷。改进方向也明确:要么让全量匹配同时产出全量位图,要么让排除运算在结果集为空时以全量位图为起点。
11.4 分页是全量排序后切片#
现在的排序算子对结果集做全量排序,分页算子再切片。
flowchart LR
HIT["命中结果集
如 1 万条"] --> SORT["全量排序
O(n log n)"]
SORT --> SLICE["按 from / size 切片
如取 20 条"]
SLICE --> OUT["返回当前页"]
在课程这种数据量下,这笔开销可以接受,但有两个问题:
- 深分页成本没有改善。
from = 10000, size = 20时仍然要全量排序,可用户只要 20 条; - 每次翻页都重复排序。同一排序条件下连续翻页,重复做的全是同样的工作。
标准解法是 Top-K 堆:维护一个大小为 size 的堆,把复杂度从 O(n log n) 降到 O(n log k)。实现上只需要给排序算子加"只保留前 K 个"的能力,不过要小心总数的语义:返回给业务的命中总数仍应该是完整命中数,不能是截断后的数量。
11.5 其他可优化点#
| 模块 | 现状 | 影响与建议 |
|---|---|---|
| 文档写入 | 更新单个文档时遍历所有字段索引删除旧值 | 成本与字段总数成正比;批量场景可按"字段是否变化"筛选 |
| 字段删除 | 位图空了会删除哈希节点,但编码字典条目永久保留 | 枚举字段合理;高基数字段会持续膨胀 |
| 多字段排序 | 每次比较都重新拼接查询键并查字典 | 比较次数 O(n log n),常数偏大;可预先算好键并缓存到文档上 |
| 排序算子 | 隐含要求排序字段非空,空切片会越界 | 目前依赖上层保证;建议在算子构造时校验 |
| 全量匹配 | 走无序遍历并把文档列表直接置位 | 与"按 id 排序可跳过排序"的假设不一致(见 5.1),两条路径的有序性前提应统一 |
| 文档读取 | 位图中的 ID 在文档集合中缺失时静默跳过 | 位图与文档一旦不同步,命中数与结果长度会对不上且无提示 |
| 服务单例 | 包级单例在未初始化时返回空值 | 调用方解引用会 panic,初始化顺序成为隐式契约;可用显式依赖注入替代 |
| 任务代码 | 残留未被调用的调试函数 | 清理项,不影响功能 |
11.6 取舍一览#
| 设计决策 | 换来了什么 | 付出了什么 |
|---|---|---|
| 全内存 + 位图运算 | 毫秒级筛选,无外部依赖 | 数据规模受单机内存限制,无水平扩展 |
| 只给 range 字段建跳表 | 按查询模式分配存储成本 | 字段类型声明必须准确,否则查询能力缺失 |
| 分层跨越位图 | 范围查询免于逐值遍历 | 每层维护位图,写入需更新 O(log n) 个位图 |
| 字符串统一编码为整数 | 位图化 + 字典翻译一次完成 | 字典只增不减,高基数字段会膨胀 |
| 位图与文档分离(惰性物化) | 筛选阶段零文档解引用,GC 压力低 | 算子需正确维护位图/文档/状态的一致性 |
| 快照存原始文档而非索引 | 格式与实现解耦,可独立分析 | 导入需重建索引;元数据不全会导致索引能力缺失 |
| 配置轮询而非事件监听 | 实现简单,无长连接维护成本 | 配置生效有 12s / 76s 延迟 |
| 条件构建失败降级为警告 | 单字段不可用不拖垮整个查询 | 筛选条件可能被静默丢弃,结果错误且难以发现 |
| 统计任务全或无 | 不会写出半截脏数据 | 一条坏数据导致整轮统计作废,时效性受损 |
| Etcd 主备选举 | 写职责唯一,避免快照写冲突 | 引入 Etcd 依赖,故障切换有最长 30 秒窗口 |
十二、小结#
1. 位图是多条件筛选的天然数据结构。 数据能全部放进内存时,把"筛选"翻译成位图的交并差,就是这套引擎性能的根本来源。布尔查询的语义跟集合运算天然对应,所以实现里几乎看不到"遍历、判断、收集"这种传统写法。
2. 位图与文档分离是架构里最巧的一处。 两者在同一个数据结构里分开存放,筛选阶段全程只动位图,到最后才批量取文档。这一个决定同时降低了 CPU 和 GC 的压力,也让"筛十万取二十"这种场景变得很廉价。
3. 索引结构按查询模式分配成本。 只有需要范围查询的字段,才付有序结构(跳表 + 分层位图)的代价,其余字段只留哈希倒排。这种"差异化存储"的思路,比"所有字段一视同仁建全套索引"务实得多。
4. 快照存原始数据而不是存索引结构,是一笔很划算的买卖。 它让存储格式与索引实现解耦,还让快照本身成为可以被独立消费的数据源(统计任务直接读快照,不需要起引擎)。代价是恢复时要重建索引,不过对分钟级的恢复窗口来说,这个代价值得。
5. 配置驱动的边界划得很清楚。 品牌差异、字段映射、排序权重全在配置里,执行引擎完全不知道业务方是谁。同一套代码能服务多个品牌和业务线,靠的就是这一点;所有"会被复用"的引擎类项目,大概也都该这么做。
6. 容错设计需要分级。 对"用户明确表达的约束",异常应该用失败来表达,而不是用降级;降级只留给"可选的增强能力"。这可能是整套设计里最值得记住的一条经验。
要说这套引擎体现了什么工程思路,那就是:先承认场景的边界(数据量可控、单机内存足够),然后在这个边界内把性能和简洁性做到极致——而不是为了应对想象中的规模,提前把系统搞复杂。它当然替代不了 Elasticsearch,但在自己被设计的那一类场景里,它更快、更简单,也更便宜。