↓ 跳过正文
  1. 文章/

自研内存检索引擎的架构与设计取舍

haifeiWu
作者
haifeiWu
北京后端开发者,热爱分布式系统、中间件与 AI 工程,正在深入学习 Go 与 Rust。
目录

做检索,第一反应通常是 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")]

这套结构里有两个设计点想单独说一下:

  1. 查询处理与存储解耦,读取路径无锁。查询全程只读内存索引,只有写入(文档变更、快照导入)才加锁,所以一次耗时的快照导入不会长时间阻塞查询。
  2. 数据流入与流出分离。业务数据经 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 层该层指针跨越区间(如 10005000、60009000)内所有取值的文档并集
第 1 层该层指针跨越区间(如 10003000、40006000)内所有取值的文档并集
第 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

有两个设计细节。

  1. must_not 针对"上游已收敛的结果集"做差,不是先在全量数据上排除:布尔算子的子节点顺序是 must → filter → must_not,执行到差集时结果集已经很小,参与运算的位图规模跟着缩小。(这个顺序对正确性也有影响,见 11.3。)
  2. 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: 持续观察主的变化
设计点取值 / 做法说明
租约 TTL30 秒主实例宕机后,租约到期,其他实例在 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["返回当前页"]

在课程这种数据量下,这笔开销可以接受,但有两个问题:

  1. 深分页成本没有改善。from = 10000, size = 20 时仍然要全量排序,可用户只要 20 条;
  2. 每次翻页都重复排序。同一排序条件下连续翻页,重复做的全是同样的工作。

标准解法是 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,但在自己被设计的那一类场景里,它更快、更简单,也更便宜。

相关文章