深入 Logstash 08 - output 插件:Elasticsearch output 与批量写入
上一篇把 filter 阶段的常用加工链梳理完:mutate 改写字段、date 覆盖 @timestamp、geoip 扩展地理信息、条件块用 tag 做分支路由。这一篇进入 output 阶段最常见的目标:Elasticsearch output 插件,重点是批量写入机制、索引命名策略、重试语义和背压来源。
核心问题:Logstash 如何把一批 event 打包成 Elasticsearch bulk 请求;索引模板与 data stream 如何决定文档落在哪个物理存储;哪些失败会被重试、哪些会真的丢掉 event;ES 的慢响应如何一路传导成 input 侧的减速。
数据流:event 从 filter 出口到 ES 的路径
1 | |
图里最上面那一行是这套机制最容易记错的地方:output 并没有一组独立线程在消费 filter 的产出。同一个 pipeline worker 线程先跑完 batch 的 filter 阶段,紧接着在自己身上跑 output 阶段,两者之间没有生产者与消费者的交接。这个模型是后面"背压来源"一节成立的前提,完整论证在第 09 篇。
batch 是 pipeline 的调度单位,大小由 pipeline.batch.size 控制(默认 125),不由 elasticsearch output 插件单独决定。插件把 batch 里的 event 序列化为 bulk request body,一次性 POST 到 ES 的 /_bulk 端点。
关键对象
1 | |
_index(文档落点)、_id(幂等写的唯一键)、pipeline.batch.size(攒批上限)各有专门小节,下面依次展开。
bulk API 的攒批机制
Logstash 的 elasticsearch output 不是每收到一个 event 就发一次 HTTP 请求,而是攒够一批再发。攒批的两个边界都不在 output 插件里,而在 pipeline 调度层:
1 | |
这两个旋钮写在 logstash.yml 或 pipelines.yml 里,elasticsearch { } 块内没有对应参数。想调攒批行为就得改 pipeline 设置,改 output 配置是无效的。
flush_size 和 idle_flush_time 是 5.x 时代 output 插件自己的攒批参数,随着批处理职责下移到 pipeline 层,两者已在 8.0 被标记 obsolete 并在后续版本移除。老博客和老配置里还常见,但现在写进配置 Logstash 会报 Unknown setting 并拒绝启动,不是"不生效"而是"起不来"。
50 毫秒这个数字经常被误记成"1 秒"级别,量级差了二十倍。低流量场景下端到端延迟的下界由它决定,估算"数据最多卡在内存里多久"时用错这个数字,结论会差出一个数量级。
单个 bulk request 的字节数没有独立上限参数,ES 侧的 http.max_content_length(默认 100MB)是硬上限。按单文档平均 1-5KB 估算,即使 pipeline.batch.size 开到 2000,单次 bulk body 也只有 10MB 上下,离 100MB 很远。所以 batch.size 的实际约束几乎从来不是这个字节上限,具体取舍见下文"误解二"。
索引命名:静态、动态与日期格式
8.x 的 data_stream 默认值是 auto,配置符合 data stream 条件时会自动走 data stream 而不是具名索引。本节所有示例都显式写了 index =>,而 index 恰好是让 auto 退回具名索引写入的条件之一,所以这些例子写的确实是具名索引。判定规则见下一节。
1 | |
%{+YYYY.MM.dd} 里的 + 前缀告诉 Logstash 用 @timestamp 的值格式化日期,语法来自 Joda-Time。这就是 date filter 必须正确改写 @timestamp 的第二个理由——除了时序查询正确,还决定了文档落在哪天的索引分片里,索引命名错了等于把不同时间段的数据混存,ILM 的 rollover 和冷热分层就全部失效。
动态索引名(%{[service][name]}-logs)可以让不同来源的日志按字段值路由到不同索引,但有风险:如果字段值来自未受信的外部数据,恶意输入可以生成数以万计的索引,导致 ES cluster state 爆炸(index count 过高是 ES 的已知性能陷阱)。生产中动态索引名要在 Grok/mutate 阶段做白名单过滤或哈希归一化。
日期插值恒按 UTC 执行
上一篇强调过时区是 date filter 的生产常见坑,同一个主题在索引命名这里还有一层:%{+...} 的日期格式化恒定按 UTC 执行。StringInterpolation 里对 Joda 分支显式做了 DateTimeFormat.forPattern(...).withZone(DateTimeZone.UTC),插值用的时间实例也按 UTC 构造。这跟 date filter 里配的 timezone 无关,跟服务器本地时区也无关。timezone 决定日志里的时间字符串怎么解析成 @timestamp,插值决定 @timestamp 怎么渲染成索引名。两者在不同阶段起作用。
对 Asia/Shanghai 的团队,后果是索引在北京时间早上 8 点滚动:每天 00:00 到 08:00 的日志会落进前一天日期的索引里。排查凌晨的问题时按本地日期猜索引名,很容易翻错一个索引。
还有一个 java.time 风格的等价写法,用双花括号,同样固定 UTC:
1 | |
真要按本地时间日切,只能在 filter 阶段自己算出日期字段(比如用 ruby filter 把 @timestamp 换算到目标时区后写进 [index_date]),再用 index => "logs-%{[index_date]}" 插值。没有任何配置项能改 %{+...} 的时区。
数据流(Data Stream)与索引模板
data_stream 有三个取值:true / false / auto。7.x 默认 false,8.0 起默认 auto,不是 true。auto 的语义是条件式的:只有当配置本身与 data stream 兼容时才走 data stream,否则退回具名索引写入。
判定条件可以在插件源码 data_stream_support.rb 的 invalid_data_stream_params 里逐条读到,归纳起来是:出现 index、document_id、routing、pipeline 这几个参数,或者 action 不是 create,或者 manage_template => false,都算与 data stream 不兼容。此外官方还要求 ecs_compatibility 必须是 v1 或 v8。设成 disabled 时 data stream 无法正常工作,显式写 data_stream => true 再配上 disabled 会直接抛 ConfigurationError。
回头看本文前面的例子:索引命名那节写了 index =>,幂等写那节写了 document_id =>,两者都会让 auto 判定为不兼容,于是退回具名索引。所以那些示例的行为是确定的,只是这个"确定"来自 auto 的退回逻辑,而不是因为 data stream 默认关闭。要真正写 data stream,配置里就不能出现这些参数:
1 | |
Data Stream 背后是一组自动滚动的 backing indices,ILM/DLM 根据策略(大小、时间、文档数)触发 rollover,老 index 自动进入 warm/cold/delete 阶段。Logstash 写入方只需要知道 data stream 的名字,索引的物理管理全部由 ES 侧策略接管。
正因为滚动交给了 ES,官方对 index 参数挂了一条 WARNING:date math(也就是 %{+yyyy.MM.dd} 这类插值)不建议与 alias 或 data stream 叠加使用。两套滚动机制放在一起,rollover 的判定基准会互相干扰。选一边:要么自己用日期索引名 + ILM 管 alias,要么交给 data stream,不要在 data stream 名字里再插日期。
索引模板(Index Template)在文档写入前就决定了字段映射(mapping)和设置(settings)。如果模板缺失或映射与字段类型不符,ES 的 dynamic mapping 会自动推断类型,但推断结果不稳定——同一个字段在不同文档里可能是 keyword 也可能是 long,导致映射冲突报错。生产部署必须在第一次写入前用 PUT _index_template/ 或 Kibana Index Management 配好模板。
重试语义:哪些错误重试,哪些丢弃
这套语义在 8.1.1 做过一次显著调整,而且必须区分两个层次:HTTP 请求层(整个 bulk 请求的响应码)和文档层(bulk response items 数组里每条操作各自的状态码)。两层的处理策略完全不同,混在一起看就会得出错误结论。
HTTP 请求层的规则极其简单:bulk API 只接受 200,所有其他响应码都无限重试。429、503、5xx 是这样,413(Payload Too Large)也是这样。这里不存在"最大重试次数"这个概念,所以也不存在"重试次数用完就丢弃"。重试期间这批 event 阻塞在 output 线程里,不前进也不退回队列。
413 也在无限重试之列,这是个反直觉的后果。batch 配得过大、bulk body 撑爆 http.max_content_length 时,管道不会快速失败,而是卡在一个永远不可能成功的重试循环里,因为每次重试发的还是同一个过大的请求体。想让它落地而不是死循环,必须显式配置:
1 | |
退避节奏由两个参数控制,跟"重试几次"无关:
1 | |
文档层的规则要细一些。把两层合起来,完整的失败信号表是这样:
| 信号 | 层次 | 触发条件 | event 去向 |
|---|---|---|---|
| 200 | HTTP | bulk 请求被正常处理 | 继续逐条查 items[].status |
| 429 | HTTP | ES 写入队列已满 | 无限重试,退避 2 秒起翻倍至 64 秒 |
| 503 | HTTP | 节点不可用 | 同上 |
| 413 | HTTP | bulk body 超过 http.max_content_length |
同上,且永远不会成功;配 dlq_custom_codes => [413] 才能落 DLQ |
| 其余非 200 | HTTP | 任何其他响应码 | 同上 |
| 200 / 201 | 文档 | 该条写入成功 | 完成 |
| 400 | 文档 | 文档结构问题、索引名非法 | 进 DLQ;未开 DLQ 则记日志后丢弃 |
| 404 | 文档 | mapping error(官方 DLQ policy 的归类) | 进 DLQ;未开 DLQ 则记日志后丢弃 |
| 409 | 文档 | 版本冲突,或 create 撞上已存在的 id |
记 warning 后丢弃,不重试、不进 DLQ |
| 其余非 2xx | 文档 | 该条写入失败 | 重试 |
| action 解析越界 | 插件 | sprintf 出的 action 不在四个合法值里 |
进 DLQ;未开 DLQ 则记日志后丢弃 |
表里真正会丢数据的只有两种情形:未开 DLQ 时的 400/404,以及任何配置下的 409。这跟上一篇 filter 阶段的失败语义正好形成对照。那边一律"打个 tag、event 继续走",这边的失败是终局的。
DLQ 收的是 400 和 404 两类文档级错误,写死在插件源码 common.rb 的 DOC_DLQ_CODES = [400, 404] 里,开关是 dead_letter_queue.enable: true,落盘后可离线修复重放(第 11 篇)。上一节讲的映射冲突按官方 DLQ policy 归在 404 而不是 400,这是最容易记反的一格:直觉上"字段类型不对"像是 400 Bad Request,实际归在 404。
409 是三类文档级错误里唯一不进 DLQ 的。官方给的建议是调高 retry_on_conflict,让 ES 自己重试比插件重试更高效。
retry_on_conflict(默认 1)不是 429/503 的重试次数控制,这是一处高频误用。它是随 bulk 请求传给 ES 的参数,含义是"ES 内部对一个 update/upsert 文档重试多少次",只与 409 版本冲突相关。想调 429/503 的退避行为,改的是上面那两个 retry_*_interval。
action 参数控制写入行为:
1 | |
action 用 sprintf 动态取值时有个边界:解析出来的字符串不在 index / create / update / delete 四个合法值里,这条 event 不会被发往 ES,而是进 DLQ(未开 DLQ 则记日志后丢弃)。字段缺失或者拼错一个字母都会走到这条路上。
action => "create" 值得结合上一节的 409 语义再看一遍。用 create 做去重是常见手法,靠"id 已存在就写不进去"来防重复。但撞上已存在的 id 时 ES 返回 409,而 409 在 8.1.1 之后是记 warning 后丢弃,不重试、不进 DLQ。也就是说去重成功的那些 event 就这么消失了,日志里只有一行 warning。如果这些"重复"其实是误判(比如 id 生成逻辑撞了),数据就是静默丢的。用 create 做去重之前,先确认能接受这个后果。
at-least-once 场景下,同一条 event 可能被重复写入(pipeline 重启、网络重传)。要实现幂等写,需要给 document_id 设一个确定性的值:
1 | |
fingerprint filter 可以对指定字段列表计算哈希,写入 @metadata,再由 elasticsearch output 用作 document_id。这条链路把 at-least-once 的重复写变成幂等写,是 Logstash 场景下实现"事实上的 exactly-once"的标准手法。
这个手法有一个明确的适用边界:document_id 在 data stream 模式下不被支持,它属于让 data_stream => auto 判定为不兼容的参数之一。所以上面这套幂等写只适用于具名索引写入。走 data stream 时想去重,得靠上游(比如让源侧保证不重发)或者在 ES 侧另做处理,output 插件这一层没有对应能力。
背压来源:ES 慢响应如何传导到 input
ES bulk 请求是同步阻塞的:worker 线程发出请求后原地等响应。关键在于这个 worker 就是刚跑完 filter 阶段的那个 worker。Logstash 没有独立的 output 线程池。pipeline.workers 控制同一组线程,官方定义就是"并行执行 pipeline 的 filter 和 output 两个阶段的 worker 数"。
所以传导路径不是"output 消费不动 filter 的产出",而是同一个线程被卡住之后回不去取下一批:ES 响应慢 → 该 worker 停在 output 阶段 → 它无法返回去从队列拉新 batch → 所有 worker 陆续卡在同一位置 → 队列填满 → input 写不进队列,被迫减速。第 09 篇专门论证这个 worker 模型,并把"filter 与 output 是两个能互相消费输出的独立阶段"明确列为需要纠正的理解。
这一点也解释了为什么第 12 篇要讲 output isolator:既然 output 阻塞会连带停住同一线程的 filter,那么把不同 output 拆到不同 pipeline,就是让一个下游故障不至于拖垮全部处理能力的唯一手段。
排障顺序:遇到"input 端读取突然变慢",先看 ES 的 bulk 响应时间和 429 频率,再顺着这条链向上游追溯,最后才怀疑 input 插件自身的配置。pipeline.workers、pipeline.batch.delay、连接池大小(pool_max、pool_max_per_route)共同决定背压信号传导的速度和幅度。
连接池与 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 成功发出。
值得注意的是这个实验观察不到什么:没有"重试次数用尽"这个终点。HTTP 层失败会一直重试下去,退避间隔从 2 秒起翻倍、上到 64 秒封顶,然后维持在 64 秒。ES 一直不恢复,这批 event 就一直卡在 worker 线程里,队列随之填满,input 被反压到停止读取。想看到 event 真的被丢弃,得构造文档级错误:比如往一个 mapping 已经把某字段定为 long 的索引里写字符串,再对比开启和关闭 DLQ 两种情况下 event 的去向。
模式提炼
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 | 批次行数 |
| pipeline.batch.delay | 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 往返次数越少,吞吐在一定范围内确实上升。但上限约束不在字节数上:按单文档 1-5KB 算,2000 条也只有 10MB 上下,离 http.max_content_length 的 100MB 很远。真正的约束是两条。一条是在途 event 占用的 JVM 堆随 batch.size × pipeline.workers 线性增长;另一条是低流量时数据要等满一批、或者等 pipeline.batch.delay 到时才发出,端到端延迟随 batch 变大而上升。实践中的做法是从 250-500 起步逐步上调,每次观察 bulk 响应时间与 ES 的 indexing 压力,看到响应时间开始抬头就停。
误解三:“关掉 sniffing 就必须在 hosts 里列出所有 ES 节点”。hosts 里只需要列出能路由到 ES cluster 的入口地址,可以是负载均衡器、代理、或者几个已知数据节点。关掉 sniffing 的含义是"不自动发现其他节点",不是"只能写入 hosts 里列出的节点"——Logstash 把请求发给 hosts 里的节点,ES 集群内部路由负责把文档分发到正确的 primary shard。
误解四:“Logstash 的 at-least-once 只要开了持久队列就够了”。持久队列保证的是 event 进入队列之后不因崩溃丢失,但它管不到 output 端。真正的丢失路径全在文档级错误上:400/404 在未开 DLQ 时被记日志后丢弃,409 无论如何都是记 warning 后丢弃、既不重试也不进 DLQ。完整的 at-least-once 需要:持久队列兜住崩溃边界 + DLQ 接住 400/404 + retry_on_conflict 或去重策略处理 409 + DLQ 的离线重放能力。缺少任何一环,"at-least-once"就只是部分成立。注意这里不需要"调大重试次数"这一环——HTTP 层的重试本来就是无限的。
误解五:“data_stream 在 8.x 默认开着,所以我的配置写的就是 data stream”。默认值是 auto 而不是 true,auto 只在配置与 data stream 兼容时才生效。配了 index、document_id、routing、pipeline 任意一个,或者 action 不是 create,或者 ecs_compatibility => disabled,都会让它退回具名索引写入。判断自己到底写进了哪里,看 ES 里是否出现了 .ds- 前缀的 backing index 最直接。
误解六:“%{+YYYY.MM.dd} 会按 date filter 里配的 timezone 切索引”。日期插值恒按 UTC 执行,timezone 影响的是解析而不是渲染,机制见前文"日期插值恒按 UTC 执行"一节。这条之所以值得单列,是因为它的表现形式极容易被误判成别的问题:索引名"晚了一天"看起来像 date filter 没配对、像服务器时区不对、像 ILM 策略有问题,唯独不像插值语法的固有行为。查这类问题的第一步应该是确认 event 的 @timestamp 转成 UTC 之后是哪一天,而不是去翻时区配置。
练习
-
本地起一个 Logstash + Elasticsearch,用 stdin input 发几条 event。停掉 ES,继续向 stdin 发数据,观察日志里
Retrying failed action的时间间隔序列,确认它从 2 秒起翻倍、到 64 秒后不再增长,并且始终不放弃。让它跑够十分钟以上,确认没有任何"重试次数用尽"的日志出现。重新启动 ES,观察积压的 event 是否被成功发出。 -
在 elasticsearch output 里设置
document_id => "%{message}"(用消息内容作 ID),向 ES 发送 10 条内容相同的 event,查询 ES 确认文档数量是 1 而不是 10。解释为什么action => "index"配合固定document_id是幂等写。然后把action改成create重跑同样的 10 条,对比两点:ES 里文档数仍然是 1,但 Logstash 日志里多出 9 条 409 warning,而且这 9 条 event 既没有重试也没有进 DLQ(把dead_letter_queue.enable打开再跑一次,确认 DLQ 文件里同样没有它们)。 -
把
index改成"logs-%{+YYYY.MM.dd}",但不用 date filter 改写@timestamp,向 stdin 发一条带历史时间戳的日志文本(但不解析它)。观察文档落在哪个索引(预期是今天的索引,而不是日志内容里的历史日期)。再加上 date filter 改写@timestamp,重复实验,确认文档落到了正确的历史日期索引。最后构造一条 UTC 时间落在前一天、本地时间已是今天的 event(东八区就取本地时间 00:00 到 08:00 之间),确认它落进的是前一天的索引。
系列导航
参考资料
- Logstash Elasticsearch output 文档:https://www.elastic.co/guide/en/logstash/current/plugins-outputs-elasticsearch.html(参数全集、8.1.1 之后的 Retry policy 与 DLQ policy 两节、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 生成)
- ES output 源码
common.rb:https://github.com/logstash-plugins/logstash-output-elasticsearch/blob/main/lib/logstash/plugin_mixins/elasticsearch/common.rb(DOC_DLQ_CODES = [400, 404]、DOC_CONFLICT_CODE = 409) - ES output 源码
data_stream_support.rb:https://github.com/logstash-plugins/logstash-output-elasticsearch/blob/main/lib/logstash/outputs/elasticsearch/data_stream_support.rb(invalid_data_stream_params是auto判定的全部依据) StringInterpolation.java:https://github.com/elastic/logstash/blob/main/logstash-core/src/main/java/org/logstash/StringInterpolation.java(%{+...}与%{{...}}两种写法都硬编码 UTC)
