
先说一句大实话分布式系统的日志监控难点从来不是“日志多”而是“日志散”。一个请求从前端打到网关再到三四个微服务最后落库全程产生的日志可能分布在几十台机器上。你 troubleshooting 的时候最怕的不是没日志而是日志在但你找不到、对不上、查不动。Logstash 在这个场景里干的活就是把散落在各处的日志收进来、洗干净、送进统一的存储与检索平台。这篇文章不讲空理论只说我这些年把 Logstash 用在分布式日志监控上的实际做法包括管道怎么设计、部署怎么分工、参数怎么调、自定义插件怎么集成以及那些文档里不会写的坑。如果你是刚接触 Logstash 的运维或后端开发或者你正在搭一套日志平台但被 grok 正则和高 CPU 占用折磨过这篇文章适合你。我会尽量把每一步的“为什么”也讲清楚而不是只丢一堆配置让你抄。1. 先搞清楚需求分布式日志监控到底卡在哪里1.1 日志分散不是最大的问题看不到才是很多人一开始搭建日志系统第一反应是“把日志集中起来”。集中当然是对的但如果只做到“集中”这一步后面一定会返工。分布式系统里的日志监控真正的诉求是三件事能搜、能关联、能告警。能搜指的是你不需要 ssh 到每一台机器上 grep。能关联指的是你能用一个 traceId 把一次请求在多个服务里的日志串起来。能告警指的是当天级的大盘出现异常时系统能主动通知你而不是等用户投诉。这三件事里Logstash 负责的是前两件事的地基采集和解析。它把非结构化的日志文本变成结构化的字段写入 Elasticsearch后面 Kibana 的搜索、聚合、可视化都依赖这些字段。如果你在 Logstash 阶段没有把 traceId、服务名、日志级别、耗时这些关键字段解析出来后面做关联分析和告警规则就是空中楼阁。所以我的建议是动手配 Logstash 之前先花半天时间梳理你系统的日志规范。哪怕只是规定“所有服务必须输出 JSON 格式日志必须包含 traceId 字段”后续的解析工作都能少掉一半。这件事越早做收益越大。1.2 方案选型为什么最终落在 Logstash 上日志采集与传输的工具有不少常见的就有 Filebeat、Fluentd、Fluent Bit、Vector再加上 Logstash。“既然有 Filebeat为什么还要用 Logstash”这是我被问过很多次的问题。我的回答是选型要看你的核心瓶颈在哪。如果是单纯的“把日志文件传到 Kafka 或 Elasticsearch”Filebeat 更轻、更省资源当然是首选。但分布式系统的日志监控瓶颈几乎都在“解析加工”这层多种日志格式、多行堆栈、字段补充、时间戳统一、敏感信息脱敏。这些逻辑用 Filebeat 的 processors 写起来很别扭用 Fluentd 的插件生态也能做但遇到复杂的 grok 解析和条件分支时Logstash 的 filter 阶段要灵活得多。Logstash 的优势是“管道模型”非常直观input 收数据filter 改数据output 送数据。每个阶段都是插件插件可以组合、可以嵌套条件。这个模型的好处是你的解析逻辑是一段可读性很高的配置而不是散落在代码里的正则和一串串 sed 命令。对于需要长期维护的日志平台来说可读性就是可维护性。当然Logstash 的缺点是重JVM 常驻内存占用天然比 Go 写的 Filebeat 高。所以生产环境里我通常的做法是“Filebeat Logstash”分层Filebeat 跑在每台业务机器上做轻量采集Logstash 集中做解析和富化这样既控制资源开销又保留了解析的灵活性。2. 管道设计input、filter、output 怎么搭才不白干2.1 input 阶段从文件到网络先把数据收进来Logstash 的 input 插件决定“日志从哪里来”。分布式系统里最常见的两种来源是 beats 和 kafka。用 filebeat 作为采集端时Logstash 这边用 beats input 接收如果中间有 Kafka 缓冲就用 kafka input。input { beats { port 5044 client_inactivity_timeout 300 } kafka { bootstrap_servers kafka-1:9092,kafka-2:9092 topics [app-logs, access-logs] group_id logstash-main consumer_threads 4 codec json } }这里有几个细节值得注意。client_inactivity_timeout默认是 60 秒我在内网环境会调大到 300 秒避免偶发网络抖动导致连接频繁断开重建。Kafka input 的consumer_threads不是越大越好它和 Logstash 的 pipeline workers 要配合一般先保持和 workers 数量一致后续再根据消费延迟逐步调。还有编码问题。很多团队一开始用codec json接收结构化日志这没问题。但如果你有部分服务输出的是纯文本或混合格式不要试图用同一个 input 接收所有数据再靠 filter 猜格式。在 input 端就用 topic 或者 beats 的 tags 区分数据流是省力的第一步。2.2 filter 阶段日志解析的关键动作filter 是整个 Logstash 配置里最花心思的部分。我常用的几个插件是 dissect、grok、mutate、date、user_agent、geoip以及几个生产环境必须的辅助插件。先说说 dissect 和 grok 怎么选。很多人一上来就写 grok然后被正则折磨得怀疑人生。实际上能用 dissect 就不用 grok。dissect 是基于分隔符的固定模式抽取性能远高于 grok适合格式规整的日志。比如下面这条日志2024-06-01 12:00:01 ERROR [order-service] traceIdabc123 userId9999 调用库存服务超时用 dissect 是这个写法filter { if [service] order-service { dissect { mapping { message %{ts} %{level} [%{service}] traceId%{traceId} userId%{userId} %{msg} } } } }如果日志格式有变化或者需要匹配更复杂的模式再上 grok。比如要匹配一个更灵活的前缀grok 可以提供容错和命名捕获filter { grok { match { message %{TIMESTAMP_ISO8601:ts} %{LOGLEVEL:level} \[%{WORD:service}\] traceId%{NOTSPACE:traceId} userId%{NUMBER:userId:int} %{GREEDYDATA:msg} } overwrite true } }注意overwrite true这个参数很关键。如果 pipeline 里有多个 grok 或者日志通过 topic 混入不设置 overwrite 可能导致旧字段残留。还有tag_on_failure解析失败时 Logstash 会默认打上_grokparsefailure标签这个标签要保留后面排查解析覆盖率时很有用。解析完字段后mutate 几乎是必用的。我的习惯是统一时间格式、类型转换、清理无用字段filter { mutate { convert { userId integer } remove_field [path, host, beat, log] } date { match [ts, ISO8601] target timestamp } }date插件的作用是用日志里的业务时间替换 Logstash 接收到数据的系统时间。分布式环境里网络延迟和批处理都会让采集时间晚于实际记录时间如果不做这个替换你在 Kibana 里按时间排序会发现日志顺序是乱的。除了解析filter 阶段还适合做富化。比如根据客户端 IP 追加归属地和运营商根据 userAgent 解析浏览器和设备信息。这些富化字段在后续的访问日志分析里非常有用。但我要提醒一句geoip 这类查库操作比较耗 CPU如果日志量很大建议只在部分数据流上开启或者用 Logstash 自身的geoip_database缓存配置来加速。2.3 output 阶段写入 Elasticsearch 的细节output 通常就是写 Elasticsearch但这里有几个容易忽略的配置output { elasticsearch { hosts [http://es-1:9200, http://es-2:9200] index app-logs-%{yyyy.MM.dd} user elastic password ${ES_PWD} manage_template false template /etc/logstash/templates/app-logs-template.json template_name app-logs } }用index app-logs-%{yyyy.MM.dd}按天分索引是常见做法。如果你的日志量特别大也可以按小时分索引但要考虑到这样索引数量会膨胀集群分片管理压力变大。日志量在中等规模单日几十 GB 到几百 GB以内按天分索引最合适。manage_template false配合自定义 template 是生产上必须做的事情。默认的模板不会有针对性的字段映射比如你希望traceId走 keyworduserId走 integer默认映射猜不到。不自定义模板的后果是未来在 Kibana 里看到一堆.keyword后缀的重复字段索引变大查询变慢。还有连接 ES 的认证信息不要明文写在配置文件里。Logstash 支持从环境变量读取上面的${ES_PWD}就是这么个写法。配合 systemd 的 EnvironmentFile 或者 K8s 的 Secret可以做到密文管理。3. 分布式部署多节点采集与性能调优实录3.1 部署拓扑采集端、汇聚端怎么分工分布式系统的日志量一上来单机跑一个 Logstash 顶不住这就涉及到部署拓扑。我常用的架构是“边缘采集 中央汇聚”两层。边缘采集层用 Filebeat 或 Logstash 的轻量模式跑在应用机器上只做读取和转发不做复杂解析。中央汇聚层部署一组 Logstash 节点专门跑 filter 和 output负责把日志推给 Elasticsearch。这两层的分工背后是资源隔离的考虑边缘节点占用高会影响业务应用而集中解析需要相对高的 CPU 和内存独立部署可以单独扩容。汇聚层前面最好加一层 Kafka。业务日志先进 KafkaLogstash 从 Kafka 消费。这个设计在分布式系统里尤其重要因为采集端的生产速度和 ES 的写入能力存在天然的不匹配。比如业务高峰期日志产量是平时的十倍如果 Logstash 消费不了Kafka 可以帮你扛住积压等高峰过去再慢慢消费。如果没有 KafkaLogstash 内存队列满的时候会直接丢弃日志这个损失在排查问题时是无法接受的。拓扑如下图所示意不用画得太细重点是把角色分清楚应用机器Filebeat只负责读文件或收 syslog轻量转发。Kafka 集群日志消息的缓冲和削峰填谷。Logstash 汇聚层消费 Kafka执行 grok、dissect、mutate、date 等解析。Elasticsearch写入与检索。Kibana查询、图表、告警配置。这个拓扑在几十台到上千台机器的场景我都验证过。核心的思路是每一层只做一件自己擅长的事出了问题也好定位。3.2 内存、队列、并发这些参数到底怎么调Logstash 性能调优绕不开 pipeline 参数。先说一个常识Logstash 处理数据的模型是“输入 - 队列 - 工人”不是每来一条就立刻处理。它把数据放进内存队列由多个 pipeline workers 并发从队列里取数据执行 filter 和 output。最常用的三个参数是pipeline.workers、pipeline.batch.size、pipeline.batch.delay在logstash.yml里配置pipeline.workers: 8 pipeline.batch.size: 1250 pipeline.batch.delay: 5pipeline.workers建议先设成 CPU 核数然后观察jvm的 CPU 占用和队列积压再微调。batch.size是每个 worker 每批从队列里取多少条事件批越大吞吐越高但单批处理时间也越长遇到慢的正则时会造成明显延迟。batch.delay是批内最大等待时间单位毫秒。还有一个容易被忽略的参数是 JVM 堆内存。很多人只调 workers忘了 Logstash 的 JVM 堆大小。默认 1GB 对于生产环境远远不够。我一般设 4GB 到 8GB具体看日志量和 filter 复杂度。修改方式是在jvm.options里改-Xms和-Xmx两个值要一致避免 JVM 动态伸缩带来性能抖动。但有个冷知识Logstash 的内存占用不能全看 JVM 堆。即使你把堆设为 4GB实际常驻内存可能到 6-7GB因为这中间还有非堆区和其他开销。所以部署的时候我给 Logstash 节点的内存配置通常是“堆大小再加 30% 到 40%”比如设 4GB 堆给机器至少留 6GB 给这个进程。3.3 采集端到汇聚端的传输可靠性日志采集链路里的可靠性分为两种情况边缘到中央中央到 ES。中央到 ES 这块Logstash 的 output 插件自带重试机制ES 临时不可用时会阻塞重试不会丢数据前提是你的磁盘还有空间写持久化队列。边缘到中央如果中间是 KafkaFilebeat 端的 ack 机制和 Kafka 的副本机制已经能提供很高的可靠性。但如果你直接用 Filebeat 到 Logstash就需要关注 Logstash 的持久化队列。Logstash 默认的队列是纯内存的进程重启或崩溃后队列里积压的数据会全部丢失。生产环境下我强烈建议开启持久化队列queue.type: persisted queue.max_bytes: 4gb queue.checkpoint.writes: 1024开启后队列数据会写磁盘进程重启能恢复未处理完的数据。queue.max_bytes设多大呢我一般按“高峰期 10 分钟的日志处理量”来估算。设太大磁盘占用和启动恢复时间都变长设太小日志突发时还是会有丢弃风险。这个值需要在容量规划阶段做一次实测再确定。提示持久化队列不是万能的。如果写入 ES 的环节长期故障磁盘上的队列占满后Logstash 会进入阻断状态表现为 input 端停止消费。这不是 bug而是保护机制。运维上要配合 ES 集群的健康监控尽早发现写入故障。4. 自定义插件集成给 Logstash 补上最后一块短板4.1 为什么默认插件不够用Logstash 自带的插件覆盖了大多数场景但分布式系统里总有几个“非常规”需求是默认插件覆盖不了的。我遇到过的典型场景有几个某老系统的日志是自定义的二进制格式只有公司内部的基础库能解析某个内部系统需要把解析后的字段同步到自研的告警平台而告警平台只接受 HTTP 接口的特殊鉴权头还有为了兼容历史字段命名需要在 filter 里做非常复杂的字段映射默认 mutate 的语法写出来让人崩溃。默认插件不够用的时候有两条路一是用 Logstash 自带的rubyfilter 嵌入一段 Ruby 代码简单小逻辑可以这么干二是写一个真正的自定义插件。前者适合几十行的临时逻辑后者适合需要复用的、有完整数据流处理的场景。如果你只是需要一段可维护的复杂逻辑我也建议直接考虑自定义插件。原因很简单ruby filter 的代码嵌在配置里无法单元测试无法复用出了问 题只能靠看日志硬调。而自定义插件是一个标准的 Ruby gem 工程能写测试、能做版本管理、能在多个 Logstash 节点间分发。4.2 写一个自定义输出的完整步骤我拿一个实际案例来演示我要把 Logstash 解析后的告警事件通过 HTTP 推送到自研告警平台这个平台要求请求头带签名且签名算法是“MD5(时间戳 密钥 body)”。Logstash 自带的httpoutput 没办法定制这种签名逻辑所以我自己写了一个 output 插件。第一步创建插件骨架。Logstash 提供了一个工具logstash-plugin generate但手工建目录也行。一个 output 插件的基本结构如下logstash-output-myalert/ ├── logstash-output-myalert.gemspec ├── Gemfile ├── lib/ │ └── logstash/ │ └── outputs/ │ └── myalert.rb └── spec/ └── outputs/ └── myalert_spec.rb第二步实现核心代码。以下是简化版的lib/logstash/outputs/myalert.rb# encoding: utf-8 require logstash/outputs/base require logstash/namespace require net/http require digest/md5 require json class LogStash::Outputs::MyAlert LogStash::Outputs::Base config_name myalert concurrency :single config :url, :validate :string, :required true config :token, :validate :string, :required true public def register uri URI.parse(url) http Net::HTTP.new(uri.host, uri.port) http.use_ssl uri.scheme https http.open_timeout 5 http.read_timeout 5 end public def receive(event) payload event.to_hash body payload.to_json ts Time.now.to_i.to_s sign Digest::MD5.hexdigest(ts token body) headers { Content-Type application/json, X-Timestamp ts, X-Sign sign } request Net::HTTP::Post.new(uri.request_uri, headers) request.body body http.request(request) rescue e logger.error(Failed to send alert: #{e.message}, :event event.to_hash) end end这里说几个关键点。concurrency :single表示这个插件的并发模型是单线程对于 HTTP 推送类 output我倾向于单线程加外部重试避免并发请求打爆下游告警平台。register方法在插件启动时执行一次适合建立 HTTP 连接等初始化工作。receive(event)对每一条事件调用事件流经这里后正常的返回即可。第三步写 gemspec 和 Gemfile。这一步是为了让 Logstash 能识别这个插件作为 gem 安装# logstash-output-myalert.gemspec Gem::Specification.new do |s| s.name logstash-output-myalert s.version 0.1.0 s.platform java s.summary Custom output plugin for sending alerts to internal platform s.authors [Your Name] s.require_paths [lib] s.files Dir[lib/**/*, spec/**/*, *.gemspec] s.test_files s.files.grep(/spec/) s.add_runtime_dependency logstash-core-plugin-api, 2.1.0 s.add_runtime_dependency logstash-codec-plain s.add_development_dependency logstash-devutils end对于 Java 平台这行要注意Logstash 是基于 JRuby 的所以s.platform java不能少。logstash-core-plugin-api的版本也要和你部署的 Logstash 版本匹配我一般先查当前版本再去定直接在 Logstash 安装目录下跑bin/logstash-plugin list --verbose logstash-core-plugin-api能看当前版本。4.3 插件打包、安装与踩坑插件代码写完后把 gem 包构建出来cd logstash-output-myalert gem build logstash-output-myalert.gemspec安装到 Logstash 节点bin/logstash-plugin install logstash-output-myalert-0.1.0.gem验证安装bin/logstash-plugin list | grep myalert如果你是在多节点部署每台都要安装一次。更省事的方式是搭一个内部 gem 源或者把插件做成一个离线 gem 包配合配置管理工具批量分发。这里有个坑我必须提醒Logstash 8.x 对 JRuby 插件的兼容性要求比较严格。如果你的 Logstash 版本比较新写插件时最好用对应的 JRuby 语法不要使用太老的 Ruby 特性。另外打包出来的 gem 如果在本地开发机是 Ruby非 Java 平台安装到 Logstash 时会报 platform 不匹配这就是为什么 gemspec 里必须写s.platform java。自定义 filter 插件和 input 插件结构类似只是父类不同。filter 插件继承LogStash::Filters::Base核心方法叫filter(event)input 插件继承LogStash::Inputs::Base核心是run方法。理解了这三个基类的区别基本上就能应对绝大多数自定义需求。经验之谈不要一上来就写复杂插件。先确认默认插件和 ruby filter 是否真的搞不定。写插件的成本除了开发还有维护和兼容性测试。我曾经为了一个排序功能写了上百行插件后来发现用默认的 ruby filter 配合event.set十几行就能解决白白增加了一个 gem 的维护负担。5. 常见问题排查与实战避坑清单5.1 高频问题速查表分布式日志监控跑起来之后日常运维里最常碰到的几个问题我整理成一张表都是实际操作中踩过的。现象可能原因排查方向日志在 Kibana 里搜不到采集端没读到文件Logstash input 未匹配先看 Filebeat 日志再看 Logstash input 的统计最后查 ES 索引是否存在解析覆盖率很低大量_grokparsefailuregrok 正则写错或日志格式变了用 Kibana 搜索该 tag抽样本日志逐条对比正则Logstash 进程 CPU 飙高grok 或 dissect 解析太慢workers 太少逐步关闭 filter 看 CPU 变化调大 workers优先把 grok 改 dissect日志时间错乱没有用 date 插件重置timestamp在 filter 里加 date 插件match 业务时间字段高峰时段丢日志内存队列打满ES 写入跟不上开启持久化队列检查 output 的批量写入配置评估 ES 集群性能日志重复出现input 读取 KafKa 提交 offset 失败或 Logstash 重启后重复消费检查 Kafka 的group_id和提交策略幂等处理写入端第一行“日志搜不到”是最常见的求助问题。我的排查顺序是先看来源机器的 Filebeat 是否正常读取再看 Logstash 的 input 有没有收到数据。Logstash 提供了一个查询管道统计的 API或者直接在配置文件里阶段性打印日志来定位问题在哪一段。别一上来就查 ES 权限顺序反了会浪费很多时间。“日志时间错乱”这条我多补充一句。很多日志系统部署时没有统一业务时间和日志内容时间导致日志在检索页面上看起来是乱序的。date 插件修正了timestamp之后不只是排序变正常告警里的时间范围判断也会准确。这一步是分布式系统日志监控里的“隐形刚需”。5.2 几个值得长期维护的运维习惯日志监控系统的上线只是开始长期稳定运行靠的是几个好习惯。第一解析覆盖率要持续监控。我在 Kibana 建了一个仪表盘专门展示每天_grokparsefailure的占比超过阈值就告警。新上线的服务日志格式不规范或者老系统升级改了日志前缀这些都会导致解析覆盖率下降。没有这个监控你会在某个深夜排查问题时才发现这个月日志一直没解析成功。第二配置模板和字段映射要有版本管理。我对 Logstash 配置、ES 索引模板都纳入 Git 管理每次变更走评审。日志字段的增删会影响查询和告警如果谁都能随意改出问题是迟早的事。第三容量规划要留缓冲。日志量的增长不完全线性大促、促销、活动都可能导致高峰暴涨。我给 Logstash 节点预留了高峰期 2 倍以上处理能力Kafka 的保留期也设置了至少 3 天防止问题排查时数据早就过期被清理。第四单元测试解析规则。Logstash 从 6.x 开始提供bin/logstash -t配置检查但只检查语法不检查正则能否正确匹配。我的做法是把常用日志样本存成测试数据每次改 grok 规则后先在本地跑一遍 Logstash 管道用stdoutoutput 观察解析结果确认无误再上生产。这一步能省掉很多“配置上线后才发现解析失败”的返工。注意修改 Logstash 过滤规则时先确认改的是“解析逻辑”还是“字段语义”。解析逻辑改了之后新数据按新规则解析但历史索引里的字段结构不会变。如果你依赖历史数据做对比分析要注意新旧字段的一致性必要时在 ES 做一次 reindex。尾巴关于日志监控我最后想说的跑了几年 Logstash最大的体会是日志监控系统真正难的不是技术本身而是团队对日志规范的共识。工具再强大如果每个服务想怎么写日志就怎么写日志平台迟早变成“日志垃圾场”。反过来只要日志规范定好Logstash 的解析规则会非常稳定新增服务只需要套用统一的模板就能接入。如果你现在正要开始搭这套体系我建议从最小闭环做起先接两三个核心服务的日志把字段规范、索引模板、Kibana 仪表盘、告警规则跑通再逐步推广。别一上来就追求全量接入那样很容易被各种历史日志的“脏格式”拖垮信心。最后分享一个我常用的小技巧Logstash 的stdoutoutput 是调试神器新环境调解析规则时先把处理结果打到控制台观察再切到 Elasticsearch output。很多人上线第一天就在 filter 里留了个连自己都没跑过的 grok 正则结果数据进去全是 parsefailure排查一圈才发现是正则的问题。先在本地把管道跑通再上生产这个习惯会让你少掉很多头发。