Skip to content

Elasticsearch 写入、Refresh 与 Segment 全过程

很多人只会背“ES 是近实时搜索,写入后要 refresh 才能查到”。这句话没错,但面试和项目里真正要理解的是:写入请求从客户端到主分片、副本、内存 buffer、translog、segment、flush、merge,每一步到底做什么,为什么这样设计,不这样会有什么问题

这一页把一条商品文档写入 ES 的完整过程拆开讲清楚。

这一页解决这些问题

问题你要掌握到什么程度
写入请求怎么找到分片知道协调节点、routing、主分片、副本的关系
写入成功为什么不一定能搜到知道 buffer、refresh、segment 的关系
translog 做什么知道它用于故障恢复,不等于倒排索引
flush 和 refresh 有什么区别知道 refresh 让数据可搜,flush 做提交和裁剪 translog
update/delete 为什么贵知道 Lucene segment 不可变,只能标记删除再写新版本
商业项目怎么用知道商品改价、上下架、订单搜索同步如何设置 refresh、bulk、重试和补偿

先建立整体画面

ES 底层每个分片本质上是一个 Lucene 索引。Lucene 为了让查询快,把数据写成不可变的 segment。不可变的好处是查询稳定、缓存友好、并发简单;代价是更新不能原地改,新增数据也要经过 refresh 才能变成可搜索的 segment。

mermaid
flowchart TD
    A["客户端写入文档"] --> B["协调节点接收请求"]
    B --> C["根据 routing 定位主分片"]
    C --> D["主分片执行写入"]
    D --> E["写入内存 indexing buffer"]
    D --> F["追加 translog"]
    D --> G["复制到副本分片"]
    G --> H["副本写入 buffer 和 translog"]
    H --> I["主分片收到副本确认"]
    I --> J["返回写入成功"]
    E --> K["refresh 生成可搜索 segment"]
    F --> L["flush 提交并裁剪 translog"]
    K --> M["merge 合并小 segment"]

注意这张图里有两个时间点:

时间点含义
返回写入成功主分片和必要副本已经接收写入,数据有 translog 保护
搜索可见refresh 后新 segment 被打开,查询才能搜到

所以 ES 不是写入失败才搜不到,写入成功但 refresh 还没发生,也可能短时间搜不到

第一步:协调节点接收写入

客户端可以把请求发给集群中任意节点。接到请求的节点就是这次请求的协调节点。协调节点不一定保存这条数据,它主要负责转发和汇总结果。

写入示例:

bash
curl -X PUT "http://localhost:9200/product_search/_doc/1001" \
  -H "Content-Type: application/json" \
  -d '{
    "productId": 1001,
    "title": "无线蓝牙降噪耳机 Pro",
    "brand": "SoundMax",
    "price": 299,
    "status": "ON_SALE",
    "updatedAt": "2026-07-05T10:00:00"
  }'

协调节点要先决定这条文档应该落到哪个主分片。默认路由逻辑可以简化理解为:

text
shard = hash(_id 或 routing) % primary_shard_count

如果没有自定义 routing,通常使用 _id 计算分片。这样同一个 _id 的文档总能落到同一个主分片。

设计点为什么需要
_id 路由保证同一文档更新时能找到原来的位置
用 primary shard count 取模把数据水平拆到多个主分片
主分片数量创建后不能随便改改了取模规则,历史数据就找不到原来的分片

第二步:主分片执行写入

请求到达主分片后,主分片会做字段处理、分词、构建索引结构,并把写入记录追加到 translog。

mermaid
flowchart TD
    A["主分片收到文档"] --> B["检查 Mapping"]
    B --> C{"字段类型是否匹配"}
    C -- "不匹配" --> D["返回 Mapping 冲突"]
    C -- "匹配" --> E["text 字段分词"]
    E --> F["生成倒排索引条目"]
    B --> G["keyword、number、date 建精确索引"]
    F --> H["写入 indexing buffer"]
    G --> H
    H --> I["追加 translog"]

这里有两个非常重要的写入位置。

位置保存什么作用
indexing buffer还没形成正式 segment 的索引数据等 refresh 时生成可搜索 segment
translog原始写入操作日志节点故障后用于恢复未提交数据

可以把它类比成做账:

  1. indexing buffer 像正在整理的账本页,还没装订成册。
  2. translog 像流水小票,机器挂了可以按小票重新整理。
  3. segment 像装订好的账本,搜索可以稳定翻阅。

第三步:复制到副本

主分片写入后,会把请求复制给副本分片。副本也要执行类似的写入动作:检查、写 buffer、写 translog。

mermaid
sequenceDiagram
    participant C as Client
    participant N as Coordinating Node
    participant P as Primary Shard
    participant R as Replica Shard
    C->>N: index document
    N->>P: route to primary
    P->>P: write buffer and translog
    P->>R: replicate operation
    R->>R: write buffer and translog
    R-->>P: ack
    P-->>N: ack
    N-->>C: write success

副本的意义不只是备份,还能承担查询流量。但副本写入也会增加写入成本,所以副本数不是越多越好。

副本数好处代价
0写入快、省资源节点故障可能丢可用性,查询吞吐低
1常见生产配置,兼顾高可用和查询写入要复制一份
2 或更多读多写少、可用性要求高写放大明显,磁盘和网络压力更大

第四步:什么时候返回成功

写入是否返回成功,和 wait_for_active_shards、副本状态、主分片状态有关。初学阶段可以先抓住这条主线:主分片写入成功,并且满足需要等待的活跃分片数量后,客户端收到成功响应

但这依然不代表搜索一定能看到,因为此时数据可能还在 buffer 中。

状态数据是否安全是否可搜索
写入 buffer 和 translog,未 refresh节点正常时安全,故障可用 translog 恢复不一定
refresh 后生成 segment安全性仍依赖 translog/flush可以搜索
flush 后 commitLucene 提交点稳定,translog 可裁剪可以搜索

第五步:Refresh 让数据可搜索

refresh 会把内存里的 indexing buffer 转成一个新的 Lucene segment,并打开一个新的 searcher。这样后续搜索请求才能看到新写入的数据。

mermaid
flowchart TD
    A["写入成功"] --> B["数据在 indexing buffer"]
    B --> C["定时 refresh"]
    C --> D["生成新的 segment"]
    D --> E["打开 searcher"]
    E --> F["新数据可以被搜索"]

默认情况下,ES 常见 refresh_interval1s。这就是“近实时”的来源:通常不是毫秒级立即可见,而是有一个短暂窗口。

配置或操作效果适用场景
默认 refresh写入吞吐和可见延迟折中商品、订单、工单搜索
调大 refresh_interval减少小 segment,提高批量写入吞吐全量导入、重建索引
refresh=true请求返回前强制 refresh少量后台操作,需要马上查到
refresh=wait_for等待下一次 refresh 后返回测试或少量强可见需求

生产里不要给每条写入都加 refresh=true。这会导致大量小 segment,写入慢、查询慢、merge 压力变大。

第六步:Flush 做提交和裁剪 translog

refresh 只解决“可搜索”,不等于完成一次完整提交。flush 会触发 Lucene commit,把内存状态提交到稳定存储,并裁剪已经不需要的 translog。

mermaid
flowchart TD
    A["持续写入"] --> B["translog 越来越大"]
    B --> C["达到条件触发 flush"]
    C --> D["Lucene commit"]
    D --> E["生成新的提交点"]
    E --> F["裁剪旧 translog"]

refresh 和 flush 的区别是面试高频点:

对比项refreshflush
核心目标让新数据可搜索做持久提交并清理 translog
是否生成可搜索 segment不是它的主要目标
是否清理 translog不清理会裁剪旧 translog
频率通常更频繁通常更少
影响影响搜索可见性和小 segment 数量影响恢复成本和磁盘日志大小

第七步:Segment 为什么要 Merge

频繁 refresh 会产生很多小 segment。查询时需要从多个 segment 查,再合并结果。小 segment 太多会让查询、文件句柄、缓存、内存都变差。

merge 会把多个小 segment 合并成大 segment,并在合并过程中清理被删除标记的旧文档。

mermaid
flowchart TD
    A["segment 1"] --> D["merge"]
    B["segment 2"] --> D
    C["segment 3"] --> D
    D --> E["新的大 segment"]
    D --> F["清理删除标记文档"]
现象原理后果
写入高峰后查询也变慢小 segment 多,merge 抢 IO查询延迟抖动
更新很多导致磁盘涨旧文档先标记删除,merge 后才真正清理短时间磁盘放大
手动强制 merge可以减少 segment会消耗大量 IO,生产慎用

更新和删除为什么不是原地修改

Lucene segment 是不可变的,所以 ES 更新可以理解成:

  1. 找到旧文档。
  2. 给旧文档打删除标记。
  3. 写入新版本文档。
  4. refresh 后新版本可搜索。
  5. merge 时真正清理旧版本。
mermaid
flowchart TD
    A["update 商品价格"] --> B["旧文档打删除标记"]
    B --> C["写入新版本文档"]
    C --> D["refresh 后新版本可搜"]
    D --> E["merge 清理旧文档"]

这解释了为什么 ES 不适合做高频强一致状态表。比如秒杀库存、账户余额、支付状态,不应该靠 ES 做最终判断。

商业场景:商品改价和上下架

商品服务中,MySQL 保存事实价格和上下架状态,ES 保存搜索视图。一次商品改价通常这样走:

mermaid
flowchart TD
    A["运营修改商品价格"] --> B["商品服务事务更新 MySQL"]
    B --> C["发送商品变更消息"]
    C --> D["ES 同步消费者收到消息"]
    D --> E["按 productId 查询 MySQL 最新快照"]
    E --> F["组装搜索文档"]
    F --> G["upsert 到 ES"]
    G --> H["refresh 后搜索页可见"]

为什么消费者要查 MySQL 最新快照,而不是直接相信消息里的价格?

原因解释
防乱序旧消息后到时,查主库能拿到最新状态
降低消息体变化消息只放 ID 和版本,结构稳定
多表组装搜索文档常需要品牌、类目、标签、库存等冗余字段
失败可重放重放消息时仍能组装当前正确文档

可运行 Demo:Java Bulk 写入商品索引

下面示例用 Elasticsearch Java API Client 表达商业项目常见写法:批量 upsert 商品搜索文档。真实项目中还要加重试、死信、监控和补偿。

java
import co.elastic.clients.elasticsearch.ElasticsearchClient;
import co.elastic.clients.elasticsearch.core.BulkRequest;
import co.elastic.clients.elasticsearch.core.BulkResponse;
import co.elastic.clients.elasticsearch.core.bulk.BulkOperation;

import java.io.IOException;
import java.math.BigDecimal;
import java.util.List;

public class ProductIndexWriter {
    private final ElasticsearchClient client;

    public ProductIndexWriter(ElasticsearchClient client) {
        this.client = client;
    }

    public void upsertProducts(List<ProductDoc> docs) throws IOException {
        BulkRequest.Builder builder = new BulkRequest.Builder();

        for (ProductDoc doc : docs) {
            builder.operations(BulkOperation.of(op -> op
                .index(idx -> idx
                    .index("product_search")
                    .id(String.valueOf(doc.productId()))
                    .document(doc)
                )
            ));
        }

        BulkResponse response = client.bulk(builder.build());
        if (response.errors()) {
            response.items().forEach(item -> {
                if (item.error() != null) {
                    System.err.println("ES 写入失败 id=" + item.id()
                        + ", reason=" + item.error().reason());
                }
            });
            throw new IllegalStateException("部分商品索引写入失败,需要重试或进入死信");
        }
    }

    public record ProductDoc(
        Long productId,
        String title,
        String brand,
        BigDecimal price,
        String status,
        Long updatedAt
    ) {}
}

关键点:

  1. _id 使用业务主键,重复消费时覆盖同一文档。
  2. 批量写入不要无限大,常见从几百到几千条压测调整。
  3. 对失败项不能吞掉,要重试、死信、告警。
  4. 搜索展示可以接受短暂延迟,交易判断必须回主库。

常见坑

为什么会发生正确做法
写完马上查不到就认为失败refresh 还没发生区分写入成功和搜索可见
每次写入都 refresh=true大量小 segment普通业务用默认 refresh,批量导入调大 refresh
高频更新 ES 状态字段update 是删除加新增高频强一致状态放主库,ES 只做搜索视图
Bulk 批次过大内存、网络、失败重试成本高压测确定批次和并发
忽略 partial failureBulk 可能部分成功部分失败检查每个 item 的 error
Mapping 冲突后反复重试不可重试错误修 Mapping 或数据,再重放

排查方法

问题排查顺序
写入成功但搜不到按 ID 查文档、等 refresh、查查询条件、查同步日志
写入慢看 bulk 大小、并发、refresh、副本、merge、磁盘 IO、线程池拒绝
磁盘增长快看更新删除比例、segment merge、索引副本、保留周期
更新失败看响应错误、Mapping 冲突、磁盘只读、线程池 rejected、网络超时
搜索看到旧值查 MQ Lag、旧消息覆盖、refresh 延迟、补偿任务

常用命令:

bash
curl "http://localhost:9200/_cat/indices?v"
curl "http://localhost:9200/_cat/segments/product_search?v"
curl "http://localhost:9200/_cat/thread_pool/write?v"
curl "http://localhost:9200/product_search/_doc/1001"

面试标准回答

text
ES 写入时,请求先到协调节点,协调节点根据 _id 或 routing 定位主分片。主分片检查 Mapping、处理字段、对 text 分词,写入 indexing buffer,同时追加 translog,然后复制到副本。主分片和满足条件的副本确认后,客户端收到写入成功。

但写入成功不等于马上可搜索,因为新数据还可能在 buffer 中。refresh 会把 buffer 生成新的可搜索 segment,并打开 searcher,所以 ES 叫近实时搜索。flush 负责 Lucene commit 和裁剪 translog,merge 负责合并小 segment 并清理删除标记。

ES 更新不是原地修改,因为 Lucene segment 不可变,更新本质是旧文档打删除标记再写入新文档,后续 merge 清理旧版本。因此 ES 适合搜索视图,不适合做库存、支付、账户余额这种强一致高频状态主库。

关联知识点

知识点说明
Elasticsearch 底层原理总览写入、查询、评分、doc values
Query 与 Fetch 全过程搜索请求如何执行
倒排索引与 BM25为什么 ES 查询快、结果怎么排序
MySQL 与 ES 数据一致性商品改价、失败重试、补偿对账
性能优化与排查写入慢、查询慢、搜不到怎么排查