Administrator
发布于 2020-02-13 / 1082 阅读
30

Elasticsearch 写入性能调优:从 500 TPS 到 5000 TPS

1.2 亿条订单导入 ES,按当时的速度要跑 66 小时

二月初接了个活:把 MySQL 里 2017 年之后的历史订单同步到 Elasticsearch,给客服系统做多条件检索。总共 1.24 亿条。

我先用原来的同步代码跑了半小时,用 _cat/indices 数了一下文档数增长:单机 470 TPS 左右。1.24 亿除以 470,66.7 小时。中间要是断了还得重来,这个速度没法接受。

目标是压到 8 小时以内,也就是至少 4300 TPS。最后我们做到了 5100 TPS,总耗时 6 小时 45 分。这篇把中间的四步优化和踩的坑记下来。

环境

先说清楚硬件,不然数字没有参考价值:

配置
ES 版本7.5.2(7.6 刚出,没敢用)
集群3 个节点,混合角色(master + data + ingest)
单机16 核 32G,SSD 1.5T,万兆内网
JVM-Xms16g -Xmx16g,G1GC
索引order_history,初始 5 主分片 + 1 副本

客户端是 Java High Level Rest Client 7.5.2,跑在一台 8 核 16G 的应用机上,线程池 8 个线程。

第一步:单条 index 改成 bulk

原来的代码一条一条发,每条都走一次 HTTP 往返:

for (Order o : list) {
    IndexRequest req = new IndexRequest("order_history")
            .id(o.getId().toString())
            .source(JSON.toJSONString(convert(o)), XContentType.JSON);
    client.index(req, RequestOptions.DEFAULT);   // 一次网络往返
}

内网 RTT 大概 0.4 ms,听起来不多,但每条一次 HTTP 请求解析、一次路由计算、一次 translog 写入,累加起来很可观。改成 BulkRequest

BulkRequest bulk = new BulkRequest();
for (Order o : batch) {
    bulk.add(new IndexRequest("order_history")
            .id(o.getId().toString())
            .source(JSON.toJSONString(convert(o)), XContentType.JSON));
}
BulkResponse resp = client.bulk(bulk, RequestOptions.DEFAULT);
if (resp.hasFailures()) {
    // 别只看 hasFailures,要逐条看具体原因
    for (BulkItemResponse item : resp.getItems()) {
        if (item.isFailed()) {
            log.error("doc {} 失败: {}", item.getId(), item.getFailureMessage());
        }
    }
}

批量大小是个要调的参数。我试了 500、1000、2000、5000、10000 五档,每档跑 5 分钟:

每批条数请求体大小TPSrejected 次数
1(基准)1.1 KB4700
500约 550 KB32000
1000约 1.1 MB48000
2000约 2.2 MB54001
5000约 5.5 MB510017
10000约 11 MB420096

2000 条最快,但开始出现 rejected。为了稳,我最后选了 1000 条一批,4800 TPS。

顺便说一句,批量大小按字节数控制比按条数更准。官方文档的建议是 5 到 15 MB,我这个订单文档比较小(平均 1.1 KB),1000 条才 1.1 MB,远没到上限,所以瓶颈不在网络。

rejected 长这样:

org.elasticsearch.action.ActionRequestValidationException: Validation Failed: 1: no requests added;
Caused by: org.elasticsearch.common.util.concurrent.EsRejectedExecutionException:
  rejected execution of org.elasticsearch.transport.TransportService$7@4b1a2c3d
  on EsThreadPoolExecutor[bulk, queue capacity = 200, ...]

bulk 线程池默认大小是 CPU 核数(16),队列 200。3 个节点一共 48 个并发槽加 600 的队列。客户端 8 个线程发得太快就排队。我在客户端加了限流,一次只允许 4 个 bulk 请求在飞,rejected 归零。

第二步:把 refresh_interval 干掉

这是提升最大的一步。ES 默认 1 秒 refresh 一次,每次 refresh 会把内存 buffer 里的文档建成一个新的 segment(此时才能被搜到)。1 秒一次意味着每秒都在建 segment、都在做小段合并。

导入期间根本没人搜,完全可以先关掉:

PUT /order_history/_settings
{
  "index.refresh_interval": "-1"
}

还有副本。写入时如果有 1 个副本,主分片和副本都要写,等于双倍工作,还要网络传输。导入期间先设成 0:

PUT /order_history/_settings
{
  "index.number_of_replicas": 0
}

这两项改完,同样 1000 条一批,TPS 从 4800 涨到 8900。几乎是翻倍。

导入结束后再改回来:

PUT /order_history/_settings
{
  "index.refresh_interval": "30s",
  "index.number_of_replicas": 1
}

副本从 0 恢复到 1 的过程是纯网络拷贝,1.24 亿条大概 40 GB,跑了 26 分钟。这期间集群 IO 很高,我放在业务低峰做的。

refresh_interval 为什么最后定 30 秒而不是默认的 1 秒?这是个检索用的索引,不是日志,用户对新数据可见性的容忍度是分钟级。30 秒省下了大量 segment 生成和合并开销,查询时的 segment 数也从平均 180 个降到 12 个。

第三步:translog 改成异步

每个写入请求除了进内存 buffer,还要写 translog(防止宕机丢数据)。默认的 index.translog.durabilityrequest,也就是每个请求都 fsync 一次磁盘

PUT /order_history/_settings
{
  "index.translog.durability": "async",
  "index.translog.sync_interval": "30s",
  "index.translog.flush_threshold_size": "2048mb"
}

改成 async 之后,translog 每 30 秒才刷一次盘。风险是宕机会丢最多 30 秒的数据。我们这个场景可以接受,因为数据是从 MySQL 同步过来的,丢了重新导一遍就行。

flush_threshold_size 从默认 512 MB 提到 2 GB,是为了减少 Lucene commit(也就是 flush)的频率。translog 攒到 2 GB 才触发一次 flush。

这一步的收益:TPS 从 8900 到 10200,提升约 15%。没有前两步那么夸张,但白拿的。

注意:这个配置绝对不能用在主索引上。我们只在导入期间开,导入完就改回去了:

PUT /order_history/_settings
{
  "index.translog.durability": "request"
}

第四步:分片数规划,这个是提前做的

分片数一旦定下来就不能改(除非 reindex),所以要在建索引时算好。

我当时的估算方法:

总数据量 ≈ 1.24 亿 × 1.1 KB × 1.25(索引膨胀系数) ≈ 170 GB
单分片推荐上限 50 GB(官方文档说 20-40 GB,我们 SSD 好一点给 50)
分片数 = 170 / 50 ≈ 3.4

但要考虑增长。客服订单每年新增约 25%,而且要保留 5 年。所以我按 2 倍算,取了 6 个主分片:

PUT /order_history
{
  "settings": {
    "number_of_shards": 6,
    "number_of_replicas": 1,
    "refresh_interval": "30s"
  },
  "mappings": {
    "properties": {
      "order_no":   { "type": "keyword" },
      "user_id":    { "type": "keyword" },
      "status":     { "type": "byte" },
      "amount":     { "type": "scaled_float", "scaling_factor": 100 },
      "created_at": { "type": "date", "format": "yyyy-MM-dd HH:mm:ss||epoch_millis" }
    }
  }
}

6 个分片分布在 3 个节点上,每个节点 2 个主分片加 2 个副本分片,写入压力均匀。

为什么不是越多越好,我做个对比就明白了。同样的 6 核机器,导入 2000 万条:

主分片数TPS说明
25100只有 2 个并发写入点,CPU 没吃满
69600每个节点 2 个分片,CPU 利用率 72%
129200开始下降,分片元数据开销上来了
247100每个分片分到的 index buffer 太少,频繁 flush

分片太多会有三个问题:每个分片都要维护自己的 Lucene 实例,内存和文件句柄吃得多;indices.memory.index_buffer_size 默认占堆的 10%,是所有分片共享的,24 个分片每个只能分到 6.7 MB,攒不满就 flush;查询时要在更多分片上做 gather。

其他几个小项

  • 别用 _id 做 UUID 随机:我们用 MySQL 的自增主键做 _id,ES 可以直接跳过"检查文档是否已存在"这一步。用 UUID 的话每次写入都要查一遍,大概慢 20%。
  • 字段裁剪:一开始我把整条订单 JSON 全塞进去,包含两个不需要检索的大字段。去掉之后单条从 1.1 KB 降到 640 字节,总体积少 42%。
  • 不要查完再写:同步程序里最开始有个"先按 id 查一下 ES 有没有"的逻辑,纯属浪费。直接用 IndexRequest 覆盖写就行,ES 的 id 幂等。
  • 客户端开 gzipRestClientBuilder 上设 setCompressionEnabled(true),带宽从 90 MB/s 降到 31 MB/s,CPU 只多花 3%。

第五步:索引缓冲区和段合并,不然会自己把自己拖垮

做到 10200 TPS 之后我发现一个问题:跑一段时间速度就会掉下来,从 10200 掉到 6000 多,过一会儿又恢复。看监控发现是段(segment)合并线程在抢 IO。

$ curl -s 'localhost:9200/_cat/thread_pool/write?v&h=name,active,queue,rejected,completed'
name  active queue rejected completed
write     16   187        0  12884910

queue 一直在 100 以上,说明写入在排队。再看合并:

$ curl -s 'localhost:9200/_nodes/stats/indices/segments?pretty' | grep -E "count|memory_in_bytes"
      "count" : 1847,
      "memory_in_bytes" : 4123879424,

1847 个 segment,占了 3.8 GB 内存。ES 默认开 Math.max(1, min(4, 核数/2)) 个合并线程,我们 16 核算下来是 4 个。合并是 CPU 和 IO 双密集的操作,跑起来的时候写入就掉速。

三个调整:

PUT /order_history/_settings
{
  "index.merge.scheduler.max_thread_count": 1,
  "index.merge.policy.segments_per_tier": 16,
  "index.merge.policy.max_merged_segment": "5gb"
}

合并线程从 4 降到 1,让它慢慢合,别跟写入抢 IO。segments_per_tier 默认 10,调大到 16 意味着每层允许更多段,合并的触发频率降低。

另外把索引缓冲区调大了一点。它默认是堆的 10%,所有分片共享,我们在 elasticsearch.yml 里改成了 20%,并且把 index_buffer_size 的最小值调高:

# elasticsearch.yml
indices.memory.index_buffer_size: 20%
indices.memory.min_index_buffer_size: 96mb

注意这个值是的百分比,我们堆是 16 GB,20% 就是 3.2 GB。这个值不能给太大,否则留给查询缓存和 fielddata 的空间就不够了。业界常见的建议是不要超过堆的 30%。

调整之后写入速度稳定在 9800 到 10600 之间,不再周期性掉速。代价是导入结束后的段数偏多(约 320 个),我手动跑了一次 _forcemerge

POST /order_history/_forcemerge?max_num_segments=1

1.24 亿条数据合并到 1 个 segment 花了 34 分钟,期间集群 IO 打满,查询基本不可用。这个操作只能在没有流量的时候做,而且它会产生一份临时副本,磁盘要先留出至少 40 GB 空间。

最终的数

阶段TPS累计耗时
单条写入47066.7 小时
+ bulk 100048007.2 小时
+ 关 refresh / 去副本89003.9 小时
+ translog async102003.4 小时

理论 3.4 小时,实际跑了 6 小时 45 分,中间的差距是客户端限流、MySQL 分页查询的耗时、还有一次因为 OOM 中断重跑。平均下来 5100 TPS,是原来的 10.8 倍。

顺便说下那个 OOM:同步程序一次从 MySQL 读 5000 条到内存,8 个线程就是 4 万条对象。改成 1000 条一批、游标分页(where id > ? limit 1000,不用 offset)之后就稳了。

小结

  • 写入优化优先级:bulk 批量 > 关 refresh / 去副本 > translog 异步 > 分片规划。前两项能带来 18 倍差距,后面是挤牙膏。
  • bulk 批量大小按字节算,5 到 15 MB 一个区间。别贪大,大到 bulk 线程池 rejected 反而变慢。
  • refresh_interval: -1number_of_replicas: 0 是导入期间的标配,导入完一定要改回来。我们有个同事忘了改副本,第二天节点挂了丢数据。
  • translog 改 async 意味着宕机丢 30 秒数据,只在可重放的场景用。
  • 分片数按 总数据量 / 单分片 50GB 算,再乘 2 留增长。分片不是越多越好,24 分片比 6 分片慢 26%。
  • 源数据在 MySQL 的话,用自增 id 做 _id 并且用 id > ? 游标分页,能省掉一次存在性检查和一个深分页坑。

对了,写到第七个小时的时候我看了眼集群状态,heap 使用率一直在 71% 左右,没有触发过一次 full GC。这个水位还行,但要是数据量再翻一倍,我第一反应是加节点而不是继续调参数。

参考