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

如何通过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秒,哪怕只有少量数据,也会提交新文件。

解决方案

  1. 优先按数据量生成文件:删除或注释掉rotate.schedule.interval.ms配置,这样只有当记录数达到flush.size,或文件大小达到s3.part.size时,才会生成新文件,实现单个文件持续追加。
  2. 保留时间轮转但减少小文件:如果必须保留时间轮转策略,调大rotate.schedule.interval.ms的值(例如设为3600000即1小时),同时确保在该时间内数据量能达到flush.size或s3.part.size的阈值,避免频繁生成小文件。
  3. 检查分区粒度:若数据按高频维度(如分钟级时间)分区,每个分区会独立生成文件。可调整分区粒度为小时/天级,减少分区数量,从而减少小文件的生成。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 03:57:16