上一篇(03)讲清了 codec 作为字节与 event 边界转换器的角色。这一篇进入 input 插件层,回答一个集中问题:file input 如何用 sincedb 追踪读取位点、beats input 的 ack 机制如何实现 at-least-once、kafka input 的消费位点和消费组是怎么回事——以及这三种机制背后"拉模型 vs 推模型"的根本差异。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
┌──────────────────────────────────────────────────────────┐
│ INPUT 层 │
│ │
│ PULL 模型 PUSH / LISTEN 模型 │
│ file ──┐ beats ──┐ │
│ jdbc ──┤ codec → tcp ──┤ codec → │
│ │ event udp ──┤ event │
│ │ kafka ──┘ │
└─────────┼───────────────────────────┘ │
▼ ▼ │
┌─────────────────────────────────────────────────────────┐
│ queue (memory / persistent) │
│ queue.full → input 减速 / 停止读取(背压) │
└─────────────────────────────────────────────────────────┘

filter workers (N)

拉模型与推模型

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
2
3
4
5
sincedb 结构(每行一个文件):
<inode> <dev_major> <dev_minor> <pos>

例: 12345678 8 1 102400
含义:inode 12345678 的文件,已读到字节偏移 102400

几个核心配置项:

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 而不是文件名追踪文件。logrotateapp.log 重命名成 app.log.1 并创建新 app.log 后,旧文件的 inode 不变,Logstash 会继续读完旧文件再切到新文件,不会漏读也不会因为文件名变了就停止。

beats input 与 Lumberjack 协议

beats input 监听 TCP 端口(默认 5044),接收来自 Filebeat、Metricbeat 等 Beats 家族工具的数据。传输协议是 Lumberjack v2,基于 TCP,支持 TLS,有序列号和批量 ACK 机制。

协议工作过程:

1
2
3
4
5
6
7
8
9
Filebeat                               Logstash beats input
│ │
│── batch(seq=1..N, events) ──────────────▶│
│ │ codec → event ─▶ queue
│ │ queue 接收完成
│◀── ACK(seq=N) ──────────────────────────│
│ │
│ 收到 ACK 后 Filebeat 推进本地 registry │
│ (registry 记录读到的文件位点) │

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 => truessl_certificatessl_key 配置 TLS,生产环境应开启。
  • port 可以配置多个 beats input 监听不同端口,用于区分不同来源的 Beats 数据流。
  • 吞吐量调优主要靠 Filebeat 的 bulk_max_size(每批事件数)和 Logstash 的 pipeline.batch.size

kafka input 与消费组

kafka input 让 Logstash 作为 Kafka 消费者接入消息队列。核心配置是 topicsgroup_idbootstrap_servers

消费位点由 Kafka broker 侧的 consumer group 管理,而不是由 Logstash 本地的文件维护。消费组(group_id)下的每个分区都有一个 offset,标记"这个 partition 已消费到第几条"。Logstash kafka input 处理完一批 event(写入 Logstash 内部队列)后,提交这批消息对应的 offset 给 Kafka broker。

1
2
3
4
5
6
7
Kafka broker
topic: logs
partition 0: offset 0..1000 → Logstash consumer (group_id=ls-prod)
partition 1: offset 0..800 → Logstash consumer (group_id=ls-prod)

Logstash kafka input:
poll(timeout) → 拿一批消息 → codec → events → queue → commit offset

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
2
3
4
5
6
7
8
9
10
11
12
13
# file-input-demo.conf
input {
file {
path => "/tmp/demo.log"
start_position => "beginning"
sincedb_path => "/tmp/sincedb-demo"
codec => plain { charset => "UTF-8" }
}
}
filter { }
output {
stdout { codec => rubydebug }
}

实验步骤:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
# 1. 准备日志文件
echo "line one" >> /tmp/demo.log
echo "line two" >> /tmp/demo.log

# 2. 启动 Logstash
bin/logstash -f file-input-demo.conf

# 3. 观察输出:两行 event,message 字段分别是 "line one" 和 "line two"

# 4. 查看 sincedb 内容(启动后约 15 秒落盘)
cat /tmp/sincedb-demo
# 输出类似: 12345678 8 1 18
# 18 = "line one\n"(9) + "line two\n"(9)

# 5. 不停 Logstash,继续追加
echo "line three" >> /tmp/demo.log
# 观察 Logstash 自动读取新行并输出

把上述行为对应到内部对象:file input 维护了一个 FileWatch 组件,按 stat_interval(默认 1 秒)轮询受监控路径的 inode 和大小变化。发现新内容后通知 reader,reader 按 chunk 读取字节,交给 codec 切成 event,event 进入 queue。

模式提炼

1
2
3
4
5
6
7
8
模式:位点外置,让 input 可重启

- 把"读到哪里"的状态写到外部持久存储(sincedb 文件、Kafka offset、Filebeat registry)
- 重启后从外部状态恢复,实现"至少一次"读取
- 位点推进时机决定可靠性级别:
进内存队列后推进 → 崩溃时可能重读,下游需幂等
写入下游确认后推进 → 端到端更强,但需要源端支持回溯
- 背压由队列容量控制:队列满时 input 停止生产,压力自然上传

这个模式在 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)写入使用幂等或去重机制。三个条件缺一不可。

练习

  1. 运行本文的 file input 实验,观察 sincedb 文件的内容变化。把 sincedb_path 设成 /dev/null 后重启 Logstash,验证所有行从头重读。解释这个行为和 start_position => "beginning" 的关系。

  2. 用 Filebeat 向本地 beats input(端口 5044)发送几行日志,在 rubydebug 输出里找到 hostagentlog 等 Filebeat 自动注入的字段。思考这些字段是 Logstash 加的还是 Filebeat 加的,以及如何用 @metadata 把它们隔离出去不进入最终文档。

  3. 思考题:一个 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 演进、生态与对比

参考资料