如何通过Kafka仅向S3追加写入单个文件?
问题分析与解决
我在EKS上部署了Kafka,使用Strimzi及Schema Registry,采用Parquet格式将分区数据通过Sink写入S3。期望实现仅向单个文件追加写入,当前设置文件达到约5MB时生成新文件,但实际每次插入都会在S3生成新文件。
当前Connector配置
Sink Connector配置
s3.part.size: 5242880 flush.size: 20000 rotate.schedule.interval.ms: 10000 transforms: "InsertField" transforms.InsertField.type: "org.apache.kafka.connect.transforms.InsertField$Value" transforms.InsertField.timestamp.field: "integration_timestamp"
Source Connector配置
max.batch.size: 100000 max.queue.size: 100000 query.fetch.size: 10000 poll.interval.ms: 10000
日志提示
等待rotateIntervalMs时间后提交文件,但可用记录数少于flush.size
问题原因
日志已明确指出问题根源:rotate.schedule.interval.ms: 10000配置强制每10秒就轮转生成新文件,无论flush.size(20000条记录)或s3.part.size(5MB)的阈值是否达到。只要到了10秒,哪怕只有少量数据,也会提交新文件。
解决方案
- 优先按数据量生成文件:删除或注释掉
rotate.schedule.interval.ms配置,这样只有当记录数达到flush.size,或文件大小达到s3.part.size时,才会生成新文件,实现单个文件持续追加。 - 保留时间轮转但减少小文件:如果必须保留时间轮转策略,调大
rotate.schedule.interval.ms的值(例如设为3600000即1小时),同时确保在该时间内数据量能达到flush.size或s3.part.size的阈值,避免频繁生成小文件。 - 检查分区粒度:若数据按高频维度(如分钟级时间)分区,每个分区会独立生成文件。可调整分区粒度为小时/天级,减少分区数量,从而减少小文件的生成。
内容的提问来源于stack exchange,提问作者OdiumPura
相关产品推荐
相关产品推荐

