上一篇把 filter 阶段的常用加工链梳理完:mutate 改写字段、date 覆盖 @timestamp、geoip 扩展地理信息、条件块用 tag 做分支路由。这一篇进入 output 阶段最常见的目标:Elasticsearch output 插件,重点是批量写入机制、索引命名策略、重试语义和背压来源。

核心问题:Logstash 如何把一批 event 打包成 Elasticsearch bulk 请求;索引模板与数据流如何决定文档落在哪个物理存储;429/503 重试和 400 丢弃背后是什么语义;ES 的慢响应如何一路传导成 input 侧的减速。

数据流:event 从 filter 出口到 ES 的路径

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
filter workers


output queue (per-pipeline batch)

▼ (flush_size events or idle_flush_time elapsed)
elasticsearch output plugin

├── build bulk request body
│ action: index / create / update / delete
│ _index: logs-%{+YYYY.MM.dd}
│ _id: %{[@metadata][fingerprint]} (optional)
│ source: serialized event fields

├── HTTP POST /_bulk ──▶ Elasticsearch cluster
│ │
│ ┌──────────────────┐│
│ │ bulk response ││
│ │ per-item status ││
│ └──────────────────┘│
│ │
│◀── 429 / 503: retry │
│◀── 400: drop + DLQ │
└── 200 (partial): per-item check

output 阶段接收的是 filter worker 处理完的 event batch。batch 是 pipeline 的调度单位,大小由 pipeline.batch.size 控制(默认 125),不是由 elasticsearch output 插件单独决定的。插件把 batch 里的 event 序列化为 bulk request body,一次性 POST 到 ES 的 /_bulk 端点。

关键对象

1
2
3
4
5
6
7
8
9
10
11
12
bulk request  ── 多个操作的 NDJSON 序列:
每个操作 = 一行 action meta + 一行 source document
action 可以是 index / create / update / delete

bulk response ── per-item 结果数组:每个操作独立的 HTTP 状态码
整体请求 HTTP 200 不代表每条文档都写入成功

_index ── 文档落点:可以是静态字符串、字段插值、日期格式
_id ── 文档唯一键:不设则 ES 自动生成,设则支持幂等写

pipeline.batch.size ── filter worker 一次处理的最大 event 数
直接决定一个 bulk request 里最多有多少文档

bulk API 的攒批机制

Logstash 的 elasticsearch output 不是每收到一个 event 就发一次 HTTP 请求,而是攒够一批再发。攒批的边界有两个:

1
2
3
4
batch_size:       pipeline.batch.size(默认 125)
filter worker 攒满这么多 event 后把 batch 交给 output
idle_flush_time: elasticsearch output 的 idle_flush_time(默认 1 秒)
超过这个时间没有新 event 也强制 flush,防止低流量时数据积压

在 Logstash 8.x 的架构里,pipeline.batch.size 是控制单批大小的主旋钮。elasticsearch output 自身的 flush_size 参数在旧版本中存在,但在新版本架构里已经废弃或由 pipeline batch 取代——批次由 pipeline 调度层统一管理,output 插件只是消费整个 batch。

单个 bulk request 的大小(字节数)没有固定上限参数,但 ES 的 http.max_content_length(默认 100MB)是硬上限。实践中把 pipeline.batch.size 设为 500-2000、单文档平均 1-5KB 时,单次 bulk 一般在几 MB 以内,安全。

索引命名:静态、动态与日期格式

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
output {
elasticsearch {
hosts => ["https://es-node:9200"]
user => "logstash_writer"
password => "${LOGSTASH_ES_PASSWORD}"

# 静态索引名
# index => "my-logs"

# 按字段动态命名
# index => "%{[service][name]}-logs"

# 按事件时间分日索引(最常见)
index => "logs-%{+YYYY.MM.dd}"

# 按月分索引,减少 shard 数量
# index => "logs-%{+YYYY.MM}"
}
}

%{+YYYY.MM.dd} 里的 + 前缀告诉 Logstash 用 @timestamp 的值格式化日期,语法来自 Joda-Time。这就是 date filter 必须正确改写 @timestamp 的第二个理由——除了时序查询正确,还决定了文档落在哪天的索引分片里,索引命名错了等于把不同时间段的数据混存,ILM 的 rollover 和冷热分层就全部失效。

动态索引名(%{[service][name]}-logs)可以让不同来源的日志按字段值路由到不同索引,但有风险:如果字段值来自未受信的外部数据,恶意输入可以生成数以万计的索引,导致 ES cluster state 爆炸(index count 过高是 ES 的已知性能陷阱)。生产中动态索引名要在 Grok/mutate 阶段做白名单过滤或哈希归一化。

数据流(Data Stream)与索引模板

Logstash 8.x 起,elasticsearch output 默认启用 Data Stream 写入模式(data_stream => true),不再直接写具名索引,而是写 data stream:

1
2
3
4
5
6
7
8
9
10
output {
elasticsearch {
hosts => ["https://es-node:9200"]
data_stream => true
data_stream_type => "logs"
data_stream_dataset => "%{[service][name]}"
data_stream_namespace => "production"
# 生成的 data stream 名:logs-<dataset>-<namespace>
}
}

Data Stream 背后是一组自动滚动的 backing indices,ILM/DLM 根据策略(大小、时间、文档数)触发 rollover,老 index 自动进入 warm/cold/delete 阶段。Logstash 写入方只需要知道 data stream 的名字,索引的物理管理全部由 ES 侧策略接管。

索引模板(Index Template)在文档写入前就决定了字段映射(mapping)和设置(settings)。如果模板缺失或映射与字段类型不符,ES 的 dynamic mapping 会自动推断类型,但推断结果不稳定——同一个字段在不同文档里可能是 keyword 也可能是 long,导致映射冲突报错。生产部署必须在第一次写入前用 PUT _index_template/ 或 Kibana Index Management 配好模板。

重试语义:哪些错误重试,哪些丢弃

ES bulk response 里每个操作有独立的状态码:

1
2
3
4
5
201 / 200  ── 写入成功
429 ── Too Many Requests:ES 的写入队列已满,需要退让重试
503 ── Service Unavailable:节点不可用,短暂重试
400 ── Bad Request:文档格式或映射冲突,无法写入,视为不可恢复
5xx(其他)── 服务端错误,可重试

Logstash 的 elasticsearch output 对 429 和 503 做指数退避重试(默认重试次数由 retry_on_conflict 等参数控制,整体重试逻辑由 plugin 内部管理),重试期间这批 event 阻塞在 output 插件里,不会前进也不会后退到队列。

400 是不可重试错误——字段映射冲突、索引名非法、文档结构问题。默认行为是记录错误日志后丢弃这条 event。如果开启了 DLQ(dead_letter_queue.enable: true),400 导致失败的 event 会被写入死信队列,供后续离线修复(第 11 篇)。

action 参数控制写入行为:

1
2
3
4
5
6
7
8
# 默认 index:有则覆盖,无则创建(幂等写需要 document_id)
action => "index"

# create:文档存在则报错 409(防止覆盖)
action => "create"

# update:只更新现有文档(结合 doc_as_upsert)
action => "update"

at-least-once 场景下,同一条 event 可能被重复写入(pipeline 重启、网络重传)。要实现幂等写,需要给 document_id 设一个确定性的值:

1
2
3
4
5
6
7
output {
elasticsearch {
document_id => "%{[@metadata][fingerprint]}"
action => "index"
# ES 的 index 操作是幂等的:相同 _id 的写入会覆盖而非追加
}
}

fingerprint filter 可以对指定字段列表计算哈希,写入 @metadata,再由 elasticsearch output 用作 document_id。这条链路把 at-least-once 的重复写变成幂等写,是 Logstash 场景下实现"事实上的 exactly-once"的标准手法。

背压来源:ES 慢响应如何传导到 input

output 端的 ES bulk 请求是同步阻塞的(每个 output worker 线程发出请求后等待响应)。当 ES 响应变慢时:

1
2
3
4
5
6
7
8
ES 慢响应
→ elasticsearch output 阻塞等待
→ output 无法消费 filter worker 的输出
→ filter worker 的 batch 无法交给 output,停止处理新 batch
→ filter worker 停止从队列拉取新 event
→ 队列积压
→ input 无法把新 event 放入队列(队列满或背压信号)
→ input 减速甚至暂停读取

这条背压链路是 Logstash 流量控制的主干。pipeline.batch.delay(等待攒批的最大时间)、ES 连接池大小(pool_maxpool_max_per_route)、output worker 线程数(pipeline.output.workers)共同决定了背压信号从 ES 传导到 input 的速度和幅度。

理解背压的意义在于:当遇到"input 端读取突然变慢"的现象时,不要先查 input 插件的配置,而是先看 ES 的 bulk 响应时间和 429 频率,再顺着背压链路向上游追溯。

连接池与 sniffing

1
2
3
4
5
6
7
8
output {
elasticsearch {
hosts => ["es-node1:9200", "es-node2:9200", "es-node3:9200"]
sniffing => false # 默认 false;true 时自动发现集群所有节点
pool_max => 1000 # HTTP 连接池最大连接数
pool_max_per_route => 100
}
}

sniffing => true 让 Logstash 在启动时调用 ES 的 /_nodes API 获取所有数据节点地址,动态维护连接列表。优点是不需要在 hosts 里枚举所有节点,自动应对节点增减;缺点是如果 ES 节点绑定的 publish_address 是内网地址而 Logstash 在另一网段,sniffing 发现的地址可能不可达。sniffing 在容器/云环境里默认关闭并保持关闭,用 hosts 列出负载均衡地址或所有已知节点。

实验:观察 bulk 写入与重试行为

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
# logstash.conf
input {
stdin {}
}

filter {
mutate {
add_field => { "env" => "test" }
}
}

output {
elasticsearch {
hosts => ["http://localhost:9200"]
index => "logstash-test-%{+YYYY.MM.dd}"
action => "index"
}
stdout { codec => rubydebug }
}

把 ES 暂时停掉或改为不可达地址,输入几行文本,观察 Logstash 的日志输出:会看到 Retrying failed action 类型的日志,包含 HTTP 状态码和退避间隔。恢复 ES 后,Logstash 会把积压的重试 event 成功发出。这个实验直接观察到重试机制——积压的 event 留在 output 线程里等待,不会返还到队列,也不会被丢弃(直到超过最大重试次数)。

模式提炼

1
2
3
4
5
6
7
8
9
模式:批量 + 幂等 + 背压感知 output

- 以 pipeline batch 为单位攒批发送,摊薄 HTTP 往返开销
- 区分可重试错误(429/503)和不可重试错误(400),
前者退避重试,后者交给 DLQ 处理,不能一律重试也不能一律丢弃
- 用确定性 document_id 把 at-least-once 的重复写变成幂等覆盖
- ES 的慢响应会通过背压链路反向传导到 input,
output 的阻塞时间是整条 pipeline 吞吐的上限
- 索引命名要与 @timestamp 的事件时间绑定,不能用处理时间替代

工程迁移表

Elasticsearch output 概念 Kafka 生态对应 Flink 对应 通用 ETL 对应
bulk request(攒批 POST) Producer batch + linger.ms checkpoint 触发 sink flush 批量 INSERT / COPY
pipeline.batch.size batch.size + buffer.memory sink buffer size 批次行数
idle_flush_time linger.ms sink idle timeout 超时 flush
429 退避重试 Producer retry + max.block.ms sink 重试策略 写入重试 + 退避
400 → DLQ 解析失败 → dead letter topic side output error stream 错误记录表
document_id 幂等写 Producer idempotence + exactly-once upsert sink MERGE / ON CONFLICT DO UPDATE
背压传导链 Producer block → Consumer slow watermark/checkpoint 背压 上游限速
sniffing Kafka 的 metadata.fetch.timeout 动态 JobManager 发现 连接池健康检查
index template Kafka topic schema registry Flink table DDL 目标表 DDL
data stream + ILM Kafka topic retention Flink state TTL 分区裁剪 + 归档

常见误解

误解一:“bulk response 返回 HTTP 200 说明所有文档都写成功了”。bulk API 的 HTTP 200 只代表请求本身被服务端接收并处理,响应体里的 items 数组里每个操作有独立的 status。只有遍历 items 检查到没有 error 字段,才能确认所有文档写入成功。Logstash 的 elasticsearch output 插件内部会做这个检查,不是说把 HTTP 状态码当结论。

误解二:“pipeline.batch.size 越大写入越快”。batch size 越大,单次 bulk 请求的文档数越多,HTTP 往返次数越少,吞吐在一定范围内确实上升。但 batch size 过大有两个反效果:单次 bulk body 体积超过 ES 的 http.max_content_length 会报 413;filter worker 需要先攒满 batch 才交给 output,低流量时 idle_flush_time 之前数据卡在内存,端到端延迟上升。实践中从 250-500 开始调,观察 bulk 响应时间和 ES indexing 压力,不要直接设到几千。

误解三:“关掉 sniffing 就必须在 hosts 里列出所有 ES 节点”。hosts 里只需要列出能路由到 ES cluster 的入口地址,可以是负载均衡器、代理、或者几个已知数据节点。关掉 sniffing 的含义是"不自动发现其他节点",不是"只能写入 hosts 里列出的节点"——Logstash 把请求发给 hosts 里的节点,ES 集群内部路由负责把文档分发到正确的 primary shard。

误解四:“Logstash 的 at-least-once 只要开了持久队列就够了”。持久队列保证的是 event 进入队列之后不因崩溃丢失,但 output 端的写入失败(网络中断、ES 重启)如果超过了 output 的最大重试次数,event 仍然会丢失。完整的 at-least-once 需要:持久队列兜住崩溃边界 + output 端足够的重试次数/窗口 + DLQ 接住最终失败的 event + DLQ 的离线重放能力。缺少任何一环,"at-least-once"就只是部分成立。

练习

  1. 本地起一个 Logstash + Elasticsearch,用 stdin input 发几条 event。停掉 ES,继续向 stdin 发数据,观察 Logstash 日志里的重试行为;重新启动 ES,观察积压的 event 是否被成功发出。

  2. 在 elasticsearch output 里设置 document_id => "%{message}"(用消息内容作 ID),向 ES 发送 10 条内容相同的 event,查询 ES 确认文档数量是 1 而不是 10。解释为什么 action => "index" 配合固定 document_id 是幂等写。

  3. index 改成 "logs-%{+YYYY.MM.dd}",但不用 date filter 改写 @timestamp,向 stdin 发一条带历史时间戳的日志文本(但不解析它)。观察文档落在哪个索引(预期是今天的索引,而不是日志内容里的历史日期)。再加上 date filter 改写 @timestamp,重复实验,确认文档落到了正确的历史日期索引。

系列导航

序号 主题 状态
00 导读:核心对象是 event,骨架是三段管道 已发布
01 Logstash 架构: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-14 运维、监控与调优
15-17 演进、生态与对比

参考资料