深入 Logstash 04 - input 插件:拉取、监听与 Beats 接入
上一篇(03)讲清了 codec 作为字节与 event 边界转换器的角色。这一篇进入 input 插件层,回答一个集中问题:file input 如何用 sincedb 追踪读取位点、beats input 的 ack 机制如何实现 at-least-once、kafka input 的消费位点和消费组是怎么回事——以及这三种机制背后"拉模型 vs 推模型"的根本差异。
1 | |
拉模型与推模型
input 插件在获取数据时有两种根本不同的方向:
拉模型(pull):Logstash 主动去源端取数据。file input 定期轮询文件的 inode 和字节偏移;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 文件,默认写在 $HOME/.sincedb_*)。下次启动或文件被轮转时,Logstash 从 sincedb 读出上次停在哪里,从该偏移继续读,而不是从头重读。
1 | |
几个核心配置项:
start_position:控制 Logstash 第一次看到某个文件时从哪里开始。end(默认)表示只读新追加的内容,适合尾随实时日志;beginning 表示从文件开头读,适合处理历史文件。注意"第一次"是以 sincedb 里没有这个 inode 的记录为准,不是以文件是否存在为准。
sincedb_path:sincedb 文件的路径。设成 /dev/null 可以让每次启动都从 start_position 指定的位置开始——这在开发调试时有用,生产环境不应这样做。
sincedb_write_interval:多久把位点刷一次盘,默认 15 秒。这是 at-least-once 的代价:最坏情况下崩溃可能丢失 15 秒内已读但未落盘 sincedb 的位点,重启后会重复读这段内容。
文件轮转处理:file input 按 inode 而不是文件名追踪文件。logrotate 把 app.log 重命名成 app.log.1 并创建新 app.log 后,旧文件的 inode 不变,Logstash 会继续读完旧文件再切到新文件,不会漏读也不会因为文件名变了就停止。
beats input 与 Lumberjack 协议
beats input 监听 TCP 端口(默认 5044),接收来自 Filebeat、Metricbeat 等 Beats 家族工具的数据。传输协议是 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。
几个实践要点:
- beats input 支持
ssl => true加ssl_certificate、ssl_key配置 TLS,生产环境应开启。 port可以配置多个 beats input 监听不同端口,用于区分不同来源的 Beats 数据流。- 吞吐量调优主要靠 Filebeat 的
bulk_max_size(每批事件数)和 Logstash 的pipeline.batch.size。
kafka input 与消费组
kafka input 让 Logstash 作为 Kafka 消费者接入消息队列。核心配置是 topics、group_id、bootstrap_servers。
消费位点由 Kafka broker 侧的 consumer group 管理,而不是由 Logstash 本地的文件维护。消费组(group_id)下的每个分区都有一个 offset,标记"这个 partition 已消费到第几条"。Logstash kafka input 处理完一批 event(写入 Logstash 内部队列)后,提交这批消息对应的 offset 给 Kafka broker。
1 | |
enable_auto_commit:默认 true,由 Kafka 客户端后台自动提交 offset。设成 false 并配合 commit_offsets_on_completion 可以在 event 进入队列后手动提交,提供更精确的 at-least-once 语义(但仍以进队列为边界,不是以写入 ES 为边界)。
consumer_threads:kafka input 内部的消费者线程数,决定可并行消费的分区数量。设置超过 topic 分区数无效,多余的线程空转。
decorate_events:设为 true 时,kafka input 把 topic、partition、offset 等元数据写入 event 的 @metadata[kafka] 字段,可用于下游按 topic 路由或做精确去重。
最小实验: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 可以保证不漏行”。file input 依赖 inode 追踪,在某些文件系统(NFS、某些容器 overlay fs)上 inode 行为与本地 ext4/xfs 不同,可能导致 sincedb 失效。此外 sincedb 落盘有间隔,崩溃后可能重读最近一段内容(重复而非漏读)。不保证精确一次。
误解二:“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后重启 Logstash,验证所有行从头重读。解释这个行为和start_position => "beginning"的关系。 -
用 Filebeat 向本地 beats input(端口 5044)发送几行日志,在
rubydebug输出里找到host、agent、log等 Filebeat 自动注入的字段。思考这些字段是 Logstash 加的还是 Filebeat 加的,以及如何用@metadata把它们隔离出去不进入最终文档。 -
思考题:一个 Logstash 管道从 Kafka 读取,使用内存队列,output 写 Elasticsearch。Logstash 进程在 filter 处理完、output 写入前崩溃,会发生什么?换成持久队列后答案如何变化?这两种场景下 Kafka 的 offset 分别停在哪里?
系列导航
| 序号 | 主题 | 状态 |
|---|---|---|
| 00 | 导读:核心对象是 event,骨架是三段管道 | 已发布 |
| 01 | Logstash 架构:JRuby、JVM 与 pipeline 的运行形态 | 已发布 |
| 02 | event 模型:@timestamp、@metadata 与字段引用 |
已发布 |
| 03 | codec:字节流与 event 的边界转换 | 已发布 |
| 04 | input 插件:拉取、监听与 Beats 接入 | 本篇 |
| 05 | Grok 的本质:命名正则加预定义 pattern | 下一篇 |
| 06 | dissect 与结构化 filter:放弃回溯换吞吐 | |
| 07 | date、mutate、geoip:常用 filter 的精确用法 | |
| 08 | Elasticsearch output:bulk 写入与索引路由 | |
| 09-12 | 管道执行与可靠性(worker·batch·背压 / 队列 / DLQ / Multiple Pipelines) | |
| 13-14 | 运维、监控与调优 | |
| 15-17 | 演进、生态与对比 |
参考资料
- 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 语义边界)
