写入性能问题往往在数据量涨上来之后才暴露:昨天还能跑满的导入任务,今天开始报 429 Too Many Requests;bulk 批次调大之后吞吐没升反降;明明只写主分片,CPU 却先打满。这些现象背后是同一件事——写入链路有多个环节,任何一环过载都会成为瓶颈,而盲目调参通常只是把瓶颈从一个环节推到另一个环节。本文按「链路 → 批次 → 持久化 → 分片 → 背压」的顺序讲清批量写入的调优方法。
1. 写入路径与瓶颈定位
1.1 一次写入的完整链路
一条文档从客户端到可被搜索,经过这些环节:协调节点接收 bulk 请求 → 按 _id 或 _routing 算出目标分片 → 转发到主分片所在节点 → 主分片写内存缓冲与 translog → 复制到副本分片 → 副本确认后主分片返回 → 定期 refresh 生成新段 → 段累积到阈值后 merge。
任何一环变慢,上游都会感知为写入变慢。定位瓶颈就是找出这一环。
1.2 三类瓶颈
| 瓶颈类型 | 表现 | 常见原因 |
|---|---|---|
| CPU | 节点 CPU 打满,merge 线程高 | 分词、merge、脚本 |
| IO | 磁盘 util 高,fsync 延迟大 | translog fsync、段合并 |
| 线程池 | 队列满,出现 429 | 并发过高、单批过大 |
| 内存 | 堆压力大,GC 频繁 | 批量缓冲、段元数据 |
先看 _nodes/stats 的 CPU 与 IO,再看 _cat/thread_pool 的队列与拒绝数,最后看 _nodes/hot_threads 定位具体线程。
1.3 先测量再调优
调优前必须先建立基线。用固定数据集跑一次导入,记录吞吐(docs/s)、耗时、CPU、磁盘 util、拒绝次数。没有基线,任何「优化」都无法证明有效。基线数据也让后续每次改动可归因。
2. bulk 的大小与并发
2.1 bulk 请求格式
bulk 用 NDJSON,每两行一组(动作 + 文档),最后必须换行:
curl -X POST "localhost:9200/_bulk?pretty" \
-H "Content-Type: application/x-ndjson" -d'
{"index": {"_index": "orders", "_id": "1"}}
{"order_id": "1", "amount": 100, "created_at": "2026-01-01T00:00:00Z"}
{"index": {"_index": "orders", "_id": "2"}}
{"order_id": "2", "amount": 200, "created_at": "2026-01-01T00:00:01Z"}
'
漏掉末尾换行会报 The bulk request must be terminated by a newline,这是最常见的格式错误。
2.2 单批大小的取舍
批次太小,网络与请求开销占比高;批次太大,单个请求的协调与内存占用高,失败重试的代价也大。经验区间:
| 单批文档数 | 单批体积 | 适用 |
|---|---|---|
| 100~500 | < 1 MB | 小文档、低延迟要求 |
| 1000~5000 | 5~15 MB | 通用批量导入 |
| 5000~10000 | 15~50 MB | 大文档、高吞吐 |
一个可靠的判据是「单批体积控制在 5~15 MB」:文档小就多装几条,文档大就少装几条。按条数一刀切容易在小文档场景浪费、大文档场景爆内存。
2.3 并发度
并发不是越高越好。当服务端线程池队列开始堆积,再加并发只会增加拒绝。推荐做法是固定并发、用「在途请求数」控制:
import concurrent.futures
def index_batch(client, index, docs):
body = build_ndjson(index, docs)
return client.bulk(body=body)
with concurrent.futures.ThreadPoolExecutor(max_workers=8) as pool:
futures = [pool.submit(index_batch, client, "orders", chunk)
for chunk in chunks(docs, 2000)]
for f in concurrent.futures.as_completed(futures):
handle(f.result())
max_workers 从 4~8 起步,观察拒绝率再决定是否上调。写入线程池 write 的队列默认是 200(每节点),并发乘以分片数一旦远超它,拒绝就不可避免。
2.4 客户端选择与自动批量
官方客户端提供 bulk helper,能自动按体积或条数切批并处理重试:
from elasticsearch.helpers import parallel_bulk
for ok, item in parallel_bulk(client, actions, index="orders",
chunk_size=2000, thread_count=4,
raise_on_error=False):
if not ok:
dead_letter.append(item)
raise_on_error=False 让单条失败不中断整批,失败项落到死信队列后续处理。自己写批量逻辑时要实现同样的语义,否则一条脏数据会让整批回滚。
3. refresh 与 translog 策略
3.1 refresh 的代价
refresh 让内存缓冲的数据生成可搜索的新段。默认 refresh_interval: 1s,对搜索友好,对写入昂贵:每次 refresh 都会产生小段,增加 merge 压力。批量导入期把它调大:
curl -X PUT "localhost:9200/orders/_settings" -H "Content-Type: application/json" -d'
{ "index": { "refresh_interval": "30s" } }'
导入完成后恢复 1s。搜索时效性要求不高的场景,长期用 30s 也是合理选择。
3.2 写入期关掉 refresh
导入大批量数据时,可以临时关闭自动 refresh,手工在结束时触发一次:
curl -X PUT "localhost:9200/orders/_settings" -H "Content-Type: application/json" -d'
{ "index": { "refresh_interval": "-1" } }'
# 导入完成后
curl -X POST "localhost:9200/orders/_refresh"
curl -X PUT "localhost:9200/orders/_settings" -H "Content-Type: application/json" -d'
{ "index": { "refresh_interval": "1s" } }'
-1 表示不自动 refresh。这一步常能带来 20%~30% 的吞吐提升,代价是导入期间数据不可搜索。
3.3 translog 的持久化策略
translog 保证未刷盘的写入在节点崩溃后可恢复。默认每个请求都 fsync,这是可靠性与吞吐的权衡点:
curl -X PUT "localhost:9200/orders/_settings" -H "Content-Type: application/json" -d'
{
"index": {
"translog.durability": "async",
"translog.sync_interval": "30s"
}
}'
async 表示由后台按 sync_interval 批量刷盘。代价是节点在两次 fsync 之间崩溃会丢失这期间的数据。只应在可重建的数据(如日志、可重放的导入)上使用,订单等关键数据保持 request。
3.4 副本数与刷盘
导入期可以把副本数临时设为 0,导完再调回:
curl -X PUT "localhost:9200/orders/_settings" -H "Content-Type: application/json" -d'
{ "index": { "number_of_replicas": 0 } }'
省掉副本写入与复制确认,吞吐能显著提升。但导入期间集群没有冗余,节点宕机会丢数据;而且副本数调回时会触发全量分片复制,产生额外 IO。是否值得取决于数据能否重放。
4. 分片与路由
4.1 分片数对写入的影响
分片是写入并行的单位:分片太少,单分片成为串行瓶颈;分片太多,每分片都是独立的 Lucene 实例,段合并与元数据开销累加。经验值是「每 GB 堆内存对应 2025 个分片」,单分片体积控制在 1050 GB。
4.2 routing 打散
默认按 _id 哈希路由,能均匀打散。但如果业务自定义了 _routing,且取值集中(如按天分区、按租户),会造成热点分片:
curl -X POST "localhost:9200/_bulk" -H "Content-Type: application/x-ndjson" -d'
{"index": {"_index": "orders", "_routing": "tenant-1"}}
{"order_id": "1", "amount": 100}
'
所有 tenant-1 的文档落到同一分片。若某租户写入量远超其他,该分片所在节点会先饱和。解决办法是给 routing 值加后缀打散,或改用 _id 默认路由。
4.3 写入确认与 wait_for_active_shards
默认写入需要主分片活跃即可返回(wait_for_active_shards: 1)。要提高写入安全性可以设为 all,但会显著降低吞吐:
curl -X POST "localhost:9200/orders/_doc/1?wait_for_active_shards=all" \
-H "Content-Type: application/json" -d'
{ "order_id": "1", "amount": 100 }'
导入期建议用默认值,导完再依赖副本与快照兜底。
4.4 版本冲突与幂等
bulk 里可以指定版本号做乐观并发控制:
{"index": {"_index": "orders", "_id": "1", "version": 2, "version_type": "external"}}
{"order_id": "1", "amount": 150}
版本不匹配时该条返回 version_conflict_engine_exception,整批的其他条仍会执行。做幂等重放时用 version_type: external 加上游版本号,能避免重复写入导致的数据错乱。
5. 拒绝与背压
5.1 429 的成因
es_rejected_execution_exception 来自线程池队列满。写入线程池 write 的队列默认 200,当到达速率超过处理速率,队列填满后新请求被拒。这是 Elasticsearch 的背压信号,不是 bug:它在告诉上游「我处理不过来了」。
5.2 查看线程池状态
curl -s "localhost:9200/_cat/thread_pool/write?v&h=node_name,active,queue,rejected,completed"
queue 长期大于 0 说明已经在排队,rejected 持续增长说明过载。此时加并发只会让情况更糟。
5.3 客户端重试与退避
对 429 必须重试,但要带退避,否则重试风暴会加剧过载:
import random, time
def bulk_with_backoff(client, body, max_retries=5):
for attempt in range(max_retries):
resp = client.bulk(body=body)
if not resp.get("errors"):
return resp
retry_items = [i for i in resp["items"]
if i.get("index", {}).get("status") == 429]
if not retry_items:
return resp
sleep = min(2 ** attempt + random.random(), 30)
time.sleep(sleep)
raise RuntimeError("bulk failed after retries")
指数退避加随机抖动,避免多个客户端同时重试形成尖峰。注意只重试被拒的条目,不要整批重发。
5.4 削峰与限流
背压的根本解法是让到达速率匹配处理能力。三种手段:
| 手段 | 说明 |
|---|---|
| 客户端限速 | 控制每秒提交的文档数 |
| 队列缓冲 | 上游用 Kafka 削峰,消费端按能力拉取 |
| 扩容 | 加节点或加分片,提高处理能力 |
对导入任务,用消息队列做缓冲是最稳的方案:上游按业务速率生产,下游消费端按集群能力消费,天然形成背压闭环。
6. 监控与排错
6.1 关键指标
| 指标 | 位置 | 含义 |
|---|---|---|
indexing.index_total | _nodes/stats | 累计写入文档数 |
indexing.index_time_in_millis | _nodes/stats | 累计写入耗时 |
indexing.throttle_time_in_millis | _nodes/stats | merge 限流时间 |
merges.current | _nodes/stats | 进行中的 merge 数 |
write.rejected | _cat/thread_pool | 写入拒绝数 |
translog.operations | _nodes/stats | translog 操作数 |
throttle_time_in_millis 持续增长说明段合并跟不上写入,需要减少 refresh 频率或降低写入速率。
6.2 常见症状与对策
| 症状 | 原因 | 对策 |
|---|---|---|
| 429 持续 | 并发过高、队列满 | 降并发、加退避 |
| 吞吐不升反降 | 单批过大、GC 频繁 | 减小批次体积 |
| CPU 打满 | 分词或 merge 重 | 简化分析器、调大 refresh |
| 磁盘 util 高 | translog fsync 频繁 | 可重放数据改 async |
| 单节点饱和 | routing 热点 | 打散 routing 或加分片 |
| 段数量暴涨 | refresh 过频 | 调大 refresh_interval |
6.3 排错顺序
遇到写入慢,按这个顺序排查:先看 _cat/thread_pool 是否有拒绝;再看 _nodes/stats 的 CPU 与 IO 是否饱和;再看 _cat/indices 的段数量与 merge 状态;最后看 _cat/shards 是否分片倾斜。多数问题在这四步内能定位。
7. 总结
| 环节 | 要点 |
|---|---|
| 批次大小 | 按体积 5~15 MB 控制,不按条数一刀切 |
| 并发度 | 固定并发 + 在途控制,观察拒绝率再调 |
| refresh | 导入期调大或关闭,结束恢复 |
| translog | 关键数据保持 request,可重放数据用 async |
| 分片路由 | 避免 routing 热点,分片数按堆内存规划 |
| 背压 | 429 是信号,退避重试 + 队列削峰 |
| 监控 | 盯 rejected、throttle_time、merge 数 |
| 原则 | 先建基线,单变量调整,逐项验证 |
批量写入调优没有通用最优解,它取决于文档大小、集群规格与可靠性要求。可靠的路径是:建立基线、定位瓶颈、单变量调整、验证收益,再进入下一轮。分片与分配策略可阅读《路由与分片分配》,整体性能手段可阅读《性能调优与缓存策略》。
延伸阅读
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。