Skip to content

Elasticsearch 深度解析 ​

#组件 · #Elasticsearch · #搜索引擎 · #倒排索引 · #全文检索

从集群架构到倒排索引,从写入链路到查询 DSL,覆盖 ES 核心设计思想与性能优化。

一、架构与进程模型 ​

1.1 分布式架构 ​

                            ┌──────────────────┐
                            │  Coordinator Node │  ← 任意节点都可接收请求
                            │  (协调节点)        │
                            └───┬──────────┬───┘
                    ┌───────────┘          └───────────┐
                    ▼                                   ▼
        ┌──────────────────────┐          ┌──────────────────────┐
        │   Data Node 1        │          │   Data Node 2        │
        │  ┌────┐ ┌────┐      │          │  ┌────┐ ┌────┐      │
        │  │ P0 │ │ R1 │      │          │  │ P0 │ │ R1 │      │
        │  │    │ │    │      │          │  │    │ │    │      │
        │  └────┘ └────┘      │          │  └────┘ └────┘      │
        │  ┌────┐             │          │  ┌────┐             │
        │  │ R2 │             │          │  │ P2 │             │
        │  └────┘             │          │  └────┘             │
        └──────────────────────┘          └──────────────────────┘

节点角色(7.x 后简化):

角色职责
Master集群管理(创建/删除索引、分片分配),需单数
Data存储数据、执行 CRUD、搜索、聚合
Ingest预处理管道(grok、date、rename 等)
Coordinating路由请求、合并结果(默认所有节点都是)
ML机器学习作业(X-Pack)

1.2 线程模型 ​

┌────────────────────────────────────────────────────────┐
│                   ES Node (JVM 进程)                    │
│                                                         │
│  ┌──────────────────────────────────────────────────┐  │
│  │         ThreadPool (线程池, 按操作隔离)             │  │
│  │                                                    │  │
│  │  ┌──────────┐ ┌──────────┐ ┌──────────┐          │  │
│  │  │ search   │ │ write    │ │ get      │  ...     │  │
│  │  │ 线程池   │ │ 线程池   │ │ 线程池   │          │  │
│  │  │ (读)     │ │ (写)     │ │ (单查)   │          │  │
│  │  └──────────┘ └──────────┘ └──────────┘          │  │
│  │  ┌──────────┐ ┌──────────┐ ┌──────────┐          │  │
│  │  │ bulk     │ │ management│ │ refresh  │          │  │
│  │  │ 线程池   │ │ 线程池    │ │ 线程池    │          │  │
│  │  └──────────┘ └──────────┘ └──────────┘          │  │
│  └──────────────────────────────────────────────────┘  │
│                                                         │
│  JVM Heap: Lucene segments + Fielddata + 查询缓存       │
│  Off-Heap: mmap 文件映射 (操作系统 Page Cache)          │
└────────────────────────────────────────────────────────┘

核心线程池:

线程池类型队列说明
searchfixed无界CPU 核心数 × 3/2 + 1
writefixed200单文档写入,CPU 核心数
bulkfixed50(可配)批量写入,CPU 核心数
getfixed1000按 ID 查询
refreshscaling—刷新 segment
flushscaling—flush translog

被拒绝时:队列满 → EsRejectedExecutionException


二、核心概念 ​

Index (索引)
 ├── Shard 0 (主分片, P0)
 │    ├── Segment 1  ← 不可变倒排索引文件 (磁盘)
 │    ├── Segment 2
 │    └── Segment N
 ├── Shard 1 (主分片, P1)
 └── Shard N (主分片, PN)

每个主分片可以有 N 个副本 (Replica),分布在不同节点。
概念类比说明
Index数据库一组相关文档的逻辑集合
Shard分表一个 Lucene 实例,水平拆分数据
Replica从库主分片的副本,提供高可用和读扩展
Document行JSON 格式的单条数据
Field列文档中的字段
MappingSchema字段类型定义,支持动态映射

三、倒排索引 ​

3.1 原理 ​

正排索引 (文档 → 词):
Doc1: "Elasticsearch is fast"
Doc2: "Redis is fast"
Doc3: "Elasticsearch is powerful"

倒排索引 (词 → 文档):
┌──────────────┬──────────────────────┐
│ Term         │ Posting List (文档列表)│
├──────────────┼──────────────────────┤
│ Elasticsearch│ [Doc1, Doc3]          │
│ fast         │ [Doc1, Doc2]          │
│ is           │ [Doc1, Doc2, Doc3]    │
│ powerful     │ [Doc3]                │
│ Redis        │ [Doc2]                │
└──────────────┴──────────────────────┘

3.2 Term Dictionary → Term Index ​

查询 "fast":
  Term Index (FST 有限状态转换器, 内存)
      ↓ 快速定位
  Term Dictionary (所有 term 排序列表, 磁盘)
      ↓ block 内二分查找
  Posting List (文档 ID 列表, 磁盘)
      ↓
  Doc Values / Stored Fields → 取出字段值

数据结构链:

结构存储位置作用
Term Index (FST)堆内存前缀压缩+Trie,快速定位 Term 在字典中的 block
Term Dictionary磁盘所有 Term 有序存储,block 内二分查找
Posting List磁盘每个 Term 的文档 ID 列表 (Frame of Reference + Roaring Bitmap 压缩)
Doc Values磁盘 (列存)排序/聚合/脚本时用,正排索引
Stored Fields磁盘 (行存)原始 JSON _source

3.3 Posting List 压缩 ​

文档 ID 列表 (未压缩):
[1, 3, 5, 8, 10, 12, 15]

Frame of Reference (FOR 编码):
1. 计算增量: [1, 2, 2, 3, 2, 2, 3]   ← 值更小
2. 按 block=128 分块
3. 每块取最小值 bit 长度,统一编码

Roaring Bitmap (稀疏/密集混合):
< 4096 个 docID → 数组
≥ 4096 个 docID → Bitmap (65536 bit = 8KB)

四、写入链路 ​

4.1 写流程图 ​

Client
  │ POST /index/_doc
  ▼
┌──────────────┐
│ Coordinating  │ 路由: shard = hash(_id) % num_primary_shards
│    Node       │
└──────┬───────┘
       │ 转发到主分片所在节点
       ▼
┌──────────────────────────────────────────────┐
│              主分片写入流程                     │
│                                               │
│  1. 写入 Translog (WAL, 顺序写, crash-safe)     │
│  2. 写入 In-Memory Buffer                     │
│  3. Refresh (默认 1s): Buffer → Segment        │
│     ↓ 此时数据可搜索                            │
│  4. Flush (默认 30m/512MB): Segment 持久化      │
│     + Translog 截断                             │
└──────────────────────────────────────────────┘

4.2 Refresh vs Flush vs Translog ​

In-Memory Buffer (不可搜索)
    │  refresh (1s 一次, 轻量)
    ▼
Segment (可搜索, OS Cache)
    │  fsync + commit point
    ▼
Disk (持久化)
操作触发条件作用代价
Refresh1s / index.refresh_intervalBuffer → Segment,数据可搜索轻量,生成新 segment
Flush30m / Translog 512MBSegment fsync + 写 commit point + 截断 Translog较重
Translog每次写WAL 日志,crash 后回放恢复顺序写,极快

4.3 Segment 合并 ​

碎片化 Segments:
[ S1 ][ S2 ][ S3 ][ S4 ][ S5 ]  ← 读需跨 segment

合并后:
[         S12_merged         ][ S3 ][ S45_merged ]

策略: TieredMergePolicy (分层)
- 大小相近的 segment 合并
- 合并过程消耗 CPU + 磁盘 I/O

五、查询与评分 ​

5.1 查询上下文 vs 过滤上下文 ​

json
{
  "query": {
    "bool": {
      "must": [     ← 查询上下文 (算分)
        { "match": { "title": "elasticsearch" } }
      ],
      "filter": [   ← 过滤上下文 (不算分, 可缓存)
        { "term": { "status": "published" } },
        { "range": { "date": { "gte": "2024-01-01" } } }
      ]
    }
  }
}
上下文是否算分是否缓存适用
Query (must/should)✅❌全文搜索
Filter (filter)❌✅精确匹配、范围、黑白名单

5.2 BM25 评分(7.0+ 默认) ​

BM25(D, Q) = Σ IDF(qi) × TF(qi, D)

IDF(qi) = ln(1 + (N - n(qi) + 0.5) / (n(qi) + 0.5))
                        ↑ N=总文档数, n(qi)=含qi的文档数

TF(qi, D) = f(qi,D) / (f(qi,D) + k1 × (1 - b + b × |D|/avgdl))
                        ↑ k1=1.2 (词频饱和度), b=0.75 (长度归一化)

比 TF-IDF 更好的原因:词频饱和(对出现多次的词不过度加权)+ 长度归一化(短文档不天然吃亏)。

5.3 查询执行流程 ​

协调节点
  │ Query Phase: 广播到所有相关分片
  │   ┌─────────────┐
  ▼   │ Shard 1: 查 │ → 返回 top N 的 docID + score
  │   │ Shard 2: 查 │ → 返回 top N 的 docID + score
  │   │ Shard 3: 查 │ → 返回 top N 的 docID + score
  │   └─────────────┘
  │ 协调节点合并排序 (选出全局 top N)
  │
  │ Fetch Phase: 按 docID 精确获取
  │   ┌─────────────┐
  ▼   │ Shard 1:    │ → 返回 _source 等字段
  │   │ Shard 2:    │ → 返回 _source 等字段
  │   └─────────────┘
  ▼
合并返回 Client

六、聚合分析 ​

6.1 聚合类型 ​

json
{
  "aggs": {
    "by_category": {           ← Bucket 聚合 (分桶)
      "terms": { "field": "category" },
      "aggs": {
        "avg_price": {         ← Metric 聚合 (计算)
          "avg": { "field": "price" }
        },
        "sales_over_time": {   ← Pipeline 聚合 (对聚合结果再计算)
          "date_histogram": { "field": "date", "interval": "month" }
        }
      }
    }
  }
}
类型示例说明
Bucketterms, range, date_histogram分组
Metricavg, sum, max, cardinality计算
Pipelinederivative, moving_avg二次计算
Matrixmatrix_stats矩阵运算

6.2 聚合的内存陷阱 ​

聚合依赖 Doc Values (列式存储, 磁盘)
大量 terms 聚合 → Fielddata (堆内存!) → OOM 风险

安全配置:
indices.breaker.fielddata.limit: 40%   ← 断路器
indices.breaker.request.limit: 60%

七、分词器 ​

7.1 组成 ​

Char Filter → Tokenizer → Token Filter
(字符预处理)  (分词)      (词处理)

示例: "I'm learning <b>Elasticsearch</b>"

1. HTML Strip Char Filter → "I'm learning Elasticsearch"
2. Standard Tokenizer    → [I'm, learning, Elasticsearch]
3. Lowercase Token Filter → [i'm, learning, elasticsearch]

7.2 常用分词器 ​

分词器适用示例输入 → 输出
standard通用英文The 2 QUICK Brown-Foxes → [the, 2, quick, brown, foxes]
ik_smart中文粗粒度南京市长江大桥 → [南京市, 长江大桥]
ik_max_word中文细粒度南京市长江大桥 → [南京市, 南京, 市长, 长江大桥, 长江, 大桥]
whitespace空白分割按空格分
keyword不分词整体作为一个 term

八、性能优化 ​

8.1 写入优化 ​

优化项方法
Bulk 批量5-15MB/批,而非单条发
调大 refresh_interval-1(索引完改回) 或 30s
副本数设为 0写完再改回 number_of_replicas: 1
关闭 swapbootstrap.memory_lock: true
增大 translog flushindex.translog.flush_threshold_size: 1gb

8.2 查询优化 ​

优化项原因
filter 替代 must精确匹配不用算分,且可缓存
限制 from+size 深度分页用 search_after 替代
避免 * 前缀通配查所有 term 极慢
_source 过滤_source: ["title"] 只返回需要的字段
routing指定路由减少扫描分片数

8.3 分片规划 ​

参数建议
单分片大小10-50GB(日志类可到 50-100GB)
单节点分片数堆内存每 GB ≤ 20 个分片
JVM 堆不超过 32GB(压缩指针),≤ 物理内存 50%

8.4 冷热分离 ​

Hot 节点 (SSD, 近期数据)  ← 高频写入/查询
  │ ILM 策略自动迁移
  ▼
Warm 节点 (HDD, 历史数据) ← 低频查询, 可缩副本
  │
  ▼
Cold/Delete → 冻结或删除

九、集群监控与运维 ​

关键 API ​

bash
# 集群健康
GET _cluster/health

# 节点状态
GET _nodes/stats

# 索引状态
GET _cat/indices?v&s=store.size:desc

# 分片分布
GET _cat/shards?v

# 慢查询日志
GET _nodes/stats/thread_pool?pretty

# 热点线程
GET _nodes/hot_threads

常见问题排查 ​

问题可能原因排查方式
写入拒绝bulk 队列满_cat/thread_pool 看 bulk rejected
查询慢深度分页/大聚合Profile API + slowlog
集群 RED主分片丢失_cat/shards 找 UNASSIGNED
OOM聚合 fielddata 爆炸断路器日志 + 设为 true 禁用 fielddata
节点离线GC/网络/磁盘_nodes/hot_threads + GC 日志

工程实践:ES 写入链路与倒排索引结构 ​

1. 写入链路全景 ​

mermaid
flowchart TD
    A["Client: POST /index/_doc"] --> B["Coordinating Node<br/>路由: hash(_id) % shard_count"]
    B --> C["Primary Shard"]
    C --> D["1. Write Translog<br/>(同步写入, 保证不丢)"]
    D --> E["2. Index Buffer<br/>(内存, 1s refresh → Segment)"]
    E --> F["3. 写入 Lucene Segment<br/>(磁盘, 不可变)"]
    C --> G["4. Replicate to Replica"]
    G --> H["Replica Shard"]
    H --> I["同步完成 → 返回 200"]

    subgraph "后台 Merge"
        J["多个小 Segment"] -->|"ES/Lucene Merge"| K["合并为大 Segment"]
        K --> L["删除旧 Segment + 清理 Translog"]
    end

各阶段延迟:

阶段典型耗时可控参数
Translog 写入~1ms磁盘 fsync 频率
Index Buffer → Segment~1srefresh_interval
Replication~1-5ms网络 RTT
Merge后台merge.policy.max_merged_segment=5gb

2. 倒排索引结构 ​

mermaid
flowchart LR
    subgraph "Term Dictionary (FST)"
        T1["elasticsearch"]
        T2["index"]
        T3["search"]
    end

    subgraph "Posting List for 'search'"
        P1["DocID=1<br/>TF=3, Positions: [7,15,22]"]
        P2["DocID=5<br/>TF=1, Positions: [42]"]
        P3["DocID=9<br/>TF=2, Positions: [3,18]"]
    end

    T3 --> P1
    T3 --> P2
    T3 --> P3

    subgraph "Skip List (加速合并)"
        S1["Level 0: Doc 1 → 5 → 9 → ..."]
        S2["Level 1: Doc 1 → 9 → ..."]
    end
text
倒排索引的数据结构:
  Term Dictionary: FST (Finite State Transducer)
    - 共享前缀压缩
    - 存储 Term → 在 Posting List 文件中的偏移

  Posting List: FOR (Frame of Reference) 编码
    - 存 DocID delta (差值压缩)
    - 存 TF (Term Frequency) 和 Position

  Skip List: 加速 AND/OR 合并
    - 跳过大段不相关的 DocID

3. 常见问题深度排查 ​

Shard 过多 ​

text
问题: 单节点 3000+ shard, 重启要 1 小时

原因:
  - 每个 shard = 一个 Lucene 实例 (独立的 JVM 堆外内存 + 线程)
  - shard 太多 → JVM 内存碎片 → GC 频繁
  - 1000 shard 大概 10GB 堆外内存开销

建议:
  - 单节点 shard 数 < 1000
  - 单 shard 大小 10-50GB 最佳
  - 使用 _shrink API 减少 shard 数

Segment Merge 打满 IO ​

text
现象: 磁盘 util 100%, 查询全部超时

原因:
  - 大量小 Segment 触发 Merge
  - Merge 是大 IO 操作 (读多个小 segment + 写一个大 segment)
  - 与其他查询竞争 IO

排查:
  GET _cat/segments?v  # 看每个 shard 的 segment 数

修复:
  PUT _cluster/settings
  {
    "transient": {
      "indices.store.throttle.max_bytes_per_sec": "20mb"  # 限制 merge IO
    }
  }

Deep Pagination 问题 ​

text
问题: GET /index/_search { "from": 100000, "size": 100 } 超时

原因:
  - ES 需要从每个 shard 取 from+size 条数据 (100100 条)
  - 在 coordinating node 合并后丢弃前 100000 条
  - 越翻页越慢

修复:
  1. 用 search_after (推荐):
     GET /index/_search { "size": 100, "search_after": [last_sort_value] }

  2. 用 scroll (批量导出):
     POST /index/_search?scroll=1m { "size": 1000 }

  3. 不要用 from/size 翻到 10000+

4. ES 排障命令速查 ​

bash
# 集群健康
GET _cluster/health

# 未分配 shard
GET _cat/shards?v&h=index,shard,prirep,state,unassigned.reason

# 节点热线程
GET _nodes/hot_threads?interval=1s

# pending tasks (堆积)
GET _cluster/pending_tasks

# 磁盘使用
GET _cat/allocation?v

# 慢日志
GET /_slowlog/index?pretty

# 查看 fielddata 内存
GET _cat/fielddata?v

参考 ​

批注模式

💬 文章评论

暂无评论,来说点什么吧 👇

编程学习笔记