数据库基础体系 · 第 43/139 篇。文章以各产品官方稳定版本的公开语义为准;示例会明确引擎、事务与部署边界。

Elasticsearch 数据写入与运维:Bulk、Ingest、ILM、快照和升级

Elasticsearch 的数据写入与运维,不能只理解为“调用一个写入 API”。一条数据从客户端进入集群,通常会经历以下阶段:

客户端
  └─ Bulk 请求
       └─ Ingest Pipeline
            └─ 路由到目标分片
                 └─ 主分片写入
                      └─ 副本分片复制
                           └─ Refresh 后可搜索
                                └─ Translog 持久化
                                     └─ ILM 生命周期管理
                                          └─ Snapshot 备份

这些阶段解决的问题不同:

  • Bulk:降低大量写入时的协议和请求开销。
  • Ingest Pipeline:在写入前执行字段补充、清洗、解析和转换。
  • ILM(Index Lifecycle Management):根据索引年龄、大小或其他条件推进索引生命周期。
  • Snapshot:备份和恢复集群数据及部分集群元数据。
  • 升级:在版本变化、插件变化或架构变化时保持数据和服务可用。

它们之间存在依赖关系。例如:

  • Ingest Pipeline 处理的是“文档进入索引之前”的数据流;
  • ILM 管理的是索引,而不是单个文档;
  • Snapshot 保存的是索引分片及相关元数据,而不是某次 Bulk 请求;
  • 升级前需要先确认 Mapping、Pipeline、ILM、快照仓库和插件是否兼容。

本文中的示例以 Elasticsearch 官方稳定版本公开语义为准,假设使用单节点或小型集群进行演示。示例中的 localhost:9200、认证方式和节点数量需要根据实际部署替换。生产环境还应启用 TLS、认证和最小权限控制。


一、写入前必须明确的边界:文档、索引、分片和可见性

1.1 文档写入不是数据库事务

Elasticsearch 的基本写入单位是文档,例如:

{
  "@timestamp": "2025-01-01T12:00:00Z",
  "service": "payment",
  "level": "error",
  "message": "timeout"
}

文档属于某个索引,索引又被划分为多个主分片。一个文档通过路由值选择目标主分片:

shard=hash(routing)modPshard = hash(routing) \bmod P

其中:

  • routing 默认是文档 _id
  • P 是主分片数量;
  • hash 是 Elasticsearch 使用的路由哈希函数。

如果指定了自定义 routing,那么同一个 routing 值的文档会被路由到同一个主分片。这可以保证相关文档集中,但也可能造成热点分片。

Elasticsearch 不提供跨多个文档、多个分片的通用 ACID 事务。一次 Bulk 请求也不是一个整体事务:

  • 某些 item 成功;
  • 某些 item 可能因 Mapping 冲突、版本冲突、超时或权限错误失败;
  • 成功的 item 不会因为其他 item 失败而自动回滚。

因此,不能把“Bulk 请求返回 HTTP 200”理解为“所有文档都写入成功”。

1.2 写入成功和搜索可见不是同一件事

写入过程至少涉及两个概念:

  1. 持久化:数据写入事务日志(translog)并按配置持久化。
  2. 搜索可见:数据经过 refresh 后进入可搜索的 Lucene segment。

默认情况下,Elasticsearch 采用周期性 refresh,因此一个刚刚成功写入的文档可能不会立即被普通搜索查到。

如果业务需要等待文档可搜索,可以使用:

POST logs-2025.01.01/_doc/1?refresh=wait_for
{
  "message": "payment timeout"
}

refresh=wait_for 会等待下一次 refresh,而不是为每个文档强制执行一次 refresh。强制 refresh=true 会增加刷新开销,不应作为批量写入的默认设置。

在写入和检索之间存在以下典型状态:

请求被接受
  ↓
主分片写入成功
  ↓
副本满足写入确认条件
  ↓
客户端收到成功响应
  ↓
refresh
  ↓
搜索可见

写入响应成功但立即搜索不到,并不必然表示数据丢失。


二、Bulk API:批量写入的协议、并发和失败处理

2.1 Bulk 使用 NDJSON,而不是一个 JSON 数组

Bulk API 的请求体是 NDJSON(Newline Delimited JSON)。每个操作通常由两行组成:

操作元数据
文档内容

一个完整示例:

curl -X POST 'http://localhost:9200/_bulk' \
  -H 'Content-Type: application/x-ndjson' \
  --data-binary $'{"index":{"_index":"logs-2025.01.01","_id":"1"}}\n{"@timestamp":"2025-01-01T12:00:00Z","service":"payment","level":"error","message":"timeout"}\n{"create":{"_index":"logs-2025.01.01","_id":"2"}}\n{"@timestamp":"2025-01-01T12:00:01Z","service":"order","level":"warn","message":"slow response"}\n'

最后必须包含换行符。--data-binary 比普通 -d 更适合传递原始换行格式。

常见操作包括:

  • index:写入文档;指定 _id 时,已存在文档会被覆盖。
  • create:只创建新文档;如果 _id 已存在,则该 item 失败。
  • update:部分更新或执行脚本。
  • delete:删除文档。

update 示例:

curl -X POST 'http://localhost:9200/_bulk' \
  -H 'Content-Type: application/x-ndjson' \
  --data-binary $'{"update":{"_index":"users","_id":"u-1"}}\n{"doc":{"last_login":"2025-01-01T12:00:00Z"},"doc_as_upsert":false}\n'

其中:

  • 第一行声明目标索引和文档 ID;
  • 第二行是 update 请求体;
  • doc 表示部分字段更新;
  • doc_as_upsert=false 表示文档不存在时不自动创建。

2.2 Bulk 响应必须逐项检查

Bulk 请求通常返回 HTTP 200,但响应中有顶层字段:

{
  "took": 12,
  "errors": true,
  "items": [
    {
      "index": {
        "_index": "logs-2025.01.01",
        "_id": "1",
        "status": 201,
        "result": "created"
      }
    },
    {
      "index": {
        "_index": "logs-2025.01.01",
        "_id": "2",
        "status": 400,
        "error": {
          "type": "mapper_parsing_exception",
          "reason": "failed to parse field"
        }
      }
    }
  ]
}

正确的判断方式是:

  1. 检查 HTTP 层是否返回错误;
  2. 检查响应中的 errors
  3. 遍历 items
  4. 根据每个 item 的状态码和错误类型决定是否重试。

错误处理不能简单地“重发整个 Bulk 请求”。因为其中一部分 item 可能已经成功:

第一次请求:
  item A 成功
  item B 临时失败
  item C 成功

错误重试整个请求:
  item A 可能被覆盖或产生版本冲突
  item B 重试
  item C 可能被重复处理

更稳妥的做法是只提取失败 item,并根据操作语义重试。

2.3 重试要求幂等性

所谓幂等,是指同一个操作重复执行,最终状态与执行一次相同。

以下操作通常更容易设计为幂等:

{"index":{"_index":"events","_id":"event-1001"}}
{"event_id":"event-1001","status":"paid"}

同一个 _id 重试时,最终得到同一份文档。

create、脚本更新和自动生成 _id 需要更谨慎:

  • create 重试可能得到 409 version_conflict_engine_exception
  • 使用自动 _id 时,每次重试可能产生新文档;
  • ctx._source.count += params.delta 这样的脚本重试可能造成重复累加;
  • 如果事件本身带有唯一业务 ID,可以将其作为 _id,或者在应用侧保存处理状态。

对于临时故障,可以采用指数退避:

第 1 次重试:100 ms
第 2 次重试:200 ms
第 3 次重试:400 ms
……

但重试次数必须有限,并区分错误类型:

  • 429:通常表示资源暂时不足,可以退避重试;
  • 502503、连接超时:可能是临时故障;
  • 400 Mapping 解析错误:重试不会解决问题;
  • 403 权限错误:应修正权限;
  • 409:需要根据版本控制或业务幂等逻辑处理。

2.4 Bulk 大小不是越大越好

Bulk 可以减少 HTTP 请求次数,但单个请求过大也会带来问题:

  • 占用更多客户端、网络和节点内存;
  • 增加请求排队和单次失败的重试成本;
  • 可能触发 HTTP 请求大小限制或代理限制;
  • 让单个请求长时间占用写入线程池;
  • 造成 GC 压力。

批量大小应根据以下因素压测确定:

  • 文档平均大小;
  • 字段数量和分析成本;
  • 主分片数量;
  • 副本数量;
  • ingest processor 复杂度;
  • 集群磁盘和 CPU;
  • 客户端并发数。

更合理的控制方式是同时限制:

每批文档数量
每批字节数
客户端并发 Bulk 数量
单个请求超时时间
失败重试队列长度

客户端可以根据 429、请求延迟和失败率动态降低并发,而不是无限增加并发。

2.5 refreshwait_for_active_shards 和版本控制

Bulk 支持多个影响写入行为的参数。

refresh

POST _bulk?refresh=wait_for

要求该 Bulk 中受影响的分片在 refresh 后再返回。它适合测试或少量需要立即检索的数据,不适合高吞吐写入的默认模式。

wait_for_active_shards

POST _bulk?wait_for_active_shards=all

它要求指定数量的活动分片副本满足条件后再确认写入。值为 all 时,要求主分片及其所有副本处于可用状态。

这会提高写入确认要求,但也意味着:

  • 某个副本故障时,写入可能等待或失败;
  • 集群恢复期间吞吐可能下降;
  • 它不是跨文档事务。

外部版本号和并发控制

如果数据来自外部系统,可以使用外部版本控制,让 Elasticsearch 根据版本号拒绝旧事件。使用时必须明确版本单调性和重试行为,否则乱序事件可能被拒绝。

对于读取后修改的场景,if_seq_noif_primary_term 可以实现乐观并发控制:

PUT products/_doc/p-1?if_seq_no=7&if_primary_term=2
{
  "price": 99.9
}

如果文档已经被其他请求修改,序列号不匹配,写入会失败而不是覆盖新数据。


三、Ingest Pipeline:写入前的数据处理链

3.1 Pipeline 的位置和作用

Ingest Pipeline 是 Elasticsearch 节点在写入索引之前执行的一组处理器(processor)。

简化数据流如下:

Bulk item
  ↓
选择 pipeline
  ↓
处理器按顺序执行
  ↓
处理成功:继续索引
处理失败:执行 on_failure 或返回错误
  ↓
路由到目标分片
  ↓
写入主分片和副本

Pipeline 处理的是写入请求中的文档 _source 及相关字段。它不是查询时的转换,也不是已存在索引数据的自动回填机制。

常见处理器包括:

  • set:设置字段;
  • rename:重命名字段;
  • remove:删除字段;
  • grok:从字符串中提取结构化字段;
  • date:解析日期;
  • convert:类型转换;
  • json:解析 JSON 字符串;
  • script:执行 Painless 脚本;
  • pipeline:调用另一个 pipeline。

3.2 创建一个可验证的 Pipeline

下面的 Pipeline 从日志消息中提取状态码,并标记数据来源:

curl -X PUT 'http://localhost:9200/_ingest/pipeline/logs-normalize' \
  -H 'Content-Type: application/json' \
  -d '{
    "description": "Normalize application logs",
    "processors": [
      {
        "set": {
          "field": "event.ingested_by",
          "value": "logs-normalize"
        }
      },
      {
        "grok": {
          "field": "message",
          "patterns": [
            "%{GREEDYDATA:message_text} status=%{NUMBER:http.status_code:int}"
          ],
          "ignore_missing": true
        }
      }
    ],
    "on_failure": [
      {
        "set": {
          "field": "error.pipeline",
          "value": "logs-normalize"
        }
      }
    ]
  }'

输入:

{
  "@timestamp": "2025-01-01T12:00:00Z",
  "message": "request failed status=503"
}

使用 _simulate 验证:

curl -X POST 'http://localhost:9200/_ingest/pipeline/logs-normalize/_simulate' \
  -H 'Content-Type: application/json' \
  -d '{
    "docs": [
      {
        "_source": {
          "@timestamp": "2025-01-01T12:00:00Z",
          "message": "request failed status=503"
        }
      }
    ]
  }'

预期结果中会出现类似字段:

{
  "event": {
    "ingested_by": "logs-normalize"
  },
  "http": {
    "status": {
      "code": 503
    }
  }
}

_simulate 的价值在于:

  • 不需要真正写入索引;
  • 可以验证字段路径;
  • 可以观察处理器执行后的 _source
  • 可以检查 on_failure 的行为。

3.3 Pipeline 的选择规则

写入时可以显式指定:

POST logs-2025.01.01/_doc/1?pipeline=logs-normalize

Bulk 中也可以指定:

POST _bulk?pipeline=logs-normalize

还可以在索引设置中配置默认 Pipeline 和最终 Pipeline。需要区分两者:

  • Default pipeline:当请求没有显式指定 pipeline 时使用;
  • Final pipeline:作为最终处理步骤执行,不能被普通请求通过选择其他 pipeline 替代。

显式指定 pipeline=_none 可以绕过默认 pipeline,但不能用来绕过 final pipeline。

这带来一个重要边界:如果 final pipeline 中有强制字段清洗,那么运维人员执行回灌、重建索引或修复数据时,必须确认这条 Pipeline 是否适用于所有输入数据。

3.4 处理失败、部分成功和死信

Pipeline 处理器可能因为以下原因失败:

  • 输入字段不存在;
  • 类型不是预期类型;
  • Grok 模式无法匹配;
  • 脚本访问了不存在的路径;
  • 日期格式不匹配;
  • Painless 脚本超时或被限制。

如果没有 on_failure,该文档通常会失败,Bulk 只会在对应 item 中返回错误。

如果设置了 on_failure,可以补充错误字段:

{
  "processors": [
    {
      "date": {
        "field": "@timestamp",
        "formats": ["ISO8601"]
      }
    }
  ],
  "on_failure": [
    {
      "set": {
        "field": "error.message",
        "copy_from": "_ingest.on_failure_message"
      }
    },
    {
      "set": {
        "field": "error.pipeline_failed",
        "value": true
      }
    }
  ]
}

失败文档有两种常见策略:

  1. 让写入失败,由客户端记录并进入重试或死信队列;
  2. 将错误信息写入专门的失败索引。

第二种方式需要谨慎:如果错误文档再次使用同一条失败 Pipeline,可能形成循环。失败索引通常应使用不同的写入路径或明确指定绕过相关 Pipeline。

3.5 Pipeline 的性能和可重复执行

Pipeline 运行在集群节点上,因此复杂处理会消耗集群 CPU,而不是免费地消耗客户端资源。

特别需要关注:

  • grok 模式匹配成本;
  • Painless 脚本循环;
  • 大字段 JSON 解析;
  • 外部资源访问型处理器;
  • 多层 pipeline 嵌套。

Pipeline 不是天然幂等。例如:

{
  "script": {
    "source": "ctx.retry_count += 1"
  }
}

同一个文档重试时会重复增加计数。若写入系统允许重试,Pipeline 应尽可能基于输入字段计算结果,而不是无条件累加状态。


四、Mapping 与写入管道的关系

Ingest Pipeline 负责“如何改造文档”,Mapping 负责“改造后的字段如何被索引”。

例如 Pipeline 产生:

{
  "http": {
    "status": {
      "code": 503
    }
  }
}

如果 Mapping 将 http.status.code 定义为 integer,它可以进行数值范围查询和聚合:

{
  "query": {
    "range": {
      "http.status.code": {
        "gte": 500
      }
    }
  }
}

如果该字段被动态映射成 keyword,虽然仍能进行精确匹配,但数值范围语义和聚合行为会不同。

因此,写入链路的顺序通常是:

原始文档
  ↓
Pipeline 解析和补充字段
  ↓
Dynamic Mapping 或显式 Mapping
  ↓
倒排索引、Doc Values、Stored Fields

一个典型索引模板可以预先声明字段:

curl -X PUT 'http://localhost:9200/_index_template/logs-template' \
  -H 'Content-Type: application/json' \
  -d '{
    "index_patterns": ["logs-*"],
    "priority": 100,
    "template": {
      "settings": {
        "index.default_pipeline": "logs-normalize"
      },
      "mappings": {
        "dynamic": "strict",
        "properties": {
          "@timestamp": {
            "type": "date"
          },
          "service": {
            "type": "keyword"
          },
          "message": {
            "type": "text"
          },
          "http": {
            "properties": {
              "status": {
                "properties": {
                  "code": {
                    "type": "integer"
                  }
                }
              }
            }
          }
        }
      }
    }
  }'

dynamic: strict 的效果是:出现未声明字段时拒绝写入,而不是自动扩展 Mapping。它能防止拼写错误或恶意字段造成 Mapping 爆炸,但会提高数据生产端对字段演进的要求。

如果 Mapping 类型已经确定,不能通过修改 Mapping 把已有 text 字段直接改成 keyword。通常需要:

  1. 创建新的索引;
  2. 使用正确的 Mapping;
  3. 通过 _reindex 配合 Pipeline 迁移;
  4. 验证后切换别名。

五、索引生命周期管理(ILM)

5.1 ILM 管理的是索引状态

ILM(Index Lifecycle Management)根据索引生命周期策略推进索引状态。它管理的是索引或数据流 backing index,不是单个文档。

典型阶段包括:

hot → warm → cold → frozen → delete

阶段名称表达的是生命周期意图:

  • hot:持续写入、查询频繁;
  • warm:基本不再写入,但仍可能查询;
  • cold:查询较少,通常更重视成本;
  • frozen:极少访问,可使用更低成本的存储形式;
  • delete:到期删除。

阶段本身不等于自动把数据移动到另一台机器或另一套存储。具体动作取决于配置,例如分配到带有特定属性的节点、执行 rollover、force merge 或删除索引。

5.2 ILM Policy 的基本结构

下面是一个简化策略:

curl -X PUT 'http://localhost:9200/_ilm/policy/logs-policy' \
  -H 'Content-Type: application/json' \
  -d '{
    "policy": {
      "phases": {
        "hot": {
          "actions": {
            "rollover": {
              "max_primary_shard_size": "30gb",
              "max_age": "1d"
            }
          }
        },
        "warm": {
          "min_age": "7d",
          "actions": {
            "forcemerge": {
              "max_num_segments": 1
            },
            "allocate": {
              "number_of_replicas": 1
            }
          }
        },
        "delete": {
          "min_age": "30d",
          "actions": {
            "delete": {}
          }
        }
      }
    }
  }'

这里的语义是:

  1. 索引进入 hot 阶段;
  2. 当索引年龄达到一天,或主分片大小达到 30 GB 时满足 rollover 条件;
  3. 索引进入 warm 后,执行 force merge,并设置副本数;
  4. 到达 30 天后删除。

不同阶段的动作受索引状态约束。例如 force merge 会产生新的 segment,并消耗大量磁盘和 I/O。它通常应在索引不再写入后执行,而不是对活跃写入索引频繁执行。

5.3 Rollover 需要写入别名或数据流

Rollover 的核心是:当当前写入索引满足条件时,创建下一个索引,并把写入目标切换过去。

使用写入别名

创建第一个索引:

curl -X PUT 'http://localhost:9200/logs-000001' \
  -H 'Content-Type: application/json' \
  -d '{
    "aliases": {
      "logs-write": {
        "is_write_index": true
      }
    },
    "settings": {
      "index.lifecycle.name": "logs-policy",
      "index.lifecycle.rollover_alias": "logs-write",
      "number_of_shards": 1,
      "number_of_replicas": 1
    },
    "mappings": {
      "properties": {
        "@timestamp": {"type": "date"},
        "message": {"type": "text"},
        "service": {"type": "keyword"}
      }
    }
  }'

之后写入别名:

curl -X POST 'http://localhost:9200/logs-write/_doc' \
  -H 'Content-Type: application/json' \
  -d '{
    "@timestamp": "2025-01-01T12:00:00Z",
    "service": "payment",
    "message": "request failed"
  }'

客户端不需要知道当前实际索引是 logs-000001 还是后续的 logs-000002

当 rollover 条件满足后,ILM 创建新索引,并将:

logs-write → logs-000001

切换为:

logs-write → logs-000002

旧索引不再是写入索引,但仍可通过模式 logs-* 查询。

使用数据流

对于时间序列数据,数据流通常比手动管理索引别名更适合。数据流包含:

  • 一个逻辑数据流名称;
  • 一个写入 backing index;
  • 若干历史 backing index;
  • 与数据流关联的索引模板。

写入数据流时,文档通常必须包含时间字段,模板需要定义 data_stream,并配置匹配的数据流索引模板。数据流的 rollover 和 backing index 管理由系统协同完成,客户端只写入数据流名称。

5.4 ILM 的执行不是实时定时器

ILM 由集群中的生命周期协调机制异步执行。满足条件后,不代表动作在同一毫秒发生。

查看策略:

curl 'http://localhost:9200/_ilm/policy/logs-policy?pretty'

查看索引当前生命周期状态:

curl 'http://localhost:9200/logs-000001/_ilm/explain?pretty'

结果中通常可以看到:

  • 当前 policy;
  • 当前 phase;
  • 当前 action;
  • 当前 step;
  • 是否出错;
  • 错误信息;
  • 是否正在等待某个条件。

诊断思路是按状态机逐步检查:

是否已关联 ILM policy?
  ↓
是否进入预期 phase?
  ↓
当前 action 是什么?
  ↓
当前 step 是否失败?
  ↓
失败原因是磁盘、分配、权限、别名还是条件未满足?

不要直接反复调用 retry 或手工修改索引状态。应先检查错误原因。例如:

  • rollover alias 未配置或不是写索引;
  • 分片大小条件尚未达到;
  • 目标节点缺少满足分配条件的资源;
  • 磁盘水位阻止分片分配;
  • force merge 仍在执行;
  • policy 在升级后引用了不兼容配置。

5.5 ILM 与数据保留的实际含义

“保留 30 天”并不一定等于“每条文档精确保留 30 × 24 小时”。

ILM 通常以索引年龄和索引级动作推进。一个索引可能因为 rollover 条件、写入速率和创建时间不同,导致其中的文档实际年龄分布不同。

例如:

logs-000001 创建于 1 月 1 日
1 月 1 日到 1 月 10 日持续写入
1 月 10 日 rollover

如果 delete 阶段按索引年龄计算,那么整个 logs-000001 可能在 30 天后被删除,但其中最晚写入的文档只保存了约 21 天。

如果业务要求“每条事件至少保存 30 天”,需要选择合适的 rollover 周期和删除策略,或者以数据流和索引边界为基础计算保留窗口,而不能只看 policy 中的数字。


六、Snapshot:备份、恢复与灾难边界

6.1 Snapshot 不是导出 JSON

Snapshot 是 Elasticsearch 的分布式备份机制,保存索引分片数据以及相关元数据。它与 _search 导出或逐条 Bulk 导出不同:

  • Snapshot 保留 Lucene segment 级数据;
  • 后续 Snapshot 通常复用仓库中已经存在的 segment;
  • 备份和恢复效率取决于数据变化、仓库和网络;
  • Snapshot 不等于实时复制;
  • Snapshot 不会自动把数据恢复到另一套集群,恢复仍需执行 restore。

Snapshot repository 可以使用 Elasticsearch 支持的仓库类型,例如共享文件系统、对象存储等。仓库必须满足集群所有相关节点的访问和权限要求。不能只在某一个节点本地配置一个目录,然后认为集群已经具备可靠备份。

6.2 注册仓库和验证

以共享文件系统为例,首先需要在节点配置允许的仓库路径,并确保所有相关节点可以访问。然后注册仓库:

curl -X PUT 'http://localhost:9200/_snapshot/backup-repo' \
  -H 'Content-Type: application/json' \
  -d '{
    "type": "fs",
    "settings": {
      "location": "/srv/elasticsearch-backups",
      "compress": true
    }
  }'

验证仓库:

curl -X POST 'http://localhost:9200/_snapshot/backup-repo/_verify?pretty'

如果验证失败,应检查:

  • 节点是否都能访问该路径;
  • Elasticsearch 进程用户是否有读写权限;
  • path.repo 是否配置正确;
  • 仓库是否实际挂载;
  • 对象存储凭证和网络策略是否正确。

6.3 创建和查看 Snapshot

创建 Snapshot:

curl -X PUT 'http://localhost:9200/_snapshot/backup-repo/snapshot-2025-01-01?wait_for_completion=false' \
  -H 'Content-Type: application/json' \
  -d '{
    "indices": "logs-*",
    "include_global_state": false,
    "feature_states": []
  }'

wait_for_completion=false 会立即返回任务信息,创建过程在后台执行。

查看 Snapshot:

curl 'http://localhost:9200/_snapshot/backup-repo/snapshot-2025-01-01?pretty'

查看进行中的任务:

curl 'http://localhost:9200/_snapshot/_status?pretty'

Snapshot 期间应关注:

  • 快照是否完成;
  • 是否有分片失败;
  • 仓库空间是否足够;
  • 节点磁盘和网络负载;
  • 是否影响线上查询和写入延迟。

“创建请求成功”不等于“备份已经可恢复”。必须检查 Snapshot 的最终状态,并定期进行恢复演练。

6.4 Snapshot 的一致性和失败处理

Snapshot 针对集群中的索引状态生成一致的备份视图。正在写入的索引也可以被 Snapshot,但写入并不会因此自动停止。快照捕获的是特定时刻的数据状态,之后的新写入不会出现在已经完成的快照中。

如果某些分片快照失败:

  • Snapshot 可能以部分失败状态结束;
  • 该快照不能被当作完整备份;
  • 应查看失败分片及具体原因;
  • 修复仓库、磁盘、权限或分片状态后重新创建。

不能只根据仓库中“存在文件”判断备份有效。

6.5 恢复 Snapshot

恢复前先查看内容:

curl -X POST 'http://localhost:9200/_snapshot/backup-repo/snapshot-2025-01-01/_restore?wait_for_completion=false' \
  -H 'Content-Type: application/json' \
  -d '{
    "indices": "logs-*",
    "include_global_state": false,
    "rename_pattern": "logs-(.+)",
    "rename_replacement": "restored-logs-$1"
  }'

这里使用重命名避免与现有索引冲突。

恢复时需要注意:

  • 目标索引不能与现有索引直接冲突;
  • 现有同名索引可能需要关闭、删除或使用重命名;
  • 恢复后分片需要重新分配和恢复;
  • 索引模板、Pipeline、ILM policy 等集群级配置不一定随 include_global_state=false 一起恢复;
  • 恢复到不同版本时必须遵守官方支持的 Snapshot 兼容范围;
  • 插件字段、分析器和自定义模块必须在目标集群可用。

通常不建议在生产集群中无条件恢复 global_state。全局状态可能覆盖或引入集群级模板、持久化设置、别名和其他配置。更安全的做法是明确选择需要恢复的索引和功能状态,并单独检查集群配置。

6.6 SLM 与备份验证

SLM(Snapshot Lifecycle Management)用于自动创建和删除 Snapshot。它适合将备份策略固化为计划任务,例如:

每天创建一次快照
保留最近若干份
定期删除过期快照

但自动化调度不能替代恢复验证。至少需要验证:

  1. 快照任务是否按计划执行;
  2. 最终状态是否成功;
  3. 仓库容量是否接近上限;
  4. 最近快照是否包含关键索引;
  5. 在隔离环境能否实际恢复;
  6. 恢复后的 Mapping、别名、Pipeline 和查询是否可用。

如果 Snapshot 仓库与生产集群处于同一个故障域,例如同一磁盘阵列或同一可用区,那么生产故障可能同时破坏数据和备份。备份的可靠性取决于仓库的独立性,而不只是 API 返回成功。


七、升级:集群版本、数据格式和回滚边界

7.1 升级前必须区分几类兼容性

Elasticsearch 升级至少涉及以下层次:

  1. 节点协议兼容性:不同版本节点能否暂时组成混合版本集群。
  2. 索引数据兼容性:旧版本创建的索引能否被新版本读取。
  3. API 兼容性:客户端请求、响应字段和废弃参数是否仍可用。
  4. 插件兼容性:分析器、脚本扩展、存储插件和安全插件是否匹配。
  5. 配置兼容性:已废弃或删除的配置是否仍被接受。
  6. 业务语义兼容性:排序、分析器、查询、聚合或默认值是否发生变化。

不能因为“节点能启动”就认为升级成功。

7.2 升级前检查清单对应的原因

升级前应完成:

确认当前版本和目标版本
  ↓
阅读目标版本升级说明和重大变更
  ↓
处理废弃 API、配置和插件
  ↓
确认集群健康和分片分配
  ↓
确认 Snapshot 可用且完成恢复演练
  ↓
在测试环境执行完整升级
  ↓
安排生产升级窗口和监控

检查集群健康:

curl 'http://localhost:9200/_cluster/health?pretty'

查看节点版本:

curl 'http://localhost:9200/_cat/nodes?v&h=name,version,roles,heap.percent,ram.percent,disk.used_percent'

查看分片状态:

curl 'http://localhost:9200/_cat/shards?v'

升级前通常希望:

  • 没有长期未分配的主分片;
  • 磁盘没有触发高水位;
  • 集群没有持续的写入拒绝;
  • 没有未处理的 ILM 错误;
  • Snapshot 已经成功完成;
  • 客户端可以兼容目标版本;
  • 所有插件都有对应版本。

健康状态为 yellow 不一定阻止升级,但必须知道原因。如果 yellow 是副本长期无法分配,升级期间可能进一步放大恢复风险。

7.3 滚动升级和全量停机升级

滚动升级是逐个停止、升级、重启节点,确保集群在过程中仍有足够节点提供服务。

它要求:

  • 当前版本到目标版本存在官方支持的混合版本路径;
  • 先升级符合要求的节点角色;
  • 升级过程中遵守官方规定的顺序;
  • 升级前避免同时停止多个关键节点;
  • 确认每个节点重新加入集群并稳定后再继续。

不能自行假设任意两个版本都支持滚动升级,尤其不能把跨多个大版本直接混滚。具体路径必须以目标版本官方升级文档为准。

全量停机升级是停止整个集群,升级所有节点后再启动。它的可用性更差,但过程简单,适合:

  • 不支持滚动升级的版本跨度;
  • 单节点部署;
  • 需要重建节点或同时调整操作系统;
  • 已明确接受停机窗口的环境。

无论采用哪种方式,都应先验证快照。Elasticsearch 通常不支持将数据目录从新版本直接降级回旧版本运行,因此“回滚”不能简单理解为安装旧二进制并重新启动原数据目录。

7.4 升级中的分片和写入行为

节点重启会触发:

节点离开集群
  ↓
主分片或副本变为未分配
  ↓
集群重新分配
  ↓
副本恢复、translog 重放或 segment 恢复
  ↓
节点重新稳定

如果在副本还未恢复时继续停止其他节点,就可能失去可用副本,甚至影响主分片可用性。

升级时应观察:

curl 'http://localhost:9200/_cluster/health?level=indices&pretty'
curl 'http://localhost:9200/_cat/recovery?v'
curl 'http://localhost:9200/_cluster/allocation/explain?pretty' \
  -H 'Content-Type: application/json' \
  -d '{
    "index": "logs-2025.01.01",
    "shard": 0,
    "primary": false
  }'

_cluster/allocation/explain 用于解释某个分片为什么不能分配。常见原因包括:

  • 节点磁盘超过水位;
  • 节点角色或分配过滤条件不匹配;
  • 副本数量超过可用节点;
  • 集群设置禁止分配;
  • 节点版本或功能不满足恢复条件;
  • 分片所在节点尚未加入集群。

滚动升级过程中,如果业务允许,可以暂时降低写入速率,减少恢复压力。但不要用关闭分配、关闭副本或强制重路由来掩盖实际故障。

7.5 升级后的验证

节点全部升级后,需要验证的不只是版本号:

curl 'http://localhost:9200/_cat/nodes?v&h=name,version,roles'
curl 'http://localhost:9200/_cluster/health?wait_for_status=yellow&pretty'
curl 'http://localhost:9200/_ilm/status?pretty'
curl 'http://localhost:9200/_ingest/pipeline/logs-normalize?pretty'

还应执行端到端验证:

  1. 写入一条经过 Pipeline 的文档;
  2. 检查 Pipeline 生成的字段;
  3. 从写入别名或数据流查询;
  4. 验证 Bulk 的成功和失败 item;
  5. 验证 ILM 当前 step;
  6. 验证 Snapshot 仓库仍可访问;
  7. 检查客户端日志和废弃 API;
  8. 检查查询延迟、写入拒绝、GC 和磁盘使用率。

升级后不应立即删除旧快照。至少在业务验证完成、恢复路径确认之前,应保留可用的回退依据。


八、一个完整的写入链路示例

下面把索引模板、Pipeline、Bulk 和生命周期策略串起来。

8.1 创建 Pipeline

curl -X PUT 'http://localhost:9200/_ingest/pipeline/app-events' \
  -H 'Content-Type: application/json' \
  -d '{
    "processors": [
      {
        "set": {
          "field": "event.ingested",
          "value": true
        }
      },
      {
        "date": {
          "field": "@timestamp",
          "formats": ["ISO8601"],
          "ignore_failure": false
        }
      }
    ]
  }'

8.2 创建 ILM 策略

curl -X PUT 'http://localhost:9200/_ilm/policy/app-events-policy' \
  -H 'Content-Type: application/json' \
  -d '{
    "policy": {
      "phases": {
        "hot": {
          "actions": {
            "rollover": {
              "max_age": "1d",
              "max_primary_shard_size": "20gb"
            }
          }
        },
        "delete": {
          "min_age": "14d",
          "actions": {
            "delete": {}
          }
        }
      }
    }
  }'

8.3 创建索引模板

curl -X PUT 'http://localhost:9200/_index_template/app-events-template' \
  -H 'Content-Type: application/json' \
  -d '{
    "index_patterns": ["app-events-*"],
    "priority": 100,
    "template": {
      "settings": {
        "index.default_pipeline": "app-events",
        "index.lifecycle.name": "app-events-policy",
        "index.lifecycle.rollover_alias": "app-events-write",
        "number_of_shards": 1,
        "number_of_replicas": 1
      },
      "mappings": {
        "dynamic": "strict",
        "properties": {
          "@timestamp": {"type": "date"},
          "service": {"type": "keyword"},
          "message": {"type": "text"},
          "event": {
            "properties": {
              "ingested": {"type": "boolean"}
            }
          }
        }
      }
    }
  }'

8.4 初始化第一个索引和写入别名

curl -X PUT 'http://localhost:9200/app-events-000001' \
  -H 'Content-Type: application/json' \
  -d '{
    "aliases": {
      "app-events-write": {
        "is_write_index": true
      }
    }
  }'

之所以还需要显式创建第一个索引,是因为 rollover 需要一个初始写入目标。后续索引会由 ILM 根据模板和策略创建。

8.5 使用 Bulk 写入

curl -X POST 'http://localhost:9200/app-events-write/_bulk?refresh=wait_for' \
  -H 'Content-Type: application/x-ndjson' \
  --data-binary $'{"index":{"_id":"evt-1"}}\n{"@timestamp":"2025-01-01T12:00:00Z","service":"payment","message":"request failed"}\n{"index":{"_id":"evt-2"}}\n{"@timestamp":"2025-01-01T12:00:01Z","service":"order","message":"request completed"}\n'

处理过程是:

  1. 请求目标是 app-events-write 别名;
  2. 别名将请求路由到当前 is_write_index=true 的索引;
  3. 默认 Pipeline app-events 执行;
  4. event.ingested 被加入文档;
  5. Mapping 校验字段类型;
  6. 文档写入主分片并复制到副本;
  7. refresh=wait_for 等待可搜索;
  8. Bulk 响应逐项返回结果。

检查生命周期:

curl 'http://localhost:9200/app-events-000001/_ilm/explain?pretty'

检查最终 Mapping:

curl 'http://localhost:9200/app-events-000001/_mapping?pretty'

九、常见失败表现与诊断路径

9.1 Bulk 返回 200,但数据不完整

首先检查:

"errors": true

然后逐项查看:

  • status
  • error.type
  • error.reason
  • 文档 _id
  • 目标 _index

不要只看 HTTP 状态码,也不要只查看 Bulk 响应最后一项。

9.2 mapper_parsing_exception

常见原因:

第一次写入:
  http.status.code = 503
  → 动态映射为数字

后来写入:
  http.status.code = "unknown"
  → 无法解析为数字

解决方式通常是:

  • 在索引创建前显式定义 Mapping;
  • 在 Pipeline 中统一类型转换;
  • 对异常值使用 null 或转入失败索引;
  • 不能直接修改已有字段类型,需要重建索引。

9.3 Pipeline 没有执行

检查顺序:

  1. 请求是否显式指定了其他 Pipeline;
  2. 是否只配置了 default pipeline;
  3. 是否通过 pipeline=_none 绕过了 default pipeline;
  4. final pipeline 是否配置;
  5. Pipeline 是否存在;
  6. _simulate 是否能复现;
  7. Bulk 中是否每个 item 都指向了预期索引。

9.4 ILM 长时间停在某个 step

执行:

curl 'http://localhost:9200/index-name/_ilm/explain?pretty'

重点查看当前 action 和 step。若涉及分片分配,再检查:

curl 'http://localhost:9200/_cluster/health?pretty'
curl 'http://localhost:9200/_cat/allocation?v'

若涉及 rollover,应检查:

  • 写入别名是否存在;
  • 是否存在唯一的写索引;
  • 索引是否使用正确的 rollover_alias
  • rollover 条件是否真正满足。

9.5 Snapshot 成功但恢复失败

可能原因包括:

  • 目标集群版本不在支持范围内;
  • 目标集群缺少必要插件;
  • 同名索引已经存在;
  • 恢复了不适用的全局状态;
  • 仓库权限变化;
  • 节点磁盘不足;
  • 分片无法分配到目标节点。

恢复验证必须在与生产不同的环境中进行,否则不能证明灾难恢复路径有效。

9.6 升级后集群无法恢复到绿色

先确认是否只是副本未分配,再执行分配解释:

curl -X GET 'http://localhost:9200/_cluster/allocation/explain?pretty' \
  -H 'Content-Type: application/json' \
  -d '{
    "index": "app-events-000001",
    "shard": 0,
    "primary": false
  }'

如果主分片正常、所有数据可读,而副本因单节点部署无法分配,集群可能是 yellow;如果主分片未分配,则是更严重的问题。不要用增加副本或强制重分配替代对磁盘、节点角色和集群设置的检查。


十、如何组合这些机制

一个较完整的生产数据管道通常具有如下职责边界:

应用或消息消费者
  ├─ 负责业务事件 ID
  ├─ 控制 Bulk 批量和并发
  ├─ 逐 item 解析响应
  ├─ 处理可重试和不可重试错误
  └─ 维护死信或补偿队列

Elasticsearch Ingest
  ├─ 负责轻量字段清洗
  ├─ 负责统一时间和结构
  └─ 不承担跨文档事务

Index Template / Mapping
  ├─ 固定字段类型
  ├─ 约束动态字段
  └─ 控制分析器和索引结构

ILM
  ├─ 控制 rollover
  ├─ 控制索引阶段
  ├─ 控制副本和分配
  └─ 控制到期删除

Snapshot / SLM
  ├─ 负责灾难恢复依据
  ├─ 保存索引和选定元数据
  └─ 通过恢复演练验证可用性

升级流程
  ├─ 验证版本和插件兼容性
  ├─ 控制分片恢复风险
  ├─ 验证 API 和数据管道
  └─ 保留可恢复的旧快照

其中最容易被混淆的是三个“成功”:

  • Bulk 成功:某个 item 被接受或写入成功;
  • ILM 成功:索引生命周期状态机完成了某一步;
  • Snapshot 成功:指定数据已成功保存到仓库并可被识别。

它们分别发生在写入、索引管理和备份三个不同层面。只有将每一层的状态、错误和恢复路径分别监控,才能判断 Elasticsearch 集群是否真正处于可用状态。


系列导航与关联阅读

官方资料

本文依据数据库官方文档重新梳理;正文、示例与生产检查清单由 WR BLOG 编写。