深入 Logstash 08 - output 插件:Elasticsearch output 与批量写入
上一篇把 filter 阶段的常用加工链梳理完:mutate 改写字段、date 覆盖 @timestamp、geoip 扩展地理信息、条件块用 tag 做分支路由。这一篇进入 output 阶段最常见的目标:Elasticsearch output 插件,重点是批量写入机制、索引命名策略、重试语义和背压来源。
核心问题:Logstash 如何把一批 event 打包成 Elasticsearch bulk 请求;索引模板与数据流如何决定文档落在哪个物理存储;429/503 重试和 400 丢弃背后是什么语义;ES 的慢响应如何一路传导成 input 侧的减速。
数据流:event 从 filter 出口到 ES 的路径
1 | |
output 阶段接收的是 filter worker 处理完的 event batch。batch 是 pipeline 的调度单位,大小由 pipeline.batch.size 控制(默认 125),不是由 elasticsearch output 插件单独决定的。插件把 batch 里的 event 序列化为 bulk request body,一次性 POST 到 ES 的 /_bulk 端点。
关键对象
1 | |
bulk API 的攒批机制
Logstash 的 elasticsearch output 不是每收到一个 event 就发一次 HTTP 请求,而是攒够一批再发。攒批的边界有两个:
1 | |
在 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 | |
%{+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 | |
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 | |
Logstash 的 elasticsearch output 对 429 和 503 做指数退避重试(默认重试次数由 retry_on_conflict 等参数控制,整体重试逻辑由 plugin 内部管理),重试期间这批 event 阻塞在 output 插件里,不会前进也不会后退到队列。
400 是不可重试错误——字段映射冲突、索引名非法、文档结构问题。默认行为是记录错误日志后丢弃这条 event。如果开启了 DLQ(dead_letter_queue.enable: true),400 导致失败的 event 会被写入死信队列,供后续离线修复(第 11 篇)。
action 参数控制写入行为:
1 | |
at-least-once 场景下,同一条 event 可能被重复写入(pipeline 重启、网络重传)。要实现幂等写,需要给 document_id 设一个确定性的值:
1 | |
fingerprint filter 可以对指定字段列表计算哈希,写入 @metadata,再由 elasticsearch output 用作 document_id。这条链路把 at-least-once 的重复写变成幂等写,是 Logstash 场景下实现"事实上的 exactly-once"的标准手法。
背压来源:ES 慢响应如何传导到 input
output 端的 ES bulk 请求是同步阻塞的(每个 output worker 线程发出请求后等待响应)。当 ES 响应变慢时:
1 | |
这条背压链路是 Logstash 流量控制的主干。pipeline.batch.delay(等待攒批的最大时间)、ES 连接池大小(pool_max、pool_max_per_route)、output worker 线程数(pipeline.output.workers)共同决定了背压信号从 ES 传导到 input 的速度和幅度。
理解背压的意义在于:当遇到"input 端读取突然变慢"的现象时,不要先查 input 插件的配置,而是先看 ES 的 bulk 响应时间和 429 频率,再顺着背压链路向上游追溯。
连接池与 sniffing
1 | |
sniffing => true 让 Logstash 在启动时调用 ES 的 /_nodes API 获取所有数据节点地址,动态维护连接列表。优点是不需要在 hosts 里枚举所有节点,自动应对节点增减;缺点是如果 ES 节点绑定的 publish_address 是内网地址而 Logstash 在另一网段,sniffing 发现的地址可能不可达。sniffing 在容器/云环境里默认关闭并保持关闭,用 hosts 列出负载均衡地址或所有已知节点。
实验:观察 bulk 写入与重试行为
1 | |
把 ES 暂时停掉或改为不可达地址,输入几行文本,观察 Logstash 的日志输出:会看到 Retrying failed action 类型的日志,包含 HTTP 状态码和退避间隔。恢复 ES 后,Logstash 会把积压的重试 event 成功发出。这个实验直接观察到重试机制——积压的 event 留在 output 线程里等待,不会返还到队列,也不会被丢弃(直到超过最大重试次数)。
模式提炼
1 | |
工程迁移表
| 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"就只是部分成立。
练习
-
本地起一个 Logstash + Elasticsearch,用 stdin input 发几条 event。停掉 ES,继续向 stdin 发数据,观察 Logstash 日志里的重试行为;重新启动 ES,观察积压的 event 是否被成功发出。
-
在 elasticsearch output 里设置
document_id => "%{message}"(用消息内容作 ID),向 ES 发送 10 条内容相同的 event,查询 ES 确认文档数量是 1 而不是 10。解释为什么action => "index"配合固定document_id是幂等写。 -
把
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 | 演进、生态与对比 |
参考资料
- Logstash Elasticsearch output 文档:https://www.elastic.co/guide/en/logstash/current/plugins-outputs-elasticsearch.html(参数全集与 data stream 配置)
- Elasticsearch Bulk API 文档:https://www.elastic.co/guide/en/elasticsearch/reference/current/docs-bulk.html(请求格式与 per-item 响应结构)
- Elasticsearch Index Templates 文档:https://www.elastic.co/guide/en/elasticsearch/reference/current/index-templates.html(映射与设置的预先声明)
- Elasticsearch Data Streams 文档:https://www.elastic.co/guide/en/elasticsearch/reference/current/data-streams.html(data stream 写入模式与 ILM 集成)
- Logstash pipeline 配置参考:https://www.elastic.co/guide/en/logstash/current/logstash-settings-file.html(pipeline.batch.size、pipeline.workers 等调优参数)
- Logstash fingerprint filter 文档:https://www.elastic.co/guide/en/logstash/current/plugins-filters-fingerprint.html(document_id 幂等写的 fingerprint 生成)
