MySQL 到 Elasticsearch 的一致性:Outbox、Revision、幂等写入与 Repair Job
生产系统如何把 MySQL 业务事实可靠投影到 Elasticsearch,处理双写失败、乱序覆盖、删除复活、重试、修复和 Alias 迁移。
在知识库、Wiki、RAG、商品搜索或内容平台中,常见的数据分工是:
MySQL
→ 保存页面、权限、目录、关系、任务、版本等业务事实
Elasticsearch
→ 保存面向 BM25、向量检索、过滤和聚合的搜索投影这个边界是合理的,但它带来一个比“怎么写 ES”更难的问题:
MySQL 已经提交,ES 写入失败怎么办?两次更新乱序到达时,旧数据会不会覆盖新数据?删除事件丢失后,搜索结果会不会出现幽灵文档?Mapping 升级期间,全量回填和线上增量如何安全衔接?
如果应用只是先写 MySQL,再同步调用 Elasticsearch,它实现的不是一致性,而是一次没有恢复协议的双写。
本文给出一套可以落地的完整方案:
MySQL Source of Truth
+
Transactional Outbox
+
单调 Source Revision
+
确定性全量投影
+
Elasticsearch External Version
+
Tombstone
+
重试与 Dead Letter
+
Repair Job
+
物理索引 v1/v2 与稳定 Alias资料快照:本文依据 Elastic 与 Debezium 官方文档核查,截止日期为 2026-08-05。SQL、Go 代码和阈值是工程示例,落地时应按真实 MySQL、Elasticsearch、消息系统和客户端版本验证。
一句话结论
不要追求 MySQL 与 Elasticsearch 之间不存在的分布式强事务。应该建立五条可持续验证的不变量:
- MySQL 是唯一业务事实;
- 每次业务变化都有一个不会丢失的同步任务;
- ES 只接受比当前版本更新的投影;
- 删除、重试和乱序不会让旧状态复活;
- 整个搜索索引可以从事实库重新生成。
1. Elasticsearch 是搜索投影,不是第二份事实库
以 Wiki 页面为例,MySQL 可以保存完整业务结构:
{
"id": "page_001",
"title": "齐夏人物分析",
"revision": 42,
"relations": {
"RELATED_CHARACTER": [
{
"entity_id": "character_xxx",
"entity_type": "character",
"title": "齐夏"
}
]
}
}Elasticsearch 只保存搜索需要的稳定投影:
{
"id": "page_001",
"title": "齐夏人物分析",
"relations": {
"RELATED_CHARACTER": [
"character:character_xxx"
]
},
"relation_target_ids": [
"character:character_xxx"
],
"source_revision": 42,
"projection_schema_version": 2,
"is_deleted": false
}两边字段不需要完全相同。真正需要建立的是可追溯关系:
ES 文档必须能够证明:
它来自哪个 MySQL 实体、
哪个业务 revision、
哪一版投影规则。最关键的边不是 MySQL 到 ES,而是:
业务行与 Outbox
在同一个 MySQL 本地事务中提交。只要“以后必须同步”这件事不会丢,ES 暂时失败就只是延迟,不再是永久缺口。
2. 直接双写为什么会留下数据洞
最常见的代码大致如下:
func UpdatePage(ctx context.Context, page Page) error {
if err := mysql.UpdatePage(ctx, page); err != nil {
return err
}
if err := elastic.IndexPage(ctx, page); err != nil {
return err
}
return nil
}代码很短,但存在多个不可消除的失败窗口。
| 场景 | MySQL | Elasticsearch | 最终问题 |
|---|---|---|---|
| DB 提交后进程崩溃 | 新 | 旧 | 搜索长期缺失更新 |
| ES 成功后 DB 事务回滚 | 旧 | 新 | 搜索出现不存在的状态 |
| revision 42 先于 41 写入 | 42 | 最终被 41 覆盖 | 索引数据倒退 |
| 删除提交后 ES 删除失败 | 已删除 | 仍存在 | 幽灵文档 |
| 请求超时但 ES 实际成功 | 新 | 可能已新 | 无法判断是否应该重试 |
需要特别记住:
客户端收到 timeout
不等于
Elasticsearch 一定没有完成写入。没有事件 ID、revision 和幂等协议时,调用方无法区分:
- 第一次根本没有执行;
- 第一次已经成功,只是响应丢失;
- 第一次只完成了部分批次;
- 第二次携带的投影已经过时。
为什么分布式事务通常不是答案
搜索投影通常不值得进入跨系统强事务:
- Elasticsearch 不是关系数据库事务参与者;
- 搜索可见性还受到 refresh 影响;
- 业务数据与搜索文档本来就是不同模型;
- 搜索索引应该允许异步重建;
- Mapping 升级仍然需要新索引和迁移;
- 强事务并不能替代重试、修复和可观测性。
更合理的目标是:
业务事务原子
+
索引最终一致
+
失败可恢复
+
延迟可观测
+
整库可重建3. Transactional Outbox:把同步任务写进业务事务
Outbox Pattern 的核心不是“再建一张队列表”,而是把业务变化和待发布事件放进同一个本地事务。
3.1 业务实体需要单调 Revision
CREATE TABLE wiki_page (
id VARCHAR(64) NOT NULL,
title VARCHAR(512) NOT NULL,
content MEDIUMTEXT NOT NULL,
relations_json JSON NOT NULL,
revision BIGINT UNSIGNED NOT NULL,
deleted_at DATETIME(6) NULL,
updated_at DATETIME(6) NOT NULL,
PRIMARY KEY (id)
);revision 是实体状态的单调逻辑时钟:
创建 revision = 1
第一次修改 revision = 2
第二次修改 revision = 3
删除 revision = 4
恢复 revision = 5只要某次变化会影响搜索结果,就应该递增 revision,包括:
- 标题或正文修改;
- 关系变化;
- 权限变化;
- 启停状态变化;
- 删除与恢复;
- 影响过滤或排序的元数据变化。
3.2 Polling Outbox 表
CREATE TABLE search_outbox (
id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT,
event_id CHAR(36) NOT NULL,
aggregate_type VARCHAR(64) NOT NULL,
aggregate_id VARCHAR(128) NOT NULL,
operation VARCHAR(16) NOT NULL,
source_revision BIGINT UNSIGNED NOT NULL,
status VARCHAR(16) NOT NULL DEFAULT 'PENDING',
available_at DATETIME(6) NOT NULL,
attempts INT UNSIGNED NOT NULL DEFAULT 0,
locked_by VARCHAR(128) NULL,
locked_at DATETIME(6) NULL,
last_error TEXT NULL,
created_at DATETIME(6) NOT NULL,
updated_at DATETIME(6) NOT NULL,
PRIMARY KEY (id),
UNIQUE KEY uk_search_outbox_event (event_id),
KEY idx_search_outbox_poll (status, available_at, id),
KEY idx_search_outbox_entity (
aggregate_type,
aggregate_id,
source_revision
)
);operation 至少可以包含:
UPSERT
TOMBSTONE
REBUILD3.3 同一个事务写业务行与 Outbox
func UpdatePage(
ctx context.Context,
cmd UpdatePageCommand,
) error {
return mysql.WithTx(ctx, func(tx *sql.Tx) error {
page, err := pageRepo.GetForUpdate(ctx, tx, cmd.PageID)
if err != nil {
return err
}
page.Apply(cmd)
page.Revision++
if err := pageRepo.Save(ctx, tx, page); err != nil {
return err
}
event := SearchOutboxEvent{
EventID: uuid.NewString(),
AggregateType: "wiki_page",
AggregateID: page.ID,
Operation: "UPSERT",
SourceRevision: page.Revision,
AvailableAt: time.Now(),
}
return outboxRepo.Insert(ctx, tx, event)
})
}事务结果只有两种:
业务行和 Outbox 都提交
或
业务行和 Outbox 都回滚这堵住了最危险的永久缺口:业务已经成功,但系统里根本不存在后续同步任务。
4. Polling Outbox 与 CDC Outbox
Outbox 常见两种执行方式。
4.1 应用 Worker 轮询
Worker 可以使用租约或 SKIP LOCKED 领取任务:
SELECT *
FROM search_outbox
WHERE status = 'PENDING'
AND available_at <= NOW(6)
ORDER BY id
LIMIT 100
FOR UPDATE SKIP LOCKED;优点:
- 不依赖 Kafka Connect 或 CDC 平台;
- 重试、DLQ 和人工操作容易按业务定制;
- 适合中小规模和已有任务框架的系统。
需要自己处理:
- Claim 与租约;
- Worker 崩溃后的过期恢复;
- 指数退避;
- Outbox 清理;
- 批量公平性;
- 热点实体串行化。
4.2 Debezium CDC
另一种做法是:
应用只向 Outbox 追加事件
→ Debezium 读取 MySQL Binlog
→ Outbox Event Router 转换事件
→ Kafka
→ ES Projection ConsumerDebezium 默认 Outbox 模型强调:
唯一事件 ID
聚合类型
聚合 ID
事件类型
Payload聚合 ID 可以作为消息 Key,让同一实体尽量进入同一 Partition。
优点:
- 业务请求不需要同步发送 Kafka;
- 适合已有 Kafka 与 CDC 基础设施的团队;
- 同一业务事件可以服务多个下游。
但 CDC 不会自动解决:
- 消费者幂等;
- ES 乱序写入;
- 删除复活;
- DLQ 与人工重放;
- Mapping 升级;
- 历史漂移修复。
4.3 选择建议
| 情况 | 更合适的方案 |
|---|---|
| 单一 ES 下游,规模有限 | Polling Outbox |
| 已有 Kafka 与 Debezium | CDC Outbox |
| 任务需要复杂业务重试 | Polling,或 CDC 后接任务控制表 |
| 只打算请求后开协程写 ES | 不建议,进程退出就会丢任务 |
5. Thin Event、Fat Event 与状态收敛
Outbox 事件可以只保存实体引用,也可以保存完整快照。
5.1 Thin Event
{
"event_id": "evt-001",
"aggregate_type": "wiki_page",
"aggregate_id": "page_001",
"operation": "UPSERT",
"source_revision": 42
}Worker 收到事件后重新读取 MySQL 当前状态,再生成投影。
优点:
- 事件小;
- 多个旧事件可以自然合并到最新状态;
- 投影逻辑集中;
- 修复业务数据后可直接重建。
风险:
- 处理时需要额外读 DB;
- 硬删除后可能无法读取实体;
- 如果业务表不保留删除状态,需要额外 Tombstone 表。
5.2 Fat Event
{
"event_id": "evt-001",
"aggregate_id": "page_001",
"source_revision": 42,
"operation": "UPSERT",
"payload": {
"title": "齐夏人物分析",
"relation_target_ids": [
"character:character_xxx"
]
}
}优点:
- 可以重放历史快照;
- 消费者不必重新读 DB;
- 删除事件可以携带完整信息。
风险:
- Payload 大;
- 事件 Schema 与投影代码强耦合;
- 可能把敏感数据带入消息系统;
- 旧事件可能长期保存过时投影。
5.3 推荐折中
搜索投影通常适合:
Thin Event
+
业务行软删除或独立 Tombstone
+
Worker 读取当前最新事实事件表达的是:
这个实体至少发生过一次需要重新投影的变化。Worker 的目标不是逐条重演所有中间状态,而是让 ES 收敛到当前状态。
6. 使用唯一的确定性 Projector
收到 UPSERT 后,不建议让不同模块分别 Patch ES 的一部分字段。
更稳妥的做法是读取当前完整事实,并重建整个搜索文档:
func BuildWikiProjection(page Page) WikiDocument {
relations, targets := ProjectRelations(page.Relations)
return WikiDocument{
ID: page.ID,
Title: page.Title,
Content: page.Content,
Relations: relations,
RelationTargetIDs: targets,
SourceRevision: page.Revision,
ProjectionSchemaVersion: 2,
IsDeleted: page.DeletedAt.Valid,
}
}确定性意味着:
相同业务事实
+
相同投影代码
=
完全相同的搜索文档需要统一:
- 数组去重和排序;
- 空值处理;
- 时间格式;
- Map Key;
- ID 规范;
- 标题和关系引用的生成规则;
- 哪些字段进入 Projection Hash。
6.1 为什么全量替换通常优于 Partial Update
零散 Patch 容易导致:
- 不同 Worker 互相覆盖;
- 旧字段无法被清除;
- 新投影字段只出现在新文档;
- Repair 无法证明文档是否完整;
- Mapping 升级后历史文档长期缺字段;
- 多套投影逻辑逐渐分叉。
Elasticsearch Update API 最终仍需要读取 _source、合并或执行脚本,并重新索引文档。MySQL 已经保存完整事实时,直接生成完整投影通常更简单。
7. Source Revision 与 External Version
仅有 Outbox 仍然不能阻止消息乱序。
revision 41 → Worker B 处理较慢
revision 42 → Worker A 处理较快
42 先写入 ES
41 后写入 ES如果没有版本约束,最终索引会倒退到 41。
7.1 文档字段
{
"mappings": {
"dynamic": "strict",
"properties": {
"id": {
"type": "keyword"
},
"source_revision": {
"type": "long"
},
"projection_schema_version": {
"type": "integer"
},
"projection_hash": {
"type": "keyword"
},
"is_deleted": {
"type": "boolean"
}
}
}
}source_revision 用于:
- 判断索引是否落后;
- Repair 对比;
- 线上排障;
- 审计搜索结果来自哪个业务版本。
但只把 revision 写进 _source 还不够,写入协议也要使用它。
7.2 使用 External Versioning
PUT kb_gateway_wiki/_doc/page_001
?version=42
&version_type=external
&require_alias=true正文:
{
"id": "page_001",
"title": "齐夏人物分析",
"source_revision": 42,
"projection_schema_version": 2,
"is_deleted": false
}version_type=external 的核心语义是:
请求版本严格大于当前版本
→ 接受写入
请求版本小于或等于当前版本
→ 返回版本冲突因此:
ES 当前版本 42
旧事件尝试写入 41
→ 409 Version Conflict
→ 文档不会倒退7.3 Bulk 同样要携带版本
下面使用 json 代码块表示 NDJSON 请求,每行仍然是一条独立记录:
{ "index": { "_index": "kb_gateway_wiki", "_id": "page_001", "version": 42, "version_type": "external" } }
{ "id": "page_001", "source_revision": 42, "is_deleted": false }
{ "index": { "_index": "kb_gateway_wiki", "_id": "page_002", "version": 18, "version_type": "external" } }
{ "id": "page_002", "source_revision": 18, "is_deleted": false }Bulk HTTP 返回成功不代表每个 Item 都成功。必须逐项检查:
errors
items[*].status
items[*].error至少区分:
- 写入成功;
- 旧版本冲突;
- Mapping 错误;
- 429 或临时不可用;
- 权限错误;
- 文档过大。
7.4 几种并发控制机制的边界
| 机制 | 版本来源 | 典型用途 |
|---|---|---|
external | MySQL 等外部事实库 | DB 到 ES 异步投影 |
external_gte | 外部事实库 | 允许相等版本覆盖,必须谨慎 |
if_seq_no 与 if_primary_term | Elasticsearch | 基于已读 ES 版本做条件更新 |
retry_on_conflict | Elasticsearch | Update 脚本冲突重试 |
MySQL 投影场景通常优先使用 external。
7.5 重复事件
相同 revision 重试时,严格 external version 会返回冲突。
推荐同时使用:
event_id做消费去重;- 投影保持确定性;
- 相同或更高版本冲突分类为“已应用或已被新版本覆盖”;
- 必要时读取
source_revision验证; - 不要为了消除 409 就默认改成
external_gte。
8. Worker 应读取当前状态
对于 Thin Event,Worker 可以这样收敛:
旧事件不一定要写旧状态。只要业务行已经到 42,Worker 可以直接投影 42。
func HandleSearchEvent(
ctx context.Context,
event SearchOutboxEvent,
) error {
page, err := pageRepo.GetIncludingDeleted(
ctx,
event.AggregateID,
)
if err != nil {
return fmt.Errorf("load page: %w", err)
}
if page.Revision < event.SourceRevision {
return fmt.Errorf(
"source revision moved backwards: row=%d event=%d",
page.Revision,
event.SourceRevision,
)
}
doc := BuildWikiProjection(page)
err = searchRepo.IndexExternalVersion(
ctx,
"kb_gateway_wiki",
page.ID,
page.Revision,
doc,
)
if err == nil {
return nil
}
if errors.Is(err, ErrVersionConflict) {
return verifyAlreadyApplied(
ctx,
page.ID,
page.Revision,
)
}
return err
}8.1 合并同一实体的任务
短时间内产生 41、42、43 时,不一定需要做三次完整投影。
可以:
- 按 aggregate ID 合并待处理事件;
- 只处理最大 revision;
- 其他任务标记 SUPERSEDED;
- Kafka 使用 aggregate ID 作为 Key;
- 消费者仍然读取当前最新事实。
任务合并只是吞吐优化,不能代替 external version。重试和旧批次仍然可能乱序。
9. 删除与 Tombstone
删除是最容易被低估的状态。
假设 revision 42 的页面被删除,MySQL 更新到 revision 43。一个延迟很久的 revision 41 UPSERT 仍可能在之后到达。
9.1 立即硬删除的隐患
Delete API 可以携带 external version,但 Elasticsearch 只在有限时间内保留已删除文档的版本信息。index.gc_deletes 的默认值是 60 秒。
因此下面的假设不安全:
只要 Delete 携带 external version,
任何时间到达的旧 UPSERT 都永远不能复活。极度延迟的旧事件可能在删除版本记录被回收后重新创建文档。
9.2 推荐:先写 Tombstone
MySQL:
revision: 42 → 43
deleted_at: null → 删除时间ES 写入版本 43 的 Tombstone:
{
"id": "page_001",
"source_revision": 43,
"projection_schema_version": 2,
"is_deleted": true
}所有正常搜索统一过滤:
{
"term": {
"is_deleted": false
}
}Tombstone 至少保留到:
最大消息延迟
+
最大重试窗口
+
DLQ 人工重放窗口
+
安全缓冲之后再由清理任务硬删除。
9.3 删除策略对比
| 策略 | 优点 | 风险 |
|---|---|---|
| 立即硬删除 | 索引干净 | 长延迟旧事件可能复活 |
| 永久 Tombstone | 版本保护简单 | 占空间,查询必须过滤 |
| 延迟清理 Tombstone | 兼顾安全与空间 | 需要定义 replay horizon |
有重试、DLQ 和人工重放时,通常选择延迟清理。
9.4 Repair 仍要发现 Orphan
无论采用哪种删除策略,都要定期识别:
MySQL 已删除或不存在
但
ES 仍然 is_deleted=false删除不能只依赖一次消息。
10. 重试、退避和 Dead Letter
不同错误需要不同动作。
| 错误类型 | 例子 | 处理 |
|---|---|---|
| 成功 | 200 或 201 | DONE |
| 已收敛 | 旧 revision 冲突 | 验证后 DONE 或 SUPERSEDED |
| 临时错误 | 429、503、连接超时 | 指数退避 |
| 数据错误 | strict mapping 拒绝、类型错误 | DEAD 并告警 |
| 配置错误 | Alias 不存在、无权限 | 暂停相关任务 |
| 源数据异常 | 业务行不存在 | 补 Tombstone 或 DEAD |
10.1 指数退避
func NextRetry(attempt int) time.Duration {
base := time.Second
maxDelay := 30 * time.Minute
shift := min(attempt, 10)
delay := base * time.Duration(1<<shift)
if delay > maxDelay {
delay = maxDelay
}
jitterLimit := max(delay/5, time.Millisecond)
jitter := time.Duration(
rand.Int63n(int64(jitterLimit)),
)
return delay + jitter
}需要明确:
- 最大尝试次数;
- 最大总重试时长;
- 单次超时;
- 单实体并发策略;
- 限流和熔断;
- Dead Letter 保留期;
- 人工负责人。
10.2 DLQ 不是垃圾桶
Dead Letter 至少记录:
event_id
aggregate_id
source_revision
operation
first_failed_at
last_failed_at
attempts
last_error_code
last_error_message
projection_schema_version还要提供:
- 单条重放;
- 批量重放;
- 按错误类型统计;
- 查看源业务行;
- Mapping 修复后重试;
- 标记人工解决;
- 记录谁执行了重放。
没有重放入口的 DLQ,只是延迟暴露的数据丢失。
11. Repair Job:最终一致性的保险
Outbox 解决正常变化不丢任务。Repair 解决历史 Bug、手工改库、索引误删、过期 DLQ、早期系统缺陷和灾难恢复后的漂移。
11.1 漂移类型
| 类型 | 判断 |
|---|---|
| Missing | MySQL 有活动实体,ES 没有 |
| Stale | MySQL revision 大于 ES revision |
| Ahead | ES revision 大于 MySQL revision |
| Orphan | ES 有活动文档,MySQL 已删除或不存在 |
| Schema Stale | projection schema 不是当前版本 |
| Hash Mismatch | revision 相同,但规范化投影不同 |
11.2 Projection Hash
func ProjectionHash(doc WikiDocument) (string, error) {
normalized := NormalizeProjection(doc)
payload, err := json.Marshal(normalized)
if err != nil {
return "", err
}
sum := sha256.Sum256(payload)
return hex.EncodeToString(sum[:]), nil
}不要把这些内容放进 Hash:
- 每次重试都会变化的时间;
- Worker ID;
- Trace ID;
- 随机值;
- 不稳定的 Map 迭代顺序。
11.3 Repair 必须复用 Projector
正确流程:
Repair 发现漂移
→ 插入 REBUILD Outbox
→ 走同一个 Projector
→ 走同一套 external version 写入不要维护:
在线 Worker 一套投影
Repair 脚本另一套投影
迁移工具第三套投影三套代码最终会产生三种索引语义。
11.4 Projection State
CREATE TABLE search_projection_state (
aggregate_type VARCHAR(64) NOT NULL,
aggregate_id VARCHAR(128) NOT NULL,
source_revision BIGINT UNSIGNED NOT NULL,
indexed_revision BIGINT UNSIGNED NULL,
projection_schema_version INT UNSIGNED NULL,
projection_hash CHAR(64) NULL,
last_success_at DATETIME(6) NULL,
last_error_code VARCHAR(64) NULL,
last_error_at DATETIME(6) NULL,
PRIMARY KEY (
aggregate_type,
aggregate_id
)
);它是同步控制面,不是新的业务事实库。用途包括:
- 查询投影 Lag;
- 定位最后成功时间;
- 支持运维重试;
- 降低频繁读取 ES revision 的成本;
- 生成一致性报表。
仍然需要定期抽查真实 ES,避免控制面自己也发生漂移。
11.5 扫描策略
大数据量下可以组合:
- 按更新时间增量扫描;
- 按实体 ID Hash 分桶;
- 每天扫描部分分区;
- 每周或每月完整审计;
- ES 使用
search_after; - 对高风险知识库提高频率;
- Outbox Lag 或 DLQ 异常时主动触发 Repair。
12. 三种版本不能混用
至少需要区分:
source_revision
→ 某个业务实体变化了多少次
projection_schema_version
→ 投影代码与字段契约是哪一版
physical_index_version
→ 当前物理索引是 v1、v2 还是 v3例如:
page revision = 42
projection schema = 2
physical index = kb_gateway_wiki_v3业务数据没变,投影逻辑仍可能升级:
schema 1:只写 relations
schema 2:增加 relation_target_ids
schema 3:改变内容拼接规则此时不应该伪造业务 revision。应该:
- 增加 projection schema version;
- 新建物理索引或触发 REBUILD;
- 重新生成搜索投影;
- 用 Repair 识别旧 Schema 文档。
13. v1/v2 重建与增量追平
假设当前:
kb_gateway_wiki Alias
→ kb_gateway_wiki_v1需要升级 Mapping 或投影逻辑。
13.1 为什么记录水位
全量扫描期间,线上仍会更新:
T0 记录 Outbox 水位 W0
T1 开始全量扫描
T2 page_001 更新到 revision 43
T3 全量扫描可能读到 42 或 43
T4 扫描结束
T5 回放 W0 之后的增量External Version 使全量与增量可以安全竞争:
- 全量写 42,增量写 43,最终 43;
- 全量已经写 43,增量再写 43,识别为重复;
- 全量写删除前状态,增量写 Tombstone 44,最终删除;
- 增量乱序,更高 revision 获胜。
13.2 从旧 ES Reindex 还是从 MySQL 重建
| 情况 | 建议 |
|---|---|
只改变物理设置,旧 _source 完整可信 | 可以考虑 Reindex |
| Mapping 改变但文档结构不变 | Reindex 或 DB 重建 |
| 投影逻辑、权限、关系或 Chunk 结构变化 | 从事实库重建 |
| 旧索引可能已经漂移 | 从事实库重建 |
| 向量需要重新计算 | 走完整 Embedding 流程 |
Reindex 不会自动复制源索引的 Mapping、Shard、Replica 和 Template。目标索引必须提前创建。
13.3 Alias 原子切换
{
"actions": [
{
"remove": {
"index": "kb_gateway_wiki_v1",
"alias": "kb_gateway_wiki"
}
},
{
"add": {
"index": "kb_gateway_wiki_v2",
"alias": "kb_gateway_wiki",
"is_write_index": true
}
}
]
}Aliases API 可以在一个原子操作中完成多项动作,应用不会经历 Alias 为空的中间状态。
所有写入建议使用:
require_alias=true这样代码误写物理索引名时会失败,而不是继续把数据写回旧 v1。
14. Read-your-write 与新鲜度契约
最终一致系统不能假装搜索结果与数据库实时完全一致。
应该明确:
正常索引延迟目标
严重 Lag 告警阈值
强一致读回 MySQL 的场景
允许搜索延迟的场景用户刚保存页面后,可以采用:
- 直接展示保存接口返回的 MySQL 实体;
- 前端短时间将新实体与搜索结果合并;
- 对少量关键写入使用
refresh=wait_for; - 返回 revision,查询 Projection State;
- 强一致页面直接读取 MySQL。
搜索结果也可以返回:
{
"id": "page_001",
"source_revision": 42,
"projection_schema_version": 2
}这样可以判断结果是否低于用户刚保存的 revision。
15. 可观测性
15.1 Outbox
outbox_pending_count
outbox_oldest_pending_age_seconds
outbox_processing_count
outbox_retry_rate
outbox_dead_count
outbox_claim_timeout_count15.2 ES 写入
es_bulk_item_success_rate
es_bulk_item_429_rate
es_mapping_error_rate
es_version_conflict_rate
es_write_latency_p50
es_write_latency_p95
es_write_latency_p99
es_tombstone_count15.3 一致性
projection_revision_lag
projection_missing_count
projection_stale_count
projection_orphan_count
projection_hash_mismatch_count
projection_schema_stale_count15.4 重建与 Alias
reindex_scanned
reindex_indexed
reindex_failed
reindex_incremental_lag
alias_current_physical_index
v1_v2_query_diff_rate15.5 Trace
{
"event_id": "evt-001",
"aggregate_type": "wiki_page",
"aggregate_id": "page_001",
"event_revision": 41,
"loaded_revision": 42,
"projection_schema_version": 2,
"target_alias": "kb_gateway_wiki",
"operation": "UPSERT",
"attempt": 3,
"result": "indexed"
}这些字段让团队能够解释“为什么没更新”,而不是只看到模糊的 index failed。
16. 常见反模式
反模式一:业务请求直接双写 MySQL 和 ES
任意步骤失败都会留下没有持久恢复任务的窗口。
反模式二:先发送消息,再提交数据库
消息可能先被消费,而业务事务最终回滚。Outbox 必须与业务行同事务提交。
反模式三:事件没有 Revision
系统只能知道“发生过变化”,不能判断先后、乱序和漂移。
反模式四:只把 Revision 放进 _source
没有 external version 时,旧事件仍能覆盖新文档。
反模式五:默认使用 external_gte
相等版本覆盖可能掩盖投影不确定性。优先用事件去重和确定性投影处理重复。
反模式六:删除后立即硬删
长延迟旧事件可能在删除版本记录回收后重新创建文档。
反模式七:Repair 直接维护另一套写 ES 代码
线上、修复和重建会逐渐产生不同语义。
反模式八:Bulk HTTP 成功就认为整批成功
Bulk 必须检查每一个 Item。
反模式九:应用可以写物理索引名
Alias 切换后,旧实例可能继续写回 v1。
反模式十:Reindex 只复制旧 ES
旧索引已经漂移或缺字段时,只会把旧问题复制到新索引。
17. 分阶段落地
P0:堵住永久丢失和乱序覆盖
- 业务实体增加单调 revision;
- 业务变化与 Outbox 同事务提交;
- Worker 使用唯一确定性 Projector;
- ES 保存 source revision;
- 写入使用 external version;
- Bulk 逐 Item 判断;
- 失败可重试并进入 DLQ;
- 所有写入只认 Alias。
P1:补删除、修复和观测
- 删除产生新 revision;
- 定义 Tombstone 与 replay horizon;
- 建立 Projection State;
- 建立 Missing、Stale、Orphan Repair;
- 增加 Projection Hash;
- 建立 Outbox Lag 与 Mapping Error 告警;
- 提供 DLQ 人工重放。
P2:支持安全的 Schema 演进
- 物理索引使用 v1、v2;
- 保存 projection schema version;
- 建立全量扫描与增量水位回放;
- 建立 v1/v2 查询 Diff;
- Alias 原子切换与回滚;
- 延迟删除旧索引。
18. 推荐生产架构
| 组件 | 负责什么 | 不负责什么 |
|---|---|---|
| Domain Service | 业务校验、revision、本地事务 | 同步等待 ES 完成 |
| MySQL | 业务事实、Outbox、删除和审计 | 全文与向量检索 |
| Outbox 或 Kafka | 不丢失变化、传递任务 | 定义搜索文档结构 |
| Projector | 从事实生成完整 ES 文档 | 保存业务唯一事实 |
| Elasticsearch | 搜索、过滤、聚合、向量 | 关系事务和级联事实 |
| Repair Job | 发现并修复漂移 | 维护另一套投影规则 |
| Alias | 隔离物理索引版本 | 自动保证数据一致 |
总结
MySQL 与 Elasticsearch 的一致性,不是靠“多重试几次”实现的,而是靠一套明确协议:
业务变化
→ MySQL 事务内递增 revision 并写 Outbox
→ Worker 读取当前事实
→ 唯一 Projector 生成完整文档
→ external version 写入稳定 Alias
→ 重试与 DLQ 保证失败可恢复
→ Tombstone 阻止旧事件复活
→ Repair 持续发现漂移
→ v1/v2 与 Alias 支持整库重建最重要的工程判断是:
最终一致不是相信消息最终会成功,而是让每个中间状态都可识别、每个失败都可重放、每个文档都能证明来源版本,并且整个索引可以从事实库重新生成。
继续阅读:
- Elasticsearch Mapping 设计:object、flattened、nested 到关系投影
- RAG 的真实架构
- Metadata Filtering:RAG 检索里的隐形控制面
- RAG 可观测性
- RAG 生产上线检查清单
官方资料
- Debezium:Outbox Event Router
- Elastic:Index API 与 External Versioning
- Elastic:Optimistic Concurrency Control
- Elastic:Bulk API
- Elastic:Delete API
- Elastic:Update API
- Elastic:Reindex API
- Elastic:Aliases
- Elastic:General Index Settings
讨论
继续讨论这篇笔记
有问题、补充案例或不同观点,可以通过 GitHub Discussions 继续交流。
Elasticsearch Mapping 设计:object、flattened、nested 到关系投影
从 Elasticsearch 如何索引 JSON 出发,系统理解 object、flattened、nested、keyword、dynamic strict 与 mapping explosion,并用 Wiki relations 场景完成选型。
RAG Freshness 生产实战:时间语义、增量索引、版本治理与过期答案拦截
从 Event Time、Valid Time、CDC、Outbox、External Version、Tombstone、Watermark、缓存失效、时间感知检索到 Freshness SLO,构建能够解释“数据何时发生、何时可见、答案依据哪个版本”的生产级 RAG 新鲜度系统。