ELK 组件部署(Elasticsearch / Logstash / Kibana / Filebeat)

这篇把 ELK 四个组件的单机/集群独立部署过程整理在一起:怎么装、配置文件里哪几项必须改、起不来时先看什么,以及 Logstash 各类插件的用法。四个组件装完之后就能拼出各种日志链路。

如果要看的是 Kubernetes 场景下的日志方案(sidecar 采集、基于 Helm 的 EFK、以及 Filebeat + Kafka + Logstash + ES + Kibana 的 EFLK 链路),那部分内容在 Kubernetes 日志收集 里,本文不重复。

总览与数据流向

先交代环境和四个组件各自的位置,后面每一节的配置都基于这套主机名。

环境信息

复用已有的 hadoop 完全分布式集群,三个节点:

1
2
3
192.168.2.241 hadoop01
192.168.2.242 hadoop02
192.168.2.243 hadoop03

组件版本统一用 8.2.0,安装目录统一放在 /opt/bigdata/<组件名>/,并用软链接 current 指向具体版本目录,后续升级只需要换软链接。

各组件的分工

  • Elasticsearch:实时的分布式搜索和分析引擎,用于全文搜索、结构化搜索及分析,是整条链路的存储与检索层。三个节点都装,组成集群
  • Kibana:Elasticsearch 的可视化前端,负责检索、图表和索引管理,只需要装一台(这里放在 hadoop03)
  • Logstash:负责接收数据、解析过滤转换、再输出数据,是链路中的重量级处理环节
  • Filebeat:轻量级日志采集器,装在产生日志的机器上,只负责把日志读出来发走,资源占用远低于 Logstash

常见的组合方式有两种:日志量不大时 Filebeat → Logstash → Elasticsearch → Kibana;日志量大或者需要削峰时,中间加一层 Kafka,Filebeat → Kafka → Logstash → Elasticsearch → Kibana。本文按 Elasticsearch、Kibana、Logstash、Filebeat 的顺序逐个装,最后做一次最小链路的联调。

Elasticsearch

三个节点都要装。官网下载地址 https://www.elastic.co/downloads/elasticsearch

解压与用户准备

Elasticsearch 默认不允许用 root 启动,先建一个专用用户:

1
2
3
4
5
6
7
8
9
useradd elasticsearch
wget https://artifacts.elastic.co/downloads/elasticsearch/elasticsearch-8.2.0-linux-x86_64.tar.gz

mkdir -p /opt/bigdata/elasticsearch
tar -zxf elasticsearch-8.2.0-linux-x86_64.tar.gz -C /opt/bigdata/elasticsearch
cd /opt/bigdata/elasticsearch/
ln -s elasticsearch-8.2.0 current

chown -R elasticsearch:elasticsearch /opt/bigdata/elasticsearch

配置文件

/opt/bigdata/elasticsearch/current/config/elasticsearch.yml,其中 node.name 每个节点不同,discovery.seed_hostscluster.initial_master_nodes 三个节点写一样:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
cluster.name: my-application
node.name: hadoop01 # 按需修改
path.data: /data1/elasticsearch,/data2/elasticsearch
path.logs: /opt/bigdata/elasticsearch/current/logs
bootstrap.memory_lock: true
network.host: 0.0.0.0
http.port: 9200
discovery.seed_hosts: ["hadoop01", "hadoop02", "hadoop03"]
cluster.initial_master_nodes: ["hadoop01", "hadoop02", "hadoop03"]
# 首次启动必须把全部 master-eligible 节点都列上:
# 只列 2 台时初始投票集就是 2,bootstrap 期 quorum = 2/2,任一台没起来就选不出 master,容错为 0。
# 这个参数只在集群第一次启动时生效,集群成型后应该删掉。
xpack.security.enabled: false
xpack.security.transport.ssl.enabled: false

⚠️ 注:8.x 默认开启 xpack 安全认证,这里为了实验方便直接关掉了,仅适用于隔离的内网环境。生产集群应保留认证与传输加密,Kibana、Logstash、Filebeat 侧相应配置账号密码或 API Key。

系统调优

这一步不做,Elasticsearch 很可能直接启动失败。

/etc/sysctl.conf

1
2
fs.file-max=655360
vm.max_map_count = 262144

写完文件记得让它生效,否则内核参数还是默认值,bootstrap check 照样会拦下启动:

1
2
sysctl -p            # 或者临时生效:sysctl -w vm.max_map_count=262144
sysctl vm.max_map_count # 确认一下
  1. fs.file-max 是系统最大打开文件描述符数,建议 655360 或更高
  2. vm.max_map_count 限制单个进程能持有的内存映射区(VMA)数量,跟线程没有关系。Elasticsearch 要抬高它是因为用 mmapfs 映射 Lucene 段文件,要求至少 262144

/etc/security/limits.conf 添加如下内容。其中 memlock unlimited 对应配置文件里的 bootstrap.memory_lock: true,不放开这一项,内存锁定会失败并导致启动中断:

1
2
3
4
5
6
* soft nproc 20480
* hard nproc 20480
* soft nofile 65536
* hard nofile 65536
* soft memlock unlimited
* hard memlock unlimited

/etc/security/limits.d/20-nproc.conf 会覆盖上面的 nproc 设置,一并改掉,或者直接把这个文件删掉:

1
* soft nproc 20480

堆和 page cache 怎么分

ES 的内存分配有一条和其他 Java 服务不太一样的原则:不要把内存都给堆。它底层是 Lucene,段文件靠 mmap 映射后由操作系统的 page cache 承载,检索性能很大程度上取决于有多少段文件能常驻内存。堆开得太大,留给 page cache 的就少了,反而更慢。

所以常规做法是一半给堆、一半留给操作系统,并且堆有一个硬上限:

  • 堆必须小于 32GB 左右的压缩指针(compressed oops)阈值。堆一旦到达这个阈值,JVM 就关掉对象指针压缩,对象头变大,实际能装的对象反而比 31GB 还少。
  • 官方的保守建议是不超过 26GB;确认开启了 zero-based oops(启动日志里能看到 heap address: ..., zero based Compressed Oops)才可以到 30~31GB。
  • -Xms-Xmx 必须设成相同值,避免运行期扩堆造成停顿,也避免和下面的内存锁定打架。

配套还有一项文中容易漏掉的:bootstrap.memory_lock: true 需要系统侧同时放开 memlock 限制,否则 ES 启动时锁定失败会直接中断。

1
2
3
4
5
6
7
8
9
10
11
# /etc/security/limits.conf
* soft memlock unlimited
* hard memlock unlimited

# 用 systemd 启动的话 limits.conf 不生效,要改 unit
# /etc/systemd/system/elasticsearch.service.d/override.conf
[Service]
LimitMEMLOCK=infinity

# 启动后确认
curl "localhost:9200/_nodes?filter_path=**.mlockall&pretty" # 应为 true

锁定内存的目的是防止堆被换出到 swap——堆页一旦进 swap,GC 扫描时要把它们逐页换回来,一次 Full GC 能卡到分钟级。这也是为什么前面要求关掉 swap。

堆内存在 /opt/bigdata/elasticsearch/current/config/jvm.options 里设置,一般取可用内存的一半,且必须小于 32G(堆一旦到 32G 附近就会关掉压缩指针,-Xmx32g 正好落在失效那一侧,此时对象头变大、实际能装的对象比 31G 还少)。官方保守值是不超过 26G,确认开启了 zero-based oops 才可以到 30~31G。剩下的内存会被 Lucene 用作文件系统缓存 —— Lucene 是一个开源全文检索工具包,Elasticsearch 底层就是基于它实现的:

1
2
3
-Xms4g

-Xmx4g

启动与集群验证

用 elasticsearch 用户启动,-d 表示后台运行,起不来就看 path.logs 下的日志:

1
2
cd /opt/bigdata/elasticsearch/current
bin/elasticsearch -d

访问任意节点的 9200 端口,能返回版本信息就算成功:

1
curl hadoop01:9200
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
{
"name" : "hadoop01",
"cluster_name" : "my-application",
"cluster_uuid" : "nbaivjfyTOmZD9uC02G7mw",
"version" : {
"number" : "8.2.0",
"build_flavor" : "default",
"build_type" : "tar",
"build_hash" : "b174af62e8dd9f4ac4d25875e9381ffe2b9282c5",
"build_date" : "2022-04-20T10:35:10.180408517Z",
"build_snapshot" : false,
"lucene_version" : "9.1.0",
"minimum_wire_compatibility_version" : "7.17.0",
"minimum_index_compatibility_version" : "7.0.0"
},
"tagline" : "You Know, for Search"
}

Kibana

Kibana 是 Elasticsearch 的检索与可视化界面,本身不存数据,只要能连上 ES 就行,所以只在 hadoop03 上装一台。官网下载地址 https://www.elastic.co/downloads/kibana

安装与配置

1
2
3
4
5
6
wget https://artifacts.elastic.co/downloads/kibana/kibana-8.2.0-linux-x86_64.tar.gz

mkdir -p /opt/bigdata/kibana
tar -zxf kibana-8.2.0-linux-x86_64.tar.gz -C /opt/bigdata/kibana
cd /opt/bigdata/kibana/
ln -s kibana-8.2.0 current

修改 /opt/bigdata/kibana/current/config/kibana.yml

1
2
3
server.port: 5601
server.host: "hadoop03"
elasticsearch.url: "hadoop01:9200"

⚠️ 注:elasticsearch.url 是 6.x 的写法,7.x 起改成了列表形式并且必须带协议头,8.2 上应写作 elasticsearch.hosts: ["http://hadoop01:9200"],同时建议把三个节点都列上以便故障时自动切换。

浏览器直连 ES 的场景下(例如一些自建面板),需要在 hadoop01 的 /opt/bigdata/elasticsearch/current/config/elasticsearch.yml 里放开跨域:

1
2
http.cors.enabled: true
http.cors.allow-origin: "*"

启动

1
2
cd /opt/bigdata/kibana/current/
nohup bin/kibana --allow-root &

启动后访问 http://hadoop03:5601 ,页面里的 Discover 用来查日志,Stack Management 里管理索引。

Logstash

Logstash 实现的功能分为接收数据、解析过滤并转换数据、输出数据三部分,分别对应三类插件:

  1. input 插件,必选
  2. filter 插件,可选
  3. output 插件,必选

配置文件就是这三段的组合:

1
2
3
4
5
6
7
8
9
input {
输入插件
}
filter {
过滤匹配插件
}
output {
输出插件
}

更准确地说,Logstash 不只是一个 input → filter → output 的数据流,而是一个 input → decode → filter → encode → output 的数据流,中间的编解码由 codec 插件完成。

安装与最小示例

官网下载地址 https://www.elastic.co/downloads/logstash 。所有节点用 root 用户安装:

1
2
3
4
5
6
wget https://artifacts.elastic.co/downloads/logstash/logstash-8.2.0-linux-x86_64.tar.gz

mkdir -p /opt/bigdata/logstash
tar -zxf logstash-8.2.0-linux-x86_64.tar.gz -C /opt/bigdata/logstash
cd /opt/bigdata/logstash/
ln -s logstash-8.2.0 current

先用命令行参数跑一个读标准输入、打印到标准输出的最小例子:

1
2
cd /opt/bigdata/logstash/current
bin/logstash -e 'input{stdin{}} output{stdout{codec=>rubydebug}}'

输入 123 # 输入 后得到:

1
2
3
4
5
6
7
8
9
10
11
{
"message" => "123 # 输入",
"host" => {
"hostname" => "hadoop01"
},
"event" => {
"original" => "123 # 输入"
},
"@timestamp" => 2022-05-22T11:14:26.267984Z,
"@version" => "1"
}

三个要点:-e 表示直接执行后面的配置;input 选了 stdinoutput 选了 stdout,其中 codec 是插件,用来指定输出格式,rubydebug 是专门用来做测试的格式,在终端输出的是 Ruby Hash 形式("key" => value)而不是 JSON——下面的示例本身就不是合法 JSON。真要输出 JSON 用 codec => json

同样的内容写成配置文件 /opt/bigdata/logstash/current/logstash-simple.conf

1
2
3
4
5
6
input {
stdin { }
}
output {
stdout { codec => rubydebug }
}

-f 指定配置文件启动,输出与上面一致:

1
bin/logstash -f logstash-simple.conf

input 插件

从文件读取数据,start_position => "beginning" 表示从文件开头读:

1
2
3
4
5
6
7
8
9
10
11
12
input {
file {
path => ["/var/log/secure"]
type => "system"
start_position => "beginning"
}
}
output {
stdout{
codec=>rubydebug
}
}

从标准输入读取,同时给事件加字段和标签:

1
2
3
4
5
6
7
8
9
10
11
12
input{
stdin{
add_field=>{"key"=>"ok"}
tags=>["add field"]
type=>"mytype"
}
}
output {
stdout{
codec=>rubydebug
}
}

输入 hello world,可以看到 typetagskey 都进到事件里了:

1
2
3
4
5
6
7
8
9
10
11
12
13
{
"host" =>{
"hostname" => "hadoop02"
},
"message" => "hello world",
"type" => "mytype",
"tags" => [
[0] "add field"
],
"key" => "ok",
"@timestamp" => 2022-05-22T11:25:56.419Z,
"@version" => "1"
}

从网络读取 TCP 数据,顺便用 grok 解析 syslog 行,配置文件 /opt/bigdata/logstash/current/logstash-tcp.conf

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
input {
tcp {
port => "5044"
}
}
filter {
grok {
match => { "message" => "%{SYSLOGLINE}" }
}
}
output {
stdout{
codec=>rubydebug
}
}

启动后在另一个终端用 nc 把日志灌进去:

1
2
3
4
bin/logstash -f logstash-tcp.conf

# 另一个终端
nc 192.168.2.241 5044 < /var/log/secure

结果节选,可以看到 timestampprocess 等字段已经被 SYSLOGLINE 拆出来了:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
{
"event" => {
"original" => "May 22 07:49:28 hadoop01 sshd[3106]: Disconnected from 192.168.2.243 port 41710"
},
"timestamp" => "May 22 07:49:28",
"process" => {
"pid" => 3106,
"name" => "sshd"
},
"@version" => "1",
"host" => {
"hostname" => "hadoop01"
},
"message" => [
[0] "May 22 07:49:28 hadoop01 sshd[3106]: Disconnected from 192.168.2.243 port 41710",
[1] "Disconnected from 192.168.2.243 port 41710"
],
"@timestamp" => 2022-05-23T01:55:22.020103Z
}

编码插件(Codec)

编码插件用于在输入或输出时处理不同类型的数据,前面用到的 rubydebug 就是其中一个。常见的格式有 plainjsonjson_lines 等。

plain 直接输出原始文本:

1
2
3
4
5
6
7
8
9
input{
stdin {
}
}
output{
stdout {
codec => "plain"
}
}
1
2
hello world # 输入
2022-05-23T02:03:30.140800Z {hostname=hadoop01} hello world # 输入

json 输出压缩成一行的 JSON:

1
2
3
4
5
6
7
8
9
input {
stdin {
}
}
output {
stdout {
codec => json
}
}
1
2
hello world # 输入
{"message":"hello world # 输入","@version":"1","host":{"hostname":"hadoop02"},"event":{"original":"hello world # 输入"},"@timestamp":"2022-05-23T02:04:49.242024Z"}

过滤器插件(Filter)

丰富的过滤器插件是 Logstash 功能强大的重要因素。名字叫过滤器,实际提供的不只是过滤能力,还可以对进入的原始数据做复杂的逻辑处理,甚至往后续流程里添加新的事件。

以这样一行访问日志为例:

1
192.168.2.241 [16/Jun/2021:16:24:19 +0800] "GET / HTTP/1.1" 403 5039

最常用的是三个插件配合:grok 做正则捕获拆字段、date 做时间处理、mutate 做字段改名与类型转换:

1
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
input {
stdin {}
}
filter {
grok {
match => { "message" => "%{IP:clientip}\ \[%{HTTPDATE:timestamp}\]\ %{QS:referrer}\ %{NUMBER:response}\ %{NUMBER:bytes}" } # %{语法: 语义},以以上格式收集数据
remove_field => [ "message", "event" ] # 删除掉 message 字段
}
date {
match => ["timestamp", "dd/MMM/yyyy:HH:mm:ss Z"] # 收集然后转存到 @timestamp 字段里
}
mutate {
rename => { "response" => "response_new" } # 重命名字段
# 同一个 mutate 内部的执行顺序是插件硬编码的(coerce → rename → update → replace → convert → gsub …),
# 不按书写顺序。rename 先跑完,再 convert "response" 时这个字段已经不存在了,转换是空操作,
# 输出里就会看到带引号的 "response_new" => "403" 而不是 403.0。所以要转的是新名字:
convert => { "response_new" => "float" } # 将字段类型修改为 float
gsub => ["referrer","\"",""] # 将 referrer 字段中所有 "\" 字符替换为 ""。
remove_field => ["timestamp"] # timestamp 是 grok 从日志行里抠出来的日志时间,已经被 date 过滤器解析进 @timestamp 了,删它只为去冗余(承载采集时间的是 @timestamp
split => ["clientip", "."] # 将 ip 以 . 分为列表
}
}
output {
stdout {
codec => "rubydebug"
}
}

转换后的结果,注意 @timestamp 已经变成日志里的时间而不是采集时间:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
{
"host" => {
"hostname" => "hadoop03"
},
"@timestamp" => 2021-06-16T08:24:19Z,
"referrer" => "GET / HTTP/1.1",
"response_new" => "403",
"clientip" => [
[0] "192",
[1] "168",
[2] "2",
[3] "241"
],
"@version" => "1",
"bytes" => "5039"
}

输出插件(Output)

常用的输出有这么几类:stdout 一般只用来调试;file 把日志写到磁盘文件;elasticsearch 把数据发给 ES,便于高效查询和长期保存;此外还支持 Nagios、HDFS、Email、Exec 等。

输出到标准输出,也就是前面一直在用的模式:

1
2
3
4
5
6
7
8
9
input{
stdin {
}
}
output {
stdout {
codec => rubydebug
}
}

输出到文件,路径里可以直接用时间和字段做变量:

1
2
3
4
5
6
7
8
9
input{
stdin {
}
}
output {
file {
path => "/data/log/%{+yyyy-MM-dd}/%{host}_%{+HH}.log"
}
}

启动后输入 hello world,落盘内容如下。这里 %{host} 取到的是一个对象,所以文件名有点难看,实际使用时建议写成具体的子字段(如 %{[host][hostname]}):

1
2
$ cat /data/log/2022-05-23/\{\"hostname\"\:\"hadoop03\"\}_02.log
{"host":{"hostname":"hadoop03"},"message":"hello world","@version":"1","@timestamp":"2022-05-23T02:20:57.849569Z","event":{"original":"hello world"}}

输出到 Elasticsearch 的写法见后面的「联调验证」。

Filebeat

Filebeat 装在需要采集日志的机器上,这里三个节点都用 root 用户安装。官网下载地址 https://www.elastic.co/cn/downloads/beats/filebeat

解压安装

1
2
3
4
5
6
wget https://artifacts.elastic.co/downloads/beats/filebeat/filebeat-8.2.0-linux-x86_64.tar.gz

mkdir -p /opt/bigdata/filebeat
tar -zxf filebeat-8.2.0-linux-x86_64.tar.gz -C /opt/bigdata/filebeat
cd /opt/bigdata/filebeat/
ln -s filebeat-8.2.0-linux-x86_64 current

输出到 Kafka

集群里已经装了 Kafka,所以让 Filebeat 直接把日志投到 Kafka。

先想清楚一件事:你希望投到 Kafka 的消息长什么样。 下面这份配置里有几组选项其实是互相抵消的,先决定形态能省很多来回:

  • codec.format.string: '%{[message]}' 的意思是”只发原始日志行”。一旦写了这个,前面所有 add_*_metadata processor 生成的字段、以及 drop_fields 删掉的字段,统统作废——因为输出根本不带它们。想保留结构化字段就别配 codec.format.string,让它默认发 JSON。
  • 反过来,如果确实只要原始行,那么 add_host_metadataadd_docker_metadata 这些就该一并去掉,它们白白消耗 CPU 和内存。
  • drop_fields 里删 hostecsinput 之前想一下:Logstash 侧还需不需要靠 host.hostname 区分来源?删早了后面就补不回来。
  • setup.template.settingssetup.kibanaoutput 是 Kafka 时完全不生效。索引模板注册和 Kibana 加载走的是 ES output 那条路,Kafka 输出时这两段是空配置,留着只会让人误以为模板已经建好了。真要注册模板得临时切到 ES output 跑一次 filebeat setup,或者干脆在 Logstash/ES 侧手工建。

一句话:要原始行就砍掉所有 processor;要结构化就砍掉 codec.format.string 两者都留着,实际生效的永远是后者,前面的配置全是自我安慰。修改 /opt/bigdata/filebeat/current/filebeat.yml,采集登录日志 /var/log/secure

1
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
filebeat.inputs:
- type: log
enabled: true
paths:
- /var/log/secure # 收集登录日志
fields:
log_topic: omessages
filebeat.config.modules:
path: ${path.config}/modules.d/*.yml
reload.enabled: false
setup.template.settings:
index.number_of_shards: 1
name: "hadoop01" # 按需修改
setup.kibana:
output.kafka:
enabled: true
hosts: ["hadoop01:9092", "hadoop02:9092", "hadoop03:9092"]
version: "0.10"
topic: 'my_test'
codec.format.string: '%{[message]}' # 输出原始格式, 删除则输出 json 处理后
partition.round_robin:
reachable_only: true
worker: 2
required_acks: 1
compression: gzip
max_message_bytes: 10000000
logging.level: debug
processors:
- add_host_metadata:
when.not.contains.tags: forwarded
- add_cloud_metadata: ~
- add_docker_metadata: ~
- add_kubernetes_metadata: ~
- drop_fields: # 删除的字样
fields: ["input", "host", "agent.type", "agent.ephemeral_id", "agent.id", "agent.version", "ecs"]

几个容易踩的点:

  1. topic 写死成了 my_testfields.log_topic 并没有生效。要按来源分 topic,就把它写成 topic: '%{[fields][log_topic]}'
  2. drop_fields 用来裁掉 Filebeat 自动附加的一大堆元数据字段,日志量大时能省不少带宽和存储
  3. codec.format.string 决定投出去的是原始行还是 JSON,二者对下游 Logstash 的解析方式影响很大
  4. logging.level: debug 只适合调试期,稳定后改回 info

启动与在 Kafka 侧验证

用 root 用户启动,-e 表示日志打到标准错误、-c 指定配置文件:

1
2
cd /opt/bigdata/filebeat/current
nohup ./filebeat -e -c filebeat.yml &

改完配置要重启 Filebeat,然后 ssh 连一次主机制造新日志。在 Kafka 侧消费对应 topic:

1
2
cd /opt/bigdata/kafka/current/bin
./kafka-console-consumer.sh --bootstrap-server hadoop01:9092,hadoop02:9092,hadoop03:9092 --topic my_test --from-beginning

默认配置下消息是完整的 JSON,字段非常多:

1
{"@timestamp":"2022-05-22T08:20:32.262Z","@metadata":{"beat":"filebeat","type":"_doc","version":"8.2.0"},"ecs":{"version":"8.0.0"},"log":{"offset":3204,"file":{"path":"/var/log/secure"}},"message":"May 22 04:20:30 hadoop02 sshd[18047]: pam_systemd(sshd:session): Failed to release session: Interrupted system call","input":{"type":"log"},"host":{"containerized":false,"ip":["192.168.2.242","fe80::ec97:d991:4336:2e98","fe80::f0df:f765:7f99:9634"],"mac":["00:0c:29:68:79:09"],"hostname":"hadoop02","name":"hadoop02","architecture":"x86_64","os":{"type":"linux","platform":"centos","version":"7 (Core)","family":"redhat","name":"CentOS Linux","kernel":"3.10.0-1160.el7.x86_64","codename":"Core"},"id":"0988a88e747e428dbcf4fdc212a6c1ac"},"agent":{"ephemeral_id":"4f44629c-d5e9-4ae4-a5fc-6f96df866dfe","id":"5d7e5f81-16e5-4863-b736-89f6873105ec","name":"hadoop02","type":"filebeat","version":"8.2.0"}}

加上 drop_fields 之后,只剩下关心的字段:

1
{"@timestamp":"2022-05-22T08:32:49.556Z","@metadata":{"beat":"filebeat","type":"_doc","version":"8.2.0"},"log":{"file":{"path":"/var/log/secure"},"offset":5550},"message":"May 22 04:32:49 hadoop02 sshd[23225]: Received disconnect from 192.168.2.242 port 55948:11: disconnected by user","agent":{"name":"hadoop02"}}

再加上 codec.format.string: '%{[message]}',投出去的就是原始日志行:

1
2
May 22 04:43:41 hadoop01 sshd[3106]: Accepted publickey for root from 192.168.2.243 port 41710 ssh2: RSA SHA256:F1RBzp64noGdTdwWX8w+PYfi0zs8ifzkv+etLAOaCJQ
May 22 04:43:41 hadoop01 sshd[3106]: pam_unix(sshd:session): session opened for user root by (uid=0)

联调验证

四个组件都起来之后,用最短的一条链路验证一遍:Logstash 读本机日志文件,处理后写入 Elasticsearch,再到 Kibana 里查。

新建 /opt/bigdata/logstash/current/secure_into_es.conf

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
input {
file {
path => ["/var/log/secure"]
start_position => "beginning"
}
}
filter {
grok {
match => { "message" => "%{SYSLOGLINE}" }
overwrite => ["message"] # SYSLOGLINE 内部又捕获了一个 message,不加这句会变成数组
}
date {
match => ["timestamp", "MMM d HH:mm:ss", "MMM d HH:mm:ss"]
target => "@timestamp" # 把日志自己的时间写进 @timestamp
}
}
output {
elasticsearch {
hosts => ["hadoop01:9200","hadoop02:9200","hadoop03:9200"]
index => "secure-%{+yyyy.MM.dd}"
}
}

这份配置里有三个地方是踩过坑之后才补上的,单独说一下。

一、date 过滤器不能省。 前面第 4 节专门强调过要把日志时间解析进 @timestamp,但落地示例里如果只有 grok 没有 date,@timestamp 就会是 Logstash 读到这一行的时刻。配上 start_position => "beginning" 之后,一个存了半年的 /var/log/secure 会被整批打上”导入那一分钟”的时间戳——在 Discover 里按时间轴看全糊成一根竖线,正是前文想避免的结果。

二、索引名用 %{+YYYY.MM.dd} 会在跨年那几天出错。 Joda 时间格式里大写 YYYYweek-year(ISO 周所属的年),不是日历年。12 月最后几天如果属于下一年的第 1 周,YYYY 就会给出下一年:

日期 yyyy.MM.dd YYYY.MM.dd
2025-12-28 2025.12.28 2025.12.28
2025-12-29 2025.12.29 2026.12.29

于是每年年底会凭空多出一批索引名跳到明年的数据,按 secure-2025.* 查就漏掉了。小写 yyyy 才是日历年。同一篇里第 522 行 file output 用的正是小写,两处本来就不一致——统一成小写即可。

三、%{SYSLOGLINE} 会把 message 变成数组。 这个内置 pattern 内部自己又捕获了一个名叫 message 的字段,而 grok 对已存在的字段是追加而不是覆盖。结果 message 从字符串变成了两元素数组(原始整行 + 解析出的正文),写进 ES 后成了多值字段,影响检索和高亮,而且报错信息里看不出原因,很容易以为是自己 pattern 写错了。加 overwrite => ["message"] 让它覆盖即可。

启动 Logstash:

1
2
cd /opt/bigdata/logstash/current
nohup bin/logstash -f secure_into_es.conf &

先在 ES 侧确认索引已经建出来、文档数在涨:

1
2
curl "hadoop01:9200/_cat/indices?v"
curl "hadoop01:9200/secure-*/_search?size=1&pretty"

再到 Kibana(http://hadoop03:5601 )里为 secure-* 建一个 data view,就能在 Discover 中按时间和字段检索了。

需要中间加 Kafka 缓冲的完整链路(Filebeat → Kafka → Logstash → Elasticsearch → Kibana),包括 Logstash 侧的 kafka input 写法和 Kibana 建索引的界面步骤,见 Kubernetes 日志收集 一文的 EFLK 部分。