深入 Logstash 04 - input 插件:拉取、监听与 Beats 接入
input 层能给出多强的可靠性,由一件事决定:"读到哪里"这个状态存放在哪里、在什么时刻推进。file input 把它写进本地 sincedb 文件;kafka input 交给 broker 侧的 consumer group;beats input 自己不存,靠 Filebeat 的 registry 配合一次 ACK 来推进。三种介质完全不同,推进时刻却停在同一条线上:event 进入 Logstash 内部队列,而不是写进下游成功之后。
这一篇沿这条线索展开三件事:sincedb 到底记了哪几列、beats 的 ACK 语义边界在哪里、kafka 的 offset 提交时机由哪个参数控制。顺带回答一个更靠前的问题——拉模型和推模型的分野,落到可靠性上究竟差在哪。
1 | |
拉模型与推模型
input 插件在获取数据时有两种根本不同的方向:
拉模型(pull):Logstash 主动去源端取数据。file input 定期轮询受监控路径下文件的标识与大小变化;jdbc input 按 schedule 执行 SQL,把查询结果转成 event。源端不感知 Logstash 的存在,Logstash 决定何时取、取多少。
推模型(push / listen):Logstash 监听端口等待数据进来。beats input 在 TCP 端口上运行 Lumberjack 协议服务端;tcp/udp input 监听对应端口;kafka input 看起来像"主动连 Kafka",但从消息流动角度更接近推模型——Kafka broker 持续推送新消息给消费者,Logstash 只需维持消费组成员身份。
这个区分在可靠性上有直接后果:拉模型的源端通常无法感知消费方是否处理完毕,位点管理由 Logstash 自己维护;推模型里(以 beats 为典型)源端等待消费方的确认,确认收到后才推进位点或删除暂存数据。
file input 与 sincedb
file input 是最常见的 pull 插件。它监控一组文件路径(支持 glob),检测到新内容后读取并逐行生成 event。
关键内部机制是 sincedb。Logstash 把每个被跟踪文件的标识(inode、主设备号、次设备号)连同当前字节偏移记录到一个状态文件里。下次启动或文件被轮转时,从 sincedb 读出上次停在哪里,从该偏移继续读,而不是从头重读。
sincedb_path 没有默认值。不显式配置时,sincedb 落在 <path.data>/plugins/inputs/file 下,文件名由 path 的 glob 推导,所以改了 glob 的写法就等于换了一个新 sincedb,已有位点全部作废。官方还给了一条硬约束:每个 file input 必须用各自独立的 sincedb_path,多个 input 共用同一个路径会互相覆盖对方的位点。
sincedb 的列数不是固定的,< v5.0.0 是四列,之后是五列或六列:
1 | |
第五列的浮点时间戳不是给人看的,它驱动 sincedb_clean_after(默认 "2 weeks",传数字时按天解释):一个被跟踪文件在这段时间内没有任何变化,它的 sincedb 记录就过期、不再持久化。这个机制存在的理由是文件系统会复用 inode,而 file input 没有可靠办法判断某个 inode 是否已经换成了新内容。代价是记录过期后如果这个文件又被发现,它会被当成从未见过的文件从头读一遍——刷盘窗口最多让你重读十几秒的内容,记录过期让你重读整个文件。
几个核心配置项:
start_position:控制 Logstash 第一次看到某个文件时从哪里开始。end(默认)表示只读新追加的内容,适合尾随实时日志;beginning 表示从文件开头读,适合处理历史文件。"第一次"是以 sincedb 里没有这个文件标识的记录为准,不是以文件是否存在为准。mode => "read" 时本参数被直接忽略,下面单独说。
sincedb_path:sincedb 文件的路径。设成 /dev/null 可以让每次启动都从 start_position 指定的位置开始,这在开发调试时有用,生产环境不应这样做。
sincedb_write_interval:多久把位点刷一次盘,默认 15 秒。这是 at-least-once 的代价:最坏情况下崩溃可能丢失 15 秒内已读但未落盘的位点,重启后会重复读这段内容。
ecs_compatibility:决定 file input 注入的来源元数据落在哪里。ECS 关闭时是扁平的 host 和 path,开启 v1/v8 后变成 [host][name] 和 [log][file][path]。本系列的版本前提与这个开关的取值约定见 01 篇,本文的实验输出以关闭态的字段名为准。
tail 模式与 read 模式
mode 是 file input 的顶层开关,默认 tail,另一个取值是 read。上面讲的位点推进、轮转检测、start_position 都是 tail 模式的行为:文件被当成没有终点的流,EOF 不代表结束,插件始终假设后面还会有内容。
read 模式把每个文件当成内容已经完整的有限流,EOF 在这里有明确语义——最后一行不需要等分隔符就能成 event,文件读完即关闭并移出活动窗口。这个语义差异解锁了三件 tail 模式做不到的事:处理 gzip 压缩文件、用 file_completed_action(delete / log / log_and_delete,默认 delete)在读完后删掉或记录文件路径、用 ignore_older 跳过修改时间过旧的文件。
代价是两个参数在 read 模式下被直接忽略,配了也不报错:start_position(read 模式永远从头读)和 close_older(EOF 即关闭,不再按"最后一次读取距今超过默认 1 小时"来关)。批量导入历史日志是 read 的场景,尾随线上日志是 tail 的场景。
文件轮转的两种方式
file input 按文件标识(inode 加设备号)而不是文件名追踪文件。两种常见轮转方式下,正确的配法不同,配错的失效方向也不同。
logrotate 默认的 rename 方式下,app.log 被改名成 app.log.1,随后创建新的 app.log。旧文件的标识不变,file input 会检测到这个标识现在对应新路径,把内部状态一起搬过去,旧内容不重读、改名后文件上的新增内容继续读。官方对此有一个前置要求:path 里必须同时包含轮转前和轮转后两个路径,例如 path => ["/var/log/app.log", "/var/log/app.log.1"]。只写 app.log 时,改名后的文件不在监控范围内,而写日志的进程在收到 reopen 信号之前往往还在往旧句柄里追加,这段内容就漏了。
copytruncate 的坑在另一头。原文件先被复制出一个副本,然后原文件被截断到零长度。截断会被检测到,该标识的位点重置为零、继续从头读新内容,这部分是自动的。需要人工保证的是副本路径不要进 path 的 glob:副本对 file input 是一个全新的文件标识,一旦被发现就会被当作新文件整份读一遍。
所以 rename 的失效模式是 path 配窄了漏读,copytruncate 的失效模式是 path 配宽了重复读。本文后面那个实验用的 path => "/tmp/demo.log" 只满足 copytruncate 的要求。
beats input 与 Lumberjack 协议
beats input 在一个 TCP 端口上接收来自 Filebeat、Metricbeat 等 Beats 家族工具的数据。port 是必填项,没有默认值;5044 是社区惯例而不是插件默认值,来源是 Filebeat 输出端的默认目标端口。传输协议是 Lumberjack v2,基于 TCP,支持 TLS,有序列号和批量 ACK 机制。
协议工作过程:
1 | |
ACK 由 Logstash beats input 在 event 进入队列之后才发送,而不是在 event 被 filter 处理完或写入 Elasticsearch 之后。这意味着 ACK 的语义是"进入 Logstash 内部队列",不是"端到端已处理完"。如果使用内存队列,Logstash 崩溃后队列内容丢失,但 Filebeat 已经前进了 registry,不会重发——这是内存队列在 beats 接入时无法提供 at-least-once 的根本原因。
换成持久队列后,event 进队列即落盘,ACK 发出时的状态有盘上记录。Logstash 重启后从持久队列恢复,加上 Filebeat 在没收到 ACK 的情况下会重发,两端结合才构成端到端的 at-least-once。
几个实践要点:
- TLS 用
ssl_enabled => true打开,再配ssl_certificate和ssl_key,生产环境应开启。 - 可以配置多个 beats input 监听不同端口,用于区分不同来源的 Beats 数据流。
- 吞吐量调优主要靠 Filebeat 的
bulk_max_size(每批事件数)和 Logstash 的pipeline.batch.size。
TLS 那一条有个存量配置的坑值得单独交代。这批参数在插件 7.0.0 做过一次改名,旧名字不是"弃用后仍可用",而是被列进 obsolete 清单:配置里带任意一个,插件启动阶段直接失败。它们通常成组出现在同一份旧配置里,迁移时会一起炸。
| 7.0.0 起失效的旧名 | 替代 |
|---|---|
ssl |
ssl_enabled |
ssl_verify_mode |
ssl_client_authentication(取值 none / optional / required) |
ssl_peer_metadata |
enrich(把 ssl_peer_metadata 作为一项加进列表) |
tls_min_version / tls_max_version |
ssl_supported_protocols(直接列协议版本,不再给上下界) |
cipher_suites |
ssl_cipher_suites |
kafka input 与消费组
kafka input 让 Logstash 作为 Kafka 消费者接入消息队列。核心配置是 topics、group_id、bootstrap_servers。
消费位点由 Kafka broker 侧的 consumer group 管理,而不是由 Logstash 本地的文件维护。消费组(group_id)下的每个分区都有一个 offset,标记"这个 partition 已消费到第几条"。至于 Logstash 在什么时刻把 offset 提交回 broker,由 enable_auto_commit 决定。
1 | |
enable_auto_commit:默认 true,此时提交由 Kafka 客户端在后台按 auto_commit_interval_ms(默认 5000)定时执行,与这批消息有没有进 Logstash 队列无关。改成 false 之后提交时机变成"每批数据写入内存队列或持久队列之后同步提交",不需要额外参数配套。
这个参数是本篇那条线索最干净的例子:位点介质没变(都在 broker 侧),可靠性级别却变了。默认值下崩溃可能丢掉"已提交但还没进队列"的那 5 秒消息;改成 false 后这个窗口消失,边界收紧到进队列,但仍不是"写入 ES 成功"。
consumer_threads:kafka input 内部的消费者线程数,决定可并行消费的分区数量。设置超过 topic 分区数无效,多余的线程空转。
decorate_events:控制是否把 Kafka 侧的元数据附到 event 上。8.x 是三值枚举,默认 none(不附加)、basic(附加 topic、consumer_group、partition、offset、key、timestamp)、extended(在 basic 之上再附加 record headers,仅限 UTF-8 编码的 header 值)。旧写法 true 和 false 只是 basic 和 none 的 deprecated 别名。
元数据落在 [@metadata][kafka] 下面,引用时用嵌套括号逐层写:
1 | |
写成 @metadata[kafka] 是非法的字段引用语法,02 篇讲过字段引用只有 [a][b] 一种写法。另一个容易忽略的点是 @metadata 下的字段在 output 阶段不会被发出,想让 topic 名进入 Elasticsearch 文档,必须像上面那样用 mutate 拷到普通字段上。
最小实验:file input 观察 sincedb
用一个最小配置观察 file input 的行为:
1 | |
实验步骤:
1 | |
把上述行为对应到内部对象:file input 维护了一个 FileWatch 组件,按 stat_interval(默认 1 秒)轮询受监控路径的 inode 和大小变化。发现新内容后通知 reader,reader 按 chunk 读取字节,交给 codec 切成 event,event 进入 queue。
模式提炼
1 | |
这个模式在 Kafka(offset commit)、Flink(checkpoint + source offset)、Spark Structured Streaming(offsets 写入 checkpoint 目录)里都有对应。
工程迁移表
| Logstash 概念 | Kafka 消费者对应 | Flink source 对应 | 通用 ETL 对应 |
|---|---|---|---|
| sincedb(file input 位点) | consumer offset | source state in checkpoint | watermark 文件 / cursor 表 |
| beats input ACK | producer ack + broker 确认 | checkpoint barrier | 消息确认回执 |
| kafka input consumer group | consumer group | source parallelism | 分片读取 |
| start_position=beginning | earliest offset reset | savepoint 回溯 | 全量重跑 |
| input 停止读取(背压) | consumer 减速 / 暂停分区 | source rate limiting | ETL 节流 |
| sincedb_write_interval | auto.commit.interval.ms | checkpoint interval | cursor 持久化频率 |
常见误解
误解一:“file input 可以保证不漏行”。漏读有两条独立成因。一条是文件标识本身不可靠:inode + 设备号在 NFS、部分容器 overlay fs 上的行为与本地 ext4/xfs 不同,重挂载还会改变设备号,sincedb 记录对不上。另一条是 path 没覆盖轮转后的路径,rename 之后那段追加内容根本不在监控范围内。加上 sincedb 落盘有间隔,实际语义是至少一次而不是精确一次。
误解二:“sincedb 里有这个文件的记录,它就永远不会被重头读一遍”。三种情况会让记录失效:sincedb_clean_after(默认两周)到期后记录不再持久化;path 的 glob 写法改了,sincedb 文件名跟着变,等于换了一份全新状态;sincedb_path => "/dev/null" 每次启动都是空状态。前两种在生产里都不需要人为操作就会发生,而后果都是同一份内容再进一次管道。
误解三:“beats input 收到数据就意味着 Logstash 已处理完”。ACK 仅表示 event 进入了 Logstash 的内部队列,不代表 filter 处理完或 output 已写入下游。下游写入失败时,event 仍可能在 Logstash 内部丢失(内存队列崩溃)或重发(持久队列重启)。
误解四:“kafka input 的 consumer_threads 设越大越快”。kafka input 的有效并发受 topic 分区数限制,consumer_threads 超过分区数后,多余线程空转不会提速。调大 consumer_threads 要同步调大 topic 分区数,同时注意 Logstash 的 pipeline.workers 也要跟上,否则 kafka input 产的快但 filter 消的慢,队列积压。
误解五:“只要用了 beats input,就自动实现了端到端不丢数据”。端到端 at-least-once 需要三件事同时满足:Filebeat 未收到 ACK 时会重发(Lumberjack 协议保证)、Logstash 持久队列兜底崩溃后重启、output 下游(如 Elasticsearch)写入使用幂等或去重机制。三个条件缺一不可。
练习
-
运行本文的 file input 实验,数一下 sincedb 实际有几列,确认第五列是不是一个浮点时间戳。再把
sincedb_path设成/dev/null重启,验证所有行从头重读,解释这个行为和start_position => "beginning"的关系。 -
把同一份配置的
mode改成read,同时把start_position留在end,观察文件是否仍然从头被读完一遍。然后加上file_completed_action => "log"和file_completed_log_path,确认读完之后被记录的是文件路径而不是内容。 -
用 Filebeat 向本地 beats input 发送几行日志(
port记得显式配),在rubydebug输出里找到host、agent、log等字段。判断每个字段是 Filebeat 注入的还是 beats input 的enrich加的,再试着用 mutate 把其中一个搬进[@metadata],验证它在 output 阶段消失了。 -
思考题:一个 Logstash 管道从 Kafka 读取,使用内存队列,output 写 Elasticsearch。Logstash 进程在 filter 处理完、output 写入前崩溃,会发生什么?
enable_auto_commit取true和false时,Kafka 的 offset 分别停在哪里?换成持久队列后这两个答案又如何变化?
系列导航
参考资料
- Logstash file input 文档:https://www.elastic.co/guide/en/logstash/current/plugins-inputs-file.html(sincedb、start_position 参数)
- Logstash beats input 文档:https://www.elastic.co/guide/en/logstash/current/plugins-inputs-beats.html(Lumberjack 协议、ACK 语义)
- Logstash kafka input 文档:https://www.elastic.co/guide/en/logstash/current/plugins-inputs-kafka.html(consumer group、offset 管理)
- Filebeat 文档 - How Filebeat works:https://www.elastic.co/guide/en/beats/filebeat/current/how-filebeat-works.html(registry 文件与 ACK 关系)
- Lumberjack v2 协议规范:https://github.com/elastic/logstash-forwarder/blob/master/PROTOCOL.md(序列号与批量 ACK)
- Logstash 持久队列文档:https://www.elastic.co/guide/en/logstash/current/persistent-queues.html(at-least-once 语义边界)
