Elasticsearch Ingest Pipeline 是什么?它和 Logstash 有什么区别?
简化版
Ingest Pipeline 是 ES 写入前的数据处理链,可以用 processors 做字段解析、重命名、转换、补充等操作。它适合轻量预处理;复杂采集、队列缓冲和多输出分发通常更适合 Logstash 或专门的数据管道。
详细版
数据进入 ES 前,经常要清洗格式、补字段、解析日志。Ingest Pipeline 运行在 ingest node 上,在文档写入索引前执行处理器。
- 常见 processor 有
set、rename、remove、grok、date、geoip、script。 - 可在 index request 或 index template 中指定默认 pipeline。
- 适合轻量字段转换、日志解析、统一补充元数据。
- 复杂逻辑会增加 ingest node CPU 压力。
- Logstash 更像外部数据处理系统,能力更重,也更适合缓冲和多输出。
完整版教学
一、为什么写入前要做数据清洗
ES 对 mapping 很敏感,脏数据一旦写进去,字段类型可能被错误推断,后续修复就要 reindex。日志场景里,原始 message 往往是一整行字符串,但查询需要结构化字段,例如 status、latency、ip、traceId。Ingest Pipeline 的目标就是在写入前把文档整理成适合索引和查询的形状。它把一部分清洗逻辑放到 ES 入口处,减少应用端重复处理。
PUT _ingest/pipeline/access-log-pipeline
{
"processors": [
{ "set": { "field": "service", "value": "gateway" } },
{ "date": { "field": "time", "formats": ["ISO8601"] } }
]
}
记忆钩子:Ingest Pipeline 是 ES 门口的“质检流水线”,文档过线后才真正入库。
二、processor 是怎么串起来的
Pipeline 由一组 processors 按顺序执行。前一个 processor 产出的字段可以被后一个 processor 使用,例如先 grok 解析 message,再 date 解析时间,再 remove 删除临时字段。顺序错误会导致字段不存在或解析失败。生产里要给失败路径配置 on_failure,否则一条脏日志可能导致写入失败或丢失关键错误信息。
{
"processors": [
{ "grok": { "field": "message", "patterns": ["%{IP:client_ip} %{NUMBER:status:int}"] } },
{ "remove": { "field": "message" } }
],
"on_failure": [
{ "set": { "field": "ingest_error", "value": "{{ _ingest.on_failure_message }}" } }
]
}
三、带数字看 ingest node 压力
假设每秒写入 20000 条日志,每条都要执行 grok 和 geoip。即使每条处理只花 0.2ms CPU,总 CPU 需求也达到 4 秒 CPU/秒,约等于持续占满 4 个 CPU 核。复杂脚本、正则和外部库处理会让成本更高。因此 Ingest Pipeline 的处理逻辑越重,越要关注 ingest node 的 CPU、队列和写入延迟。
20000 docs/s × 0.2ms = 4000ms CPU/s
约等于 4 个 CPU 核持续工作
四、它和 Logstash 的边界
Ingest Pipeline 运行在 ES 集群内部,部署简单,适合轻量处理。Logstash 是独立数据处理组件,有输入、过滤、输出插件生态,可以接 Kafka、文件、数据库,也能做缓冲、多输出、复杂过滤。区别可以理解为:Ingest Pipeline 更靠近 ES 写入入口,Logstash 更像外部 ETL 管道。如果清洗逻辑会明显拖累 ES,或者需要可靠缓冲,应该考虑外置管道。
| 维度 | Ingest Pipeline | Logstash |
|---|---|---|
| 部署位置 | ES ingest node | 独立组件 |
| 适合逻辑 | 轻量转换 | 复杂 ETL |
| 缓冲能力 | 较弱 | 更强 |
| 多输出 | 不擅长 | 擅长 |
五、默认 pipeline 和模板怎么配合
为了避免写入方漏传 pipeline,可以在索引设置或模板里配置默认 pipeline。这样匹配某类索引的新文档都会经过统一处理,尤其适合日志索引和数据流。模板、ILM、pipeline 通常一起出现:模板定义 mapping 和默认 pipeline,ILM 控制 rollover,pipeline 保证字段结构统一。这个组合能让数据接入更稳定。
PUT _index_template/logs-template
{
"index_patterns": ["logs-*"],
"template": {
"settings": { "index.default_pipeline": "access-log-pipeline" }
}
}
六、如何测试和排查 pipeline
上线前应该用 _simulate API 测试 pipeline,输入样例文档,观察输出字段是否符合预期。上线后如果写入失败,要查看 bulk 响应里的 item 错误、pipeline 的 on_failure 字段和 ingest node 指标。不要直接在生产大流量上试复杂 grok,正则写错可能把写入链路拖慢。稳定做法是先用小样本模拟,再灰度应用到新索引。
POST _ingest/pipeline/access-log-pipeline/_simulate
{
"docs": [
{ "_source": { "message": "127.0.0.1 200", "time": "2026-07-30T09:00:00Z" } }
]
}
七、常见误区与追问
- 误区:Ingest Pipeline 不会影响 ES 性能。 它运行在 ingest node 上,复杂处理会消耗 CPU 并拖慢写入。
- 误区:有 Pipeline 就不需要 Logstash。 Pipeline 适合轻量处理,复杂采集、缓冲和多输出仍可能需要 Logstash。
- 误区:解析失败可以忽略。 失败路径不处理会导致数据丢失、写入失败或字段缺失。
- 追问:如何让所有新日志都走 pipeline? 在索引模板中设置
index.default_pipeline。 - 追问:上线前怎么验证? 使用
_simulateAPI 输入样例文档,检查输出和失败路径。
八、加强记忆
Ingest Pipeline 记成“ES 写入口的轻量清洗线”。它能做 set、rename、grok、date、geoip 等处理,但处理越复杂越吃 ingest node。面试回答时把 processor 顺序、on_failure、默认 pipeline、和 Logstash 的边界讲清,就能体现真正的工程判断。