如何搭建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
相关产品推荐
相关产品推荐

