You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何搭建Kafka输入S3输出的Logstash管道实现单消息存为自定义命名独立S3文件

解决方案

你的问题核心是Logstash S3输出插件默认会聚合多个事件写入同一个文件,仅调整Kafka输入侧的拉取参数无法解决该问题,需从字段解析、S3输出旋转规则、文件名规则三个维度调整配置,具体如下:

配置调整要点

  • 补充Kafka消息解码配置,将消息内容解析为Logstash结构化字段,才能提取指定属性作为文件名
  • 调整S3输出的文件旋转策略,强制单条事件生成独立文件
  • 配置S3对象名规则,使用占位符引用事件字段作为文件名

完整参考配置

input {
  kafka {
    bootstrap_servers => "mykafkaserver:9092"
    topics => "document"
    group_id => "xLogAna1"
    auto_offset_reset => "earliest"
    # 若Kafka消息为JSON格式,添加JSON解码,可根据实际消息格式更换对应解码器
    codec => json
  }
}

output {
  s3{
    access_key_id => "XXXXXXXXXXXXXXX"
    secret_access_key => "SSSSSSSSSSSSSSS"
    region => "eu-west-1"
    bucket => "<my-documnt-bucket>"
    codec => "plain"
    # 强制按大小触发文件旋转
    rotation_strategy => "size"
    # 文件大小超过1字节即触发上传
    size_file => 1
    # 禁用批量写入
    batch => false
    # 配置文件名,%{字段名}会自动读取事件对应字段的值,示例为用document_id作为文件名
    # 如需防止重名,可拼接Kafka元数据字段,例如:"%{document_id}_%{[@metadata][kafka][offset]}.txt"
    object_name => "%{document_id}.txt"
  }
}

注意事项

如果需要引用的属性是消息中的嵌套字段,使用%{[一级字段名][二级字段名]}的格式填写占位符即可。


内容的提问来源于stack exchange,提问作者Amit Meena

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.07 14:30:01