input 层能给出多强的可靠性,由一件事决定:"读到哪里"这个状态存放在哪里、在什么时刻推进。file input 把它写进本地 sincedb 文件;kafka input 交给 broker 侧的 consumer group;beats input 自己不存,靠 Filebeat 的 registry 配合一次 ACK 来推进。三种介质完全不同,推进时刻却停在同一条线上:event 进入 Logstash 内部队列,而不是写进下游成功之后。

这一篇沿这条线索展开三件事:sincedb 到底记了哪几列、beats 的 ACK 语义边界在哪里、kafka 的 offset 提交时机由哪个参数控制。顺带回答一个更靠前的问题——拉模型和推模型的分野,落到可靠性上究竟差在哪。

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

pipeline workers (N):filter + output

拉模型与推模型

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

例: 12345678 8 1 102400 1786000000.123 /var/log/app.log
含义:该文件标识已读到字节偏移 102400,
最后活跃时间是那个浮点时间戳,上次匹配到的路径是 /var/log/app.log

第五列的浮点时间戳不是给人看的,它驱动 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 关闭时是扁平的 hostpath,开启 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_actiondelete / 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
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。

几个实践要点:

  • TLS 用 ssl_enabled => true 打开,再配 ssl_certificatessl_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 消费者接入消息队列。核心配置是 topicsgroup_idbootstrap_servers

消费位点由 Kafka broker 侧的 consumer group 管理,而不是由 Logstash 本地的文件维护。消费组(group_id)下的每个分区都有一个 offset,标记"这个 partition 已消费到第几条"。至于 Logstash 在什么时刻把 offset 提交回 broker,由 enable_auto_commit 决定。

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 客户端在后台按 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 值)。旧写法 truefalse 只是 basicnone 的 deprecated 别名。

元数据落在 [@metadata][kafka] 下面,引用时用嵌套括号逐层写:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
input {
kafka {
bootstrap_servers => "kafka-1:9092"
topics => ["logs"]
group_id => "ls-prod"
decorate_events => "basic"
}
}
filter {
# 正确:逐层括号
if [@metadata][kafka][partition] == 0 { ... }

# @metadata 不随 output 发出,要进最终文档必须显式拷出来
mutate {
add_field => { "kafka_topic" => "%{[@metadata][kafka][topic]}" }
}
}

写成 @metadata[kafka] 是非法的字段引用语法,02 篇讲过字段引用只有 [a][b] 一种写法。另一个容易忽略的点是 @metadata 下的字段在 output 阶段不会被发出,想让 topic 名进入 Elasticsearch 文档,必须像上面那样用 mutate 拷到普通字段上。

最小实验: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
18
# 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 1786000000.123 /tmp/demo.log
# 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
9
10
11
12
13
14
模式:位点外置,让 input 可重启

- 把"读到哪里"的状态写到外部持久存储(sincedb 文件、Kafka offset、Filebeat registry)
- 重启后从外部状态恢复,实现"至少一次"读取
- 位点推进时机决定可靠性级别,而这个时机通常是某个参数的默认值,不是引擎的固有性质:
按时钟定时推进 → 崩溃丢掉一个时间窗口的进度
file : sincedb_write_interval 默认 15 秒
kafka : enable_auto_commit 默认 true + auto_commit_interval_ms 默认 5000
进队列后同步推进 → 时间窗口消失,边界收紧到"进了 Logstash"
kafka : enable_auto_commit => false
beats : ACK 在 event 进队列后发出
写入下游确认后推进 → 端到端更强,但要求源端支持回溯
Logstash input 层不提供这一档,只能靠源端(如 Kafka)自己保留
- 背压由队列容量控制:队列满时 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 可以保证不漏行”。漏读有两条独立成因。一条是文件标识本身不可靠: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)写入使用幂等或去重机制。三个条件缺一不可。

练习

  1. 运行本文的 file input 实验,数一下 sincedb 实际有几列,确认第五列是不是一个浮点时间戳。再把 sincedb_path 设成 /dev/null 重启,验证所有行从头重读,解释这个行为和 start_position => "beginning" 的关系。

  2. 把同一份配置的 mode 改成 read,同时把 start_position 留在 end,观察文件是否仍然从头被读完一遍。然后加上 file_completed_action => "log"file_completed_log_path,确认读完之后被记录的是文件路径而不是内容。

  3. 用 Filebeat 向本地 beats input 发送几行日志(port 记得显式配),在 rubydebug 输出里找到 hostagentlog 等字段。判断每个字段是 Filebeat 注入的还是 beats input 的 enrich 加的,再试着用 mutate 把其中一个搬进 [@metadata],验证它在 output 阶段消失了。

  4. 思考题:一个 Logstash 管道从 Kafka 读取,使用内存队列,output 写 Elasticsearch。Logstash 进程在 filter 处理完、output 写入前崩溃,会发生什么?enable_auto_committruefalse 时,Kafka 的 offset 分别停在哪里?换成持久队列后这两个答案又如何变化?

系列导航

序号 主题
00 导读:核心对象是 event,骨架是三段管道
01 架构:JRuby、JVM 与 pipeline 的运行形态
02 event 模型:@timestamp、@metadata 与字段引用
03 codec:字节流与 event 的边界转换
04 input 插件:拉取、监听与 Beats 接入(本篇)
05 Grok 的本质:命名正则加预定义 pattern
06 dissect 与结构化 filter:放弃回溯换吞吐
07 常用 filter 组合:mutate、date、geoip 与条件
08 output 插件:Elasticsearch output 与批量写入
09 pipeline 执行模型:worker、batch 与背压
10 内存队列 vs 持久队列:可靠性的分界线
11 死信队列(DLQ):无法处理的 event 去哪
12 Multiple Pipelines 与 pipeline-to-pipeline
13 监控:Node Stats API、hot threads 与瓶颈定位
14 性能调优:JVM heap、批处理与持久队列磁盘
15 Logstash vs Beats vs Ingest Pipeline:该用谁
16 Logstash vs Fluentd vs Vector:日志管道的三种取舍
17 Logstash 的演进与 Elastic Agent 的冲击

参考资料