Skip to content

日志流水线:Filebeat + Logstash + Elasticsearch

基于 ELK Stack 8.x · 核于 2026-08

速查

  • 完整流水线日志文件 → Filebeat(tail 采集)→ Logstash(grok/mutate 解析)→ Elasticsearch(存+倒排索引)→ Kibana(Discover/Dashboard)
  • Filebeat 核心概念:①input/prospector(采集源,日志文件路径);②harvester(每个文件一个,tail 跟踪新行);③spooler(批量打包 event);④registry(记录每文件读到哪,崩溃续传);⑤output(发到 Logstash/ES/Kafka)。
  • Logstash 管道三段:①input(beats/kafka/file/syslog);②filter(grok 正则/mutate 改字段/date 解析时间/json/geoip/ruby);③output(elasticsearch/kafka/file/stdout)。
  • grok:基于正则的「模式匹配」——预置模式(TIMESTAMP_ISO8601/LOGLEVEL/IP/NUMBER)+ 自定义,把非结构化日志行解析成结构化字段。
  • mutate:字段操作——add_field/rename/replace/remove_field/convert(类型转换)/uppercase/lowercase。
  • ES index 模型:按时间建索引(app-logs-2026.08.09),便于按天 ILM(Index Lifecycle Management)滚动删除。
  • mapping:字段类型——text(分词,全文检索)/ keyword(不分词,精确匹配与聚合)/ long/date/boolean。错配类型会导致检索失效(如把 IP 设成 text 无法精确匹配)。
  • Ingest Node:ES 内置的轻量解析 pipeline,可替代简单 Logstash(grok/date/convert),减少组件——小规模推荐。
  • ILM(Index Lifecycle Management):自动管理索引生命周期——hot(写入)/ warm(查询)/ cold(少查)/ delete(删除),控制存储成本。

一、Filebeat:轻量日志采集

Filebeat 是 Go 写的轻量 agent,部署在每台产生日志的机器:

yaml
# filebeat.yml
filebeat.inputs:
  - type: log
    enabled: true
    paths:
      - /var/log/myapp/*.log          # 采集路径
    fields:                            # 附加字段
      env: prod
      service: order
    fields_under_root: true
    multiline.pattern: '^\d{4}-\d{2}-\d{2}'  # 多行合并(异常栈)
    multiline.negate: true
    multiline.match: after

output.logstash:
  hosts: ["logstash:5044"]             # 发 Logstash
  # 或直接发 ES:output.elasticsearch.hosts: ["es:9200"]

registry.file: /var/lib/filebeat/registry  # 记录读取进度
  • harvester + spooler:每个文件起一个 harvester(goroutine)tail 新行,spooler 把事件批量发 output。
  • registry 续传:每读完一段记录 file + offset,崩溃重启从上次 offset 续读——不丢不重。
  • multiline:异常栈是多行的(Exception 行后跟多个 at 行),用 multiline 把多行合并成一个 event(以时间戳开头的行算新 event)。
  • Module:Filebeat 自带 Nginx/Apache/MySQL/Redis/Kafka 等日志的 Module,预置 grok + Kibana Dashboard,开箱即用。

二、Logstash:filter 管道详解

Logstash 的 filter 是日志加工的核心:

grok:正则解析非结构化日志

ruby
filter {
  grok {
    match => {
      # 把日志行解析成 timestamp/level/service/message 字段
      "message" => "%{TIMESTAMP_ISO8601:timestamp} %{LOGLEVEL:level} \[%{DATA:service}\] %{GREEDYDATA:msg}"
    }
    tag_on_failure => ["_grokparsefailure"]
  }
}
  • 预置模式TIMESTAMP_ISO8601LOGLEVELIPNUMBERWORDDATA(非贪婪)、GREEDYDATA(贪婪)。
  • 自定义模式(?<fieldname>pattern),复杂场景写自己的正则。
  • 代价:grok 是正则,CPU 密集——日志量大时 Logstash 资源消耗明显,要监控解析延迟。

mutate:字段操作

ruby
filter {
  mutate {
    convert => { "status" => "integer" }     # 类型转换(聚合需要)
    rename => { "clientip" => "client_ip" }  # 重命名
    remove_field => ["message"]              # 删原始字段省空间
    add_field => { "env" => "prod" }         # 加字段
  }
}

date:时间字段解析

ruby
filter {
  date {
    match => ["timestamp", "yyyy-MM-dd HH:mm:ss", "ISO8601"]
    target => "@timestamp"                   # 写入 @timestamp(ES 时间字段)
  }
}
  • ES 默认按 @timestamp 排序与按时间检索,所以要把日志原始时间解析到这个字段。

geoip:IP 转地理

ruby
filter {
  geoip { source => "client_ip" }            # 自动加 geoip.country_name/geoip.city_name
}
  • 用于在 Kibana 画地理分布地图(用户访问来源)。

三、Elasticsearch:存储与索引

index 按时间滚动

app-logs-2026.08.07    app-logs-2026.08.08    app-logs-2026.08.09
   (已删)                 (warm, 可查)            (hot, 写入中)
  • 按天建索引:Logstash output 用 index => "app-logs-%{+YYYY.MM.dd}",每天一个索引。
  • ILM:定义策略——hot(保留 N 天写入)→ warm(迁到慢盘)→ cold(再久)→ delete(自动删)。控制总存储成本。
  • index template:预定义每个 app-logs-* 索引的 mapping(字段类型)与 settings(分片数)。

mapping:字段类型决定检索行为

json
{
  "properties": {
    "@timestamp": { "type": "date" },
    "level":   { "type": "keyword" },          // 精确匹配 + 聚合
    "message": { "type": "text" },             // 分词全文检索
    "status":  { "type": "integer" },          // 数值聚合
    "client_ip": { "type": "ip" }              // IP 类型(范围查询)
  }
}
  • text vs keyword:text 会分词(适合全文检索 message),keyword 不分词(适合精确匹配 level + 聚合统计)。错配会让检索失效(把 level 设成 text 会让精确匹配失败)。
  • dynamic mapping:ES 默认自动推断字段类型,但容易猜错(把数字 ID 当 keyword)——生产建议显式 mapping。

聚合(aggregation)

json
// 按服务分组统计错误数
{
  "size": 0,
  "query": { "match": { "level": "error" } },
  "aggs": {
    "by_service": {
      "terms": { "field": "service", "size": 10 },
      "aggs": {
        "error_over_time": {
          "date_histogram": { "field": "@timestamp", "fixed_interval": "1h" }
        }
      }
    }
  }
}
  • 类似 SQL 的 GROUP BY service + 时间分桶——支持 terms/histogram/date_histogram/avg/sum/cardinality 等。

四、Kibana:检索与可视化

Discover:日志检索主界面

  • KQL 查询level:error and service:"order service" and @timestamp > "now-1h"
  • 字段过滤:点字段值快速过滤(如点 level 的 error 值加筛选条件)。
  • 时间轴下钻:时间分布柱状图,拖选时间段缩小范围。

Dashboard:多 Visualize 大盘

  • 类似 Grafana 的 Dashboard——多个 Visualize 聚合展示(错误数趋势、各服务占比、地理分布)。
  • 共享时间范围与过滤,支持联动。

Visualize 类型

类型用途
Lens拖拽式可视化(推荐)
Vertical Bar柱状图(趋势)
Pie饼图(占比)
Line折线图(时序)
Data Table表格
Metric单个数字
Map地理分布
TSVB时序聚合
Vega自定义高级可视化

五、Ingest Node:轻量解析替代 Logstash

ES 内置的 Ingest Node 可做简单解析,省去 Logstash:

json
// 定义 pipeline(grok + add field)
PUT _ingest/pipeline/app-logs
{
  "processors": [
    { "grok": {
        "field": "message",
        "patterns": ["%{TIMESTAMP_ISO8601:timestamp} %{LOGLEVEL:level} %{GREEDYDATA:msg}"]
    }},
    { "date": { "field": "timestamp", "formats": ["ISO8601"] } },
    { "remove": { "field": "message" } }
  ]
}

// Filebeat 直接发 ES 时指定 pipeline
// output.elasticsearch.pipeline: app-logs
  • 何时用 Ingest Node:日志结构简单、解析规则固定 → Ingest Node 省组件;解析复杂(多源、多变体)→ Logstash 更灵活。
  • 代价:Ingest Node 在 ES 集群内跑,消耗 ES 资源——日志量大时 Logstash 独立扩展更稳。

下一步

掌握了日志流水线后,下一步看与 Grafana LGTM 对比——ELK vs Loki 的全文检索 vs 省存储取舍、ELK vs Prometheus 的日志 vs 指标边界、选型决策矩阵。