从零构建 Logstash 日志监控闭环:Fluent Bit 旁路采集与 ECS 规范化落地#
引言:谁来监控日志中枢?#
在现代云原生与安全大数据架构中,Logstash 往往承担着多源异构日志汇聚、清洗、规整并输入 OpenSearch 的核心角色。然而,“灯下黑” 是许多平台架构师面临的共同痛点——我们用 Logstash 监控一切,却唯独没有监控 Logstash 自身。
当某个数据管道因为不合理的正则发生死锁、当上游网络抖动触发反压(Backpressure)、甚至当组件自身因 OOM 崩溃时,监控往往陷入盲区。
本方案旨在从零到一构建一套生产级的旁路监控闭环:通过轻量级 Agent 旁路采集 Logstash 自身的 JSON 日志,利用 OpenSearch Ingest Pipeline 将其对齐至工业级 ECS(Elastic Common Schema)规范,并最终落地于 OpenSearch 原生的 Data Stream(数据流)架构中。
一、 轻量级旁路采集#
让 Logstash 将自身的日志发给自己,在逻辑上是一个 “循环论证” 的悖论:一旦 Logstash 假死,它就永远无法上报自己的故障。因此,生产环境必须坚持 旁路轻量化采集原则。
1.1 架构流程图#
flowchart LR
subgraph LogSource["📦 日志来源"]
LS["Logstash 实例\n/var/log/logstash/logstash-json.log\n/var/log/logstash/logstash-slowlog.json"]
end
subgraph Collection["🚚 旁路采集层 (Fluent Bit)"]
FB["Fluent Bit Agent\nTail Input → Filter → OpenSearch Output"]
end
subgraph Storage["💾 存储与处理层 (OpenSearch)"]
direction TB
DS["Data Stream\nlog_siem_logstash_log"]
Tpl["Index Template\n绑定 Ingest Pipeline"]
Pipe["Ingest Pipeline\n字段重映射 / ECS 对齐"]
BI["Back-Index\n.ds-log_siem_logstash_log-000001"]
end
subgraph Verify["🔍 验证层"]
OD["OpenSearch Dashboards\nDiscover 查询验证"]
end
LS -->|" Tail 文件尾随 "| FB
FB -->|" HTTP Bulk Create "| DS
DS -->|" 自动匹配 "| Tpl
Tpl -->|" 触发 "| Pipe
Pipe --> BI
BI --> OD1.2 旁路采集核心原则#
| 原则 | 说明 |
|---|---|
| 零性能污染 | Logstash 本身是重度消耗 CPU 和内存的 Java 进程。旁路采集必须采用 C 语言编写、高并发、极低内存占用的轻量级 Agent(如 Fluent Bit)进行日志监听,避免引入额外的资源竞争。 |
| 免多行解析 | 传统的 Java 堆栈日志(Stack Trace)由于多行断裂,需要 Agent 进行复杂的正则合并,这在生产环境中极易漏判且消耗 CPU。本方案要求 Logstash 必须原生以 纯 JSON 格式 向磁盘滚动输出标准日志与慢日志,实现开箱即用的结构化采集。 |
| 断点续传 | Fluent Bit 的 Tail 插件通过 SQLite 数据库将文件偏移量持久化到本地磁盘,Agent 重启后能够从上次中断的位置继续采集,保障日志零丢失。 |
1.3 Fluent Bit 数据流处理#
flowchart LR
subgraph Input["📥 输入 (Input)"]
Tail["Tail Plugin\n监控文件偏移 (Checkpoint)"]
end
subgraph Filter["🔧 过滤 (Filter)"]
Modify["Modify\n注入静态元数据"]
Nest["Nest\n字段嵌套整理"]
end
subgraph Buffer["💾 缓冲与背压"]
Mem["内存缓冲区\n限制 50MB"]
Disk["磁盘辅助队列\n防止内存溢出"]
end
subgraph Output["📤 输出 (Output)"]
OS["OpenSearch Plugin\nBulk API 异步批量写入"]
end
Tail --> Modify --> Nest --> Mem
Mem -->|" 内存写满时 "| Disk
Mem & Disk --> Output二、 Logstash 日志输出配置#
2.1 为什么需要 JSON 格式日志#
Logstash 默认输出的日志为纯文本格式,包含多行堆栈跟踪信息,采集时需要复杂的正则表达式进行行合并与字段提取。这种方式存在以下问题:
- CPU 开销大:正则匹配在多行场景下消耗大量计算资源
- 容易漏判:复杂的堆栈格式可能导致正则匹配失败
- 维护成本高:每次 Logstash 版本升级可能改变日志格式
本方案要求 Logstash 原生输出 结构化 JSON 日志,每条日志为一行完整的 JSON 对象,Fluent Bit 可以直接解析,无需额外处理。
2.2 Logstash JSON 日志配置#
在 logstash.yml 中配置日志输出格式:
# 日志路径
path.logs: /var/log/logstash
# 日志级别(生产环境建议 info)
log.level: info
# 日志格式:json(关键配置)
log.format: json配置生效后,Logstash 将生成以下日志文件:
| 文件 | 说明 | 内容 |
|---|---|---|
/var/log/logstash/logstash-json.log |
标准运行时日志 | 管道状态、插件事件、JVM 信息、错误堆栈等 |
/var/log/logstash/logstash-slowlog.json |
慢日志 | 处理耗时超过阈值的事件记录 |
2.3 JSON 日志样例#
一条典型的 Logstash JSON 日志如下:
{
"level": "ERROR",
"loggerName": "logstash.codecs.json",
"timeMillis": 1782311195698,
"thread": "76f91f03e1c5829e7064a5b55f616fa1a4ffb616712ab0ad025e8909025df8b4-beatsHandler[T#3]",
"pipeline.id": "winlogbeat",
"plugin.id": "76f91f03e1c5829e7064a5b55f616fa1a4ffb616712ab0ad025e8909025df8b4",
"logEvent": {
"message": "JSON parse error, original data now in message field",
"message": "Unrecognized token 'An': was expecting (JSON String, Number, Array, Object or token 'null', 'true' or 'false')\n at [Source: REDACTED (`StreamReadFeature.INCLUDE_SOURCE_IN_LOCATION` disabled); line: 1, column: 1]",
"exception": "LogStash::Json::ParserError",
"data": "An account was logged off.\n\nSubject:\n\tSecurity ID:\t\tN-1-5-20\n\tAccount Name:\t\tXXXXXX$\n\tAccount Domain:\t\tXXXXXXXX\n\tLogon ID:\t\t0x377640AE6\n\nLogon Type:\t\t\t3\n\nThis event is generated when a logon session is destroyed. It may be positively correlated with a logon event using the Logon ID value. Logon IDs are only unique between reboots on the same computer."
}
}三、 Fluent Bit 采集配置详解#
在前一节中,我们完成了 Logstash 的 JSON 日志配置。接下来,需要部署 Fluent Bit Agent 来旁路采集这些日志。
3.1 安装 Fluent Bit#
# 1. 添加 Fluent Bit 官方源
sudo sh -c 'curl https://packages.fluentbit.io/fluentbit.key | gpg --dearmor > /usr/share/keyrings/fluentbit-keyring.gpg'
codename=$(grep -oP '(?<=VERSION_CODENAME=).*' /etc/os-release 2>/dev/null || lsb_release -cs 2>/dev/null)
echo "deb [signed-by=/usr/share/keyrings/fluentbit-keyring.gpg] https://packages.fluentbit.io/ubuntu/$codename $codename main" | sudo tee /etc/apt/sources.list.d/fluent-bit.list
# 2. 安装 Fluent Bit(推荐 5.0.7)
sudo apt-get update
sudo apt-get install fluent-bit=5.0.7
# 3. 验证安装
/opt/fluent-bit/bin/fluent-bit --version3.2 创建sqlite数据库路径#
sudo mkdir -p /var/lib/fluent-bit/3.3 修改 /lib/systemd/system/fluent-bit.service#
[Unit]
Description=Fluent Bit
Documentation=https://docs.fluentbit.io/manual/
Requires=network.target
After=network.target
[Service]
Type=simple
EnvironmentFile=-/etc/sysconfig/fluent-bit
EnvironmentFile=-/etc/default/fluent-bit
Environment=FB_VERSION=5.0.7 # 注入Fluent Bit 版本环境变量
Environment=SERVICE_NAME=xxxxx # 注入服务名称环境变量
Environment=SERVICE_IP=xxxxxx # 注入服务 IP 环境变量
ExecStart=/opt/fluent-bit/bin/fluent-bit -c /etc/fluent-bit/fluent-bit-logstash.yaml # 指定 Logstash日志采集配置文件
Restart=always
[Install]
WantedBy=multi-user.target3.4 Logstash日志采集配置文件#
创建 /etc/fluent-bit/fluent-bit-logstash.yaml:
---
service:
http_server: "on"
http_listen: "0.0.0.0"
http_port: 2020
Health_Check: "on"
flush: 1
daemon: off
log_level: info
parsers_file: parsers.conf
plugins_file: plugins.conf
pipeline:
inputs:
- name: tail
path: /var/log/logstash/logstash-json.log
tag: test
db: /var/lib/fluent-bit/logstash_tail.db
path_key: logPath
key: message
filters:
- name: modify
match: '*'
Add:
- agent_name fluent-bit
- agent_version ${FB_VERSION}
- agent_host ${HOSTNAME}
- service_ip ${SERVICE_IP}
- service_name ${SERVICE_NAME}
- service_host ${HOSTNAME}
Rename:
- date @timestamp
- name: nest
match: '*'
operation: nest
wildcard:
- agent_name
- agent_version
- agent_host
nest_under: agent
remove_prefix: agent_
- name: nest
match: '*'
operation: nest
wildcard:
- service_ip
- service_name
- service_host
nest_under: service
remove_prefix: service_
- name: modify
match: '*'
outputs:
- name: opensearch
match: '*'
host: 192.168.100.57
port: 9200
index: log_siem_logstash_log
write_operation: create
suppress_type_name: on
generate_id: on
http_user: xxxxx
http_passwd: '*******'
tls: on
tls.verify: off四、 Logstash日志ECS标准化#
完成 Fluent Bit 采集配置后,接下来需要将采集到的日志对齐至 ECS(Elastic Common Schema)规范,以便于后续的查询、分析和可视化。
4.1 组件化模板#


- Index settings & Advanced settings
{
"index.number_of_replicas": "0",
"index.default_pipeline": "logstash_log_default",
"index.refresh_interval": "5s",
"index.number_of_shards": "1"
}| 参数 | 值 | 说明 |
|---|---|---|
number_of_shards |
1 | 监控日志单分片即可满足查询需求 |
number_of_replicas |
0 | 单节点的话需要设置副本为0 |
refresh_interval |
5s | 降低刷新频率,提升写入性能 |
index.default_pipeline |
logstash_log_default | 绑定 ingest pipeline |
- Index mappings
{
"properties": {
"logstash": {
"type": "object",
"properties": {
"log": {
"type": "object",
"properties": {
"plugin_id": {
"type": "keyword"
},
"log_path": {
"type": "keyword"
},
"module": {
"type": "keyword"
},
"pipeline_id": {
"type": "keyword"
},
"log_event": {
"type": "flat_object"
},
"thread": {
"fields": {
"text": {
"type": "text",
"fields": {
"keyword": {
"ignore_above": 256,
"type": "keyword"
}
}
}
},
"type": "keyword"
}
}
}
}
},
"agent": {
"type": "object",
"properties": {
"host": {
"type": "keyword"
},
"name": {
"type": "keyword"
},
"version": {
"type": "keyword"
}
}
},
"@timestamp": {
"type": "date"
},
"ecs": {
"type": "object",
"properties": {
"version": {
"type": "keyword"
}
}
},
"log": {
"type": "object",
"properties": {
"file": {
"type": "object",
"properties": {
"path": {
"type": "keyword"
}
}
},
"level": {
"type": "keyword"
}
}
},
"data_stream": {
"type": "object",
"properties": {
"namespace": {
"value": "siem-mon",
"type": "constant_keyword"
},
"type": {
"value": "logs",
"type": "constant_keyword"
},
"dataset": {
"value": "logstash.log",
"type": "constant_keyword"
}
}
},
"host": {
"type": "object",
"properties": {
"ip": {
"type": "ip"
},
"name": {
"type": "keyword"
}
}
},
"message": {
"type": "match_only_text"
},
"event": {
"type": "object",
"properties": {
"ingested": {
"type": "date"
},
"original": {
"type": "keyword"
},
"kind": {
"type": "keyword"
},
"created": {
"type": "date"
},
"category": {
"type": "keyword"
},
"type": {
"type": "keyword"
}
}
}
}
}4.2 OpenSearch Ingest Pipeline#
创建 Logstash 日志数据处理 Ingest Pipeline
PUT _ingest/pipeline/logstash_log_default
{
"description": "Pipeline for parsing logstash node logs",
"processors": [
{
"set": {
"tag": "set_ecs_version_f5923549",
"field": "ecs.version",
"value": "8.17.0"
}
},
{
"set": {
"tag": "set_event_category",
"field": "event.category",
"value": "process"
}
},
{
"set": {
"field": "event.kind",
"value": "event"
}
},
{
"set": {
"tag": "set_event_ingested",
"field": "event.ingested",
"value": "{{_ingest.timestamp}}"
}
},
{
"rename": {
"tag": "rename_message_to_event_original",
"field": "message",
"target_field": "event.original",
"override_target": true,
"ignore_missing": true
}
},
{
"json": {
"tag": "parse_logstash_log",
"field": "event.original",
"target_field": "logstash.log",
"if": "ctx.event?.original != null"
}
},
{
"convert": {
"tag": "convert_logstash_log_timeMillis_to_string",
"field": "logstash.log.timeMillis",
"type": "string"
}
},
{
"date": {
"tag": "parse_logstash_log_timeMillis",
"field": "logstash.log.timeMillis",
"formats": [
"UNIX_MS"
],
"target_field": "@timestamp"
}
},
{
"rename": {
"tag": "rename_logstash_log_loggerName_to_logstash_log_module",
"field": "logstash.log.loggerName",
"target_field": "logstash.log.module",
"ignore_missing": true
}
},
{
"rename": {
"tag": "rename_logstash_log_logEvent_message_to_message",
"field": "logstash.log.logEvent.message",
"target_field": "message",
"ignore_missing": true
}
},
{
"rename": {
"tag": "rename_logstash_log_logEvent_to_logstash_log_log_event",
"field": "logstash.log.logEvent",
"target_field": "logstash.log.log_event",
"ignore_missing": true
}
},
{
"rename": {
"tag": "rename_logstash_log_level_to_log_level",
"field": "logstash.log.level",
"target_field": "log.level",
"ignore_missing": true
}
},
{
"rename": {
"tag": "rename_logPath_to_log_file_path",
"field": "logPath",
"target_field": "log.file.path",
"ignore_missing": true
}
},
{
"script": {
"tag": "convert_logstash_log_log_event_action_to_string",
"description": "Convert logstash.log.log_event.action elements to string.",
"if": "ctx?.logstash?.log?.log_event?.action instanceof List",
"lang": "painless",
"source": """def items = [];
ctx.logstash.log.log_event.action.forEach(v -> {
items.add(v.toString());
});
ctx.logstash.log.log_event.action = items;
"""
}
},
{
"script": {
"tag": "set_event_type",
"lang": "painless",
"source": """def errorLevels = ["ERROR", "FATAL"]; if (ctx?.log?.level != null) {
if (errorLevels.contains(ctx.log.level)) {
ctx.event.type = ["error"];
} else {
ctx.event.type = ["info"];
}
}"""
}
},
{
"copy": {
"tag": "copy_timestamp_to_event_created",
"source_field": "@timestamp",
"target_field": "event.created",
"override_target": true,
"ignore_missing": true
}
},
{
"copy": {
"tag": "copy_service_ip_to_host_ip",
"source_field": "service.ip",
"target_field": "host.ip",
"override_target": true,
"ignore_missing": true
}
},
{
"copy": {
"tag": "copy_agent_host_to_host_name",
"source_field": "agent.host",
"target_field": "host.name",
"override_target": true,
"ignore_missing": true
}
}
],
"on_failure": [
{
"append": {
"field": "error.message",
"value": "Processor \"{{{ _ingest.on_failure_processor_type }}}\" with tag \"{{{ _ingest.on_failure_processor_tag }}}\" in pipeline \"{{{ _ingest.on_failure_pipeline }}}\" failed with message \"{{{ _ingest.on_failure_message }}}\""
}
},
{
"set": {
"field": "event.kind",
"value": "pipeline_error"
}
},
{
"append": {
"field": "tags",
"value": "preserve_original_event",
"allow_duplicates": false
}
}
]
}4.3 索引模板#

五、Logstash日志采集与验证#
完成 ECS 标准化配置后,启动 Fluent Bit 采集服务,日志将根据已有索引模板自动创建 DataStream,并调用 Ingest Pipeline 进行数据处理。
flowchart LR
subgraph Raw["原始 JSON 日志"]
A["@timestamp\nlevel\nloggerName\nmessage\nlogEvent\ntimeMillis\nthread"]
end
subgraph FluentBit["Fluent Bit Filter"]
B["Modify: 注入 agent.name\nagent.version"]
end
subgraph Pipeline["OpenSearch Ingest Pipeline"]
C["Rename: 字段重映射"]
D["Convert: 类型转换"]
E["Lowercase: 级别标准化"]
end
subgraph ECS["ECS 标准化输出 DataStream"]
F["@timestamp → date\nlog.level → keyword\nlog.logger → keyword\nmessage → text"]
end
A --> B --> C --> D --> E --> F5.1 启动日志采集服务#
# 启动 Fluent Bit
sudo systemctl start fluent-bit
# 查看状态
sudo systemctl status fluent-bit
# 查看日志
sudo journalctl -u fluent-bit -f --since "10 min ago"
# 健康检查
curl http://127.0.0.1:2020/5.2 查看 DataStream#
如果一切正常,此时应该会自动创建 DataStream

5.3 在 Dashboards 中创建 Index Pattern#



5.4 在 Dashboards 中 Discover 查看日志数据#

六、 为什么选择 Data Stream 而非传统 Index?#
监控日志作为典型的时序数据(Time-series Data),传统的 “每日新建索引” 模式维护成本高,且不利于字段类型的动态演进。本方案完全拥抱 OpenSearch 原生的 Data Stream(数据流) 特性。
Data Stream 简化并规范了这一流程,强制实施最适合时序数据的架构设计:**专为 append-only 数据设计,并确保每个文档包含时间戳字段 **。
Data Stream 内部结构#
一个 Data Stream 在内部由多个 backing index 组成:
- Search 请求:自动路由到所有 backing index,实现统一查询
- Indexing 请求:自动路由到最新的 write backing index
- ISM 策略:可自动化处理索引 rollover 或删除
传统 Index 与 Data Stream 核心对比#
| 对比维度 | 传统 Index + Alias | Data Stream |
|---|---|---|
| 写入模式 | 可更新/删除任意文档,存在误操作风险 | 强制 append-only,防止误修改历史数据 |
| 内部结构 | 单索引或多个独立索引 | 自动管理多个 backing index,对外统一命名 |
| 查询路由 | 需使用通配符 logs-* 匹配多个索引 |
查询自动路由到所有 backing index |
| 时间戳要求 | 无强制要求 | 强制每个文档包含 @timestamp 字段 |
| 数据安全性 | 可能意外更新/删除历史数据 | append-only 机制保障历史数据不可变 |
Logstash 日志场景的适配性#
Logstash 日志是典型的时序数据,其特征与 Data Stream 的设计目标完美契合:
- 只增不改:日志产生后无需更新或删除
- 持续增长:日志数据量随时间快速增长
- 按时间查询:90% 以上的查询基于时间范围过滤
- 需要自动化运维:希望自动处理索引轮转和旧数据清理
七、 ISM 自动化生命周期(Index State Management)#
为应对海量监控日志带来的容量压力,设计符合工业规范的滚动退役策略。 具体策略可根据自己的实际业务场景设置即可,下面为举例说明。
flowchart LR
subgraph Hot["🔥 Hot 阶段 (1-3天)"]
H1["高频写入与检索"]
H2["NVMe SSD 存储"]
H3["1 个副本"]
H4["Rollover: 50GB 或 1天"]
end
subgraph Warm["⚡ Warm 阶段 (4-7天)"]
W1["只读索引"]
W2["Force Merge 合并段"]
W3["普通 HDD 存储"]
end
subgraph Cold["❄️ Cold 阶段 (8-30天)"]
C1["副本数降为 0"]
C2["深度压缩"]
C3["低频安全审计"]
end
subgraph Delete["🗑️ Delete 阶段 (>30天)"]
D1["自动物理删除"]
end
Hot --> Warm --> Cold --> Delete7.1 创建ISM#
PUT _plugins/_ism/policies/logstash_log_policy
{
"policy": {
"policy_id": "logstash_log_policy",
"description": "Logstash 监控日志生命周期管理",
"default_state": "hot",
"states": [
{
"name": "hot",
"actions": [
{
"retry": {
"count": 3,
"backoff": "exponential",
"delay": "1m"
},
"allocation": {
"require": {
"temp": "hot"
},
"include": {},
"exclude": {},
"wait_for": false
}
},
{
"retry": {
"count": 3,
"backoff": "exponential",
"delay": "1m"
},
"rollover": {
"min_size": "50gb",
"min_index_age": "1d",
"copy_alias": false
}
}
],
"transitions": [
{
"state_name": "warm",
"conditions": {
"min_index_age": "1d"
}
}
]
},
{
"name": "warm",
"actions": [
{
"retry": {
"count": 3,
"backoff": "exponential",
"delay": "1m"
},
"allocation": {
"require": {
"temp": "warm"
},
"include": {},
"exclude": {},
"wait_for": false
}
},
{
"retry": {
"count": 3,
"backoff": "exponential",
"delay": "1m"
},
"force_merge": {
"max_num_segments": 1
}
}
],
"transitions": [
{
"state_name": "cold",
"conditions": {
"min_index_age": "7d"
}
}
]
},
{
"name": "cold",
"actions": [
{
"retry": {
"count": 3,
"backoff": "exponential",
"delay": "1m"
},
"allocation": {
"require": {
"temp": "cold"
},
"include": {},
"exclude": {},
"wait_for": false
}
},
{
"retry": {
"count": 3,
"backoff": "exponential",
"delay": "1m"
},
"replica_count": {
"number_of_replicas": 0
}
}
],
"transitions": [
{
"state_name": "delete",
"conditions": {
"min_index_age": "30d"
}
}
]
},
{
"name": "delete",
"actions": [
{
"delete": {}
}
]
}
]
}
}7.2 绑定索引模板#

#
结语#
通过本篇方案,我们完成了 Logstash 日志监控体系中最核心的 “基础设施铺设”。至此,Logstash 的自身日志已经源源不断地以标准 ECS 规范、极低的资源损耗落地于 OpenSearch 数据流中。
这套方案的核心价值在于:
- 旁路轻量化:采用 Fluent Bit 替代 Logstash 自采集,避免循环依赖和性能污染
- 结构化采集:强制 Logstash 输出
JSON日志,彻底告别多行正则解析- 标准化存储:通过 Ingest Pipeline 对齐 ECS 规范,为后续的可视化、告警、审计打下坚实基础
- 自动化运维:Data Stream + ISM 实现索引全生命周期自动化管理,降低运维负担
下一步行动建议:
在后续实践中,我们可以基于这套标准的 ECS 数据流,进一步构建自定义告警规则、日志分析 Dashboard,以及自动化故障诊断流程。