ES Shard Runbook
Architecture + shard scaling

Elasticsearch 架构与 Shard 扩容方案

这是一份面向生产实施的单页手册:先把 Elasticsearch 从集群、节点、分片、Lucene segment 到读写链路讲清楚,再给出 shard 扩容的决策树、命令模板、验证项与回滚路径。

静态 primary shard 数只能在索引创建时设定;已建索引不能直接修改。
动态 replica shard 数可在线调整,用于读扩展与故障冗余。
10-50GB 单 shard 常用目标区间,同时关注每 shard 文档数上限。
3 路径 rollover 面向未来,split 快速增 primary,reindex 最通用。
01 / distributed architecture

Elasticsearch 的架构骨架

Elasticsearch 是把 Lucene 索引分布式化后的搜索与分析系统。真正承载数据的是 shard; cluster、node、role、routing、allocation 都是在围绕 shard 的生命周期做组织。

Cluster

集群是状态与调度边界

一个 cluster 由多个互联 node 组成,共享 cluster state。cluster state 记录节点成员、 index metadata、mapping、setting、alias、shard 分配等信息,由 elected master 负责发布更新。

  • 客户端可连接任意节点,不必知道每个 shard 的物理位置。
  • 节点加入、离开、故障时,集群根据分配规则自动迁移或恢复 shard。
  • 稳定 master 很重要,因为建索引、改 mapping、rollover、恢复等都依赖 cluster state。
Node roles

节点角色决定职责边界

节点可以同时承担多个角色,也可以被专用化。大型集群通常会把 master、data、ingest、ML、 transform、coordinating 等职责拆开,避免控制面与数据面互相抢资源。

  • master-eligible:参与选主,维护集群元数据与 shard 分配决策。
  • data / data_hot / data_warm:保存 shard,执行 CRUD、search、aggregation。
  • ingest:执行 pipeline 预处理,适合高写入链路前置转换。
  • coordinating:每个节点天然具备,负责 scatter/gather 与 bulk 分发。
Shard

分片是横向扩展单位

每个 index 会被切成一个或多个 primary shard。每个 primary 可有若干 replica。 primary 负责写入主路径,replica 提供冗余,也能承担读请求。

  • 一个 shard 本质上是一个完整 Lucene index。
  • primary shard 数决定写入与数据分布的并行度。
  • replica shard 数决定故障容忍度与读吞吐潜力。
Client / Kibana
REST API、Bulk API、Search API、管理 API
Coordinating
接收请求,路由到 shard,汇总多 shard 结果
Cluster state
节点、索引、mapping、alias、routing、allocation 的共享元数据
Data nodes
承载 primary / replica shard,执行搜索、聚合、写入与恢复
Lucene
segment、inverted index、doc values、stored fields、merge
Storage
磁盘、水位线、文件句柄、I/O、snapshot repository
运行时心智模型

ES 的“分布式”主要解决四件事

  1. 把数据切开:document 通过 routing 进入某个 primary shard。
  2. 把副本散开:primary 与 replica 尽量不在同一节点,结合 zone awareness 提高可用性。
  3. 把请求扇出:一次 search 可能需要访问多个 shard,coordinating node 做汇总。
  4. 把故障收敛:节点故障后 replica 可提升为 primary,缺失副本由集群恢复。
设计关键点 shard 不是越多越好。更多 shard 会提高并行度和可分布性,但也增加 heap、文件句柄、 cluster state、查询 fan-out、恢复和 merge 成本。
02 / data model

从 Index 到 Segment

shard 扩容不是只改一个数字。要理解为什么要 split、reindex 或 rollover,需要先知道 index、mapping、routing、Lucene segment 之间的关系。

Index / data stream

逻辑数据集合

index 是查询、写入和管理的逻辑入口;日志、指标这类时间序列数据更推荐用 data stream, 由多个 backing index 组成,并通过 rollover 持续生成新的写入索引。

静态内容型数据常见形态:一个业务 alias 指向一个或多个 index。时间序列数据常见形态: data stream 或 alias + rollover。

Mapping / analyzer

字段与倒排结构

mapping 决定字段类型;analyzer 决定 text 字段如何分词。mapping 维度越大、segment 越多,每个 segment 为字段元数据付出的 heap 成本也越高。

扩容前同步 mapping 很关键。reindex 到新 index 时,目标 index 不会自动继承你想要的 shard 数、analyzer 或 dynamic template。

Lucene segment

不可变的小索引文件

写入先进入内存 buffer 与 translog,refresh 后产生可搜索 segment;segment 后续被 merge 成更大的 segment。删除并非立刻释放磁盘,而是在 merge 后回收。

扩容期间需要预留磁盘,因为 split 可能硬链接或复制 segment;reindex 则会同时保存新旧两份数据。

03 / request path

写入与搜索链路

shard 扩容的效果取决于瓶颈在哪:写入瓶颈多看 primary shard 与 ingest/bulk; 查询瓶颈多看 shard fan-out、replica、cache、聚合字段和 coordinating node。

Write path

写入:routing 到 primary,再复制到 replica

  1. 客户端 bulk/index 请求到任意节点。
  2. coordinating node 根据 index、document id 或自定义 routing 计算目标 primary shard。
  3. primary shard 写入、更新 translog,随后把操作复制到 replica。
  4. refresh 后文档对搜索可见;flush 会生成 commit point 并清理 translog。
扩容含义 增加 replica 不会增加单个 index 的写入分区数;要提升长期写入分布,通常需要更多 primary shard 的新索引,或对满足条件的索引执行 split。
Search path

搜索:scatter 到 shard,再 gather 汇总

  1. coordinating node 解析查询,决定要访问哪些 shard copy。
  2. query phase 在各 shard 本地执行,返回排序候选与聚合局部结果。
  3. fetch phase 拉取命中文档内容,coordinating node 合并排序、聚合、分页结果。
  4. 如果每次请求打到过多 shard,延迟会被 fan-out 和 reduce 放大。
读扩容优先级 查询压力大时,先确认查询是否跨了过多 index/shard;再评估增加 replica、加 data/coordinating node、收窄时间范围、改聚合字段或优化 mapping。
04 / shard sizing

Shard 容量设计方法

扩容前先算目标,而不是只凭“慢了”就加 shard。目标 shard 数来自数据量、增长率、查询并发、 节点数、可用区、副本数和单 shard 目标大小。

经验区间

单 shard 目标

常用 starting point 是每个 primary shard 约 10GB 到 50GB,或者不超过 约 200M documents。实际值要通过压测和生产指标校准:日志类写多读少、商品搜索类读多聚合多, 合理 shard 大小可能不同。

目标 primary 数
ceil(目标热数据量 / 目标单 primary shard 大小)
总 shard copy 数
primary_shards * (replicas + 1)
磁盘需求
源数据 * (replicas + 1) / 目标磁盘水位利用率
split 目标
target_shards = source_shards * N,且必须整除 number_of_routing_shards
示例

从症状反推方案

假设某业务写入索引当前 5 primary、1 replica,热数据已到 1.8TB,单 shard 约 360GB, 写入延迟和 recovery 时间都很高。目标单 primary 约 45GB,则目标 primary 大约是 40 个。

  • 如果原索引创建时 routing shards 支持 5 -> 10 -> 20 -> 40,优先评估 split。
  • 如果不满足 split 条件,创建 v2 index 后 reindex,并通过 alias 切换。
  • 如果只需要未来数据变好,更新 template 后 rollover 新写入索引即可。
  • 如果主要是读压力,先加 replica 与数据节点,不要误把 replica 当成写分区扩容。
现象 优先判断 更可能的动作
单 shard 过大,恢复慢,merge 压力高 primary shard 数不足或 rollover 周期过长 split / reindex / 缩短 rollover 周期
查询 QPS 高,CPU 高但写入正常 replica 不足、查询 fan-out 过大、聚合字段昂贵 加 replica、加 data/coordinating node、收窄查询范围
写入线程池排队,bulk 延迟上升 primary 分区并行度、ingest pipeline、磁盘 I/O 新索引更多 primary、优化 bulk、拆 ingest
cluster state 发布慢,master 压力高 index/shard/field 数过多 减少小索引小 shard,合并索引,治理 mapping
05 / decision tree

Shard 扩容方案选择

“扩容 shard”至少有四种含义:加节点、加 replica、让未来新索引有更多 primary、 或迁移/拆分已有索引。先选路径,再动手。

A. 加节点 + rebalance

适用于节点磁盘、CPU、I/O 不足,但单 index 的 shard 数设计还合理。新增 data node 后, ES 会在同 tier 内重新平衡 shard。

  • 风险低
  • 不改变索引结构
  • 解决资源不足,不解决单 shard 过大

B. 增加 replica

适用于读多写少或读 QPS 不够。replica 可在线调整,能提高读并发和容灾,但写入仍由 primary 分区决定。

  • 在线可调
  • 增加磁盘和写复制成本
  • 不能解决 primary shard 太大

C. Rollover 新索引

适用于时间序列或可按版本切换的写入链路。更新 template 的 primary 数后 rollover, 未来数据进入新 shard 布局,历史数据不迁移。

  • 停机影响最小
  • 历史大 shard 仍存在
  • 推荐配合 ILM 自动化

D. Split / Reindex

适用于已有索引必须改变 primary shard 数。split 更快但限制多;reindex 最通用但成本最高, 需要别名切换和回滚窗口。

  • 真正改变历史数据分布
  • 需要预留磁盘与变更窗口
  • 要提前设计 alias/双写/停写策略
本方案建议的默认路径 如果业务 index 已经过大且必须让历史数据也重新分布:先判断能否 split;能 split 就走 “短暂停写 + split + alias 切换”;不能 split 就走 “v2 index + reindex + final delta + alias 原子切换”。如果是日志/指标新数据,优先 rollover,而不是重搬全部历史。
06 / implementation

具体实施步骤

下面的命令以 Kibana Dev Tools 或 curl 思路书写。变量请按实际环境替换: my-index-v1my-index-v2my-write-alias

0.1

确认目标与变更窗口

明确本次要解决的是写入分区不足、读吞吐不足、单 shard 过大、节点容量不足,还是未来索引规划。 对 split/reindex,提前确认业务可接受的停写窗口,或准备双写/队列缓冲。

变量约定
INDEX=my-index-v1
NEW_INDEX=my-index-v2
ALIAS=my-write-alias
TARGET_PRIMARY_SHARDS=40
TARGET_REPLICAS=1
0.2

采集当前状态

记录集群健康、shard 大小、节点磁盘、mapping/settings、alias、写入速率和查询延迟。 这些数据既用于算目标,也用于变更后验收。

Kibana Dev Tools
GET _cluster/health?pretty
GET _cat/nodes?v&h=name,roles,heap.percent,cpu,load_1m,disk.used_percent,disk.avail
GET _cat/indices/my-index-v1?v&h=index,pri,rep,docs.count,store.size,pri.store.size
GET _cat/shards/my-index-v1?v&h=index,shard,prirep,state,docs,store,node
GET my-index-v1/_settings?include_defaults=true
GET my-index-v1/_mapping
GET _alias/my-write-alias
0.3

做 snapshot 与容量门禁

split 至少要能容纳目标索引恢复过程中的额外空间;reindex 会同时占用新旧两份数据。 如果磁盘接近 high watermark,先扩节点或释放空间。

容量与备份检查
GET _snapshot/_all
PUT _snapshot/prod-repo/pre_shard_scale_20260703?wait_for_completion=false
{
  "indices": "my-index-v1",
  "ignore_unavailable": false,
  "include_global_state": false
}

GET _cat/recovery/my-index-v1?v
GET _cluster/settings?include_defaults=true&filter_path=**.disk.watermark*
1.1

适用条件

选择 rollover 的前提是业务可以接受“历史索引保持旧 shard 布局,未来写入进入新布局”。 这是时间序列数据最自然的路径,推荐用 data stream 或 ILM 自动化。

1.2

更新 index template

把新 primary shard 数写入模板,确保下一代 backing index 或新 index 采用新布局。

更新模板示例
PUT _index_template/logs-template
{
  "index_patterns": ["logs-prod-*"],
  "template": {
    "settings": {
      "index.number_of_shards": 12,
      "index.number_of_replicas": 1,
      "index.refresh_interval": "5s"
    },
    "mappings": {
      "dynamic": false
    }
  },
  "priority": 200
}
1.3

执行 rollover

对 data stream 或 write alias 执行 rollover。可先 dry run 检查条件,再正式切换新写入索引。

rollover
POST my-write-alias/_rollover?dry_run=true
{
  "conditions": {
    "max_primary_shard_size": "45gb",
    "max_age": "7d"
  }
}

POST my-write-alias/_rollover
{
  "conditions": {
    "max_primary_shard_size": "45gb",
    "max_age": "7d"
  }
}
2.1

判断是否能 split

split 会基于源 index 创建一个 primary shard 更多的新 index。它比 reindex 快,但有硬约束: 源 index 必须 green,目标 index 不存在,目标 primary 数必须是源 primary 的倍数, 且要能整除源 index 的 index.number_of_routing_shards

检查 routing shard 与健康状态
GET _cluster/health/my-index-v1?wait_for_status=green&timeout=60s
GET my-index-v1/_settings?filter_path=*.settings.index.number_of_shards,*.settings.index.number_of_routing_shards
HEAD my-index-v2
2.2

进入变更窗口并阻断写入

split 要求源索引只读。应用侧需要暂停写入、切到队列缓冲,或确保没有写请求落到源 index。 对 write alias 场景,先停写再加写阻塞,避免静默失败。

设置只读
PUT my-index-v1/_settings
{
  "settings": {
    "index.blocks.write": true
  }
}
2.3

执行 split

目标 index 可覆盖部分 settings。扩容期间建议先把 target replica 设为 0,加速恢复; 完成验证后再恢复 replica。

split index
POST my-index-v1/_split/my-index-v2?wait_for_active_shards=1
{
  "settings": {
    "index.number_of_shards": 40,
    "index.number_of_replicas": 0,
    "index.blocks.write": null
  },
  "aliases": {
    "my-read-alias": {}
  }
}

GET _cat/recovery/my-index-v2?v
GET _cluster/health/my-index-v2?wait_for_status=green&timeout=30m
2.4

原子切换 alias

通过 _aliases 一次性从旧 index 移除写 alias,并把 alias 加到新 index,避免读写方看到半切换状态。

alias cutover
POST _aliases
{
  "actions": [
    { "remove": { "index": "my-index-v1", "alias": "my-write-alias" } },
    { "add":    { "index": "my-index-v2", "alias": "my-write-alias", "is_write_index": true } }
  ]
}

PUT my-index-v2/_settings
{
  "settings": {
    "index.number_of_replicas": 1,
    "index.refresh_interval": "1s"
  }
}
3.1

创建目标 index

reindex 不会替你配置目标 index。必须先创建好 settings、mapping、analyzer、routing 规则和目标 primary shard 数。导入期间可临时关闭 refresh、设 replica 为 0。

create target index
PUT my-index-v2
{
  "settings": {
    "index.number_of_shards": 40,
    "index.number_of_replicas": 0,
    "index.refresh_interval": "-1"
  },
  "mappings": {
    "dynamic": false,
    "properties": {
      "id": { "type": "keyword" },
      "title": { "type": "text" },
      "updated_at": { "type": "date" }
    }
  }
}
3.2

执行全量 reindex

wait_for_completion=false 返回 task id,避免长请求被中断。 大索引建议使用 slices 并限速,防止把生产写入和搜索打爆。

async reindex
POST _reindex?wait_for_completion=false&slices=auto&requests_per_second=5000
{
  "source": {
    "index": "my-index-v1"
  },
  "dest": {
    "index": "my-index-v2",
    "op_type": "create"
  },
  "conflicts": "proceed"
}

GET _tasks/<task_id>
POST _reindex/<task_id>/_rethrottle?requests_per_second=2000
3.3

处理增量数据

零停机迁移需要应用双写、CDC、消息队列回放,或有可靠的 updated_at 增量条件。没有这些机制时,采用短暂停写:停消费者/写入口,记录时间点,做最后一轮增量同步。

final delta 示例
POST _reindex?wait_for_completion=true&conflicts=proceed
{
  "source": {
    "index": "my-index-v1",
    "query": {
      "range": {
        "updated_at": { "gte": "2026-07-03T08:00:00Z" }
      }
    }
  },
  "dest": {
    "index": "my-index-v2"
  }
}
3.4

刷新、验证、切 alias

完成同步后恢复 refresh 与 replica,等待 green,再做 count、抽样查询、业务回归。 最后用 alias 原子切换。

validate and cutover
PUT my-index-v2/_settings
{
  "settings": {
    "index.refresh_interval": "1s",
    "index.number_of_replicas": 1
  }
}

POST my-index-v2/_refresh
GET _cluster/health/my-index-v2?wait_for_status=green&timeout=30m
GET my-index-v1/_count
GET my-index-v2/_count

POST _aliases
{
  "actions": [
    { "remove": { "index": "my-index-v1", "alias": "my-write-alias" } },
    { "add":    { "index": "my-index-v2", "alias": "my-write-alias", "is_write_index": true } }
  ]
}
4.1

适用条件

如果瓶颈是搜索并发、节点可用性或跨 AZ 冗余,而不是 primary shard 太大,可以在线增加 replica。前提是有足够 data node 和磁盘容纳新增副本。

4.2

在线调整副本数

调整后观察 shard recovery、搜索延迟、CPU 与磁盘水位。写入压力很高时,新增 replica 会增加复制成本,需谨慎。

replica update
PUT my-index-v1/_settings
{
  "settings": {
    "index.number_of_replicas": 2
  }
}

GET _cat/recovery/my-index-v1?v
GET _cluster/health/my-index-v1?wait_for_status=green&timeout=30m
GET _cat/shards/my-index-v1?v&h=index,shard,prirep,state,store,node
07 / verification

验收、监控与回滚

扩容不是 alias 切完就结束。保留旧 index 的只读窗口,持续观察恢复、延迟、错误率和业务指标, 才算完成。

数据一致性

对比 _count、按主键抽样、多条件查询结果、业务核心聚合。 如果存在更新/删除,需要验证 final delta 或双写回放是否覆盖。

集群健康

确认目标 index green,无 unassigned shard;恢复结束后观察 pending tasks、master CPU、 GC、磁盘水位和 thread pool rejected。

性能基线

对比扩容前后的 P50/P95/P99 查询延迟、bulk indexing latency、refresh/merge time、 search queue、indexing pressure。

回滚窗口

旧 index 至少保留 24 到 72 小时且设为只读。若新 index 异常,先停写,再用 _aliases 原子切回旧 index。

rollback command

alias 回滚模板

切回旧索引
POST _aliases
{
  "actions": [
    { "remove": { "index": "my-index-v2", "alias": "my-write-alias" } },
    { "add":    { "index": "my-index-v1", "alias": "my-write-alias", "is_write_index": true } }
  ]
}

回滚前确认旧 index 是否还包含切换后的写入。如果切换后新 index 已接受写入,回滚前要先决定 是否把这段数据回补到旧 index,避免数据丢失。

do not

常见误区

  • 不要试图在线修改已有 index 的 index.number_of_shards
  • 不要在磁盘接近 high watermark 时启动 split/reindex。
  • 不要只看总 shard 数;还要看单节点 shard 分布、单 shard 大小和 mapping 字段数。
  • 不要把 replica 当成写入 primary 分区扩容;它主要提升读能力和冗余。
  • 不要在没有 snapshot、无回滚 alias、无增量策略的情况下迁移生产大索引。
sources

参考来源

页面内容按 Elasticsearch 官方文档与 API 文档整理,命令模板需结合实际版本、权限与业务写入链路再落地。