Chico Notes
LLM Wiki / RAG

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 之间不存在的分布式强事务。应该建立五条可持续验证的不变量:

  1. MySQL 是唯一业务事实;
  2. 每次业务变化都有一个不会丢失的同步任务;
  3. ES 只接受比当前版本更新的投影;
  4. 删除、重试和乱序不会让旧状态复活;
  5. 整个搜索索引可以从事实库重新生成。

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
}

代码很短,但存在多个不可消除的失败窗口。

场景MySQLElasticsearch最终问题
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
REBUILD

3.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 Consumer

Debezium 默认 Outbox 模型强调:

唯一事件 ID
聚合类型
聚合 ID
事件类型
Payload

聚合 ID 可以作为消息 Key,让同一实体尽量进入同一 Partition。

优点:

  • 业务请求不需要同步发送 Kafka;
  • 适合已有 Kafka 与 CDC 基础设施的团队;
  • 同一业务事件可以服务多个下游。

但 CDC 不会自动解决:

  • 消费者幂等;
  • ES 乱序写入;
  • 删除复活;
  • DLQ 与人工重放;
  • Mapping 升级;
  • 历史漂移修复。

4.3 选择建议

情况更合适的方案
单一 ES 下游,规模有限Polling Outbox
已有 Kafka 与 DebeziumCDC 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 几种并发控制机制的边界

机制版本来源典型用途
externalMySQL 等外部事实库DB 到 ES 异步投影
external_gte外部事实库允许相等版本覆盖,必须谨慎
if_seq_noif_primary_termElasticsearch基于已读 ES 版本做条件更新
retry_on_conflictElasticsearchUpdate 脚本冲突重试

MySQL 投影场景通常优先使用 external

7.5 重复事件

相同 revision 重试时,严格 external version 会返回冲突。

推荐同时使用:

  1. event_id 做消费去重;
  2. 投影保持确定性;
  3. 相同或更高版本冲突分类为“已应用或已被新版本覆盖”;
  4. 必要时读取 source_revision 验证;
  5. 不要为了消除 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 或 201DONE
已收敛旧 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 漂移类型

类型判断
MissingMySQL 有活动实体,ES 没有
StaleMySQL revision 大于 ES revision
AheadES revision 大于 MySQL revision
OrphanES 有活动文档,MySQL 已删除或不存在
Schema Staleprojection schema 不是当前版本
Hash Mismatchrevision 相同,但规范化投影不同

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 的场景
允许搜索延迟的场景

用户刚保存页面后,可以采用:

  1. 直接展示保存接口返回的 MySQL 实体;
  2. 前端短时间将新实体与搜索结果合并;
  3. 对少量关键写入使用 refresh=wait_for
  4. 返回 revision,查询 Projection State;
  5. 强一致页面直接读取 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_count

15.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_count

15.3 一致性

projection_revision_lag
projection_missing_count
projection_stale_count
projection_orphan_count
projection_hash_mismatch_count
projection_schema_stale_count

15.4 重建与 Alias

reindex_scanned
reindex_indexed
reindex_failed
reindex_incremental_lag
alias_current_physical_index
v1_v2_query_diff_rate

15.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 支持整库重建

最重要的工程判断是:

最终一致不是相信消息最终会成功,而是让每个中间状态都可识别、每个失败都可重放、每个文档都能证明来源版本,并且整个索引可以从事实库重新生成。

继续阅读:

官方资料

讨论

继续讨论这篇笔记

有问题、补充案例或不同观点,可以通过 GitHub Discussions 继续交流。

On this page

一句话结论1. Elasticsearch 是搜索投影,不是第二份事实库2. 直接双写为什么会留下数据洞为什么分布式事务通常不是答案3. Transactional Outbox:把同步任务写进业务事务3.1 业务实体需要单调 Revision3.2 Polling Outbox 表3.3 同一个事务写业务行与 Outbox4. Polling Outbox 与 CDC Outbox4.1 应用 Worker 轮询4.2 Debezium CDC4.3 选择建议5. Thin Event、Fat Event 与状态收敛5.1 Thin Event5.2 Fat Event5.3 推荐折中6. 使用唯一的确定性 Projector6.1 为什么全量替换通常优于 Partial Update7. Source Revision 与 External Version7.1 文档字段7.2 使用 External Versioning7.3 Bulk 同样要携带版本7.4 几种并发控制机制的边界7.5 重复事件8. Worker 应读取当前状态8.1 合并同一实体的任务9. 删除与 Tombstone9.1 立即硬删除的隐患9.2 推荐:先写 Tombstone9.3 删除策略对比9.4 Repair 仍要发现 Orphan10. 重试、退避和 Dead Letter10.1 指数退避10.2 DLQ 不是垃圾桶11. Repair Job:最终一致性的保险11.1 漂移类型11.2 Projection Hash11.3 Repair 必须复用 Projector11.4 Projection State11.5 扫描策略12. 三种版本不能混用13. v1/v2 重建与增量追平13.1 为什么记录水位13.2 从旧 ES Reindex 还是从 MySQL 重建13.3 Alias 原子切换14. Read-your-write 与新鲜度契约15. 可观测性15.1 Outbox15.2 ES 写入15.3 一致性15.4 重建与 Alias15.5 Trace16. 常见反模式反模式一:业务请求直接双写 MySQL 和 ES反模式二:先发送消息,再提交数据库反模式三:事件没有 Revision反模式四:只把 Revision 放进 _source反模式五:默认使用 external_gte反模式六:删除后立即硬删反模式七:Repair 直接维护另一套写 ES 代码反模式八:Bulk HTTP 成功就认为整批成功反模式九:应用可以写物理索引名反模式十:Reindex 只复制旧 ES17. 分阶段落地P0:堵住永久丢失和乱序覆盖P1:补删除、修复和观测P2:支持安全的 Schema 演进18. 推荐生产架构总结官方资料