Elasticsearch 数据迁移方案对比:esm vs Logstash
Elasticsearch 集群在扩容、机房迁移、冷热分离等场景下,通常需要将数据从一个集群搬迁到另一个。同集群内的索引重建可以用 reindex,备份恢复可以用 snapshot,但这两者在跨集群或跨机房场景下存在明显限制:reindex 要求网络可达且不能跨大版本,snapshot 要求共享存储或版本严格一致。真正的跨集群数据迁移,通常需要在 esm(轻量同步工具)和 Logstash(ETL 管道)之间做选型。
版本说明
本文原写于 2022-03,2026-09 重写。
esm 工具的基本用法未变;新集群若已升级到 8.x 并开启安全认证,迁移命令需额外补充认证参数(-x/-y 的 URL 中携带用户名密码或改用 --auth),迁移逻辑本身不受影响。
# 1. 四种迁移方案对比
| 工具 | 适用场景 | 限制 |
|---|---|---|
| esm | 索引级迁移、增量同步(按时间戳) | 不处理 mapping 变更、无复杂转换 |
| Logstash | ETL 场景、需字段转换、多目标输出 | 需配置 input/filter/output、较重 |
| reindex | 同集群或远程集群索引重建 | 需集群间网络可达、无法跨版本 |
| snapshot | 集群级冷备/跨机房恢复 | 需共享存储(S3/NFS/HDFS)或手动拷贝、版本必须一致 |
选型速判:
- 单纯搬迁索引数据、不需要字段转换 → esm
- 需要字段清洗/转换、输出到多个目标(如 Kafka + ES)→ Logstash
- 同集群内重建索引、mapping 要调整 → reindex API
- 整集群备份、灾难恢复 → snapshot
# 2. 迁移前的准备
无论选择哪种工具,以下准备工作必须完成:
# 2.1 目标端先建好索引与 mapping
esm 和 Logstash 都只迁移文档数据,不迁移索引的 mapping 配置。如果目标索引不存在,ES 会根据第一条文档自动推断 mapping,极易与源端不一致。迁移前务必在目标端手动创建索引,并确保 mapping 与源端兼容:
# 从源端获取 mapping
GET source_index/_mapping
# 在目标端创建同样结构的索引(示例)
PUT target_index
{
"mappings": {
"properties": {
"field1": { "type": "text" },
"field2": { "type": "date" }
}
}
}
2
3
4
5
6
7
8
9
10
11
12
13
# 2.2 写入期优化:关闭副本与刷新
迁移期间大规模写入数据时,副本分片和频繁刷新是主要性能瓶颈。在目标索引上临时调整以下设置,可大幅提升写入吞吐:
# 迁移前:禁副本、关刷新
PUT target_index/_settings
{
"index": {
"number_of_replicas": 0,
"refresh_interval": "-1"
}
}
2
3
4
5
6
7
8
number_of_replicas: 0:副本在写入期是纯开销,迁完再加回来refresh_interval: -1:禁用自动刷新,避免生成过多小 segment
# 2.3 迁完后恢复设置
数据迁移完成后,恢复正常的副本和刷新策略,并等待分片分配完成:
# 恢复设置
PUT target_index/_settings
{
"index": {
"number_of_replicas": 1,
"refresh_interval": "1s"
}
}
# 确认集群状态转绿(副本分配完成)
GET _cluster/health/target_index
2
3
4
5
6
7
8
9
10
11
必须等到 _cluster/health 返回 "status": "green" 才算迁移完成。如果停留在 yellow,说明副本尚未分配完成,此时若发生节点故障可能导致数据丢失。
# 2.4 跨大版本迁移的 mapping 兼容性
跨大版本(如 6.x → 8.x)时,以下 mapping 特性已被移除或变更,必须在目标端手动调整:
| 旧版本特性 | 新版本状态 | 处理建议 |
|---|---|---|
_type 字段 | 8.x 已移除 | 目标 mapping 中不再声明 _type,文档直接写入 _doc |
string 类型 | 5.x 起拆分为 text/keyword | 源端 string 需映射为目标端的 text(全文检索)或 keyword(精确匹配) |
index_options: offsets | 部分版本不支持 | 如需高亮,改用 index_options: positions |
如果不处理这些不兼容项,批量写入时会整批失败,错误信息类似 mapper_parsing_exception。
# 3. 方案一:esm 轻量同步
esm(Elasticsearch Migration Tool)基于 scroll 与 bulk 实现,适合索引级别的全量或增量迁移。
# 3.1 安装获取
# 从 GitHub Release 下载(以 v0.x 为例)
wget https://github.com/medcl/esm/releases/download/v0.7.0/esm-0.7.0-linux-amd64.tar.gz
tar -xzf esm-0.7.0-linux-amd64.tar.gz
mv esm-0.7.0-linux-amd64 /usr/local/bin/esm
# 验证
esm --help
2
3
4
5
6
7
# 3.2 增量同步脚本
逻辑:每3分钟同步一次 10 分钟之前的数据(留 7 分钟重叠窗口防丢数据)
vim /usr/local/bin/esm.sh
#!/bin/bash
while [ true ]; do
now=$(date +%s000)
ten_min_ago=$(($now - 600000))
echo $now $ten_min_ago
/usr/bin/esm -s http://192.0.2.11:9200 -x "myindex" -d http://192.0.2.12:9200 -w 10 -t 5m -b 10 -q "updatetime:[$ten_min_ago TO $now]"
sleep 180
done
2
3
4
5
6
7
8
# 3.3 参数解释
| 参数 | 说明 |
|---|---|
-s | 源 ES 地址 |
-d | 目标 ES 地址 |
-x | 源索引名(支持通配符) |
-y | 目标索引名(省略则与源同名) |
-w | worker 并发数 |
-t | scroll 超时时间 |
-b | bulk 批量大小 |
-q | 查询条件(用于增量同步的时间范围过滤) |
# 3.4 Systemd 服务配置
vim /usr/lib/systemd/system/esm.service
[Unit]
Description=esm
After=network.target
[Service]
Type=simple
PIDFile=/var/run/esm.pid
WorkingDirectory=/usr/local/
ExecStart=/usr/local/bin/esm.sh
ExecReload=/bin/kill -s HUP $MAINPID
ExecStop=/bin/kill -s QUIT $MAINPID
PrivateTmp=true
Restart=on-failure
RestartSec=5
LimitNOFile=65536
[Install]
WantedBy=multi-user.target
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
# 3.5 启动与验证
systemctl daemon-reload
systemctl start esm.service
systemctl status esm.service -l
# 验证任务进度(查看 esm 日志)
tail -f /var/log/esm.log
# 验证目标索引文档数
GET target_index/_count
2
3
4
5
6
7
8
9
# 4. 方案二:Logstash ETL 迁移
Logstash 适合需要字段转换/清洗、多目标输出的场景。以下三个配置样例覆盖常见迁移需求。
# 4.1 基础 ES→ES 迁移配置
input {
elasticsearch {
hosts => ["192.0.2.11:9200"]
index => "my_test_*"
user => "elastic"
password => "xxxxx"
scroll => "10m"
size => 2000
docinfo => true
docinfo_target => "[@metadata][doc]"
}
}
output {
elasticsearch {
hosts => ["192.0.2.12:9200"]
user => "elastic"
password => "xxxxx"
index => "%{[@metadata][doc][_index]}"
document_id => "%{[@metadata][doc][_id]}"
}
}
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
为什么 docinfo + document_id 是必需的:
docinfo => true让 Logstash 读取每条文档的元数据(_index,_id,_type等)docinfo_target => "[@metadata][doc]"将元数据存入事件的 metadata 字段,避免污染原始文档内容document_id => "%{[@metadata][doc][_id]}"在输出时使用原文档的_id,确保幂等写入:同一文档多次迁移不会产生重复数据,ES 会根据_id做更新而非插入
如果不设置 document_id,ES 会自动生成新的 _id,迁移变成纯追加——同一批数据重跑一次,目标端就多一份,无法幂等。
# 4.2 迁移到 Kafka 的配置(多目标输出场景)
input {
elasticsearch {
hosts => ["http://192.0.2.11:19200"]
index => "user_login_log_2024_*"
user => "${ES_USER}"
password => "${ES_PASSWORD}"
query => '{
"query": {
"bool": {
"must": [
{
"range": {
"@timestamp": {
"gte": "2024-01-01T00:00:00.000Z"
}
}
}
]
}
}
}'
size => 1000
scroll => "5m"
}
}
filter {
mutate {
remove_field => [ "@version", "@timestamp" ]
}
}
output {
kafka {
codec => json
bootstrap_servers => "192.0.2.13:9092,192.0.2.14:9092,192.0.2.15:9092,192.0.2.16:9092,192.0.2.17:9092"
topic_id => "user_login_log_topic"
}
}
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
# 4.3 动态索引名迁移(Ruby filter)
以下配置通过动态生成索引名,将数据从 mytestindex_2024_11 迁移到按 tag_id 分片的多个目标索引(如 mytestindex_5282_2024_11):
input {
elasticsearch {
user => "elastic"
password => "xxxx"
hosts => ["192.0.2.18:9200"]
index => "mytestindex_2024_11"
query => '{
"query": {
"bool": {
"must": [
{ "range": {
"create_time": {
"gte": "2024-11-27T00:00:00+08:00",
"lte": "2024-11-28T00:00:00+08:00"
}
}}
]
}
}
}'
docinfo => true
docinfo_target => "[@metadata][doc]"
}
}
filter {
ruby {
code => "
tag_id = event.get('tag_id')
if tag_id
event.set('dynamic_index', 'mytestindex_' + tag_id.to_s + '_2024_11')
else
event.tag('missing_tag_id')
end
"
}
}
output {
elasticsearch {
hosts => ["http://192.0.2.18:9200"]
user => "elastic"
password => "xxxx"
index => "%{dynamic_index}"
document_id => "%{[@metadata][doc][_id]}"
}
}
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
配置要点:
rubyfilter 根据字段tag_id的值动态构建目标索引名,并存入dynamic_index字段- 如果
tag_id缺失,添加missing_tag_id标签,便于后续排查 - 输出时通过
%{dynamic_index}引用动态生成的索引名
# 5. 怎么选:判据清单
按以下判据逐一检查,一般能在 30 秒内做出选型:
- 是否需要字段转换/清洗:是 → Logstash;否 → 继续
- 是否需要持续增量同步:是 → esm(支持定时脚本);Logstash 用 input 的 schedule 选项做周期重跑,但它每轮是重新执行查询,靠 document_id 幂等覆盖去重,不是真正的增量位点
- 数据量级:esm 无 filter 环节,链路更短;Logstash 每条文档都要过 filter/codec,管道越复杂开销越大——具体差多少必须在目标环境实测
- 是否跨大版本(如 6.x → 8.x):任意工具都需先在目标端手动创建兼容 mapping,不能直接迁移
- 是否需要多目标输出(ES + Kafka):Logstash 原生支持多 output,esm 只能单目标
- 运维资源:esm 可托管为 systemd 服务,Logstash 需维护 JVM 和配置文件
# 6. 坑与边界
# 6.1 esm 侧
- 时间窗口重叠:
ten_min_ago=$(($now - 600000))配合sleep 180(3分钟),实际留 7 分钟重叠窗口;若源端数据在 7 分钟内更新,可能重复同步,目标端需按_id去重或幂等写入 - mapping 变更:esm 只迁移数据,不处理 mapping;若源目标集群 mapping 不兼容,需先在目标集群创建好正确 mapping
- 认证信息:8.x 集群建议用
--auth=user:pass而非 URL 中明文携带
# 6.2 Logstash 侧
- mapping 不一致导致写入失败:Logstash 输出时会按目标索引的现有 mapping 写入,若目标 mapping 与源数据字段类型不兼容(如 string 写入 integer 字段),batch 会整批失败。迁移前务必在目标端创建与源端兼容的 mapping。
- size/scroll 调大后 JVM 压力:
size(每批文档数)和scroll(保持查询上下文的时间)增大可提升吞吐,但会线性增加 Logstash 的堆内存占用。若出现 OOM,优先降低size而非延长scroll。 - @timestamp 被 remove 后的影响:若 filter 中移除了
@timestamp字段,且 output 到 ES 的索引名包含日期通配符(如logstash-%{+YYYY.MM.dd}),索引名将无法解析为日期,导致写入失败或进入异常索引。
# 7. 验证步骤
数据迁移完成后,按以下步骤验证完整性:
步骤 1:文档数对比
# 源索引文档数
GET source_index/_count
# 目标索引文档数
GET target_index/_count
2
3
4
5
两数应大致相等(Logstash 如有 filter 丢弃数据,目标可能更少)。
步骤 2:随机抽样校验
# 随机获取 100 条文档的 _id
GET source_index/_search
{
"size": 100,
"query": {
"function_score": {
"random_score": {}
}
},
"_source": false
}
# 抽查目标端是否存在
GET target_index/_doc/<random_id>
2
3
4
5
6
7
8
9
10
11
12
13
14
抽查 10-20 条,确认关键字段内容一致。若使用 Logstash 且有 filter 转换,需对比转换后的字段值。
迁移任务快速参考
{
"_meta": {
"doc_version": "1.0",
"article_id": "es-migration-esm-logstash",
"profile_context": "elasticsearch-migration",
"last_updated": "2026-09"
},
"quick_start": {
"esm": {
"download": "wget https://github.com/medcl/esm/releases/download/v0.7.0/esm-0.7.0-linux-amd64.tar.gz",
"full_sync": "esm -s <source> -x <index> -d <dest> -w 10",
"incremental": "esm -s <source> -x <index> -d <dest> -q \"timestamp:[start TO end]\"",
"daemon": "systemctl start esm.service"
},
"logstash": {
"config_path": "/etc/logstash/conf.d/migration.conf",
"validate": "bin/logstash -f migration.conf --config.test_and_exit",
"run": "bin/logstash -f migration.conf",
"docinfo_required": true,
"document_id": "%{[@metadata][doc][_id]}"
}
},
"safety_rules": [
{"risk": "自动推断 mapping", "action": "目标端先 PUT mapping 再迁移", "tools": "both"},
{"risk": "重复数据", "action": "确保目标索引 _id 与源端一致(幂等写入)", "tools": "both"},
{"risk": "写入期性能瓶颈", "action": "迁前设置 number_of_replicas=0, refresh_interval=-1", "tools": "both"},
{"risk": "副本未分配完成", "action": "迁后等 _cluster/health status=green", "tools": "both"},
{"risk": "跨版本不兼容", "action": "目标端移除 _type, string→text/keyword", "tools": "both"},
{"risk": "esm 时间窗口重叠", "action": "留 7 分钟重叠,依赖 _id 去重", "tools": "esm"},
{"risk": "Logstash mapping 冲突导致 batch 失败", "action": "目标端预创建兼容 mapping", "tools": "logstash"},
{"risk": "Logstash JVM OOM", "action": "优先降低 size 参数,而非延长 scroll", "tools": "logstash"},
{"risk": "@timestamp 移除导致日期索引失效", "action": "移除前确认索引名不依赖 date filter", "tools": "logstash"}
],
"verification": {
"doc_count": "GET <index>/_count",
"random_sample": "GET <index>/_search {\"size\": 100, \"query\": {\"function_score\": {\"random_score\": {}}}, \"_source\": false}",
"check_replica": "GET _cluster/health/<index>",
"sample_id": "GET <index>/_doc/<id>"
}
}
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40