ElasticSearch 与搜索服务

2023-03-24T14:11:02+08:00 | 22分钟阅读 | 更新于 2026-03-24T14:11:02+08:00

@

学习目标

学完本章,你应该能够:

  1. 用自己的话解释 ElasticSearch 是什么、解决什么问题,以及它和 MySQL 在"找数据"这件事上的本质区别。
  2. 讲清楚 ES 的核心概念——索引 / 文档 / 分片 / 副本,以及倒排索引、FST 为什么能让全文检索飞快。
  3. 设计一套平台化的搜索服务:统一接入、推送接口、搜索聚合、Kafka 削峰、隔离部署,让新业务方零痛苦接入。
  4. 设计标签服务(缓存预加载、覆盖式打标签、冗余字段),并把它接入搜索、用 boost 权重影响排序。
  5. 在面试里把"通用接入解决共性问题、扩展接口解决个性问题"的平台化思维讲成一段完整故事。

前置知识(如果下面任一点生疏,先回看对应章):

  • 第02章 Gin + GORM:知道一个 HTTP 接口怎么写、DAO 怎么分层。
  • 第07章 Kafka:知道消息是怎么从生产者到消费者的(本章搜索服务写入要用到)。
  • 第05章 缓存:知道 Redis 基本用法(标签服务缓存预加载要用到)。
  • 基本的 MySQL 索引概念(会看 KEYUNIQUE KEY)。

本章你会动手做的事

  • 用 Docker 起一个单节点 ES,亲手建 user_idx 索引、写一条文档、跑一个 match 查询,感受"近实时"。
  • 给搜索服务新增一个业务方(比如"话题"),体会"通用接入 + 注册 Handler"的平台化扩展。
  • 在本地 Redis 里把标签预加载进去,再打一次覆盖式标签,验证旧标签被整体替换。

一、ElasticSearch 入门

1.1 ElasticSearch 是什么

ElasticSearch(ES)是基于 Lucene 的搜索与分析引擎,能够快速、可靠地存储、检索和分析大量数据。核心特性:

  • 高性能:基于 Lucene,全文检索表现出色
  • 近实时:文档索引到可搜索的延迟通常仅 1 秒
  • 大数据:通过分片与副本机制支持 PB 级数据
  • RESTful API:与各种编程语言轻松集成
  • 生态丰富:与 Logstash、Kibana、Beats 等组成 ELK 体系

核心理解:ES 一切能力最终都依赖两个字 —— “搜索”。无论业务检索、日志分析、可观测性,本质都是基于倒排索引的文本搜索。

一句话类比:如果说 MySQL 是一本病历柜(按挂号号找人),ES 就像一个带超级索引的图书馆——你不知道书号,只记得"书里出现过’Go’这个词",它也能在毫秒级把相关书全捞出来。后面所有概念都是为了服务"按词找文档"这一件事。

1.2 基本概念

ES 的概念分为两大类。

数据的组织方式

ES 概念类比 MySQL说明
索引(Index)同类文档的集合
文档(Document)表中的一行一条 JSON 数据
字段(Field)文档中的属性

数据的部署方式

ES 概念类比 MySQL说明
分片(Shard)分库分表数据水平切分,单节点撑不住时的水平扩展手段
副本(Replica)从库主分片的复制副本,提供高可用与读扩展

中间件通用范式:分片解决"单点存不下",副本解决"单点挂了数据丢失"。Kafka、Redis Cluster、MySQL 分库分表都是同样思路。

flowchart LR
    A[写一条文档] --> I[(索引 Index
类似一张表)] I --> S1[分片 1] I --> S2[分片 2] I --> S3[分片 3] S1 --> R1[副本] S2 --> R2[副本] S3 --> R3[副本]

这张图在讲:一个索引被切成多个分片分布在不同节点,每个分片还有副本。数据既"装得下"(分片)又"挂不了"(副本)。

1.3 倒排索引

倒排索引是 ES 的核心,理解它是理解 ES 性能的关键。

正排索引 vs 倒排索引

  • 正排:找到数据 → 知道它的属性(按 ID 查文档)
  • 倒排:从属性出发 → 找到包含该属性的数据(按关键词查文档)

生动比喻

  • 正排:HR 找上你,问你有什么技能
  • 倒排:HR 在网站输入"5年 Go"这个关键字,搜出一堆人的简历

白话直觉:想象你有一堆文章,每篇都有编号。最笨的办法是"翻开每篇文章,看里面有没有这个词"——这就是正排,慢。聪明办法是提前做一本"词 → 文章编号列表"的字典:查 Go 直接翻字典拿到 [1,2],查 Java 拿到 [3]。这本字典就是倒排索引(Inverted Index)。“倒排"的意思是——索引方向被反过来:从"文档找词"变成"词找文档”。

倒排索引的组织方式

ES 为每个字段单独维护一张"词项 → 文档列表"的表,称为 Posting List

name 字段的 Posting List:
+------+----------+
| 词项 | 文档ID   |
+------+----------+
| Tom  | [1, 3]   |
| Bob  | [2]      |
+------+----------+

desc 字段的 Posting List:
+------+----------+
| 词项 | 文档ID   |
+------+----------+
| Go   | [1, 2]   |
| Java | [3]      |
+------+----------+
  • 搜索 name = Tom → 去 name 表找
  • 搜索 desc = Go → 去 desc 表找
  • 词项必须排序,否则查询需要全表扫描
flowchart TD
    Q[搜索 desc = Go] --> D[翻 desc 字段的倒排字典]
    D --> P[Posting List: 文档 1, 2]
    P --> R[直接返回文档 1 和 2]

这张图在讲:搜索时不需要扫描文章正文,直接查"词→文档"的字典,O(log n) 拿到结果。这就是 ES 快的秘密。

1.4 FST(Finite State Transducers)

Posting List 中的词项需要在内存中组织,ES 使用 FST 这种类似前缀树(Trie)的结构,但更进一步:前缀和后缀都共享

前缀树(共享前缀):
  0 → 1 → 2(c) → 3(t)    "ct"
  0 → 1 → 2(a) → 3(t)    "cat"

FST(前后缀都共享):
  0 → 1 → 2(c/a) → 3(t)
  - 路径 0→1→3 表示 "ct"
  - 路径 0→1→2→3 表示 "cat"

背景知识:Gin 的路由树就是前缀树。FST 在前缀树基础上压缩了后缀,进一步节省内存,适合海量词项场景。

flowchart LR
    O[0] --> N1[1]
    N1 --> CA[2: c/a 分支]
    CA --> T[3: t]
    T -. 经分支2 .-> CAT[cat]
    T -. 不经分支2 .-> CT[ct]

这张图在讲:FST 把 catct 的公共部分(ct)只存一份,内存占用远小于朴素前缀树。

1.5 节点类型

一个 ES 实例就是一个节点,可以同时扮演多个角色(但生产环境建议专节点专角色):

节点类型作用
候选主节点(Master-eligible)可被选举为主节点,负责集群元数据管理
协调节点(Coordinating Node)协调请求分发与结果聚合,类似分库分表的代理
数据节点(Data Node)实际存储数据,承担读写与查询压力

生产建议:尽量不要让节点兼职,让每个节点专注于单一角色,避免数据节点的高负载影响主节点稳定性。

1.6 ES 写入流程(面试热点)

ES 的"近实时"特性源于其写入流程:

1. 文档写入 Buffer(ES 自己的内存缓冲区)
            ↓  定时 refresh(默认 1 秒)
2. 刷新到 Page Cache(操作系统缓存)→ 生成 segment
            ↓  定时 flush
3. 刷盘 + 记录 Commit Point

关键理解

  • 只有进入 Page Cache 的数据才能被搜索到,这就是"近实时"的原因(默认 1 秒延迟)
  • 真正落盘(flush)的代价很高,间隔较长
  • 如果节点在 Page Cache 阶段崩溃,数据会丢失(可通过 translog 降低风险)
flowchart TD
    W[写文档] --> B[Buffer 内存缓冲]
    B -->|每秒 refresh| P[Page Cache
生成 segment 可被搜] P -->|定时 flush| D[(磁盘 + Commit Point)] B -. 崩溃丢数据 .-> X[靠 translog 兜底]

这张图在讲:数据先进 Buffer,1 秒后进 Page Cache 就能被搜到了(这就是为什么叫"近实时"),真正落盘是后面才发生的事。

⚠️ 新手必踩的坑:为什么刚写入搜不到。很多新手写完一条文档立刻去查,发现查不到,以为是代码错了。其实是 refresh 还没触发(默认 1 秒)。测试时可以手动调用 _refresh,或者耐心等 1 秒。生产上不要为了"立刻能搜"把 refresh 间隔调成 0,那样会疯狂生成 segment,性能崩。

1.7 Docker 部署单节点 ES

# docker-compose.yml 简化配置
services:
  elasticsearch:
    image: docker.elastic.co/elasticsearch/elasticsearch:8.0.0
    environment:
      - discovery.type=single-node    # 单节点模式
      - xpack.security.enabled=false  # 禁用 xpack(生产不要禁用)
    ports:
      - "9200:9200"

1.8 ES HTTP API 基本用法

创建索引

PUT /user_idx
{
  "settings": {
    "number_of_shards": 3,    # 3 个分片
    "number_of_replicas": 2   # 2 个副本
  },
  "mappings": {
    "properties": {
      "nickname": { "type": "text" },
      "email":    { "type": "keyword" },
      "phone":    { "type": "keyword" }
    }
  }
}

写入文档

POST /user_idx/_doc
{
  "nickname": "大明",
  "email": "john@example.com",
  "phone": "13800000000"
}

查询文档

POST /user_idx/_search
{
  "query": {
    "match": { "email": "example.com" }
  }
}

1.9 字段类型:text vs keyword

ES 支持非常多的字段类型,最需要分清的是 textkeyword

类型是否分词适用场景优缺点
keyword不分词精确匹配、聚合(标签、状态码、邮箱)索引快;不支持全文搜索
text分词全文搜索(文章、描述)支持模糊、相关性;性能开销大

选择原则

  • 邮箱、手机号、状态码 → keyword(精确匹配)
  • 文章标题、内容、描述 → text(全文搜索)
  • 需要既精确匹配又全文搜索 → 用 fields 同时定义两种类型

⚠️ 新手必踩的坑:用 termtext 字段查不到。这是最高频的错误。text 字段写入时会被分词(大明 可能变成 ),而 term 查询是不分词精确匹配词项。所以 term: {nickname: "大明"}text 字段基本查不到,得用 match 或把字段设成 keyword。一句话记忆:精确匹配用 term + keyword,全文搜索用 match + text

1.10 常用查询类型

查询类型说明适用场景
Match Query全文匹配(分词后查询)文本搜索
Term Query精确匹配keyword 字段
Range Query范围查询数字、日期
Bool Query组合查询(must/should/must_not)复杂逻辑
Match Phrase Query短语匹配(保持顺序)“Go 后端” 这种短语
Prefix Query前缀查询自动补全
Wildcard Query通配符查询模糊匹配(性能差慎用)
Fuzzy Query模糊查询(编辑距离)容错搜索
Nested Query嵌套文档查询关联对象
Aggregation Query聚合查询统计、分组、求和

Range Query 操作符gt 大于 / gte 大于等于 / lt 小于 / lte 小于等于

1.11 Go 客户端基本使用

使用官方库 github.com/elastic/go-elasticsearch/v8

package es

import (
    "context"
    "fmt"
    "github.com/elastic/go-elasticsearch/v8"
)

// NewClient 初始化 ES 客户端
func NewClient(addresses []string) (*elasticsearch.Client, error) {
    // 步骤 1:构造客户端配置(可配置用户名密码、TLS 等)
    cfg := elasticsearch.Config{
        Addresses: addresses,
    }
    // 步骤 2:创建客户端
    es, err := elasticsearch.NewClient(cfg)
    if err != nil {
        return nil, fmt.Errorf("创建 ES 客户端失败: %w", err)
    }
    // 步骤 3:健康检查,确认集群可达
    resp, err := es.Info()
    if err != nil {
        return nil, err
    }
    defer resp.Body.Close()
    fmt.Println("ES 连接成功", resp)
    return es, nil
}

官方 SDK 评价:讲义中明确提到"官方 SDK 难用"——本质上就是构造 HTTP 请求,没有做太多抽象。实践中常使用第三方封装库(如 olivere/elastic),下面搜索服务部分会展示。


二、搜索服务设计与实现

2.1 搜索服务的价值

平台搜索栏的核心目的:

  1. 方便用户快速找到所需内容(提升体验)
  2. 增加广告收入(搜索结果中插广告)
  3. 控制流量分发(信息茧房的构建入口)

表面是"按相关性排序",实际上平台可以通过搜索栏控制流量分发,暴露金主内容。

白话类比:搜索栏就像一个商场的导购台。表面上是帮你找店铺(用户体验),实际上导购先推的永远是"交了钱的店铺"(流量分发 + 广告)。理解这一点,你就懂了为什么搜索服务要做成"平台"——因为它牵动利益。

2.2 设计要考虑的点

  • 理解用户意图(查询分类、意图识别)
  • 多样化搜索选项(精确 / 模糊 / 分类)
  • 结果排序优化(相关性 + 业务策略)
  • 个性化搜索(结合用户画像)
  • 隐私与数据安全
  • 高并发、高性能、高可用

2.3 功能范围

本课程实现的简化版搜索功能:

  • 搜索具体的用户
  • 搜索某篇文章
  • 提供接口允许 Article、User 服务同步数据过来

2.4 推送接口设计:统一 vs 分离

关键决策:作为一个平台方,怎么对接不同业务方?

方案优点缺点
通用接口(接收 JSON)业务方无需适配,接入快失去编译期检查
业务专属接口编译期类型安全,调用清晰每次接入需新增接口

本课程选择:分开接口(保证编译期检查)。最佳实践:同时提供两类接口,紧急需求走通用接口,常规接入走业务接口。

flowchart LR
    U[User 服务] -->|业务专属接口| S[搜索服务]
    A[Article 服务] -->|业务专属接口| S
    X[紧急新业务] -->|通用 JSON 接口| S

这张图在讲:常规业务走"专属接口"(有类型保护),救火的新业务走"通用接口"(免适配)。两类并存是平台化的常态。

2.5 用户索引定义与推送

// User 索引定义
// 隐私字段(email、phone)的处理因产品形态而异:
// - toC 公开平台:不允许通过邮箱 / 手机号搜索
// - toB 私有 IM:需要支持邮箱 / 手机号搜索
type User struct {
    ID       int64  `json:"id"`
    Nickname string `json:"nickname"` // 通常设为 text + keyword 双类型
    Email    string `json:"email"`    // keyword 类型
    Phone    string `json:"phone"`    // keyword 类型
}
// SaveUser 保存用户数据到 ES
// 关键:指定文档 ID = 用户 ID,实现 Insert Or Update 语义
// 这样用户注册、更新信息时,同一用户始终对应同一份文档
func (r *UserESRepository) SaveUser(ctx context.Context, user domain.User) error {
    // 步骤 1:把领域对象映射成 ES 文档结构
    doc := User{
        ID:       user.ID,
        Nickname: user.Nickname,
        Email:    user.Email,
        Phone:    user.Phone,
    }
    // 步骤 2:序列化成 JSON body
    body, _ := json.Marshal(doc)
    // 步骤 3:用文档 ID 作为 ES 文档的 _id,达到 Upsert 效果
    _, err := r.es.Index().
        Index("user_idx").
        Id(fmt.Sprintf("%d", user.ID)).  // 关键:指定 _id
        Body(string(body)).
        Do(ctx)
    return err
}

2.6 帖子索引定义与推送

// Article 帖子索引定义
// 注意:通常不存储作者信息(如作者名),因为搜索人名应该走用户索引
// 保留 status 字段:文章设为不可见时,不应被搜索到
type Article struct {
    ID      int64  `json:"id"`
    Title   string `json:"title"`    // text 类型,全文搜索
    Content string `json:"content"`  // text 类型
    Status  uint8  `json:"status"`   // keyword 类型,过滤已发表状态
}

2.7 搜索接口设计

不同业务方对搜索的要求差异很大(优先展示、类目、个性化),所以搜索接口很难做成统一。本课程设计一个全站模糊搜索接口,分别查询 User 和 Article,在 Service 层聚合结果。

// SearchService 在 Service 层聚合多业务搜索结果
// 大型平台的完整流程:
// 1. 输入预处理(去空格、分词、同义词替换)
// 2. 解析表达式生成搜索计划
// 3. 执行搜索计划
// 4. 根据业务规则调整最终输出
// 本课程实现简化版:分别查询 User 和 Article,聚合返回
func (s *SearchService) Search(ctx context.Context, uid int64, expression string) (SearchResult, error) {
    // 步骤 1:输入预处理:trim、切割关键字
    keywords := strings.Fields(strings.TrimSpace(expression))
    // 步骤 2:并发查询多个业务方(可使用 errgroup 加速)
    var (
        wg    sync.WaitGroup
        users []domain.User
        arts  []domain.Article
    )
    wg.Add(2)
    go func() {
        defer wg.Done()
        users, _ = s.userRepo.Search(ctx, keywords)
    }()
    go func() {
        defer wg.Done()
        arts, _ = s.articleRepo.Search(ctx, keywords)
    }()
    wg.Wait()
    // 步骤 3:聚合返回
    return SearchResult{Users: users, Articles: arts}, nil
}
flowchart TD
    Q[用户输入 expression] --> P[预处理: trim 分词]
    P --> C[并发查询]
    C --> U[查 User 索引]
    C --> A[查 Article 索引]
    U --> M[聚合结果]
    A --> M
    M --> R[返回 SearchResult]

这张图在讲:一次全站搜索 = 预处理 + 并发查多个索引 + 聚合。Service 层把"查谁、怎么聚"统一封装,业务方无感。

Article 搜索实现

// SearchArticle 利用 ES 多关键字匹配特性
// 必须过滤 status = 已发表
func (r *ArticleESRepository) SearchArticle(ctx context.Context, keywords []string) ([]domain.Article, error) {
    // 步骤 1:构造 bool 查询:must 多字段匹配 + filter 状态过滤
    query := map[string]interface{}{
        "query": map[string]interface{}{
            "bool": map[string]interface{}{
                "must": []interface{}{
                    map[string]interface{}{
                        "multi_match": map[string]interface{}{
                            "query":  strings.Join(keywords, " "),
                            "fields": []string{"title", "content"},
                        },
                    },
                },
                "filter": []interface{}{
                    map[string]interface{}{
                        "term": map[string]interface{}{"status": uint8(2)}, // 已发表
                    },
                },
            },
        },
    }
    // 步骤 2:序列化并发送给 ES
    body, _ := json.Marshal(query)
    resp, err := r.es.Search().
        Index("article_idx").
        Body(string(body)).
        Do(ctx)
    if err != nil {
        return nil, err
    }
    // 步骤 3:解析返回结果集...
    return parseArticles(resp), nil
}

2.8 proto import 路径

search.protoimport "search/v1/sync.proto";,路径由 buf 配置的 proto 根目录决定(如 webook/api/proto),合并起来就是完整路径。

2.9 搜索表达式:谁来解析?

SearchRequest 中用 expression 字段代表用户输入的关键字,是一种通用设计,参考 GitHub 搜索栏的复杂语法。

意义:让 web/bff 层与业务服务直接透传用户输入,搜索服务自己解析、自己执行。这样未来扩展搜索语法(如 lang:go stars:>1000)不需要改动上层。

2.10 主动使用消息队列削峰

性能瓶颈时引入 Kafka:

方案 A:为不同业务方定义不同的 Event,业务方发送到特定 Topic 方案 B:定义统一 Event 格式(包含 index_namedata),数据直接穿透到 DAO

业务方 → Kafka(统一 Topic) → 消费者 → Service.Save → ES
                                    ↘
                                     gRPC 接口 → Service.Save → ES

关键:消费者和 gRPC 接口最终都调用 Service 接口,保证一致性。

flowchart LR
    B[业务方] -->|同步 gRPC| S[Service.Save]
    B -. 高并发时 .-> K[Kafka 统一 Topic]
    K --> C[消费者]
    C --> S
    S --> ES[(ES 索引)]

这张图在讲:写入 ES 有两条路——同步 gRPC 直写,和高并发时走 Kafka 削峰。两条路最终都收敛到同一个 Service.Save,保证行为一致。

2.11 被动使用消息队列

业务方强势时不适配你的接口,要求你监听它们的消息:

  • 方案 A:编写不同消费者实现,直接调用 Service
  • 方案 B:编写不同消费者,统一转发到自己的统一 Topic(推荐)

如果已经设计了统一接入 Event,方案 B 更好。

2.12 高可用:隔离部署

思路 1:读写隔离——同步数据服务一个集群,查询服务一个集群

思路 2:业务隔离——首页搜索栏 / 核心页面搜索栏走一个集群,其他业务搜索走另一个集群

核心理念:保大不包小。通过分组、降级,力保核心业务可用;多个核心业务时,一个崩溃不影响其他。

flowchart LR
    W[写集群
同步数据] -. 独立 .-> ES1[(ES 写入集群)] R[读集群
用户查询] -. 独立 .-> ES2[(ES 查询集群)]

这张图在讲:读写隔离部署,写挂了不影响用户搜,搜挂了不影响数据同步。核心思路是"别让小业务的故障拖垮核心业务"。


三、标签服务

3.1 标签功能的价值

  • 提升用户体验:方便用户找到关心的内容
  • SEO 优化:为网页提供关键词
  • 社交分享:增加内容传播效果
  • 个性化推荐:通过用户对标签的点击行为推断兴趣

白话类比:标签就像给书贴的便利贴分类。“Go"“后端"“面试"这些便利贴贴上去,别人想找同类内容,直接按便利贴抽出来就行。没有标签,所有文章就是一堆没分类的纸,找起来要翻烂。

3.2 标签形态

层级标签(toB 常见):

公司标签 → 部门标签 → 小组标签 → 个人标签
个人可用 = 自己定义的 + 所在组织链上定义的全部

通用标签 vs 分组标签

  • 通用标签:可给任何资源打
  • 分组标签:只能给特定资源(如只能给联系人、只能给文章)

3.3 标签服务用例

  • 创建标签
  • 加载用户全部标签(打标签前预加载)
  • 给资源打标签(覆盖式:原 A/B/C → 新 B/C/D,直接覆盖)
  • 查看资源上的标签

大部分系统不提供删除标签、更新标签的功能,标签一旦创建即不可变。这一前提对缓存设计很关键。

3.4 表结构设计

// Tag 标签表
type Tag struct {
    ID   int64  `gorm:"primaryKey;autoIncrement"`
    UID  int64  `gorm:"index:idx_uid_name,unique,priority:1;comment:创建者"`
    Name string `gorm:"type:varchar(128);index:idx_uid_name,unique,priority:2;comment:标签名"`
    Ctime int64
    Utime int64
}
// 唯一索引 (uid, name) 防止同一用户创建重名标签

// TagBiz 资源-标签关联表
type TagBiz struct {
    ID      int64  `gorm:"primaryKey;autoIncrement"`
    UID     int64  `gorm:"index:idx_uid_biz,priority:1;comment:打标签的人"`
    Biz     string `gorm:"type:varchar(128);index:idx_uid_biz,priority:2"`
    BizID   int64  `gorm:"index:idx_uid_biz,priority:3"`
    TID     int64  `gorm:"index:idx_tid;comment:标签ID"`
    Ctime   int64
    Utime   int64
}

3.5 创建标签 + 缓存

// CreateTag 创建标签
// 标签一旦创建即不可变,所以可以直接放入缓存
func (s *TagService) CreateTag(ctx context.Context, uid int64, name string) (int64, error) {
    // 步骤 1:落库创建标签
    id, err := s.repo.CreateTag(ctx, domain.Tag{UID: uid, Name: name})
    if err != nil {
        return 0, err
    }
    // 步骤 2:假设单用户无并发操作,直接将新标签加入缓存
    _ = s.cache.Append(ctx, uid, name)
    return id, nil
}

缓存实现:使用 Redis List

// Append 将标签追加到用户标签列表
// 前提:标签不可变(无更新、无删除),所以可以放心缓存
// 使用 pipeline 提高性能;也可用 lua 或 TxPipeline
func (c *TagCache) Append(ctx context.Context, uid int64, name string) error {
    return c.client.RPush(ctx, tagKey(uid), name).Err()
}

为什么用 List? 标签一旦创建即不可变,无需更新删除,List 是天然的适合容器。如果标签可变,应改用 Hash 或 Set。

3.6 查找用户标签:缓存预加载

// PreloadUserTags 启动时批量加载所有用户标签到缓存
// 适用于 toB 场景:所有用户标签总数不大,可全部缓存
func (r *TagRepository) PreloadUserTags(ctx context.Context) error {
    // 步骤 1:分批游标,从 UID=0 开始往后扫
    const batchSize = 100
    var lastUID int64 = 0
    for {
        // 步骤 2:按批从 DB 拉取
        tags, err := r.findBatchByUID(ctx, lastUID, batchSize)
        if err != nil {
            return err
        }
        if len(tags) == 0 {
            break
        }
        // 步骤 3:按 UID 分组,pipeline 写入 Redis
        groups := groupByUID(tags)
        for uid, names := range groups {
            _ = r.cache.BatchAppend(ctx, uid, names)
        }
        lastUID = tags[len(tags)-1].UID
    }
    return nil
}

// 在应用启动时调用
func main() {
    // ... 初始化 Repository
    repo := tag.NewTagRepository(...)
    // 启动后预加载缓存
    go func() {
        if err := repo.PreloadUserTags(context.Background()); err != nil {
            log.Warn("预加载标签缓存失败", logger.Error(err))
        }
    }()
}
flowchart TD
    S[应用启动] --> L[分批扫 DB]
    L --> G[按 UID 分组]
    G --> R[Pipeline 写入 Redis]
    R --> C{还有更多?}
    C -- 是 --> L
    C -- 否 --> D[预加载完成]

这张图在讲:缓存预加载就是"启动时把全量标签分批灌进 Redis”,之后查询走内存/Redis,不再打 DB。

缓存预加载:在极端追求高性能的场景下,必须使用本地缓存预加载,否则一次缓存未命中性能就不达标。

3.7 为资源打标签(覆盖语义)

// AttachTags 给资源打标签,覆盖原有标签
func (s *TagService) AttachTags(ctx context.Context, uid int64, biz string, bizID int64, tids []int64) error {
    return s.repo.AttachTags(ctx, uid, biz, bizID, tids)
}

// 实现思路:
// 1. 先删除该 (uid, biz, biz_id) 下的所有关联
// 2. 再插入新的关联
func (r *TagRepository) AttachTags(ctx context.Context, uid int64, biz string, bizID int64, tids []int64) error {
    // 步骤 1:开事务,先删除旧关联
    return r.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
        // 1. 删除旧关联
        if err := tx.Where("uid = ? AND biz = ? AND biz_id = ?", uid, biz, bizID).
            Delete(&TagBiz{}).Error; err != nil {
            return err
        }
        // 2. 插入新关联
        if len(tids) == 0 {
            return nil
        }
        rows := make([]TagBiz, 0, len(tids))
        for _, tid := range tids {
            rows = append(rows, TagBiz{UID: uid, Biz: biz, BizID: bizID, TID: tid})
        }
        return tx.Create(&rows).Error
    })
}

⚠️ 新手必踩的坑:覆盖式打标签要放事务里。如果"删旧"成功、“插新"失败(比如中途 DB 抖动),就会出现"标签凭空消失"的中间态。务必用事务包住"先删后插”,要么全成要么全败。另外并发打同一个资源的标签时,两个事务可能互相覆盖,生产上要对 (uid, biz, biz_id) 加锁或串行化。

3.8 获取资源上的标签(GORM Preload)

// GetTags 查询资源上的标签
// 模型始终是"某人给某资源打了某标签",所以查询要带 uid
// 支持系统默认标签:用 uid = -1 表示
//   合并 uid = 123(真实用户)和 uid = -1(系统)的结果
func (r *TagRepository) GetTags(ctx context.Context, uid int64, biz string, bizID int64) ([]domain.Tag, error) {
    // 步骤 1:查关联表,并用 Preload 一次 JOIN 出标签名
    var res []TagBiz
    err := r.db.WithContext(ctx).
        Preload("Tag").  // 通过 Preload 一次 JOIN 出标签名
        Where("uid IN ? AND biz = ? AND biz_id = ?", []int64{uid, -1}, biz, bizID).
        Find(&res).Error
    if err != nil {
        return nil, err
    }
    // 步骤 2:转换为 domain.Tag
    tags := make([]domain.Tag, 0, len(res))
    for _, tb := range res {
        tags = append(tags, domain.Tag{ID: tb.Tag.ID, Name: tb.Tag.Name})
    }
    return tags, nil
}

冗余字段优化:如果在 TagBiz 中冗余 Tag.Name,就不需要 JOIN 查询。但前提是 Tag.Name 不可变。冗余字段的本质是"用空间换时间”,且只适合冗余不可变字段,否则会引发数据不一致。


四、标签接入搜索服务

4.1 标签如何影响搜索排序

默认情况下 ES 按相关性排序,但无法满足业务需求。实践中通过赋予不同字段不同权重来干预最终排序:

  • 命中标签关键字 → 高权重
  • 命中标题 → 中权重
  • 命中正文 → 低权重

最终输出还要考虑第三方充值、合作协议等业务因素。

4.2 集成搜索的方式

方案描述优点缺点
直接推送标签服务创建独立 ES 索引通用性最强搜索时需跨索引关联
二次封装文章服务封装标签,整合进文章索引单索引查询简单标签与文章耦合

本课程采用方案一:标签独立索引。原因是标签始终与个人挂钩(uid 是必填字段),不适合放进 Article 或 User 索引。

注意:这导致后续无法使用 Nested 文档解决索引关联查询。

4.3 推送标签到搜索:统一接入

借助标签服务演示通用接入

// Kafka 消息定义(与 sync.proto 一致)
type SyncDataEvent struct {
    IndexName string                 `json:"index_name"`  // 索引名
    Data      map[string]interface{} `json:"data"`        // 文档数据
}

// 推送标签
func (s *TagService) SyncTagToSearch(ctx context.Context, tag domain.Tag) error {
    // 步骤 1:组装统一格式事件(index_name + data),与搜索服务约定
    event := SyncDataEvent{
        IndexName: "tag_idx",
        Data: map[string]interface{}{
            "id":   tag.ID,
            "uid":  tag.UID,
            "name": tag.Name,
        },
    }
    // 步骤 2:序列化并发送;用 hash(uid) 作 key 保证同用户标签进同一分区,保序
    msg, _ := json.Marshal(event)
    // 同步推送(标签打标不是高并发场景,Kafka 性能足够)
    // 注意:用 hash(uid) 作为 key,保证同一用户的标签消息发到同一 partition,保持顺序
    return s.producer.SendMessage(ctx, "sync_data", fmt.Sprintf("%d", tag.UID), msg)
}

4.4 标签索引定义

// tag_idx 索引
// 标签始终与个人挂钩,不适合放进 Article 或 User 索引
type TagDoc struct {
    ID   int64  `json:"id"`
    UID  int64  `json:"uid"`
    Name string `json:"name"`  // text 类型,支持模糊搜索
}

4.5 跨索引关联:多次查询方案

ES 内置的跨文档关联方案有:

方案机制缺点
Nested(内嵌文档)文档中嵌套其他文档不适合我们的模型(标签与个人挂钩)
Parent-Child(父子文档)父子文档关联性能差,实践中少用

本课程采用方案多次查询

// 步骤 1:先查标签索引,找到匹配的 biz_id
// 步骤 2:再查 article 索引,叠加是否命中标签的条件,并赋予高权重(boost)
func (r *SearchRepository) SearchArticleWithTagBoost(ctx context.Context, keywords []string) ([]domain.Article, error) {
    // 步骤 1:查询标签索引,拿到命中关键字的 biz_id 列表
    tagBizIDs, err := r.findBizIDsByTag(ctx, keywords)
    if err != nil {
        return nil, err
    }
    // 步骤 2:查询文章索引,用 bool.should 叠加基础匹配 + 标签命中 boost
    query := map[string]interface{}{
        "query": map[string]interface{}{
            "bool": map[string]interface{}{
                "should": []interface{}{
                    // 基础查询:title + content 匹配
                    map[string]interface{}{
                        "multi_match": map[string]interface{}{
                            "query":  strings.Join(keywords, " "),
                            "fields": []string{"title", "content"},
                        },
                    },
                    // 增强查询:命中标签的文章,boost 提升权重
                    map[string]interface{}{
                        "terms": map[string]interface{}{
                            "id":    tagBizIDs,
                            "boost": 2.0,  // boost 越高越相关,默认 1.0
                        },
                    },
                },
            },
        },
    }
    body, _ := json.Marshal(query)
    resp, err := r.es.Search().Index("article_idx").Body(string(body)).Do(ctx)
    // 解析返回...
    return parseArticles(resp), nil
}
flowchart TD
    Q[搜索关键字] --> T[查 tag_idx 拿命中 biz_id]
    T --> A[查 article_idx]
    A --> B{bool.should 合并}
    B --> M[基础匹配 title/content]
    B --> H[命中标签 boost=2.0]
    H --> R[标签命中排更前]

这张图在讲:先去标签索引捞"哪些文章的标签命中了关键词”,再查文章索引时给这些文章加权(boost=2.0),于是带标签命中的文章排名更靠前。

boost 直观理解:boost 默认值 1.0,值越高越相关。如果文章标签命中搜索关键字,权重更高,排名更靠前。


五、自测题与动手练习

自测题(合上书能答出来,才算懂)

  1. 为什么叫"倒排"索引?正排索引和倒排在"查文档"时路径有什么不同?为什么词项必须排序?
  2. ES 写入后为什么默认要等约 1 秒才能搜到(近实时)?如果节点在 Page Cache 阶段崩溃,会发生什么?
  3. text 字段和 keyword 字段在分词上有什么区别?为什么用 termtext 字段常常查不到?
  4. 搜索服务为什么同时提供"业务专属接口"和"通用 JSON 接口"?两者分别适合什么场景?
  5. 标签服务"覆盖式打标签"为什么要放在事务里?如果不在事务里,可能出现什么中间态?

动手练习(建议真做一遍)

  1. 起一个 ES 并跑通基础链路:用 Docker 起单节点 ES,建 user_idx(nickname=text、email=keyword),写一条文档,分别用 match 查 nickname、term 查 email,验证"刚写入 1 秒内可能查不到"。
  2. 加一个业务方接入搜索:仿照 SaveUser,给"话题 Topic"实现索引定义 + 推送 + 搜索,体会平台化扩展。
  3. 验证标签覆盖语义与缓存:启动标签服务预加载缓存,给某资源先打 [A,B,C] 再覆盖成 [B,C,D],查库确认旧关联被整体替换,并观察 Redis 中该用户标签列表的变化。

六、本章小结

  • ES 的本质是"按词找文档":靠倒排索引(词项→Posting List)和 FST(前后缀共享压缩内存)实现高速全文检索;写入先进 Buffer,每秒 refresh 进 Page Cache 才可被搜,所以叫"近实时"。
  • text / keyword 是高频坑:text 分词、适合 match 全文搜;keyword 不分词、适合 term 精确匹配;用 term 查 text 基本查不到。
  • 搜索服务做平台化:统一接入解决共性问题(通用 JSON + Kafka 削峰),业务专属接口解决个性问题(编译期类型安全);高可用靠读写隔离、业务隔离,“保大不包小”。
  • 标签服务靠"不可变"换简单:标签创建即不可变 → 放心用 Redis List 缓存 + 启动预加载;覆盖式打标签必须事务包裹;只冗余不可变字段。
  • 标签接入搜索用多次查询 + boost:避免 Nested/Parent-Child 的性能坑,先查标签索引拿 biz_id,再查文章索引用 boost 加权提升排序。

下一章(第18章)我们转向 ELK 日志体系与 Canal 数据同步——把"系统出问题时怎么快速定位"和"数据库变更怎么实时同步到别处"这两件生产必备的本领拿下。

About Me

没什么想介绍的,一个很大众的码农…

喜欢代码,车,马,真的是 🐎

讨厌别人让我给自己的代码写注释 最厌烦别人的程序没有写注释

目标

学AI,加油!加油!