Skip to main content

从单值到多维:一次课程搜索的 ES alias + 事件驱动改造实录

· 19 min read
Kanelli
Backend & AI Engineer

一次把单值筛选的老搜索改成多维多选、又顺手把"改筛选项就要发版"这个老毛病治好的完整过程。

目录


缘起:一个只支持单值筛选的搜索

我在做一个在线课程平台。搜索接口 /study-center/search 长这样:

GET /study-center/search?keyword=&categoryId=1&tagId=2&sort=BEST_MATCH

categoryIdtagId 都是单值。前端筛选面板只能选一个课程分类 + 一个课程标签,sort 也只有一个 BEST_MATCH

这个设计上线时够用,但业务一涨起来立刻就撑不住了:

  • 产品想加"适合角色"(研发 / 产品 / 售前 / 技术管理者 / 职能人员),一个课程可能同时适合多个角色,是多对多
  • 想加"难度"(L1–L4 四个等级),要支持多选,还要支持按难度升序 / 降序排列
  • 想加"学习形式"(沙箱 / 视频 / 混合),要能筛
  • 想加"最新 / 最热 / 按难度"三种排序

所有这些都被卡在"我们的接口只吃单值参数"这个前提上。

更糟糕的是,运营侧任何一次筛选项调整——加一个角色、改一个难度名——都得发版,因为选项是硬编码在后端 Kotlin 代码里的。

我要做的事就一句话:把搜索从单值改成多维,顺便把"改筛选项要发版"这个老毛病治了。

一、老系统的四个痛点

在动手前,我把老实现的问题列了一遍:

1. 只支持单值维度 接口签名 categoryId: Long / tagId: Long,SQL 也是 WHERE category_id = ?。想改多值就要动契约。

2. 索引是硬编码的物理名 ES 里的索引名叫 course_20251029 之类,代码里 const CourseIndexName = "course_20251029" 到处引用。想改 mapping 就得停服、切代码、发版——零停机 mapping 变更不可能。

3. Query 是 printf 模板拼出来的 搜索请求体是一个 .json 模板 + fmt.Sprintf 填参数。加一个多值字段就得改模板 + 改填充逻辑,字符串拼接的世界里没有类型安全,也没法单测。

4. 筛选项写死在代码里 Kotlin StudyCenterController.getSearchConditionlistOf("研发", "产品", "售前"...)。运营想改叫法、加一项、调排序全得走发版。

一句话总结:"扩展性 = 0,运营自助 = 0"

二、目标与几个关键取舍

我给自己定了四个目标:

  1. 多维多选筛选:4 个维度(方向 / 角色 / 难度 / 学习形式)+ 3 种排序(最新 / 最热 / 按难度)
  2. 零停机 mapping 变更:以后再加字段不需要停服
  3. 运营改筛选项 30 秒内前台生效:不要再发版
  4. 全流程灰度可控,可回滚:出问题能立刻回退

顺着这些目标,有几个关键取舍很快就明确了:

要不要做独立的"索引重建 CLI 工具"? 最初的方案是写一个 rebuild-course-index CLI,把整个重建流程封装起来。但仔细想想,它无非就是三步:拿 mapping → 灌数据 → 切 alias。三个 gRPC 就能表达清楚,再包一层 CLI 反而增加运维面和权限管理成本。决定不做,运维入口收敛到管理后台索引重建页 + 三个 RPC 组合。

要不要为每种后台变更新造独立的 Event 类? 最初拆分是 CourseDictEvent / CourseCategoryEvent / CourseTagEvent 三类事件走三个 Kafka topic。但前台的课程变更已经有一个 CourseEvent + COURSE_EVENT 通道在跑。决定复用一条 Kafka 通道,把 dict / category / tag 变更都转成"最终受影响的 courseIds 列表",用同一个 CourseEvent(OPERATION_UPSERT, courseIds) 发出去——事件语义统一、消费者不用改。

老单值参数 tag_id / category_id 要不要删? 删了简洁,但 proto wire 兼容会出问题(老客户端还在生产环境跑)。决定保留但标 DEPRECATED,Go 层加一个 mergeInt64s 兼容合并 V1 单值到 V2 多值,等下个大版本再彻底删。

这三个决定后来都被证明是对的——在扩展性和运维复杂度之间,我更愿意让代码多背一点债,也不让运维多一个组件要照看。

三、ES 索引 alias 化:mapping 变更零停机的地基

老代码里满天飞的 CourseIndexName = "course_20251029",是所有零停机改造的第一道拦路石。

ES alias 的思路很直接:读写都走 alias,不写物理索引名

// 替换掉
// const CourseIndexName = "course_20251029"

const (
CourseReadAlias = "course_read"
CourseWriteAlias = "course_write"
)

// 所有搜索请求
searchReq := esapi.SearchRequest{
Index: []string{CourseReadAlias},
Body: bytes.NewReader(body),
}

// 所有写入
bulkReq := esutil.BulkIndexerConfig{
Index: CourseWriteAlias,
// ...
}

alias 挂载的运维动作只需 30 秒(业务代码上线前先做):

POST /_aliases
{
"actions": [
{"add": {"index": "course_20251029", "alias": "course_read"}},
{"add": {"index": "course_20251029", "alias": "course_write"}}
]
}

一旦这层地基铺好,未来加字段的完整流程就变成了:

1. 建新索引 course_v20260801(带新 mapping)
2. 打开双写 write alias 挂两个索引,新旧一起写
3. 全量灌数据 后台管理页触发 refresh
4. 校验 doc count 新索引 ≥ 老索引
5. 切读 read alias 从老切到新
6. 观察 3 天
7. 停双写 write alias 摘掉老索引
8. 删老索引 DELETE /course_20251029

任意一步失败都能立刻回退——把 alias 切回去就行,业务侧完全无感。这是我在这次改造里最想推荐给别人的一个模式:如果你的 ES 索引名还写在代码里,抽空把 alias 铺上,未来任何 mapping 变更都会感谢自己。

三个小防呆

真上线时踩了三个小坑,都值得说一句:

1. alias 白名单 SwitchCourseAlias 是个高危 API,一旦被误调传错 alias 名字,能把生产读流量指到一个空索引。我给它加了个白名单:

var allowedAliases = map[string]bool{
CourseReadAlias: true,
CourseWriteAlias: true,
}

if !allowedAliases[req.Alias] {
return nil, fmt.Errorf("alias %q not in allow list", req.Alias)
}

只接受这两个别名,其他一律拒绝。

2. Preflight 检测:启动就 fail-fast 服务启动时先跑一次 EnsureCourseAliases——检测 alias 三态:

  • 都在且指向合法索引 → 正常启动
  • 索引存在但 alias 没挂 → fail-fast,服务不起来
  • alias 挂了但物理索引不存在 → 同上

比"跑着跑着搜出空结果"再排查友好一万倍。

3. unmapped_type 兜底 排序字段在老索引里可能还不存在(新加的 hotrating),直接 sort 会 500。给每个 sort 都带上:

{"hot": {"order": "desc", "unmapped_type": "float"}}

字段不存在时按 0 处理,请求不再爆。

四、Query builder:把 printf 模板换成 map+json.Marshal

老 query 是个模板文件 + fmt.Sprintf

//go:embed search_course.tpl.json
var searchTpl string

body := fmt.Sprintf(searchTpl, keyword, categoryId, tagId, from, size)

问题不用多说——多值字段没法填、参数一错就是隐式 SQL/JSON 注入面、单测约等于跑一遍 ES 才能验。

改成 map + json.Marshal

func buildCourseSearchBody(p *pb.SearchCourseParam, from, size int32) ([]byte, error) {
must := []map[string]any{}

if p.Keyword != "" {
must = append(must, map[string]any{
"multi_match": map[string]any{
"query": p.Keyword,
"fields": []string{"title^3", "title.simple^2", "description", "description.simple"},
},
})
}

// 多值维度用 terms
categoryIDs := mergeInt64s(p.CategoryIds, p.CategoryId) // V1 兼容合并
if len(categoryIDs) > 0 {
must = append(must, map[string]any{
"terms": map[string]any{"category_id": categoryIDs},
})
}
if len(p.RoleCodes) > 0 {
must = append(must, map[string]any{
"terms": map[string]any{"role_codes": p.RoleCodes},
})
}
if len(p.Levels) > 0 {
must = append(must, map[string]any{
"terms": map[string]any{"level": p.Levels},
})
}

// 单值维度用 term
if p.LearningForm > 0 {
must = append(must, map[string]any{
"term": map[string]any{"learning_form": p.LearningForm},
})
}

body := map[string]any{
"from": from,
"size": size,
"query": map[string]any{
"bool": map[string]any{
"must": must,
"must_not": []map[string]any{
{"exists": map[string]any{"field": "deleteAt"}},
},
"filter": []map[string]any{
{"term": map[string]any{"publicStatus": 1}},
},
},
},
"sort": courseSortClause(p.Sort),
}

return json.Marshal(body)
}

这一改换来三个好处:

  1. 类型安全:拼 JSON 的过程走的是 Go 结构体和 map,任何字段名/类型错误编译期就报
  2. 可单测buildCourseSearchBody 是纯函数,输入参数、输出字节,各种参数组合直接 assert JSON 结构
  3. 老 V1 参数自然兼容mergeInt64s(V2 多值, V1 单值) 一行搞定

五、Refresh + Switch:把"索引重建"降级成三个 RPC

alias 铺好之后,重建索引不再是一个"神秘的运维仪式",只是三个 RPC 的组合:

GetCourseMappings → 拿最新的 mapping JSON(embed 在二进制里)
RefreshCourseDocuments → 把 DB 数据全量灌进指定索引
SwitchCourseAlias → 把 read/write alias 指向切到新索引

后台管理页把这三个动作拼成一个向导,SRE 点几下就能完成从建索引到切读到停双写的整套流程。

RefreshCourseDocuments 里有个细节值得一提——限速

const batchSize = 200
const interval = 50 * time.Millisecond

for i := 0; i < len(courses); i += batchSize {
end := min(i+batchSize, len(courses))
if err := bulkIndex(ctx, courses[i:end]); err != nil {
return err
}
time.Sleep(interval)
}

200/批 + 50ms 间隔,不追求灌数据的极限速度,为的是让 ES 集群有喘息的空间——重建期间业务照常吃流量。数万课程灌完约几分钟,业务侧完全无感。

六、事件驱动即时刷新与 external_gte 防 stale write

搜索改多维后,还剩一个体验问题:运营在管理后台改了课程分类名,前台什么时候能看到?

老实现是"等下一次全量对账"——最坏 24 小时。这在多角色 / 多难度上线之后完全不可接受,产品诉求是 ≤5s 生效。

方案是很典型的事件驱动

Kotlin 后台改数据
↓ 事务提交后(afterCommit)
CourseMetaChangeNotifier
↓ 反查关系表拿 courseIds
Kafka topic: COURSE_EVENT

Go 消费者 (backend/liteapp)
↓ 分批 200
BulkIndexCourses (write alias)

ES doc 更新

Kotlin 侧发事件的时机很关键——必须 afterCommit

TransactionSynchronizationManager.registerSynchronization(
object : TransactionSynchronization {
override fun afterCommit() {
val courseIds = relationRepo.findCourseIdsByCategoryId(categoryId)
courseIds.chunked(500).forEach {
kafkaTemplate.send(COURSE_EVENT, CourseEvent(op = UPSERT, courseIds = it))
}
}
}
)

如果在事务里就发事件,消费者可能比 DB commit 更早读到"事件已发但 DB 还没提"的错序状态,一旦按事件反查关系表会拿到旧数据。

一个坑:stale write

事件驱动跑起来之后,我遇到了一个诡异现象:运营连续快速改一个课程的分类两次,最终 ES 里保留的是第一次的分类值。

原因是两次事件被两个不同的 Go worker 消费,第二次事件的 bulk 请求先到 ES,第一次事件后到。ES 里最后写的赢——赢的是老数据。

解法是给 bulk 请求带一个基于时间戳的版本号:

version := course.LastModifiedDate.UnixMilli()
bulkItem := esutil.BulkIndexerItem{
Action: "index",
DocumentID: strconv.FormatInt(course.ID, 10),
Body: bytes.NewReader(body),
VersionType: "external_gte",
Version: &version,
}

external_gte 的语义是"本次写入的 version 必须 ≥ 当前 doc 的 version,否则拒绝"。这样即使消息乱序到达,旧版本永远无法覆盖新版本——ES 里最终保留的一定是 LastModifiedDate 最大的那一份。

被拒绝的写入 ES 会返回 409,我把它当成"正常事件"打指标而非当错误——因为它不是失败,只是 ES 帮我做了并发裁决。

七、CronJob 兜底:二级防线的价值

事件驱动很好,但绝对不能只靠事件驱动。想想它可能失败的场景:

  • Kotlin 进程重启,事件在内存队列里丢了
  • Kafka 分区暂时不可用
  • Go 消费者版本回滚,某种 event schema 处理不了
  • 极端并发下事务提交与事件发出的错序

所以第二层防线:每天 03:00 一个 CronJob,全量重刷所有课程

apiVersion: batch/v1
kind: CronJob
metadata:
name: course-reconcile-job
spec:
schedule: "0 3 * * *"
jobTemplate:
spec:
template:
spec:
containers:
- name: reconcile
image: myapp/liteapp:latest
command: ["curl"]
args:
- "-X"
- "PUT"
- "http://liteapp:8080/api/v1/liteapps/courses/reindex"

Job 内部就是 BulkIndexingCourse——遍历所有课程、按 write alias 分批灌回 ES。每次刷完打一个 course_reconcile_courses_total counter。

这套二级防线才是让我睡得着觉的关键。 事件驱动做即时性,CronJob 做最终一致性。就算所有事件消费都挂了 24 小时,第二天早上 3 点也能拉回来。

八、字典表:把"改筛选项要发版"的坑填了

这时候前台已经完全跑通了,但还有个尴尬——三组筛选项(角色 / 难度 / 学习形式)依然硬编码在 Kotlin 代码里。

翻代码看的时候我发现,之前有人建过一个 t_course_roles 独立表,但只有 DB 存在,Kotlin 侧从没读过——实际上是"死数据"。

决定引入一个通用字典表,把三组筛选项统一承载:

CREATE TABLE t_course_dict (
id BIGINT PRIMARY KEY,
dict_type VARCHAR(32) NOT NULL, -- role / level / learning_form
code VARCHAR(64) NOT NULL, -- 稳定标识
name VARCHAR(128) NOT NULL, -- 展示名
order_index INT NOT NULL DEFAULT 0,
visible TINYINT NOT NULL DEFAULT 1,
deleted_date DATETIME NULL,

UNIQUE KEY uk_type_code (dict_type, code),
INDEX idx_type_visible_order (dict_type, visible, order_index)
);

-- 保 id 迁移老表数据
INSERT INTO t_course_dict (id, dict_type, code, name, order_index)
SELECT id, 'role', code, name, order_index FROM t_course_roles;

-- 塞入 level / learning_form 字典
INSERT INTO t_course_dict (dict_type, code, name, order_index) VALUES
('level', 'L1', 'L1 入门', 1),
('level', 'L2', 'L2 进阶', 2),
('level', 'L3', 'L3 熟练', 3),
('level', 'L4', 'L4 精通', 4),
('learning_form', 'SANDBOX', '沙箱', 1),
('learning_form', 'VIDEO', '视频', 2),
('learning_form', 'MIXED', '混合', 3);

DROP TABLE t_course_roles; -- 一次到位,不留半迁移状态

Kotlin 侧 CourseDictService 加一层 60s JVM 本地 TTL 缓存(后台改动通过 evictCache() 显式失效):

@Service
class CourseDictService(private val repo: CourseDictRepository) {
private val cache = AtomicReference<Map<String, List<CourseDict>>>(emptyMap())
private var lastLoad = 0L
private val ttl = 60_000L

fun listVisibleGrouped(): Map<String, List<CourseDict>> {
val now = System.currentTimeMillis()
if (now - lastLoad > ttl) reload()
return cache.get()
}

fun evictCache() { lastLoad = 0 } // 后台 CRUD 后主动清

private fun reload() {
val all = repo.findAllByVisibleTrueOrderByOrderIndexAsc()
cache.set(all.groupBy { it.dictType })
lastLoad = System.currentTimeMillis()
}
}

StudyCenterController.getSearchCondition 里的三段硬编码就变成了一行 courseDictService.listVisibleGrouped()

再补一套后台 CRUD(字典管理 / 分类管理 / 标签管理),运营就彻底自助了——从"改一个筛选项要发版 + 上线"变成"点几下按钮,60 秒内前台生效"。

有两条 CRUD 层的硬约束值得单独说:

  • code 创建后只读:ES doc 里会引用 code,一旦允许改就是滚雪球式的数据一致性灾难。改名走"废弃老 code + 新建新 code"路径
  • dict_type 不可迁移:跨类型移动一条字典(比如把 role 挪到 level)没有合理业务场景,全量走软删 + 新建

Controller 层针对"改 code / 硬删 / 换类型"直接返回 400,比在 Service 里再检查一遍更早失败。

九、标签合并:一个容易做错的事务

后台三个模块里最不好写的是标签合并

场景:运营发现库里有 LangChain / LangChain教程 / LangChain入门 三个语义等价的标签,想合并成一个。

朴素做法是"改改课程关联",但真上手会踩至少三个坑:

坑 1:(course_id, target_tag_id) 冲突 课程 A 同时打了 LangChainLangChain教程。要合并到前者时,直接 UPDATE tag_id = target 会撞唯一键。

坑 2:used_count 累计 tag 上有个 used_count 字段用于按热度排序,合并时要正确重算,不能漏累。

坑 3:原子性 关系重定向、软删老 tag、used_count 更新必须在同一个事务里,任何一步失败要整体回滚。

我最后落到的实现是四步走:

@Transactional
fun merge(fromIds: List<Long>, toId: Long) {
require(toId !in fromIds) { "cannot merge to self" }

// Step 1: 先删可能撞唯一键的关系
// DELETE FROM t_course_tag_relations
// WHERE tag_id IN (:fromIds)
// AND course_id IN (
// SELECT course_id FROM t_course_tag_relations WHERE tag_id = :toId
// )
relationRepo.deleteConflictingTagRelations(fromIds, toId)

// Step 2: 重定向剩余关系
// UPDATE t_course_tag_relations SET tag_id = :toId WHERE tag_id IN (:fromIds)
val moved = relationRepo.redirectTagRelations(fromIds, toId)

// Step 3: 老 tag 软删(@SQLDelete 会自动改成 UPDATE deleted_date = NOW())
tagRepo.deleteAll(tagRepo.findAllById(fromIds))

// Step 4: 重算 used_count
tagRepo.setUsedCount(toId, relationRepo.countByTagId(toId))
tagRepo.setUsedCountBatch(fromIds, 0)
}

100 条关系的合并测下来 <5s,事务里的 SQL 都是走索引的批量操作。

事务提交后由 CourseMetaChangeNotifier 反查所有受影响的 courseIds 走 Kafka 通知,前台在 5 秒内看到新数据。

十、可观测:6 条告警 + 一张 Grafana + 一份 Runbook

任何一次搜索改造都必须自带可观测。我给这套系统留了三份东西:

6 条 Prometheus 告警

名字触发条件优先级
CourseAliasBrokenread alias 未指向任何索引P0
CourseRefreshErrorSpike5min 内 refresh 失败次数陡增P1
CourseRefreshDurationHighrefresh P95 耗时 > 60sP2
CourseReconcileJobFailed每日 CronJob 失败P1
CourseAdminRefreshFailedSpike后台变更→ES 失败率 > 0P2
CourseAdminAPI5xxHigh后台接口 5xx 率 > 1%P2

第一条 CourseAliasBroken 是 P0 事故告警——read alias 一旦断开,全站课程搜索出空结果,用户立刻会投诉。

一张 Grafana 主 dashboard:QPS / P95 / 错误率 / 事件消费吞吐 / refresh 时长直方图 / 老接口调用量。

一份 Runbook:每条告警的处理步骤都在 README.md 里写死——alias 断了怎么手动接回、refresh 卡了怎么定位、CronJob 失败怎么补跑,值班时不用现问。

告警指标里最重要的一个观察点是 course_alias_health gauge——0/1 值,告警不需要看曲线,看这个 gauge 就知道 read/write alias 是不是正常。

写在最后:这次改造里几条通用的教训

回头看,这次改造真正让我留下印象的不是任何一个具体技术点,而是几条工程上的通用取舍:

1. 有些复杂度是"看起来必要"的,砍掉更好 最初拆分里有一个 CLI 工具、有三个独立事件类、有一套独立的 latency histogram。等真正开始做,发现它们要么和已有能力重合,要么增加运维负担却没换来实质收益。做减法比加法难,但更值。

2. 二级防线永远值得 事件驱动是即时性,CronJob 是最终一致性——它们不冲突,是配套关系。别指望单一机制永远正确。

3. 数据模型层的关键约束要在 Controller 层就拒 code 只读、dict_type 不可换、硬删禁止——这些约束越早拒(Controller > Service > Repository > DB constraint)越好。等到 DB 抛异常时,业务侧的错误处理链路已经跑了一半。

4. Alias 是所有 ES 系统都该有的底层能力 如果你的 ES 索引名还写在代码里,抽时间把 alias 铺上。零停机 mapping 变更这个能力,未来任何一次字段扩展都会感谢当年的自己。

5. external_gte 是事件驱动 + ES 场景的标配 只要你的更新可能被并发触发(后台 + CronJob + 事件),版本号防 stale write 就是必须的。别等出了幽灵回滚才想起。

6. 硬编码的筛选选项永远是运营效率的杀手 一旦你的产品经理开始高频问"这个选项能改吗",就是该把它挪到 DB 的信号。字典表 + 60s 缓存这个模式,几乎可以套用到任何"运营可配置项"的场景。

这次改造从技术上讲不算特别新,用的都是成熟组件(ES / Kafka / MySQL / JPA / Prometheus)。它更像是一次把老搜索系统的每个短板都补齐的过程——单值→多值、硬编码→字典化、静态→事件驱动、盲跑→可观测。每一步单独看都很朴素,但摞在一起就把一个搜索系统的天花板抬了两个身位。

有时候好架构不需要多新——把该抽象的抽象好,把该兜底的兜好,把该看的指标亮出来,就够了。


技术栈:Kotlin · Spring Boot · Go · gRPC · ElasticSearch · Kafka · MySQL · Redis · Prometheus · Grafana