学习目标
学完本章,你应该能够:
- 用自己的话解释 ElasticSearch 是什么、解决什么问题,以及它和 MySQL 在"找数据"这件事上的本质区别。
- 讲清楚 ES 的核心概念——索引 / 文档 / 分片 / 副本,以及倒排索引、FST 为什么能让全文检索飞快。
- 设计一套平台化的搜索服务:统一接入、推送接口、搜索聚合、Kafka 削峰、隔离部署,让新业务方零痛苦接入。
- 设计标签服务(缓存预加载、覆盖式打标签、冗余字段),并把它接入搜索、用
boost权重影响排序。 - 在面试里把"通用接入解决共性问题、扩展接口解决个性问题"的平台化思维讲成一段完整故事。
前置知识(如果下面任一点生疏,先回看对应章):
- 第02章 Gin + GORM:知道一个 HTTP 接口怎么写、DAO 怎么分层。
- 第07章 Kafka:知道消息是怎么从生产者到消费者的(本章搜索服务写入要用到)。
- 第05章 缓存:知道 Redis 基本用法(标签服务缓存预加载要用到)。
- 基本的 MySQL 索引概念(会看
KEY、UNIQUE 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 把
cat和ct的公共部分(c→t)只存一份,内存占用远小于朴素前缀树。
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 支持非常多的字段类型,最需要分清的是 text 和 keyword。
| 类型 | 是否分词 | 适用场景 | 优缺点 |
|---|---|---|---|
| keyword | 不分词 | 精确匹配、聚合(标签、状态码、邮箱) | 索引快;不支持全文搜索 |
| text | 分词 | 全文搜索(文章、描述) | 支持模糊、相关性;性能开销大 |
选择原则:
- 邮箱、手机号、状态码 → keyword(精确匹配)
- 文章标题、内容、描述 → text(全文搜索)
- 需要既精确匹配又全文搜索 → 用
fields同时定义两种类型
⚠️ 新手必踩的坑:用
term查text字段查不到。这是最高频的错误。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 搜索服务的价值
平台搜索栏的核心目的:
- 方便用户快速找到所需内容(提升体验)
- 增加广告收入(搜索结果中插广告)
- 控制流量分发(信息茧房的构建入口)
表面是"按相关性排序",实际上平台可以通过搜索栏控制流量分发,暴露金主内容。
白话类比:搜索栏就像一个商场的导购台。表面上是帮你找店铺(用户体验),实际上导购先推的永远是"交了钱的店铺"(流量分发 + 广告)。理解这一点,你就懂了为什么搜索服务要做成"平台"——因为它牵动利益。
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.proto 中 import "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_name 和 data),数据直接穿透到 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,值越高越相关。如果文章标签命中搜索关键字,权重更高,排名更靠前。
五、自测题与动手练习
自测题(合上书能答出来,才算懂):
- 为什么叫"倒排"索引?正排索引和倒排在"查文档"时路径有什么不同?为什么词项必须排序?
- ES 写入后为什么默认要等约 1 秒才能搜到(近实时)?如果节点在 Page Cache 阶段崩溃,会发生什么?
text字段和keyword字段在分词上有什么区别?为什么用term查text字段常常查不到?- 搜索服务为什么同时提供"业务专属接口"和"通用 JSON 接口"?两者分别适合什么场景?
- 标签服务"覆盖式打标签"为什么要放在事务里?如果不在事务里,可能出现什么中间态?
动手练习(建议真做一遍):
- 起一个 ES 并跑通基础链路:用 Docker 起单节点 ES,建
user_idx(nickname=text、email=keyword),写一条文档,分别用match查 nickname、term查 email,验证"刚写入 1 秒内可能查不到"。 - 加一个业务方接入搜索:仿照
SaveUser,给"话题 Topic"实现索引定义 + 推送 + 搜索,体会平台化扩展。 - 验证标签覆盖语义与缓存:启动标签服务预加载缓存,给某资源先打
[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 数据同步——把"系统出问题时怎么快速定位"和"数据库变更怎么实时同步到别处"这两件生产必备的本领拿下。