Logstash 插件实战与 ES 集群安全加固

Logstash 插件实战与 ES 集群安全加固
[TOC]
Logstash 插件体系分析 nginx 日志
Logstash 的 ==Filter 插件==是数据处理的核心,负责将原始日志==解析成结构化字段==
- grok:正则提取任意文本 → 映射为字段
- date:把日志里的时间字符串 → 替换
@timestamp - useragent:从 UA 字符串提取设备/浏览器/OS
- geoip:把公网 IP → 国家/城市/经纬度/坐标点
grok — 基于正则的字段提取

- 链接:官方文档
- Ctrl + F 搜索
Logstash—> other versions —> filter过滤插件 —> grok字段过滤插件
通俗解释:
grep 是查找,grok 是==拿着模板套格式==(内置的正则模板)
比如 nginx 日志有一套固定格式,grok 里就有现成的 %{HTTPD_COMMONLOG} 模板,直接套上去就把各个字段拆出来了
#!/usr/bin/env python3import randomimport timefrom datetime import datetime, timedelta
# 伪造IP池——覆盖多个国家的地址段(用于 GeoIP 出图)IPS = [ "1.2.3.4", # 澳大利亚 "8.8.8.8", # 美国 "13.107.42.14", # 美国(微软) "31.13.24.12", # 爱尔兰 "47.96.0.1", # 中国杭州 "52.84.0.1", # 美国 "59.24.3.1", # 韩国 "77.88.55.88", # 俄罗斯 "103.235.46.39", # 日本 "110.242.68.66", # 中国北京 "151.101.1.1", # 美国(Fastly) "185.199.108.153", # 德国 "202.108.22.5", # 中国北京 "203.119.24.1", # 越南 "216.58.200.4", # 美国(Google)]
# 状态码权重更丰富STATUS = [200] * 60 + [201] * 5 + [204] * 3 + [301] * 8 + [302] * 5 + [304] * 5 + [400] * 3 + [401] * 2 + [403] * 3 + [404] * 8 + [500] * 3 + [502] * 2 + [503] * 2 + [504] * 1
# 更丰富的URL列表(包含多种HTTP方法)URLS = [ # GET 请求 "/", "/index.html", "/api/health", "/api/products", "/api/products/123", "/api/products/456", "/api/users", "/api/users/789", "/api/orders", "/api/orders/321", "/api/orders/654", "/about", "/contact", "/faq", "/images/logo.png", "/images/banner.jpg", "/css/main.css", "/css/theme.css", "/js/app.js", "/js/vendor.js", "/docs/api-reference", "/docs/getting-started", "/blog/post/2024/01", "/blog/post/2024/02", "/blog/post/2024/03", "/products/list", "/products/detail/1", "/products/detail/2", "/login", "/register", "/cart", "/checkout", "/payment", "/success", "/admin/dashboard", "/admin/users", "/admin/settings", "/api/v2/products", "/api/v2/users", "/api/v2/orders"]
# HTTP方法METHODS = ["GET"] * 85 + ["POST"] * 10 + ["PUT"] * 3 + ["DELETE"] * 2
# 更丰富的User-Agent池UAS = [ # Chrome - Windows 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/141.0.0.0 Safari/537.36', 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/140.0.0.0 Safari/537.36', 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/139.0.0.0 Safari/537.36', # Chrome - Mac 'Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/141.0.0.0 Safari/537.36', 'Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/140.0.0.0 Safari/537.36', # Safari - iPhone 'Mozilla/5.0 (iPhone; CPU iPhone OS 18_0 like Mac OS X) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/18.0 Mobile/15E148 Safari/604.1', 'Mozilla/5.0 (iPhone; CPU iPhone OS 17_5 like Mac OS X) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/17.0 Mobile/15E148 Safari/604.1', # Safari - iPad 'Mozilla/5.0 (iPad; CPU OS 18_0 like Mac OS X) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/18.0 Mobile/15E148 Safari/604.1', # Android Chrome 'Mozilla/5.0 (Linux; Android 14; Pixel 8 Pro) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/141.0.792.99 Mobile Safari/537.36', 'Mozilla/5.0 (Linux; Android 13; SM-S908B) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/140.0.0.0 Mobile Safari/537.36', # Firefox - Windows 'Mozilla/5.0 (Windows NT 10.0; Win64; x64; rv:135.0) Gecko/20100101 Firefox/135.0', 'Mozilla/5.0 (Windows NT 10.0; Win64; x64; rv:134.0) Gecko/20100101 Firefox/134.0', # Firefox - Mac 'Mozilla/5.0 (Macintosh; Intel Mac OS X 10.15; rv:135.0) Gecko/20100101 Firefox/135.0', # Edge 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/141.0.0.0 Safari/537.36 Edg/141.0.0.0', 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/140.0.0.0 Safari/537.36 Edg/140.0.0.0', # Opera 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/140.0.0.0 Safari/537.36 OPR/120.0.0.0', # 搜索引擎爬虫 'Mozilla/5.0 (compatible; Googlebot/2.1; +http://www.google.com/bot.html)', 'Mozilla/5.0 (compatible; Bingbot/2.0; +http://www.bing.com/bingbot.htm)', # 其他 'curl/8.5.0', 'Wget/1.21.4', 'Python-urllib/3.11',]
# Referer池REFERERS = [ "-", # 直接访问 "https://www.google.com/", "https://www.baidu.com/", "https://www.bing.com/", "https://www.facebook.com/", "https://twitter.com/", "https://www.example.com/", "https://www.example.com/blog", "https://www.example.com/products", "https://www.example.com/login", "https://www.google.com/search?q=example", "https://www.baidu.com/s?wd=example", "https://www.example.com/api/products", "https://www.example.com/docs", "https://news.ycombinator.com/", "https://www.reddit.com/r/programming/",]
def generate_log_line(): """生成一条随机日志""" ip = random.choice(IPS) method = random.choice(METHODS) url = random.choice(URLS) status = random.choice(STATUS) body_size = random.randint(50, 102400) # 50字节到100KB
# 生成随机时间(最近7天内) days_ago = random.randint(0, 7) hours_ago = random.randint(0, 23) minutes_ago = random.randint(0, 59) seconds_ago = random.randint(0, 59)
log_time = datetime.now() - timedelta( days=days_ago, hours=hours_ago, minutes=minutes_ago, seconds=seconds_ago ) time_str = log_time.strftime("%d/%b/%Y:%H:%M:%S +0800")
ua = random.choice(UAS) referer = random.choice(REFERERS)
# 构建日志行 log_line = f'{ip} - - [{time_str}] "{method} {url} HTTP/1.1" {status} {body_size} "{referer}" "{ua}"\n' return log_line
# 生成200条日志log_file = "/var/log/nginx/access.log"print(f"开始生成200条随机日志到 {log_file}...")
with open(log_file, "a") as f: for i in range(200): log_line = generate_log_line() f.write(log_line)
# 每50条显示一次进度 if (i + 1) % 50 == 0: print(f"已生成 {i + 1}/200 条...")
# 偶尔随机延迟,模拟真实日志写入 if random.random() < 0.05: # 5%的概率延迟 time.sleep(random.uniform(0.01, 0.1))
print("✅ 完成! 共生成 200 条日志")
# 显示一些统计信息print("\n📊 日志统计:")print(f" - 总条数: 200")print(f" - 唯一IP数: {len(set(IPS))}")print(f" - 唯一URL数: {len(URLS)}")print(f" - 唯一User-Agent数: {len(UAS)}")print(f" - 状态码分布: 2xx, 3xx, 4xx, 5xx")1)准备 nginx 测试日志root@Elk02 ~# vim /root/generate_nginx_log.py# 生成 200 条测试日志用于本实验root@Elk02 ~# > /var/log/nginx/access.logroot@Elk02 ~# python3 /root/generate_nginx_log.py✅ 完成! 共生成 200 条日志root@Elk02 ~# wc -l /var/log/nginx/access.log200 /var/log/nginx/access.log
2)配置 Logstash — grok filterroot@Elk03 ~# vim /etc/logstash/conf.d/04-beats-grok-es.confinput { beats { port => "6666" }}
filter { mutate { remove_field => [ "@version","agent","log","ecs","tags","input","host" ] }
grok { match => { "message" => "%{HTTPD_COMBINEDLOG}" } }# HTTPD_COMBINEDLOG 是 logstash 内置的正则模板}
output { stdout { codec => "rubydebug" }}
- 链接地址:匹配模板
root@Elk03 ~# tree /usr/share/logstash/vendor/bundle/jruby/2.5.0/gems/logstash-patterns-core-4.3.4/patterns/legacy/usr/share/logstash/vendor/bundle/jruby/2.5.0/gems/logstash-patterns-core-4.3.4/patterns/legacy├── aws├── firewalls├── haproxy├── httpd├── java.............
root@Elk03 ~# cat /usr/share/logstash/vendor/bundle/jruby/2.5.0/gems/logstash-patterns-core-4.3.4/patterns/legacy/httpdHTTPDUSER %{EMAILADDRESS}|%{USER}HTTPDERROR_DATE %{DAY} %{MONTH} %{MONTHDAY} %{TIME} %{YEAR}
# Log formatsHTTPD_COMMONLOG %{IPORHOST:clientip} %{HTTPDUSER:ident} %{HTTPDUSER:auth} \[%{HTTPDATE:timestamp}\] "(?:%{WORD:verb} %{NOTSPACE:request}(?: HTTP/%{NUMBER:httpversion})?|%{DATA:rawrequest})" (?:-|%{NUMBER:response}) (?:-|%{NUMBER:bytes})........................# Error logsHTTPD20_ERRORLOG \[%{HTTPDERROR_DATE:timestamp}\] \[%{LOGLEVEL:loglevel}\] (?:\[client %{IPORHOST:clientip}\] ){0,1}%{GREEDYDATA:message}
# DeprecatedCOMMONAPACHELOG %{HTTPD_COMMONLOG}COMBINEDAPACHELOG %{HTTPD_COMBINEDLOG} ✅3)清理并启动'清理环境'root@Elk02 ~# killall -9 filebeatroot@Elk02 ~# rm -rf /var/lib/filebeat/root@Elk03 ~# pkill -f -9 logstash# ⚠️ logstash 进程名是 java,killall 匹配不到,用pkill -f 选项(匹配整个命令行)# beats input 监听 6666 端口,数据从 Filebeat 流入,不产生 sincedb,只需清理 filebeat 端即可
4)启动 Logstashroot@Elk03 ~# logstash -rf /etc/logstash/conf.d/04-beats-grok-es.conf --config.test_and_exit# 先干跑root@Elk03 ~# logstash -rf /etc/logstash/conf.d/04-beats-grok-es.confroot@Elk03 ~# ss -ntl | grep 6666LISTEN 0 4096 *:6666 *:*
5)配置 Filebeat — 发送 nginx 日志到 Logstashroot@Elk02 ~# vim /etc/filebeat/config/07-nginx-to-logstash.yamlfilebeat.inputs:- type: filestream paths: - /var/log/nginx/access.log # 只有一个 input,不加 id 只是 WARNING,不影响
output.logstash: hosts: ["10.0.0.8:6666"]
6)启动 Filebeat,观察 Logstash 终端root@Elk02 ~# filebeat -e -c /etc/filebeat/config/07-nginx-to-logstash.yaml
7)grok 解析结果 — Logstash 终端输出root@Elk03 ~# 终端输出 (rubydebug 格式):{ "httpversion" => "1.1", "@timestamp" => 2026-07-30T01:47:02.793Z, "bytes" => "1250", "response" => "204", "referrer" => "\"https://github.com\"", "request" => "/favicon.ico", "verb" => "DELETE", "timestamp" => "29/Jul/2026:09:00:00 +0800", "clientip" => "218.177.17.249", "agent" => "\"Mozilla/5.0 (Linux; Android 14; SM-G998B) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.6099.43 Mobile Safari/537.36\"" "auth" => "-", "ident" => "-", "message" => "218.177.17.249 - - [07/Mar/2026:08:28:32 +0800] \"DELETE /favicon.ico HTTP/1.0\" 204 1250 \"https://github.com\" \"Mozilla/5.0 (Linux; Android 14; SM-G998B) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.6099.43 Mobile Safari/537.36\"",}# grok 成功将原始日志拆成了 clientip、agent、verb、request、response、bytes 等独立字段⚠️ 注意 @timestamp 显示的是处理时间(logstash 接收事件的时间),不是日志里的真实时间'下一节用 date filter 修正'date — 修正 @timestamp

官网链接:date字段过滤插件
默认 @timestamp 是 Logstash 收到事件的时间,不是日志生成的时间
这导致 Kibana 里的时间轴不对,需要 date 插件把日志里的时间覆盖到 @timestamp
1)date filter 配置root@Elk03 ~# vim /etc/logstash/conf.d/04-beats-grok-es.confinput { beats { port => "6666" }}
filter { mutate { remove_field => [ "@version","agent","log","ecs","tags","input","host" ] }
grok { match => { "message" => "%{HTTPD_COMBINEDLOG}" } }
date { # "30/Jul/2026:09:00:00 +0800" 匹配 nginx 日志中的时间格式 match => [ "timestamp", "dd/MMM/yyyy:HH:mm:ss Z" ] target => "@timestamp" # 用日志里的时间覆盖 @timestamp }}
output { stdout { codec => "rubydebug" }}| date 模式字符 | 含义 | nginx 日志中的值 |
|---|---|---|
dd | 日 (2位) | 30 |
MMM | 月 (英文缩写) | Jul |
yyyy | 年 (4位) | 2026 |
HH | 时 (24小时制) | 09 |
mm | 分 | 00 |
ss | 秒 | 00 |
Z | 时区 (+0800) | +0800 |

2)重新发送数据root@Elk02 ~# rm -rf /var/lib/filebeat/root@Elk02 ~# filebeat -e -c /etc/filebeat/config/07-nginx-to-logstash.yaml
3)观察 time 字段变化{ "@timestamp" => 2026-07-29T01:00:00.000Z,# ↑ 变成了日志里的 09:00 +0800 → 01:00 UTC ✅️# 不再是处理时的 01:47:02 "timestamp" => "29/Jul/2026:09:00:00 +0800",}✅️ @timestamp 已修正为日志真实时间useragent — 分析用户设备

- 官网链接:useragent字段过滤插件
User-Agent 字符串里嵌入了操作系统、浏览器、设备型号等信息 useragent 插件把这些信息自动提取成结构化字段,直接用于 Kibana 出图
1)useragent filter 配置root@Elk03 ~# vim /etc/logstash/conf.d/04-beats-grok-es.confinput { beats { port => "6666" }}
filter { mutate { remove_field => [ "@version","agent","log","ecs","tags","input","host" ] }
grok { match => { "message" => "%{HTTPD_COMBINEDLOG}" } }
date { match => [ "timestamp", "dd/MMM/yyyy:HH:mm:ss Z" ] target => "@timestamp" }
useragent { # 基于哪个字段提取设备信息 source => "agent" # 将分析结果存储到指定字段(不指定则散落在顶级字段) target => "kpyun-agent" }}
output { stdout { codec => "rubydebug" }}source => "message" 为什么能工作?
useragent 插件的正则在 message 全文里搜索匹配,Mozilla/5.0 这个特征足够独特,即使前面有 IP、时间戳等多余文本也不影响,所以当前写法能正常返回解析结果
但更佳实践:将 grok 的 %{HTTPD_COMMONLOG} 换成 %{HTTPD_COMBINEDLOG} (我已经自行更换过了)
后者在前者的基础上额外提取了 referrer 和 agent 两个字段
🔥 然后 source => "agent",精准命中目标字段,干净准确,少了无意义的全文扫描开销
2)重新发送root@Elk02 ~# rm -rf /var/lib/filebeat/root@Elk02 ~# filebeat -e -c /etc/filebeat/config/07-nginx-to-logstash.yaml
3)查看 useragent 解析结果{ "kpyun-agent" => { "name" => "Chrome Mobile", "major" => "141", "os" => "Android", "version" => "141.0.792.99", "os_name" => "Android", "os_full" => "Android 14", "device" => "Pixel 8 Pro", "os_version" => "14", "patch" => "792", "os_major" => "14", "minor" => "0" },}# 从 UA 字符串中解析出了:浏览器品牌+版本、操作系统、设备型号'后面 Kibana 出图时可以按 os_name、device 等维度聚合'geoip — 解析 IP 的地理位置

- 官网链接:geoip字段过滤插件
- geoip 插件需要离线数据库(GeoLite2-City.mmdb),把公网 IP 翻译成地理信息
- 下载地址:GeoLite2-City.mmdb(需自行注册)
- 解析出
geoip.location(经/纬度)后才能在地图上打点
jiuzhao@Ubuntu ES集群$ scp GeoLite2-City_20250311.tar.gz Elk03:/tmp
1)准备 GeoIP 数据库root@Elk03 ~# tar xf /tmp/GeoLite2-City_20250311.tar.gz -C /root/root@Elk03 ~# ls -lh /root/GeoLite2-City_20250311/GeoLite2-City.mmdb-rw-r--r-- 1 root root 58M GeoLite2-City.mmdb# City 数据库 58MB,包含 IP→国家/城市/经纬度信息✅️ 足够用(ASN 数据库数据量少,不需要)
2)geoip filter 配置root@Elk03 ~# vim /etc/logstash/conf.d/04-beats-grok-es.confinput { beats { port => "6666" }}
filter { mutate { remove_field => [ "@version","agent","log","ecs","tags","input","host" ] }
grok { match => { "message" => "%{HTTPD_COMBINEDLOG}" } }
date { match => [ "timestamp", "dd/MMM/yyyy:HH:mm:ss Z" ] target => "@timestamp" }
useragent { source => "agent" target => "kpyun-agent" }
geoip { # 基于哪个字段解析IP source => "clientip"
# 指定本地的数据库路径 database => "/root/GeoLite2-City_20250311/GeoLite2-City.mmdb" }}
output { stdout { codec => "rubydebug" }}
3)重新发送,观察 geoip 解析结果root@Elk02 ~# rm -rf /var/lib/filebeat/root@Elk02 ~# filebeat -e -c /etc/filebeat/config/07-nginx-to-logstash.yaml
4)geoip 解析结果{ "geoip" => { "country_code2" => "CN", "country_name" => "China", "ip" => "110.242.68.66", "country_code3" => "CN", "continent_code" => "AS", "longitude" => 113.722, "latitude" => 34.7732, "timezone" => "Asia/Shanghai", "location" => { "lon" => 113.722, "lat" => 34.7732 }, }, "clientip" => "110.242.68.66",}# geoip 把 IP 110.242.68.66 解析成了中国(China),坐标 lat=34.77 / lon=113.72# location 字段包含经纬度(geo_point),可用于 Kibana 地图打点| geoip 产出字段 | 含义 | Kibana 用途 |
|---|---|---|
country_name | 国家名 | 饼图/柱状图聚合 |
city_name | 城市名 (部分 IP 可能缺失) | 表格/词云 |
continent_code | 大洲代码 | 地域聚合 |
longitude / latitude | 浮点经/纬度 | — |
timezone | 时区 | 辅助分析 |
location.lat / location.lon | geo_point 坐标 | ==地图打点== |
ELFK 出图展示故障案例
1)将处理结果写入 ES(先不搞模板)root@Elk03 ~# vim /etc/logstash/conf.d/04-beats-grok-es.confinput { beats { port => "6666" }}
filter { # ... grok + date + useragent + geoip 全套 filter ...}
output { stdout { codec => "rubydebug" }
# 我们把数据既输出至屏幕也也输出至ES集群 elasticsearch { hosts => ["10.0.0.6:9200","10.0.0.7:9200","10.0.0.8:9200"] index => "kpyun-end-elfk-nginx-%{+YYYY.MM.dd}" }}
2)重新发送root@Elk02 ~# rm -rf /var/lib/filebeat/root@Elk02 ~# filebeat -e -c /etc/filebeat/config/07-nginx-to-logstash.yaml
3)WebUI查看http://10.0.0.6:5601/# 这些索引是按天分布的(脚本具有随机性,是近几天的访问日志)
创建索引模式 --> kpyun-end-elfk-nginx* --> 选择时间戳

# ⚠️ location(包含 lat 和 lon)是 float(浮点型),不是 geo_point(地理坐标点) 类型!# Kibana 地图需要 geo_point(地理坐标点) 才能画点标记===============================================# 同理检查 bytes 和 clientip"bytes": { "type": "text", #(应该是 integer)"clientip": { "type": "text", #(应该是 ip)# 这些字段没有显式映射,ES 自动推断的类型不精确为什么需要索引模板?
- ES 的动态映射会猜测字段类型,但猜不准
geoip.location需要geo_point(地理坐标点)才能在地图上打点bytes需要integer才能做数值聚合(否则text类型无法用在 Y 轴)clientip设为ip类型可以启用 IP 范围查询
💡 可以通过 Kibana WebUI 创建索引模板 路径:Stack Management → 索引管理 → 索引模板 → 创建模板
- 名称填写
kpyun-end - 索引模式填写
kpyun-end-elfk-nginx* - 组件模板(跳过)—> 索引设置(分片数)
"number_of_shards": 7, # 主分片数 "number_of_replicas": 0 # 副本数- 在映射中添加(映射设置)
geoip.location为地理坐标点(geo_point)bytes为数值 --> 整型(integer)clientip为IP
- 别名(跳过)—> 详情预览

# 删除旧索引 + 清除 Kibana 索引模式`WebUI 操作即可`
4)删除旧索引 + 重新采集重新采集即可看到地图正常渲染root@Elk03 ~# pkill -f -9 logstashroot@Elk03 ~# logstash -rf /etc/logstash/conf.d/04-beats-grok-es.confroot@Elk02 ~# rm -rf /var/lib/filebeat/root@Elk02 ~# filebeat -e -c /etc/filebeat/config/07-nginx-to-logstash.yaml
5)再次检查 映射(mapping)✅ 已修正root@Elk01 ~# curl -s '10.0.0.6:9200/kpyun-end-elfk-nginx-2026.07.30/_mapping?pretty' | grep -A1 'location\|bytes\|clientip' "bytes" : { "type" : "integer",✅️ bytes → integer; "clientip" : { "type" : "ip"✅️ clientip → ip; "location" : { "type" : "geo_point",✅️ geoip.location → geo_point

Logstash 多实例
Logstash 不允许两个实例共享同一个数据目录(跟 Filebeat 一样的逻辑)
必须为每个实例指定独立的 --path.data
# 默认数据目录的内容root@Elk03 ~# ls -la /usr/share/logstash/data/total 20drwxr-xr-x 4 root root 4096 Jul 30 16:29 .drwxr-xr-x 12 root root 4096 Jul 30 16:29 ..drwxr-xr-x 2 root root 4096 Jul 30 16:29 dead_letter_queue-rw-r--r-- 1 root root 0 Jul 30 16:29 .lockdrwxr-xr-x 2 root root 4096 Jul 30 16:29 queue-rw-r--r-- 1 root root 36 Jul 30 16:29 uuid# .lock 锁文件,第二个实例检测到已锁定 → 报错退出# 不指定 --path.data 默认就是这个位置必须指定 --path.data:不写就共用 /usr/share/logstash/data/,第二个实例抢不到 .lock → FATAL 退出
root@Elk03 ~# cat /etc/logstash/conf.d/01-stdin-to-stdout.confinput { stdin {}}output { stdout { # 指定输出的数据格式,如果不指定,则默认值为: rubydebug codec => "json" }}root@Elk03 ~# cp /etc/logstash/conf.d/01-stdin-to-stdout.conf /tmp/instance1.confroot@Elk03 ~# cp /etc/logstash/conf.d/01-stdin-to-stdout.conf /tmp/instance2.conf
1)启动实例1 — 默认数据目录(/usr/share/logstash/data/)root@Elk03 ~# logstash -f /tmp/instance1.conf# 不指定 --path.data,默认就是 /usr/share/logstash/data/
2)启动实例2 — 指定独立数据目录root@Elk03 ~# logstash -f /tmp/instance2.conf --path.data /tmp/logstash-02# --path.data 指定独立目录,绕开 .lock 冲突
3)验证两个实例都在运行root@Elk03 ~# ps -ef | grep [l]ogstash | wc -l2root@Elk03 ~# ps -ef | grep logstash | grep -v greproot 2483 ... Logstash -f /tmp/instance1.confroot 2531 ... Logstash -f /tmp/instance2.conf --path.data /tmp/logstash-02# 两个 logstash 进程,各自独立的数据目录 ✅️Logstash 的 if 多分支语句
核心思路:Logstash 通过 if [type] == "xxx" 实现==按日志类型==分路由
type相当于 logstash 里的标签:在 ==input 中赋上==,在 ==output 中用if做路由分发==- 每个 beats input 指定不同的
type,在 output 中用 if/else 决定写到哪个索引
1)Logstash 配置 — 3 个 beats input + if/else 输出root@Elk03 ~# vim /etc/logstash/conf.d/05-multiple_input-to-if_es.yamlinput { beats { port => "6666" type => "auth"# 这个端口收 auth 日志 }
beats { port => "7777" type => "syslog"# 这个端口收 syslog }
beats { port => "8888" type => "default"# 这个端口收 kern 日志,走 else 兜底 }}
filter { mutate { split => { "message" => " " } # 以空格为分隔符切割
add_field => { # 把时间从message拆出来 --> 独立字段date "date" => "%{[message][0]}" } }
date { # 时间格式是 ISO 8601 格式(带时区) match => [ "date", "ISO8601" ] target => "@timestamp" }
mutate { # 删掉被拆过的 `message` remove_field => [ "message","@version","agent","log","ecs","tags","input" ] }}
output { if [type] == "auth" { elasticsearch { hosts => ["10.0.0.6:9200","10.0.0.7:9200","10.0.0.8:9200"] index => "kpyun-elfk-if-auth-%{+YYYY.MM.dd}" } # 条件表达式用 == 做判断 } else if [type] == "syslog" { elasticsearch { hosts => ["10.0.0.6:9200","10.0.0.7:9200","10.0.0.8:9200"] index => "kpyun-elfk-if-syslog-%{+YYYY.MM.dd}" } } else { # 其它都进入这个索引 elasticsearch { hosts => ["10.0.0.6:9200","10.0.0.7:9200","10.0.0.8:9200"] index => "kpyun-elfk-if-default-%{+YYYY.MM.dd}" } }}🌰 split 的暗坑:message 字段类型冲突
mutate { split => { "message" => " " } }split 把 message 从字符串 "2026-07-31T18:55:01 Elk01 CRON..." 拆成了数组 ["2026-07-31T18:55:01", "Elk01", "CRON"...]
ES 已有 message → text 映射,但数组首元素是日期格式,ES 误判想重映射为 date → 400 报错
ERROR: mapper [message] cannot be changed from type [text] to [date]解决方法:提取完需要的字段后,删掉被拆过的 message,==保留 date 字段==方便查看原始时间:
mutate { remove_field => [ "message" ] }# message 拆分完成了数组,不再需要# date 字段保留,方便 Kibana 直接查看每条日志的原始时间2)启动 Logstashroot@Elk03 ~# logstash -rf /etc/logstash/conf.d/05-multiple_input-to-if_es.yaml --config.test_and_exitroot@Elk03 ~# logstash -rf /etc/logstash/conf.d/05-multiple_input-to-if_es.yamlroot@Elk03 ~# ss -lntup | grep -E '6666|7777|8888'LISTEN *:8888LISTEN *:6666LISTEN *:7777# 三个端口全部就绪 ✅️
3)Filebeat 配置 — 采集系统日志root@Elk01 ~# vim /tmp/66-auth-to-logstash.yamlfilebeat.inputs:- type: filestream id: auth-collector paths: - /var/log/auth.log* prospector.scanner.exclude_files: ['\.gz$']# 排除压缩文件(以gz结尾的文件)(正则)
output.logstash: hosts: ["10.0.0.8:6666"]# auth 日志 → 6666(type=auth)
root@Elk01 ~# vim /tmp/77-syslog-to-logstash.yamlfilebeat.inputs:- type: filestream id: syslog-collector paths: - /var/log/syslog* # 这里没有.log,syslog就是日志文件 prospector.scanner.exclude_files: ['\.gz$']
output.logstash: hosts: ["10.0.0.8:7777"]# syslog → 7777(type=syslog)
root@Elk01 ~# vim /tmp/88-kern-to-logstash.yamlfilebeat.inputs:- type: filestream id: kern-collector paths: - /var/log/kern.log* prospector.scanner.exclude_files: ['\.gz$']
output.logstash: hosts: ["10.0.0.8:8888"]# kern → 8888(type=default → 走 else 分支)
4)启动 Filebeat 多实例 — 每个实例独立 registry# 实例1:auth 日志root@Elk01 ~# rm -rf /var/lib/filebeat/root@Elk01 ~# nohup filebeat -e -c /tmp/66-auth-to-logstash.yaml > /tmp/66.log 2>&1 &
# 实例2:syslog(独立 --path.data)root@Elk01 ~# rm -rf /tmp/filebeat-syslogroot@Elk01 ~# nohup filebeat -e -c /tmp/77-syslog-to-logstash.yaml --path.data /tmp/filebeat-syslog > /tmp/77.log 2>&1 &
# 实例3:kern(独立 --path.data)root@Elk01 ~# rm -rf /tmp/filebeat-kernroot@Elk01 ~# nohup filebeat -e -c /tmp/88-kern-to-logstash.yaml --path.data /tmp/filebeat-kern > /tmp/88.log 2>&1 &
5)验证三个索引分别收到数据root@Elk01 ~# curl -s '10.0.0.6:9200/_cat/indices/kpyun-elfk-if*?v'health status index docs.countgreen open kpyun-elfk-if-auth-2026.07.30 1678green open kpyun-elfk-if-syslog-2026.07.30 27800green open kpyun-elfk-if-default-2026.07.30 14329# auth→1678 条,syslog→27800 条,kern→14329 条(走 else)✅️'三条数据流各自写入不同的索引,互不干扰'⚠️ filestream 类型的 input 必须加 id 字段
- 单 input 时不加只是 WARNING,不影响运行
- ==多 input 或多实例==时必须加,否则 Filebeat 重启后分不清哪个 cursor 对应哪个文件 → ==数据重复采集==
- 多个 Filebeat 实例也必须用不同的
--path.data隔离 registry
Logstash 的 Pipeline

- 多实例 = 多个 Java 进程,各自消耗 1G+ 内存
- Pipeline = 单个进程内跑多个逻辑管道,==省内存、好管理==
- Pipeline 配置在
pipelines.yml中定义,每个 pipeline 有独立的pipeline.id和path.config
| 运行方式 | 命令 | 场景 |
|---|---|---|
| 单配置文件 | logstash -rf xxx.conf | 简单测试 |
| Pipeline 模式 | logstash -r (无 -f) | ==生产环境== |
| 多实例 | logstash -rf a.conf + logstash -rf b.conf --path.data /tmp/data | 资源隔离 |
1)编写 pipelines.ymlroot@Elk03 ~# vim /etc/logstash/pipelines.yml# - pipeline.id: main# path.config: "/etc/logstash/conf.d/*.conf"
- pipeline.id: xixi path.config: "/etc/logstash/conf.d/02-file-to-es.conf"
- pipeline.id: haha path.config: "/etc/logstash/conf.d/03-beats-to-es.conf"📌 通配符 vs 指定文件:
- 默认写法是把 config 合成==一个 pipeline==,共用同一套线程和队列
- 一条堵了,其他全卡
💡 这里把默认行注释掉,改为拆成 xixi/haha 两个独立 pipeline,各自跑自己的,互不拖累 → ==生产环境标准做法==
2)创建符号链接 — logstash 默认找这个路径root@Elk03 ~# mkdir -p /usr/share/logstash/configroot@Elk03 ~# ln -svf /etc/logstash/pipelines.yml /usr/share/logstash/config/pipelines.yml# logstash 默认加载 /usr/share/logstash/config/pipelines.yml⚠️ 常见报错:logstash 不会自动去读 /etc/logstash/pipelines.yml
必须创建软链接指向 /usr/share/logstash/config/pipelines.yml
`把之前的 模板、索引、索引模式 都删除了`
3)启动(不加 -f,自动加载 pipelines.yml)root@Elk03 ~# logstash -r --config.test_and_exitroot@Elk03 ~# logstash -r# ...[INFO] Pipelines running {:count=>2, :running_pipelines=>[:xixi, :haha]}✅️ 两个 pipeline 并行运行在同一个进程中ES 集群启用认证功能
🤡 启用 xpack.security 前建议拍个快照 🤡 启用后所有请求都必须带认证信息,否则一律 401
1)生成证书文件root@Elk01 ~# /usr/share/elasticsearch/bin/elasticsearch-certutil cert \ -out /etc/elasticsearch/elastic-certificates.p12 \ -pass "" --days 36500# -pass "" : 空密码(生产环境建议设密码)# --days 36500 : 100 年有效期........Certificates written to /etc/elasticsearch/elastic-certificates.p12........root@Elk01 ~# ls -l /etc/elasticsearch/elastic-certificates.p12-rw------- 1 root elasticsearch 3596 Jul 30 10:05 elastic-certificates.p12# 3600 字节的 PKCS#12 证书文件'权限600' # 属主root,属组elasticsearchroot@Elk01 ~# egrep -i 'group|user' /lib/systemd/system/elasticsearch.serviceUser=elasticsearchGroup=elasticsearch# ES集群运行用户是elasticsearchroot@Elk01 ~# chmod 640 /etc/elasticsearch/elastic-certificates.p12-rw-r----- 1 root elasticsearch 3596 Aug 1 09:31 elastic-certificates.p12
2)拷贝证书到所有节点root@Elk01 ~# scp -p /etc/elasticsearch/elastic-certificates.p12 10.0.0.7:/etc/elasticsearch/-p: #(小写的p)保留源文件的权限root@Elk01 ~# scp -p /etc/elasticsearch/elastic-certificates.p12 10.0.0.8:/etc/elasticsearch/# 三台机器的证书文件必须一致root@Elk02 ~# ls -l /etc/elasticsearch/elastic-certificates.p12-rw-r----- 1 root elasticsearch 3596 Aug 1 09:31 elastic-certificates.p12# 权限、属主、属组正确
3)修改 ES 配置文件 — 三台都要加root@Elk01 ~# vim /etc/elasticsearch/elasticsearch.yml# 在文件末尾追加:xpack.security.enabled: truexpack.security.transport.ssl.enabled: true# 启用安全认证,认证的方式是证书认证xpack.security.transport.ssl.verification_mode: certificatexpack.security.transport.ssl.keystore.path: elastic-certificates.p12xpack.security.transport.ssl.truststore.path: elastic-certificates.p12# 同步到其他节点root@Elk01 ~# scp /etc/elasticsearch/elasticsearch.yml 10.0.0.7:/etc/elasticsearch/root@Elk01 ~# scp /etc/elasticsearch/elasticsearch.yml 10.0.0.8:/etc/elasticsearch/
4)三台全部重启 ESroot@Elk01 ~# systemctl restart elasticsearch.serviceroot@Elk02 ~# systemctl restart elasticsearch.serviceroot@Elk03 ~# systemctl restart elasticsearch.service
5)验证 — 没有认证直接被拒root@Elk01 ~# curl 10.0.0.6:9200/_cat/nodes?v{"error":{"type":"security_exception","reason":"missing authentication credentials ..."},"status":401}# 401 Unauthorized ✅️ 认证已生效
6)生成随机密码root@Elk01 ~# /usr/share/elasticsearch/bin/elasticsearch-setup-passwords autoPlease confirm that you would like to continue [y/N] y
Changed password for user apm_systemPASSWORD apm_system = lmpdGQWoBWR5B7PjGpVi
Changed password for user kibana_systemPASSWORD kibana_system = iYLp0QpWVR0IUnpidbi6
Changed password for user kibanaPASSWORD kibana = iYLp0QpWVR0IUnpidbi6
Changed password for user logstash_systemPASSWORD logstash_system = aotLNka89X00g2zdutTU
Changed password for user beats_systemPASSWORD beats_system = YErtC18PxGbBSb9QioRh
Changed password for user remote_monitoring_userPASSWORD remote_monitoring_user = qI2N1Q77Y3dDtDfQtknC
Changed password for user elasticPASSWORD elastic = XYVMlZSaxUnA17uYP3Ti# 把这些密码记下来!后面 Kibana / Logstash / Filebeat 都要用| 内置用户 | 用途 | 对接组件 |
|---|---|---|
elastic | 超级管理员 | 管理员登录 Kibana |
kibana_system | Kibana 内部连接 | Kibana → ES |
logstash_system | Logstash 内部连接 | Logstash → ES |
beats_system | Beats 内部连接 | Filebeat → ES |
7)验证集群 — 带认证访问root@Elk01 ~# curl -u elastic:XYVMlZSaxUnA17uYP3Ti http://10.0.0.6:9200/_cat/nodes?v -u user:pass # 🔥 专门用来传递用户名和密码ip heap.percent ram.percent cpu node.role master name10.0.0.6 44 91 26 cdfhilmrstw - Elk0110.0.0.7 41 86 26 cdfhilmrstw * Elk0210.0.0.8 40 76 21 cdfhilmrstw - Elk03# 3 节点集群,带认证正常访问 ✅️Kibana 对接 ES 认证集群
1)修改 Kibana 配置 — 添加 ES 认证信息root@Elk01 ~# vim /etc/kibana/kibana.yml# 追加:elasticsearch.username: "kibana_system"elasticsearch.password: "iYLp0QpWVR0IUnpidbi6"# ↑ 用上一步生成的 kibana_system 密码
2)重启 Kibanaroot@Elk01 ~# systemctl restart kibana.serviceroot@Elk01 ~# ss -lntup | grep 5601LISTEN 0 511 0.0.0.0:5601 0.0.0.0:*# Kibana 监听正常
3)验证 Kibana 状态root@Elk01 ~# curl -s -u elastic:XYVMlZSaxUnA17uYP3Ti 'http://10.0.0.6:5601/api/status'{"status":{"overall":{"state":"green"}}}`这次访问的是5601端口,不是9200`# Kibana 已成功对接 ES 认证集群 ✅️
4)网页验证(无痕模式)http://10.0.0.6:5601/
重置 elastic 管理员密码
两种方式:
- 方式 A:创建超级管理员用户 → 通过 API 重置 elastic 密码
- 方式 B:💡 Kibana WebUI → Stack Management → 用户 → 找到
elastic→ 修改密码- 改成简单好记的(实验环境),如
passwd
- 改成简单好记的(实验环境),如

root@Elk01 ~# curl -u elastic:123456 http://10.0.0.6:9200/_cat/nodes?vip heap.percent ram.percent node.role master name10.0.0.8 59 87 cdfhilmrstw - Elk0310.0.0.6 37 93 cdfhilmrstw - Elk0110.0.0.7 58 89 cdfhilmrstw * Elk02✅️ 新密码已生效
1)创建超级管理员用户 admin-user# 此用户是 ES 本地用户,不依赖认证服务`即使 elastic 密码忘了、认证服务挂了,这个用户照样能工作`-r superuser # 赋予超级管理员权限,权限等同于 elasticroot@Elk01 ~# /usr/share/elasticsearch/bin/elasticsearch-users useradd admin-user -p passwd -r superuserroot@Elk01 ~# /usr/share/elasticsearch/bin/elasticsearch-users listadmin-user : superuser===========================================`更改本地管理员用户密码 admin-user`/usr/share/elasticsearch/bin/elasticsearch-users passwd admin-user -p kpyun666# useradd 创建用户,passwd 改密码
2)通过 API 重置 elastic 密码`用新用户重置 elastic 密码`root@Elk01 ~# curl -s --user admin-user:passwd \ -XPUT "http://10.0.0.6:9200/_xpack/security/user/elastic/_password?pretty" \ -H 'Content-Type: application/json' \ -d '{"password": "passwd"}'# 返回 { } 表示成功(空对象 = 无报错)
3)验证新密码root@Elk01 ~# curl -s -u elastic:passwd '10.0.0.6:9200/_cat/nodes?v'ip node.role master name10.0.0.6 cdfhilmrstw - Elk0110.0.0.7 cdfhilmrstw * Elk0210.0.0.8 cdfhilmrstw - Elk03✅️ 新密码生效Filebeat 和 Logstash 对接 ES 认证集群
1)Filebeat 添加认证信息root@Elk02 ~# vim /etc/filebeat/config/tcp-to-es-auth.yamlfilebeat.inputs:- type: tcp host: "0.0.0.0:9000"
output.elasticsearch: hosts: - "http://10.0.0.6:9200" - "http://10.0.0.7:9200" - "http://10.0.0.8:9200" index: "kpyun-auth-filebeat-%{+yyyy-MM-dd}" username: "elastic" password: "passwd"# ↑ 加这两行即可
setup.ilm.enabled: falsesetup.template.name: "kpyun-auth-fb"setup.template.pattern: "kpyun-auth-filebeat*"setup.template.overwrite: falsesetup.template.settings: index.number_of_shards: 3 index.number_of_replicas: 0
2)启动 + 测试root@Elk02 ~# filebeat -e -c /etc/filebeat/config/tcp-to-es-auth.yamlroot@Elk01 ~# echo 'hello elk auth test' | nc -w 1 10.0.0.7 9000root@Elk01 ~# curl -s -u elastic:passwd '10.0.0.6:9200/kpyun-auth-filebeat-2026-08-01/_search?pretty' | grep message "message" : "hello elk auth test",✅️ Filebeat → 认证 ES 集群成功
3)Logstash 添加认证信息root@Elk03 ~# vim /etc/logstash/conf.d/99-tcp-to-es-auth.confinput { tcp { port => 9999 }}
output { elasticsearch { hosts => ["http://10.0.0.6:9200","http://10.0.0.7:9200","http://10.0.0.8:9200"] index => "kpyun-logstash-auth-%{+yyyy-MM-dd}" user => "elastic" password => "passwd"# ↑ logstash 用 user/password,不是 username/password }}
4)启动 + 测试root@Elk03 ~# logstash -rf /etc/logstash/conf.d/99-tcp-to-es-auth.confroot@Elk01 ~# echo 'logstash auth test' | nc -w 1 10.0.0.8 9999root@Elk01 ~# curl -s -u elastic:passwd '10.0.0.6:9200/_cat/indices/kpyun-logstash-auth*?v'health status index docs.countgreen open kpyun-logstash-auth-2026-07-30 1✅️ Logstash → 认证 ES 集群成功| 组件 | 认证配置字段 | 示例 |
|---|---|---|
| Filebeat | username / password | username: "elastic" |
| Logstash | user / password | user => "elastic" |
| Kibana | elasticsearch.username / elasticsearch.password | elasticsearch.username: "kibana_system" |
⚠️ Filebeat 用 username,Logstash 用 user — 两个组件的字段名不一样,容易搞混!
Kibana 实现 RBAC

==RBAC (Role-Based Access Control)==:基于角色的访问控制
- 先创建角色(Role),定义能访问哪些索引、有哪些权限
- 再创建用户(User),把用户绑定到角色
- 不同用户登录 Kibana 后,==看到的数据不一样==
| 层级 | 说明 |
|---|---|
| 用户(User) | 实际访问 Kibana/ES 的个体(如 xixi、haha、hehe) |
| 角色(Role) | 权限集合,介于用户和资源之间(如 dba_role、k8s_role、sre_role) |
| ES 集群资源【索引】 | 角色绑定的索引级权限,包含:create_doc(写入文档)、create_index(API显式创建索引)、index(写入时自动建不存在的索引)、read(读取)、all(全部) |
| Kibana 界面资源【操作面板】 | 角色绑定的功能面板权限:Enterprise Search / Analytics / Observability / Management / Security |
用户 → 被赋予角色 → 角色通过绑定关联到索引级操作权限和 Kibana 面板访问权限,从而实现细粒度访问控制
官方权限参考:Elastic 7.17 权限文档
创建角色
-
登录 Kibana → Stack Management → 安全 → 角色
-
点击 创建角色
-
填写角色名称(如
dba_role)
两类资源 -
集群权限:勾选
monitor(允许查看集群状态)(不要填all权限太大) -
运行身份权限:什么都不要填,必须空着(允许代表其他用户提交请求,开了等于权限控制白做)
-
索引权限:
- 索引模式:
kpyun-dba* - 权限:
create_index、create_doc、index、read - 同理创建
k8s_role(索引模式:kpyun-k8s*)和sre_role(索引模式:kpyun-sre*) - 索引权限给
all即可,不是集群权限
- 索引模式:
-
Kibana 权限:根据需要勾选面板访问权限(如
Analytics、Management等)dba_role可以先给小一点Analytics(Discover)、Management(开发工具)k8s_role和sre_role给 All 即可

Kibana 权限 -
点击 创建角色
创建用户并绑定角色
- Stack Management → 安全 → 用户
- 点击 创建用户
- 填写用户名(如
xixi)、全名填中文、密码(如passwd)、确认密码 - 角色分配:勾选
dba_role - 点击 创建用户
| 用户名 | 全名 | 电子邮件地址 | 角色 |
|---|---|---|---|
| xixi | 嘻嘻 | xixi@qq.com | dba_role |
| haha | 哈哈 | haha@qq.com | k8s_role |
| hehe | 呵呵 | hehe@qq.com | sre_role |
验证权限隔离
1)准备测试数据(elastic 超级管理员)# 登录 elastic → 管理 → 开发工具,三索引各写入一条初始数据
POST kpyun-dba/_doc{ "name": "孙悟空", "hobby": ["蟠桃","仙丹","紫霞仙子"]}
POST kpyun-k8s/_doc{ "name": "猪八戒", "hobby": ["吃","睡","高老庄"]}
POST kpyun-sre/_doc{ "name": "沙和尚", "hobby": ["大师兄,师傅被妖怪抓走了","二师兄,师傅被妖怪抓走了"]}
2)测试 xixi(dba_role:仅 kpyun-dba* → create_index / create_doc / index / read)# 退出 elastic → 登录 xixi → 管理 → 开发工具
# ✅️ 写入属于自己的索引(有 create_doc)POST kpyun-dba/_doc{ "name": "唐僧", "hobby": ["念经","拜佛","女儿国"]}
# ❌️ 写入 kpyun-k8s(无权限)POST kpyun-k8s/_doc{ "name": "白骨精", "hobby": ["幻术","离间","吃人"]}# security_exception: unauthorized 403
# ✅️ 读取自己的索引(有 read)GET kpyun-dba/_search
# ❌️ 删除文档(无 delete 权限)DELETE kpyun-dba/_doc/<文档ID># security_exception: unauthorized 403
3)测试 haha(k8s_role:kpyun-k8s* → all)# 退出 xixi → 登录 haha → 开发工具
# ❌️ 读别人的索引 403GET kpyun-dba/_search# security_exception: unauthorized
# ✅️ 写入自己的索引(有 all)POST kpyun-k8s/_doc{ "name": "白骨精", "hobby": ["幻术","离间","吃唐僧肉"]}
# ✅️ 读取自己的索引GET kpyun-k8s/_search
# ✅️ 删除自己的文档(有 all 包含 delete)DELETE kpyun-k8s/_doc/<文档ID>
# hehe(sre_role)同理,只对 kpyun-sre* 有 all 权限,操作同上,不再复述xixi只能碰kpyun-dba*,读/写其他索引一律 403haha/hehe给了all权限,读/写/删都不受限,但只能操作自己的索引- dba_role 额外补了
index→ Logstash 写入时自动创建日切索引 - dba_role 故意不给
delete——防止误删数据
修改角色权限:如果要让某个角色能访问所有索引,不用重新创建,直接编辑角色:
- 切回 elastic 账号
- Stack Management → 安全 → 角色 → 点击对应角色(如
dba_role) - 索引模式从
kpyun-dba*改为*→ 该角色即可访问全部索引 - 权限级别按需调整(
read或all)
课后练习
-
Logstash 监听 6666/7777/8888 端口,用 3 个不同用户写入 3 个索引
- kpyun-dba-%{+yyyy-MM-dd}
- kpyun-k8s-%{+yyyy-MM-dd}
- kpyun-sre-%{+yyyy-MM-dd}
-
发送测试数据并验证
解法一 — if 语句实现
1)编写配置文件 — 单个 Logstash 进程内用 if/else 分发root@Elk03 ~# vim /etc/logstash/conf.d/if-solution.yamlinput { tcp { port => "6666" type => "dba"# 端口 6666 的事件打上 type=dba }
tcp { port => "7777" type => "k8s" }
tcp { port => "8888" type => "sre" }}
filter { mutate { remove_field => [ "@version","agent","log","ecs","tags","input" ] }}
output { if [type] == "dba" { elasticsearch { hosts => ["10.0.0.6:9200","10.0.0.7:9200","10.0.0.8:9200"] index => "kpyun-dba-%{+yyyy-MM-dd}" user => "xixi" password => "passwd" } } else if [type] == "k8s" { elasticsearch { hosts => ["10.0.0.6:9200","10.0.0.7:9200","10.0.0.8:9200"] index => "kpyun-k8s-%{+yyyy-MM-dd}" user => "haha" password => "passwd" } } else { elasticsearch { hosts => ["10.0.0.6:9200","10.0.0.7:9200","10.0.0.8:9200"] index => "kpyun-sre-%{+yyyy-MM-dd}" user => "hehe" password => "passwd" } }}# ⚠️ if/sle 判断的是 type 字段,type 在 input 中定义# output 中指定不同的 user/password → 用不同身份写入 ES
2)启动 Logstashroot@Elk03 ~# logstash -rf /etc/logstash/conf.d/if-solution.yaml --config.test_and_exitroot@Elk03 ~# logstash -rf /etc/logstash/conf.d/if-solution.yaml
3)发送测试数据root@Elk01 ~# echo 'dba-data-111' | nc -w 1 10.0.0.8 6666root@Elk01 ~# echo 'k8s-data-222' | nc -w 1 10.0.0.8 7777root@Elk01 ~# echo 'sre-data-333' | nc -w 1 10.0.0.8 8888- ⚠️ nc 只是测试偷懒用的,生产环境应该用 Filebeat 多实例
- 每个实例采集对应日志 → output 到不同的 Logstash 端口
- 详细配置见上方「Logstash 的 if 多分支语句」一节(3 个配置文件 + 独立 —path.data)
4)验证 — 三个索引各收到一条数据root@Elk01 ~# curl -s -u elastic:passwd '10.0.0.6:9200/_cat/indices/kpyun-dba-20*,kpyun-k8s-20*,kpyun-sre-20*?v'health status index docs.countgreen open kpyun-dba-2026-08-01 1green open kpyun-k8s-2026-08-01 1green open kpyun-sre-2026-08-01 1✅️ 三个端口 → 三个索引 → 三个用户,完美隔离if 多分支 vs Pipeline:更推荐 Pipeline
| if 多分支 | Pipeline | |
|---|---|---|
| 配置文件 | 一个文件全部搞定 | 每个通道独立配置 |
| 单点风险 | 一个语法错 → 整条管道挂 | 一个通道挂 → 其他照常运行 |
| 启动方式 | logstash -rf xxx.conf | logstash -r 热加载 |
| 生产 | 简单场景凑合用 | ==推荐标准做法== |
💡 Pipeline 把一条大管道拆成多条小管道,隔离故障、独立维护,配合 pipelines.yml 管理,==一个进程跑所有通道==,比多实例省内存
解法二 — Pipeline 实现
1)编写三个独立的配置文件root@Elk03 ~# vim /etc/logstash/conf.d/pipeline-dba.yamlinput { tcp { port => "6666" } }filter { mutate { remove_field => [ "@version","agent","log","ecs","tags","input" ] } }output { elasticsearch { hosts => ["10.0.0.6:9200","10.0.0.7:9200","10.0.0.8:9200"] index => "kpyun-dba-%{+yyyy-MM-dd}" user => "xixi" password => "passwd" }}# 同理创建 k8s(7777端口,haha用户)和 sre(8888端口,hehe用户)root@Elk03 ~# vim /etc/logstash/conf.d/pipeline-k8s.yamlinput { tcp { port => "7777" } }filter { mutate { remove_field => [ "@version","agent","log","ecs","tags","input" ] } }output { elasticsearch { hosts => ["10.0.0.6:9200","10.0.0.7:9200","10.0.0.8:9200"] index => "kpyun-k8s-%{+yyyy-MM-dd}" user => "haha" password => "passwd" }}root@Elk03 ~# vim /etc/logstash/conf.d/pipeline-sre.yamlinput { tcp { port => "8888" } }filter { mutate { remove_field => [ "@version","agent","log","ecs","tags","input" ] } }output { elasticsearch { hosts => ["10.0.0.6:9200","10.0.0.7:9200","10.0.0.8:9200"] index => "kpyun-sre-%{+yyyy-MM-dd}" user => "hehe" password => "passwd" }}
2)修改 pipelines.ymlroot@Elk03 ~# vim /etc/logstash/pipelines.yml# - pipeline.id: main# path.config: "/etc/logstash/conf.d/*.conf"
- pipeline.id: xixi path.config: "/etc/logstash/conf.d/pipeline-dba.yaml"
- pipeline.id: haha path.config: "/etc/logstash/conf.d/pipeline-k8s.yaml"
- pipeline.id: hehe path.config: "/etc/logstash/conf.d/pipeline-sre.yaml"
3)启动root@Elk03 ~# logstash -rPipelines running {:count=>3, :running_pipelines=>[:xixi, :haha, :hehe]# 一个进程跑三个 pipeline ✅️文章分享
如果这篇文章对你有帮助,欢迎分享给更多人!















