← 返回题目列表

Elasticsearch Ingest Pipeline 是什么?它和 Logstash 有什么区别?

中等 第 24 / 30 题 更新于 2026/07/30
ElasticsearchIngest PipelineLogstash数据清洗

简化版

Ingest Pipeline 是 ES 写入前的数据处理链,可以用 processors 做字段解析、重命名、转换、补充等操作。它适合轻量预处理;复杂采集、队列缓冲和多输出分发通常更适合 Logstash 或专门的数据管道。

详细版

数据进入 ES 前,经常要清洗格式、补字段、解析日志。Ingest Pipeline 运行在 ingest node 上,在文档写入索引前执行处理器。

  • 常见 processor 有 setrenameremovegrokdategeoipscript
  • 可在 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 PipelineLogstash
部署位置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
  • 追问:上线前怎么验证? 使用 _simulate API 输入样例文档,检查输出和失败路径。

八、加强记忆

Ingest Pipeline 记成“ES 写入口的轻量清洗线”。它能做 set、rename、grok、date、geoip 等处理,但处理越复杂越吃 ingest node。面试回答时把 processor 顺序、on_failure、默认 pipeline、和 Logstash 的边界讲清,就能体现真正的工程判断。