这是《Elasticsearch 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 写入链路是性能和可靠性的交汇点。理解一次写入请求在协调节点、主分片、副本分片、Lucene 内存、translog、refresh 和 flush 之间的流转后,才能解释写入延迟、可见性、数据恢复和写入放大。
本章拆解完整的写入过程。
24.1 一次 Index 请求的生命周期
假设写入 orders-v3/_doc/order_10001,索引有 3 个主分片和 1 个副本:
sequenceDiagram
participant C as Client
participant N as Coordinating Node
participant P as Primary Shard
participant R as Replica Shard
C->>N: PUT /orders-v3/_doc/order_10001
N->>N: 根据路由计算目标主分片
N->>P: 转发写请求
P->>P: 校验 Mapping / 执行 pipeline
P->>P: 写 Indexing Buffer + translog
P->>R: 同步副本
R->>R: 写 Indexing Buffer + translog
R-->>P: 副本成功
P-->>N: 主副本均成功
N-->>C: 返回 2xx
默认 wait_for_active_shards=1 只要求写入前主分片可用,但请求成功还取决于 consistency 和版本实现。8.x 默认副本写入失败会导致主流程失败,具体行为要结合版本确认。
核心点:
- 协调节点负责路由,不决定文档内容;
- 主分片负责校验和写入顺序;
- 副本接收主分片转发的内容;
- 写成功不代表立刻可搜索;
- translog 保证故障恢复,refresh 决定可见性。
24.2 协调节点做什么
协调节点执行:
- 解析请求;
- 根据
_id或自定义 routing 计算分片; - 选择目标主分片所在节点;
- 转发请求;
- 汇总响应或错误。
Bulk 请求会更复杂:
POST _bulk
{"index": {"_index": "orders-v3", "_id": "order_1"}}
{"order_id": "order_1", "pay_amount": 100}
{"index": {"_index": "orders-v3", "_id": "order_2"}}
{"order_id": "order_2", "pay_amount": 200}
协调节点需要把批量项按目标分片拆分,分别发给对应主分片,再把逐项结果按原始顺序组装返回。因此:
- Bulk 中部分成功、部分失败是正常的;
- 应用必须检查每个 item 的状态;
- 请求过大可能导致协调节点内存压力;
- Bulk size 通常从 1MB 到 10MB 之间测试;
- 发生 429 时应指数退避,而不是立即重试。
24.3 主分片写入
主分片写入前会处理:
- Mapping 校验与动态映射更新;
- 自动生成
_id; - 版本检查;
- routing 校验;
- ingest pipeline;
- 脚本更新或 doc update。
更新请求会经历:
读取 _source
-> 应用 update doc 或 script
-> 标记旧文档 deleted
-> 写入新文档
因此 update 比 index 更重。它需要读旧值、计算新值,并产生删除标记。高频更新同一个文档会显著放大写入。
24.4 版本控制与乐观并发
Elasticsearch 通过 seq_no 和 primary_term 实现乐观并发控制。
PUT products-v3/_update/sku_10001?if_seq_no=10&if_primary_term=1
{
"doc": {
"price": 4999
}
}
如果读取后有其他人已经修改,写入会返回 409 冲突。应用可以选择:
- 重新读取再修改;
- 使用
retry_on_conflict; - 改成以事件时间为准的脚本判断;
- 将热点文档更新合并到外部缓冲,再低频写入。
POST products-v3/_update/sku_10001?retry_on_conflict=3
{
"script": {
"source": "ctx._source.stock += params.delta",
"lang": "painless",
"params": { "delta": -1 }
}
}
retry_on_conflict 适合幂等性较弱但冲突可控的更新。如果每次库存扣减都直接打 ES,热点 SKU 会成为瓶颈,交易扣减应以数据库或缓存为准,再异步同步到搜索。
24.5 Refresh
Refresh 把 Indexing Buffer 中的数据生成一个新 Segment,并打开可搜索。
write -> memory + translog -> refresh -> searchable segment
查看和调整刷新间隔:
GET orders-v3/_settings/index.refresh_interval
PUT orders-v3/_settings
{
"index.refresh_interval": "5s"
}
常见策略:
| 场景 | 建议值 |
|---|---|
| 常规业务索引 | 1s |
| 日志高频写入 | 5s 到 30s |
| 批量导入 | -1,导入完成后恢复 |
| 写后必须可查 | refresh=wait_for 或接口层等待 |
| 测试或演示 | refresh=true |
批量导入示例:
PUT import-v1/_settings
{
"refresh_interval": "-1",
"number_of_replicas": 0
}
导入完成后:
PUT import-v1/_settings
{
"refresh_interval": "1s",
"number_of_replicas": 1
}
不要忘记恢复副本。很多人批量导入时把副本设为 0,导入完成后忘记恢复,节点故障时数据风险会很高。
24.6 Flush
Flush 是 Lucene commit,将变更真正提交到磁盘,并生成新的 commit point。
Indexing Buffer / Segment
-> fsync
-> Lucene commit
-> 清理 translog
Flush 比 Refresh 成本高得多。它受 translog 大小和时间阈值控制:
GET orders-v3/_stats/flush,merge,refresh,indexing
如果 flush 频繁:
- translog 阈值过小;
- 磁盘写入能力不足;
- 分片数太多导致每个分片都频繁触发;
- 写入请求过小,系统调用和元数据开销占比高。
24.7 Translog 与恢复
节点重启或崩溃时,Lucene 最近 commit 之后但已对客户端确认的变更需要通过 translog 重放。
恢复流程:
打开最近 commit point
-> 重放 translog
-> 恢复到确认状态
-> 分片重新可用
因此 translog 不是可有可无的日志,而是写确认语义的一部分。
可靠性取舍:
PUT critical-index-v1/_settings
{
"index.translog.durability": "request"
}
| 场景 | 建议 |
|---|---|
| 订单、支付、审计 | request |
| 商品搜索、画像 | 可评估 async |
| 日志与指标 | 常可评估 async |
| 测试环境 | 默认或更宽松 |
即使使用 async,也要确认应用侧有幂等重试,因为网络失败后可能重复写入。
24.8 Segment Merge 与写入放大
一次业务写入可能引发多层数据搬运:
1. 写 Indexing Buffer
2. refresh 生成 Segment
3. merge 把小 Segment 合成大 Segment
4. 副本节点重复 refresh 和 merge
5. 更新场景标记 deleted,merge 时清理
这就是写入放大。表现通常是:
- 写入吞吐不稳定;
- CPU IO wait 高;
- merge 任务排队;
- 查询延迟上升;
- 磁盘使用临时上涨。
查看 merge 统计:
GET _nodes/stats/indices/merges
GET orders-v3/_stats/merges
治理方式:
- 增大刷新间隔;
- 合理批量写入;
- 避免高频小 bulk;
- 避免高频更新同一文档;
- 控制分片数;
- 使用 rollover 控制单分片大小;
- 在低峰执行 force merge;
- 确认磁盘和页缓存充足。
24.9 写入性能模型
写入吞吐由以下资源共同决定:
| 资源 | 影响 |
|---|---|
| CPU | 分词、脚本、merge、压缩 |
| 磁盘 IO | translog fsync、flush、merge |
| 页缓存 | Segment 读取和合并效率 |
| 网络 | 副本同步和客户端传输 |
| 分片数 | 并行度与扇出 |
| 请求大小 | 每项固定成本占比 |
| Mapping 复杂度 | 解析与索引成本 |
一个实用判断表:
| 现象 | 优先检查 |
|---|---|
| 单条写入很慢 | 客户端同步调用、refresh、副本 |
| Bulk 吞吐低 | bulk size、并发、mapping、分片 |
| 429 rejected | write thread pool 队列满 |
| 延迟毛刺 | refresh、flush、merge、磁盘 |
| CPU 高 | 分词、脚本、merge、查询混部 |
| IO 高 | translog、merge、磁盘能力不足 |
| 副本延迟 | 网络、恢复限流、节点负载 |
24.10 写入背压
Elasticsearch 的写入线程池是有限的。当队列满时,请求会被拒绝:
GET /_nodes/stats/thread_pool
GET /_cat/thread_pool/write?v&h=node_name,name,active,queue,rejected,completed
遇到 429 时,客户端应该:
- 记录失败批次;
- 指数退避;
- 降低并发;
- 检查是否有慢查询抢占资源;
- 确认磁盘水位和节点是否离线;
- 评估扩容或拆分集群。
不要无脑把线程池队列调大。队列变长只会增加等待时间、占用内存,并让故障雪崩更晚暴露。
24.11 幂等写入
网络超时不代表写入失败。应用必须支持幂等:
- 明确指定
_id; - 消息系统使用 offset + 批次检查点;
- 事件携带唯一 event_id;
- 重放不会产生重复数据;
- 对部分失败的 Bulk 只重试失败项。
示意流程:
poll batch
-> 按 item 写入 bulk
-> 成功项记录 offset
-> 失败项进入重试
-> 超过阈值进入死信
-> 恢复后可人工重放
如果使用自动生成 _id,重试会写入两个文档。对于日志可能可以容忍,对于订单和支付不可接受。
24.12 写入调优清单
- 批量提交,避免单文档高频写入;
- bulk size 和并发通过压测确定;
- 明确
_id与幂等策略; - 高吞吐索引增大
refresh_interval; - 导入期间临时关闭副本并恢复;
- 避免热点文档高频 update;
- 脚本简单,参数化,避免每请求编译;
- Mapping 控制字段数量和嵌套深度;
- 监控 write rejected、merge、flush、translog;
- 磁盘保留足够余量,避免触碰 flood stage;
- 写入与重查询在高峰期错峰或资源隔离;
- 429 使用退避和限流,不盲目加大客户端并发。
本章小结
写入请求经历协调路由、主分片校验、Lucene 写入、translog 持久化、副本同步、refresh 可见、flush 提交和 merge 清理。写成功和可搜索是两个不同时点,可靠性和吞吐之间也始终存在取舍。写入调优的关键是批量、幂等、控制刷新与合并、留足磁盘与页缓存,并让客户端具备背压能力。
思考题
- 为什么写入成功后立刻查询可能看不到?
- refresh、flush、translog 分别解决什么问题?
- 为什么 update 比全新 index 更昂贵?
- Bulk 返回 200 时,是否代表所有文档都写入成功?
- 如果 write 线程池 rejected 快速上升,你会怎么处理客户端和服务端?