批量写入调优与背压

讲解 Elasticsearch 批量写入调优与背压:bulk 批次大小与并发度的取舍、refresh 与 translog 的持久化策略、分片与 routing 对写入的影响、429 拒绝的成因与重试退避,以及写入链路的监控指标与排错方法,帮助导入任务在吞吐与可靠性之间找到平衡。

写入性能问题往往在数据量涨上来之后才暴露:昨天还能跑满的导入任务,今天开始报 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~50005~15 MB通用批量导入
5000~1000015~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/statsmerge 限流时间
merges.current_nodes/stats进行中的 merge 数
write.rejected_cat/thread_pool写入拒绝数
translog.operations_nodes/statstranslog 操作数

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 数
原则先建基线,单变量调整,逐项验证

批量写入调优没有通用最优解,它取决于文档大小、集群规格与可靠性要求。可靠的路径是:建立基线、定位瓶颈、单变量调整、验证收益,再进入下一轮。分片与分配策略可阅读《路由与分片分配》,整体性能手段可阅读《性能调优与缓存策略》。

延伸阅读

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「elasticsearch」更多文章

  1. 组合模板与索引生命周期
  2. 相关性调优与离线评测
  3. Elasticsearch 与 OpenSearch 对比