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。
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 还没发生,也可能短时间搜不到。
第一步:协调节点接收写入
客户端可以把请求发给集群中任意节点。接到请求的节点就是这次请求的协调节点。协调节点不一定保存这条数据,它主要负责转发和汇总结果。
写入示例:
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"
}'协调节点要先决定这条文档应该落到哪个主分片。默认路由逻辑可以简化理解为:
shard = hash(_id 或 routing) % primary_shard_count如果没有自定义 routing,通常使用 _id 计算分片。这样同一个 _id 的文档总能落到同一个主分片。
| 设计点 | 为什么需要 |
|---|---|
用 _id 路由 | 保证同一文档更新时能找到原来的位置 |
| 用 primary shard count 取模 | 把数据水平拆到多个主分片 |
| 主分片数量创建后不能随便改 | 改了取模规则,历史数据就找不到原来的分片 |
第二步:主分片执行写入
请求到达主分片后,主分片会做字段处理、分词、构建索引结构,并把写入记录追加到 translog。
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 | 原始写入操作日志 | 节点故障后用于恢复未提交数据 |
可以把它类比成做账:
- indexing buffer 像正在整理的账本页,还没装订成册。
- translog 像流水小票,机器挂了可以按小票重新整理。
- segment 像装订好的账本,搜索可以稳定翻阅。
第三步:复制到副本
主分片写入后,会把请求复制给副本分片。副本也要执行类似的写入动作:检查、写 buffer、写 translog。
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 后 commit | Lucene 提交点稳定,translog 可裁剪 | 可以搜索 |
第五步:Refresh 让数据可搜索
refresh 会把内存里的 indexing buffer 转成一个新的 Lucene segment,并打开一个新的 searcher。这样后续搜索请求才能看到新写入的数据。
flowchart TD
A["写入成功"] --> B["数据在 indexing buffer"]
B --> C["定时 refresh"]
C --> D["生成新的 segment"]
D --> E["打开 searcher"]
E --> F["新数据可以被搜索"]默认情况下,ES 常见 refresh_interval 是 1s。这就是“近实时”的来源:通常不是毫秒级立即可见,而是有一个短暂窗口。
| 配置或操作 | 效果 | 适用场景 |
|---|---|---|
| 默认 refresh | 写入吞吐和可见延迟折中 | 商品、订单、工单搜索 |
调大 refresh_interval | 减少小 segment,提高批量写入吞吐 | 全量导入、重建索引 |
refresh=true | 请求返回前强制 refresh | 少量后台操作,需要马上查到 |
refresh=wait_for | 等待下一次 refresh 后返回 | 测试或少量强可见需求 |
生产里不要给每条写入都加 refresh=true。这会导致大量小 segment,写入慢、查询慢、merge 压力变大。
第六步:Flush 做提交和裁剪 translog
refresh 只解决“可搜索”,不等于完成一次完整提交。flush 会触发 Lucene commit,把内存状态提交到稳定存储,并裁剪已经不需要的 translog。
flowchart TD
A["持续写入"] --> B["translog 越来越大"]
B --> C["达到条件触发 flush"]
C --> D["Lucene commit"]
D --> E["生成新的提交点"]
E --> F["裁剪旧 translog"]refresh 和 flush 的区别是面试高频点:
| 对比项 | refresh | flush |
|---|---|---|
| 核心目标 | 让新数据可搜索 | 做持久提交并清理 translog |
| 是否生成可搜索 segment | 是 | 不是它的主要目标 |
| 是否清理 translog | 不清理 | 会裁剪旧 translog |
| 频率 | 通常更频繁 | 通常更少 |
| 影响 | 影响搜索可见性和小 segment 数量 | 影响恢复成本和磁盘日志大小 |
第七步:Segment 为什么要 Merge
频繁 refresh 会产生很多小 segment。查询时需要从多个 segment 查,再合并结果。小 segment 太多会让查询、文件句柄、缓存、内存都变差。
merge 会把多个小 segment 合并成大 segment,并在合并过程中清理被删除标记的旧文档。
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 更新可以理解成:
- 找到旧文档。
- 给旧文档打删除标记。
- 写入新版本文档。
- refresh 后新版本可搜索。
- merge 时真正清理旧版本。
flowchart TD
A["update 商品价格"] --> B["旧文档打删除标记"]
B --> C["写入新版本文档"]
C --> D["refresh 后新版本可搜"]
D --> E["merge 清理旧文档"]这解释了为什么 ES 不适合做高频强一致状态表。比如秒杀库存、账户余额、支付状态,不应该靠 ES 做最终判断。
商业场景:商品改价和上下架
商品服务中,MySQL 保存事实价格和上下架状态,ES 保存搜索视图。一次商品改价通常这样走:
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 商品搜索文档。真实项目中还要加重试、死信、监控和补偿。
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
) {}
}关键点:
_id使用业务主键,重复消费时覆盖同一文档。- 批量写入不要无限大,常见从几百到几千条压测调整。
- 对失败项不能吞掉,要重试、死信、告警。
- 搜索展示可以接受短暂延迟,交易判断必须回主库。
常见坑
| 坑 | 为什么会发生 | 正确做法 |
|---|---|---|
| 写完马上查不到就认为失败 | refresh 还没发生 | 区分写入成功和搜索可见 |
每次写入都 refresh=true | 大量小 segment | 普通业务用默认 refresh,批量导入调大 refresh |
| 高频更新 ES 状态字段 | update 是删除加新增 | 高频强一致状态放主库,ES 只做搜索视图 |
| Bulk 批次过大 | 内存、网络、失败重试成本高 | 压测确定批次和并发 |
| 忽略 partial failure | Bulk 可能部分成功部分失败 | 检查每个 item 的 error |
| Mapping 冲突后反复重试 | 不可重试错误 | 修 Mapping 或数据,再重放 |
排查方法
| 问题 | 排查顺序 |
|---|---|
| 写入成功但搜不到 | 按 ID 查文档、等 refresh、查查询条件、查同步日志 |
| 写入慢 | 看 bulk 大小、并发、refresh、副本、merge、磁盘 IO、线程池拒绝 |
| 磁盘增长快 | 看更新删除比例、segment merge、索引副本、保留周期 |
| 更新失败 | 看响应错误、Mapping 冲突、磁盘只读、线程池 rejected、网络超时 |
| 搜索看到旧值 | 查 MQ Lag、旧消息覆盖、refresh 延迟、补偿任务 |
常用命令:
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"面试标准回答
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 数据一致性 | 商品改价、失败重试、补偿对账 |
| 性能优化与排查 | 写入慢、查询慢、搜不到怎么排查 |
