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

Kafka Connect S3 Sink基于Key的分区及数据留存配置问询

Kafka Connect S3 Sink 连接器的分区与数据保留配置

1. 基于Key对输出数据分区

可以实现。通过指定Key分区器,让S3 Sink按照消息的Key值对数据进行分区存储,配置示例如下:

# 使用Key分区器
partitioner.class=org.apache.kafka.connect.storage.KeyPartitioner
# 若消息Key为结构化数据(如Avro、JSON格式),可指定具体字段作为分区键(可选)
# partition.field.name=your_key_field_name

配置后,S3上的路径会自动按Key值生成子目录,实现按Key分区。

2. 仅保留每个Key的最新10条记录

S3 Sink连接器本身没有直接支持按Key保留最新N条记录的配置项。要实现这个需求,建议先对原Kafka Topic的数据做预处理:
通过Kafka Streams或KSQL按Key聚合,只保留每个Key的最新10条记录,再将处理后的结果发送到新Topic,最后让S3 Sink消费这个新Topic写入S3。

举个Kafka Streams的核心逻辑示例:

StreamsBuilder builder = new StreamsBuilder();
// 读取原Topic数据
KStream<String, YourDataModel> originalStream = builder.stream("source-topic");
// 按Key分组,聚合保留最新10条
KGroupedStream<String, YourDataModel> groupedStream = originalStream.groupByKey();
KTable<String, List<YourDataModel>> latest10Records = groupedStream.aggregate(
    ArrayList::new,
    (key, newRecord, recordList) -> {
        recordList.add(newRecord);
        // 超过10条则移除最早的记录
        if (recordList.size() > 10) {
            recordList.remove(0);
        }
        return recordList;
    },
    Materialized.as("latest-10-store")
);
// 将处理后的数据发送到新Topic供S3 Sink消费
latest10Records.toStream().to("processed-topic");

3. 仅保留10分钟前的数据

这个需求分两种场景处理:

  • 只写入10分钟前的旧数据:可以通过自定义Transform过滤掉10分钟内的新记录,只保留符合时间条件的数据写入S3。如果使用内置的Filter转换,可参考以下配置(注意内置Filter的表达式支持有限,复杂场景建议自定义Transform):
transforms=filterRecentData
transforms.filterRecentData.type=org.apache.kafka.connect.transforms.Filter$Value
# 条件表达式:消息中的timestamp字段小于当前时间减10分钟(需根据实际字段名调整)
transforms.filterRecentData.condition=value.timestamp < ${current_timestamp}-600000
  • 自动删除S3中超过10分钟的数据:S3 Sink本身不支持自动清理,但可以在AWS S3控制台配置生命周期规则,设置对象创建10分钟后自动删除。

4. 同时基于Key和时间周期进行分区

可以通过CompositePartitioner组合Key分区器和时间分区器,实现双重分区。配置示例如下:

# 使用复合分区器,同时启用Key和时间分区
partitioner.class=org.apache.kafka.connect.storage.CompositePartitioner
# 指定要组合的分区器
partitioner.partitioners=org.apache.kafka.connect.storage.KeyPartitioner,org.apache.kafka.connect.storage.TimeBasedPartitioner

# 时间分区相关配置
# 用消息记录的时间作为分区依据(也可以指定消息中的自定义时间字段)
timestamp.extractor=org.apache.kafka.connect.storage.RecordTimestamp
# 10分钟一个时间分区
partition.duration.ms=600000
# S3路径的时间格式
path.format='year'=YYYY/'month'=MM/'day'=dd/'hour'=HH/'minute'=mm

# Key分区相关配置(可选,若Key为结构化数据则指定字段)
# partition.field.name=your_key_field_name

配置后,S3上的路径会同时包含Key和时间维度的子目录,比如s3://bucket/path/key=xxx/year=2024/month=05/day=20/hour=14/minute=30/

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 11:15:32