分布式系统的日志向来做起来头疼,尤其是节点一多、服务一拆分,日志散落在几十台机器上,出了问题想定位简直像大海捞针。我在这块折腾了挺长时间,最后沉淀下来一套以 Logstash 为核心的日志监控方案,今天把这套实践的思路、配置、坑点一次讲清楚。这篇文章适合刚接触 Logstash 的人,也适合那些已经在用但总感觉性能不对、日志对不齐的团队,我会从架构选型讲到参数调优,再讲到如何写自己的自定义插件,尽量让你看完能直接上手。
1. 为什么分布式日志监控非 Logstash 不可
1.1 分布式系统日志监控的三个典型难题
分布式的环境里,日志监控面对的从来不是“多几台机器”这么简单。第一个难题是数据源极其分散,有容器里的 stdout、宿主机上的文件、中间件自带的事件日志、业务系统通过 TCP 推送过来的数据,这些数据格式完全不一样,有的是 JSON,有的是多行堆栈,有的干脆是乱糟糟的纯文本。第二个难题是流量不均匀,业务高峰期的日志量可能是低谷期的几十倍,如果采集链路没有缓冲和削峰能力,一到高峰期就会疯狂丢日志。第三个难题是排查链路困难,一个请求跨了三四个服务,日志分散在不同机器上,没有一个统一处理的环节,你根本没法凭一个 traceId 把整条调用链拼起来。
这三个问题单独看都不致命,但叠加在一起,就要求日志采集层必须具备三件事:多源接入、格式解析、数据规整。Logstash 恰好就是围绕这三个能力设计的。
1.2 Logstash 的定位与其他采集器的对比
很多人一上来就问,Filebeat 不是更轻量吗,为什么还要 Logstash?这里要先搞清楚定位。Filebeat 本质是一个采集器,它只负责“把文件读出来送到某个地方”,它确实很轻、资源占用少,但它没有能力做复杂的解析和规整。而 Logstash 是一个完整的处理管道,它接收数据后可以走 filter 阶段,做 grok 正则解析、JSON 解包、字段类型转换、时间戳矫正、甚至通过自定义插件跑任意逻辑。
我见过不少团队直接用 Filebeat 推数据到 Elasticsearch,一旦原始日志格式稍微特殊一点,比如多行堆栈被切断、字段类型对不上,后期再想补救就得写一堆 Painless 脚本,维护成本非常难受。而把这些脏活放在 Logstash 的 filter 阶段统一处理,主链路干净得多。
换个角度说,Logstash 的定位更像一个“数据加工厂”,它不追求极致的轻量,而是追求“任何日志进来,出去都是规整的结构化数据”这件事。在分布式场景下,这个加工能力恰恰是刚需。当然,Logstash 也不是不能做采集端,只是它更常被部署在日志汇聚层,负责把各节点的数据统一收口。
2. 整体架构设计:从采集到落地的全链路
2.1 一条日志的完整旅程:管道模型拆解
Logstash 的所有秘密都在管道模型上。管道分三个阶段:input 负责接收数据,filter 负责处理数据,output 负责输出数据。看起来很简单对吧,实际用起来你会发现这一个模型几乎能覆盖所有接入场景,而且三个阶段是独立的,你可以任意调整中间的 filter 逻辑而不影响 input 和 output。
以我们团队为例,一条 Nginx 日志的旅程是这样的:Filebeat 在节点上读文件,通过 Logstash 的 beats input 端口送入管道;进入管道后,先由 grok 插件把一行文本拆成 clientip、request、status、body_bytes_sent 这些字段;然后 date 插件从日志里提取时间戳,替换掉默认的接收时间,避免日志展示时间和实际发生时间不一致;接着 mutate 插件把 status 从字符串转成整数、body_bytes_sent 转成 long 类型;最后 output 插件把这条规整后的数据写入 Elasticsearch。从采集到落地,整条链路的数据格式在管道里经历了从“原文”到“结构体”的蜕变。
这条链路的灵活之处在于,filter 阶段可以组合任意多个插件,之间的顺序是可以调的。顺序这件事很容易被忽略,实际上却非常关键。比如先 grok 提取字段,再 mutate 转换类型,这两步顺序一旦反了,grok 匹配到的全是转换后的数据,原始信息就丢了。我自己的习惯是,先把原始字段拆干净,再统一做类型转换和标准化,最后再做数据裁剪和脱敏,这个顺序踩过坑的都知道。
2.2 部署形态选型:集中式还是边车式
Logstash 的部署形态是分布式架构里第一个要决策的事。所谓集中式,就是所有节点的日志通过采集器上传到一台或者几台中央 Logstash 集群来处理;所谓边车式,就是每个应用节点旁边都挂一个独立的 Logstash 实例,日志在本地处理完再送给存储端。
集中式的优势是维护简单,插件配置只需要在一处管理,上层做告警、做治理都方便,但劣势是节点到中央 Logstash 之间的网络链路成为瓶颈,日志量大的时候容易把网络打满。边车式的优势是处理能力随节点水平扩展,日志传输距离短、延迟低,但劣势是每个节点都要投入一份 Logstash 的资源开销,配置变更要批量推送,管理复杂度一下就上来了。
我在实际工作里更推荐一种混合方案:采集层用 Filebeat 把日志就近送到 Kafka,由 Kafka 做统一缓冲,然后再由一组 Logstash 集群从 Kafka 消费数据做解析处理。这个方案既规避了集中式 Logstash 的网络瓶颈,又不会像边车式那样在每个节点消耗大量内存,而且 Kafka 天然具备削峰填谷的能力,业务高峰期几十倍流量也不会打到 Elasticsearch 上。
2.3 引入 Kafka 做缓冲层:削峰填谷的关键
如果没有 Kafka 这一层,Logstash 的 input 很可能直接被洪水一样的日志冲垮,尤其是业务在搞大促、做秒杀的时候。我自己就经历过凌晨两点被电话叫起来,原因是日志量瞬间暴涨,Logstash 的内存被堆到上限,整个管道卡死在等待状态。
有了 Kafka 之后,这个问题从根本上被化解了。Filebeat 只需把数据送进 Kafka 就算完成任务,不需要关心下游处理能力。Logstash 的消费速度完全由自己控制,消费不过来就暂时积压在 Kafka 里,等高峰期过去再慢慢消费。Kafka 的吞吐量几乎是线性扩展的,分区数设置合理的话,日志量再大也只是加分区的事。
这里有一个很值得注意的细节:Kafka 的 topic 分区数最好跟 Logstash 的消费并发度匹配起来。Logstash 从 Kafka 消费的时候,一个分区对应一个消费线程,如果你的分区数是 3,那即便你把 pipeline.workers 调到 10,实际并行度也只有 3。所以设计 Kafka topic 的时候,分区数要预留余量,一般来说分区数不少于 Logstash 实例数的 1.5 到 2 倍,这样后续扩容 Logstash 实例时才不会因为分区数不够而卡住并行度。
3. 核心配置细节与关键参数解析
3.1 input 插件选型与配置要点
Logstash 的 input 插件非常多,但在分布式场景下,真正高频使用的其实就三个:beats、kafka、tcp。如果你在节点上部署了 Filebeat,那 Logstash 这边就用 beats input 来接收;如果你走 Kafka 缓冲层,那就用 kafka input 来消费;如果你面对的是遗留系统,对方只能通过 TCP 直接把日志发过来,那就用 tcp input。绝大多数场景跑到这三个插件就够了。
输入插件虽然不是性能瓶颈,但配置里还是有几个容易踩坑的地方。比如 tcp input 默认是按行分割数据的,如果你的日志本身就包含换行符,比如 Java 的堆栈信息,就会导致一条完整堆栈被切成了好几条。这时候要么在发送端改造成 JSON 包装后的单行,要么在 Logstash 侧配合 codec 做多行合并。再比如 kafka input 的 group_id 要保证唯一,多个 Logstash 实例如果用了同一个 group_id,它们会分摊同一个 topic 的不同分区,这是预期的;但如果你误以为不同实例会各消费一份全量数据,那就和预期完全反了,最后结果就是每份日志只被处理了一半。
input { beats { port => 5044 client_inactivity_timeout => 3600 } kafka { bootstrap_servers => "kafka1:9092,kafka2:9092" topics => ["app-log", "nginx-log"] group_id => "logstash-prod-group" codec => json } tcp { port => 5000 mode => "server" codec => line } }小提示:如果你用了多个 input,它们进来的事件会混在同一个管道里处理。后续要在 filter 里区分来源,可以给每个 input 加上 type 或者 tags 字段,比如type => "nginx",然后在 filter 用条件判断处理逻辑,这样不同来源的日志就不会被混成一套规则。
3.2 filter 阶段的解析与字段标准化
filter 阶段是 Logstash 最有价值的部分,也是最容易写烂的部分。我见过太多的配置,filter 里堆了几十个插件,每进来一条日志就把所有正则跑一遍,性能差到离谱,最后只能靠加机器硬撑。真正合理的做法是,先用条件判断把日志分流,再按各自的格式执行最小化的解析逻辑。
举个例子,如果我们同时接收 Nginx access log 和 Java application log,那 filter 的写法应该是这样:
filter { if [type] == "nginx" { grok { match => { "message" => "%{IPORHOST:clientip} %{DATA:request} %{NUMBER:status:int}" } overwrite => ["message"] } date { match => ["timestamp", "dd/MMM/yyyy:HH:mm:ss Z"] target => "@timestamp" } } else if [type] == "java" { grok { match => { "message" => "%{TIMESTAMP_ISO8601:log_timestamp} %{LOGLEVEL:level} %{DATA:class} - %{GREEDYDATA:msg}" } } multiline { pattern => "^\s+at " negate => false what => "previous" } } }这上面的代码演示了一个重要原则:先分流、再解析、同类日志才共享规则。很多新手把所有日志混在一起,试图用一条正则覆盖所有格式,结果就是匹配率低、字段残缺、解析失败还要靠 Kanban 板维护规则。
字段标准化也是一个重要习惯。比如 IP 地址、状态码、响应时间这些字段,在原始日志里都是字符串,如果不加转换,落到 Elasticsearch 里做排序和范围查询的时候就会出问题。我的做法是在 grok 提取时就指定类型,或者统一用 mutate 插件做 convert 转换。
mutate { convert => { "[status]" => "integer" } convert => { "[responsetime]" => "float" } rename => { "responsetime" => "response_time_ms" } gsub => ["message", "\"", ""] }这里要提醒一下,grok 插件里的类型转换语法是%{NUMBER:status:int},这里的int并不是 Elasticsearch 里的 integer,它只是把字符串在管道内转成数值类型,最终写到 ES 的 mapping 类型,取决于你 output 的时候怎么设定的。很多人忽略了这层区别,导致后面 mapping 不一致,这是个非常隐蔽的坑。
3.3 output 输出策略与 Elasticsearch 的高效对接
output 阶段最常见的做法是写到 Elasticsearch,但怎么写得高效、写得稳,这是个技术活。我见过很多 Logstash 进程其实就是死在这里。第一个必须注意的问题是 index 的命名。不要用一个固定 index 名,比如logstash-%{+YYYY.MM.dd}这种按天拆分的方式,才能保证后续做索引生命周期管理(ILM)时,可以按天或者按月滚动删除旧数据,不然索引会无限膨胀,最后 Elasticsearch 自己先扛不住了。
如果你已经用了 ILM,那让你在 output 里配置的 index 名称要和 ILM 策略匹配。ES 7 以上的版本推荐用ilm-rollover-alias的方式:output 写入别名,由 ILM 负责滚动。这里有一个很经典的报错,就是你在 output 里既写了index又开了ilm_enabled => true,结果数据写不进去,原因就是索引生命周期管理的策略和手动指定的 index 冲突了。解决方案是二者只保留一个,推荐用 ILM 管理滚动,output 里只写别名。
output { if [type] == "nginx" { elasticsearch { hosts => ["es1:9200", "es2:9200"] index => "nginx-log-%{+YYYY.MM.dd}" manage_template => false user => "elastic" password => "xxxx" } } else { elasticsearch { hosts => ["es1:9200", "es2:9200"] index => "app-log-%{+YYYY.MM.dd}" manage_template => false } } }manage_template => false这个配置也是值得说明的。默认情况下 Logstash 会自己创建一套模板,导致字段类型由 Logstash 说了算。我比较推荐在 ES 侧维护一份自己控制的模板,然后在 Logstash 里关掉模板管理,这样 mapping 的变更流程完全可控,不会出现 Logstash 升级后模板被重置的诡异问题。
3.4 性能调优:从内存到批量的全面检查
Logstash 性能调优,核心就两件事:JVM 堆内存和管道批量参数。前者决定了 Logstash 能吃掉多少数据,后者决定了它处理数据的节奏感。
JVM 堆内存的默认值是 1GB,这对一个处理大流量的管道来说完全不够。我一般建议堆内存设置为物理内存的一半,但不要超过 8GB,因为 JVM 在堆内存大于 8GB 之后,对象指针压缩会失效,内存翻倍了性能反而可能下降。修改方式是在jvm.options里调整-Xms和-Xmx,注意这两个值一定要一样,避免 JVM 运行时扩容导致停顿。
管道参数方面,最重要的三个是pipeline.workers、pipeline.batch.size和pipeline.batch.delay。pipeline.workers默认是 CPU 核数,它决定了 filter 阶段的并行度,但要注意它不是越大越好,因为每个 worker 的上下文切换和内存开销都不小,我实测在 8 核机器上设 6 到 8 个 worker 往往比设 16 个效果更好。pipeline.batch.size默认是 125,这个数偏保守,如果你从 Kafka 消费数据,完全可以调到 1000 甚至 2000,批量越大,单批处理效率越高,但也要看单条数据的复杂度。批量处理是有延迟的,pipeline.batch.delay默认 50ms,意思是攒够 50ms 的时间就 flush 一次,如果你的场景对实时性要求不那么苛刻,把 delay 调大到 100ms 左右,吞吐量会有明显提升。
还有一点是管道队列的设置。queue.type默认是memory,如果管道输出端出现阻塞,数据会积压到内存队列里,内存再用完了就开始丢数据。稳妥的做法是改用持久化队列persisted,设置queue.max_bytes比如 8GB,这样即使 ES 挂了几个小时,数据也能先存在磁盘队列里,等恢复后继续消费,避免丢日志。
4. 集成自定义插件:从零开发一个日志解析过滤器
4.1 什么时候需要自定义插件
Logstash 自带的 filter 插件确实很丰富,grok、mutate、date、json、csv、geoip,覆盖了绝大多数解析需求。但真实世界的日志格式总是会超出插件的表达能力。比如我们曾经遇到一种老系统产生的二进制半结构化日志,字段之间用特殊分隔符隔开,而且同一个字段在不同时期会变换含义。这种格式用 grok 写正则几乎没法维护,用 ruby filter 写内联脚本又没法做单元测试。这个时候就该考虑自定义插件了。
自定义插件最大的价值不是炫技,而是把复杂的解析逻辑封装成可维护、可测试的独立模块。它虽然是用 Ruby 写的,但你可以把它当作一个普通的工具类来理解:接收事件对象,做处理,然后交给管道继续走。一旦写好,你可以像使用内置插件一样在配置里引用它,比如filter { example { message => "%{message}" } }。
不过我也要劝一句:能用内置插件解决的事,别急着写自定义插件。自定义插件引入了代码维护成本、升级兼容成本和测试成本。只有当内置插件的组合已经非常别扭,或者性能已经无法接受时,才值得动手。我们团队的标准是:同一个解析逻辑被动复制粘贴超过三处,或者单条解析性能已经低于每秒 3000 条且用 grok 又无法优化时,才考虑自定义插件。
4.2 插件骨架与生命周期
Logstash 插件的结构其实非常规整,一个最小可运行的 filter 插件需要四个文件:一个 gemspec 文件、一个主 ruby 文件、一个版本文件和一个配置模板。插件命名有严格的规范,filter 插件必须叫logstash-filter-插件名,对应的类名是LogStash::Filters::插件名,这个映射关系不能错,否则 Logstash 加载插件时会直接报找不到。
插件的生命周期方法只有两个核心的:register和filter。register方法在 Logstash 启动时调用一次,适合做初始化工作,比如预编译正则、建立外部连接、读取配置文件。filter方法在每个事件到来时调用,输入参数是事件对象,你在这里解析、修改、增加字段,然后调用@metrics或者直接返回。事件对象有一套自己的 API,event.get("field")用来读取字段、event.set("field", value)用来设置字段、event.include?("field")用来判断字段是否存在,这套 API 跟内部 scripting 使用的方式是一样的。
# encoding: utf-8 require "logstash/filters/base" require "logstash/namespace" class LogStash::Filters::Example < LogStash::Filters::Base config_name "example" config :prefix, :validate => :string, :default => "" def register @prefix = @prefix.to_s end def filter(event) message = event.get("message") if message.start_with?(@prefix) event.set("extracted", message.sub(@prefix, "")) event.set("parse_status", "ok") else event.set("parse_status", "skipped") end @metrics.incr(:events, :processed) end end这段代码的实际功能非常简单:给带指定前缀的消息做一个提取,标记解析状态。别看它简单,它演示了自定义插件的基本骨架,后续你要做更复杂的解析,也是在这个框架上扩展而已。
4.3 开发一个多行 Stacktrace 过滤器
既然要讲实际价值,我就拿我们真实开发的一个插件举例:解析 Java 异常堆栈。Java 的堆栈日志天生是多行的,第一行是异常类型和描述,后面跟着at com.xxx.Class.method(Class.java:12)这样的调用栈,而且往往嵌套着Caused by。如果按行读取,一条完整异常会被拆成多条日志,查问题的时候根本没法定位完整堆栈。
我们开发了一个logstash-filter-stacktrace插件,它的核心逻辑是识别多行堆栈并聚合成一条完整的异常事件。这个插件内部维护了一个状态机:新事件到来时先判断是不是异常开始行,如果是,就开启一个聚合缓冲区;后续的行只要不是新的异常开始行,就持续追加到缓冲区末尾;当遇到非堆栈内容或者超时窗口到期,就把缓冲区的完整堆栈作为一个事件释放出去。
def filter(event) line = event.get("message").to_s if line =~ /^(\w+(?:\.\w+)+Exception|\w+Error):/ flush_pending if @stack_buffer @stack_buffer = [line] event.cancel elsif line =~ /^\s+at / && @stack_buffer @stack_buffer << line event.cancel elsif @stack_buffer && line.strip.empty? @stack_buffer << line event.cancel else flush_pending end end def flush_pending return unless @stack_buffer merged = @stack_buffer.join("\n") event = LogStash::Event.new("message" => merged) event.set("stacktrace", true) @output_queue << event @stack_buffer = nil end这里面一个关键细节是event.cancel方法。它告诉 Logstash 当前这个事件不要再继续往下走,因为堆栈行已经并入缓冲区了,如果这里不取消,原始堆栈行就会变成很多条残缺日志进到 ES。等到完整堆栈拼好之后,我们通过@output_queue手动投递新事件,让合并后的堆栈作为一条完整日志继续走后续管道。
这个插件用到的技术点其实不难,难点在于边界情况:堆栈里穿插了Caused by、日志里有空行、多条堆栈并发到来。我们最终采用了一个超时策略,如果堆栈行之后 5 秒没有新的堆栈行到来,就强制 flush,保证日志不会在内存里积压太久。
4.4 安装与测试的实践经验
自定义插件写好后,打包和安装的流程也是标准化的。在插件根目录下运行gem build logstash-filter-stacktrace.gemspec,会生成一个.gem文件。然后通过bin/logstash-plugin install /path/to/logstash-filter-stacktrace-0.1.0.gem把这个文件安装到 Logstash 的本地插件目录。注意,logstash-plugin 命令的路径依赖你启动 Logstash 的账户,如果你是 root 启动,就要用 root 的 Logstash 安装路径来装插件,否则会出现“插件已装但 Logstash 加载失败”的问题,这个我踩过。
装完之后先别急着上生产,用bin/logstash -f test.conf --config.test_and_exit做一次配置校验,确认插件能被正确加载。然后再用一行测试数据跑管道,观察 stdout output 打印的结果是不是符合预期。我习惯在测试配置里加一个 stdin input 和一个 stdout output,这样可以直接手动敲测试日志验证效果,比直接上生产看 index 要快得多。
还有一个建议:把插件源码纳入版本管理,用 CI 跑单元测试。Logstash 插件本质上是 Ruby gem,所以可以用 rspec 写测试,注册一个事件对象,调用filter,断言字段结果。这些测试可以在不启动 Logstash 的情况下跑,速度很快。有了测试兜底,后续改插件才敢大胆重构。
5. 常见问题与排查技巧实录
5.1 日志“神秘”丢失的常见原因
日志丢失是分布式监控里最让人抓狂的问题。现象是业务日志明显在产生,但 ES 里查不到。排查时我第一步永远不是看管道,而是看 Kafka 的消费位置。如果 Kafka consumer group 的 lag 一直增长,说明 Logstash 在消费但处理不过来;如果 lag 是 0 但 ES 没有数据,那问题出在 Logstash 到 ES 这段链路。
有一个高频原因:Logstash 的 output 写 ES 时遇到 mapping 冲突,比如同一个字段既是 string 又是 long,ES 会直接拒绝写入,而 Logstash 默认的重试策略又是有限的,重试几次失败后就直接丢弃了。解决办法是在 ES 侧检查 rejected 日志,通常能在 ES 的日志文件里看到mapper_parsing_exception字样,然后去调整索引模板,把冲突字段的类型统一,再重新投喂数据。
另一个常见原因是 sincedb 文件的混乱。Filebeat 通过 sincedb 记录每个文件的读取位置,如果你迁移了路径或者文件被轮转,sincedb 对应关系错乱,就会重复读取或者跳过文件。解决办法是清掉旧的 sincedb 文件重新读取,或者给 Filebeat 配置ignore_older来避免处理太久远的日志。
这里我把排查顺序整理成一个速查表:
| 现象 | 第一步排查 | 常见根因 | 快速处理方式 |
|---|---|---|---|
| ES 没有数据但 Kafka lag 为 0 | 检查 Logstash 日志中 output 报错 | mapping 冲突、ES 拒绝写入 | 修正索引模板后重灌数据 |
| Kafka lag 持续上涨 | 检查 Logstash 的 CPU 和堆内存 | filter 解析性能不足 | 增大 batch size、调优正则、扩实例 |
| 同一条日志出现多次 | 检查 sincedb 和消费 offset | Filebeat 重复读取、Kafka offset 重置 | 清理 sincedb 或重置 consumer group |
| 日志出现但缺字段 | 检查 filter 的匹配结果 | grok 正则不匹配 | 用 grok debugger 调试正则 |
5.2 时间戳与时区的“八小时幻觉”
日志监控里最隐蔽的问题之一就是时间不对。你看到 Kibana 上一条日志的 @timestamp 是 8 小时之前,第一反应往往是系统时钟有问题,查了一圈发现机器时间都是对的,最后才意识到是时区转换的锅。
Logstash 默认使用 UTC 时间作为 @timestamp。如果你的日志里自带时间戳,通过 date filter 解析后,默认把它当作 UTC 来存储,而你的业务系统在东八区,Kibana 展示的时候如果没做时区转换,就会看到差 8 个小时。解决方式是在 date filter 里显式指定时区:
date { match => ["log_timestamp", "yyyy-MM-dd HH:mm:ss"] timezone => "Asia/Shanghai" target => "@timestamp" }另一种情况是日志里本身没有时间戳,只有 Logstash 的接收时间。这时你看到的就是 Logstash 处理时刻的 UTC 时间,如果你希望它跟业务本地时间一致,可以在 Kibana 里设置展示时区,也可以在 filter 里用 ruby 插件强制转换,但一般推荐前者,因为 ES 里存 UTC 才是最佳实践,展示层的时区转换应该交给 Kibana。
5.3 JVM 内存与管道积压的恶性循环
Logstash 跑着跑着突然 OOM,或者 CPU 飙升到 100%,这类问题的根源往往不在 Logstash 本身,而在下游。比如 ES 集群节点挂了,output 一直写不进去,Logstash 的重试机制会不断堆积待发送数据,内存队列越来越大,最终触发堆内存溢出。
这个问题的根因是管道里没有“泄压阀”。我的建议是两层防护:第一层,使用持久化队列,设置queue.type => persisted和queue.max_bytes,让溢出数据落到磁盘而不是内存,这个配置能做到进程重启后数据不丢。第二层,给 Logstash 配上健康检查,定时调用_node/stats接口,监控jvm.mem.heap_used_percent和pipeline.workers的状态,一旦超过 85% 就触发告警,让运维提前介入,不要等到 OOM 再慌。
这里还要注意 JVM 堆内存设置的问题。如果你把-Xmx设得太接近物理内存上限,JVM 频繁进行 Full GC,停顿时间会很长,Logstash 的吞吐量反而下降。我实测过,堆内存 8GB 的情况下,把-Xms和-Xmx都设为 8GB,再配合 CMS 收集器(JDK 8 场景下),整体表现是最稳的。
5.4 反压与背压:如何让管道自适应流量
Logstash 的管道架构天然带有流量控制的机制,但很多人没用明白。output 写不动的时候,Logstash 会自动把压力传到 filter,filter 传回 input,最终 input 暂停拉取新数据。这个过程叫 backpressure。如果你的 input 是 beats,Filebeat 收到 Logstash 的背压信号后会降低读取文件的速度,避免把 Logstash 压垮。
但如果你用的是 Kafka input,背压机制就不是那么灵敏了。Logstash 从 Kafka 拉取数据是主动行为,它不会自动感知 Kafka 里积压了多少消息,只会按自己的节奏消费。正因如此,Kafka 场景下我通常用fetch_max_bytes和max_poll_records这两个参数来控制消费节奏。前者控制每次 fetch 请求的最大字节数,后者控制每次 poll 返回的最大消息条数。如果你发现消费太快导致磁盘 IO 飙升,可以适当调低这两个值,让 Logstash 慢慢消化而不是一口吃成大胖子。
另外别忘了设置consumer_threads参数,它决定了每台 Logstash 实例上单组消费者的线程数。这个参数要和 Kafka 分区数配合,不然容易造成有的分区长时间不被消费,出现 lag 不均衡的情况。我见过一个案例,Kafka 分区有 20 个,但 Logstash 实例上consumer_threads设了 5,结果 20 个分区只有 5 个被消费,其余 15 个分区的数据一直积压,业务方看到的就是日志严重延迟,排查了半天才找到这个配置的问题。
最后说几句掏心窝的话
日志监控这个事,真正难的从来不是把 Logstash 跑起来,而是把它跑得又稳又快、出了问题还能快速定位。我个人这几年的最大体会是:一定要重视管道里的可观测性。Logstash 自己就是做可观测性的工具,但它自身的运行指标很多团队反而不看。我的习惯是在上线初期就把_node/stats和_node/hot_threads这两个接口的指标接到监控系统里,每天看一次各个插件的耗时分布,这样真出问题的时候,不是从“哪里丢日志”开始猜,而是直接拿着数据说话。
另一个体会是,配置文件的版本管理一定要做。Logstash 的配置更新频繁,尤其是 filter 规则,一条正则的调整就可能引发线上日志解析失败。我建议把所有配置纳入 Git 管理,并且每个环境(测试、预发、生产)分开维护,改动必须走 MR 流程。听起来有点重,但经历过一次“生产 filter 被热更新后 grok 全挂”的事故后,你就知道这步有多重要了。
最后再分享一个小技巧:在 filter 阶段加一个[@metadata][processor]字段,标记每条日志经过的解析步骤。排查问题时,你可以在 ES 里用这个字段做过滤,快速确认日志到底是在哪个环节被处理的。这个小改动不值钱,但排查效率能提升一大截,强烈建议试试。