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

如何用Kafka-S3-Connector将AWS MSK数据按消息变量生成S3存储路径?

使用Kafka S3 Connector将AWS MSK数据按自定义路径存储到S3

核心思路

要实现bucket-name/{schema-name}/{table-name}/{year}/{month}/{date}的存储路径,需通过Single Message Transforms(SMT) 提取消息嵌套字段到顶层,结合时间分区器生成日期路径,最终用S3 Connector的路径模板拼接完整存储结构。


1. 前置准备

  • 确保MSK集群与S3存储桶网络连通(私有子网可配置S3 VPC端点,公网环境需确认路由可达)。
  • 为MSK Connect执行角色配置S3权限:s3:PutObject、s3:ListBucket、s3:GetBucketLocation。
  • 已在MSK Connect中部署Confluent S3 Sink Connector插件(或AWS官方提供的S3连接器)。

2. 关键配置说明

2.1 提取嵌套字段到顶层

消息中的schema-name、table-name和timestamp均嵌套在metadata节点下,需用SMT将这些字段提取到记录顶层,方便路径模板调用:

  • ExtractField$Value:从metadata中提取指定字段,存入顶层自定义字段。
  • TimestampConverter$Value:将提取到的字符串格式timestamp转换为Kafka Connect支持的Timestamp类型,供时间分区器使用。

2.2 路径与分区配置

  • s3.path.format:定义存储路径模板,用${字段名}引用提取后的顶层字段,时间变量${year}、${month}、${day}由时间分区器自动生成。
  • TimeBasedPartitioner:基于指定时间字段生成日期分区,需配置分区字段、时间格式及时区。

3. 完整连接器配置示例

{
  "name": "msk-s3-custom-path-sink",
  "config": {
    "connector.class": "io.confluent.connect.s3.S3SinkConnector",
    "tasks.max": "2",
    "topics": "你的MSK主题名称",
    "s3.bucket.name": "你的S3存储桶名称",
    "s3.path.format": "${schema-name}/${table-name}/${year}/${month}/${day}",
    "flush.size": "1000",
    "rotate.interval.ms": "3600000",
    "format.class": "io.confluent.connect.s3.format.json.JsonFormat",
    "partitioner.class": "io.confluent.connect.storage.partitioner.TimeBasedPartitioner",
    "partitioner.field": "event_timestamp",
    "timestamp.extractor": "RecordField",
    "timezone": "UTC",
    "partition.duration.ms": "86400000",
    "transforms": "extractSchema,extractTable,extractTimestamp,convertTimestamp",
    "transforms.extractSchema.type": "org.apache.kafka.connect.transforms.ExtractField$Value",
    "transforms.extractSchema.field": "metadata.schema-name",
    "transforms.extractSchema.target.field": "schema-name",
    "transforms.extractTable.type": "org.apache.kafka.connect.transforms.ExtractField$Value",
    "transforms.extractTable.field": "metadata.table-name",
    "transforms.extractTable.target.field": "table-name",
    "transforms.extractTimestamp.type": "org.apache.kafka.connect.transforms.ExtractField$Value",
    "transforms.extractTimestamp.field": "metadata.timestamp",
    "transforms.extractTimestamp.target.field": "event_timestamp_str",
    "transforms.convertTimestamp.type": "org.apache.kafka.connect.transforms.TimestampConverter$Value",
    "transforms.convertTimestamp.field": "event_timestamp_str",
    "transforms.convertTimestamp.target.type": "Timestamp",
    "transforms.convertTimestamp.format": "yyyy-MM-dd'T'HH:mm:ss.SSS'Z'"
  }
}

配置项解释

  • transforms:依次执行字段提取与格式转换:
    1. extractSchema:提取metadata.schema-name到顶层schema-name字段
    2. extractTable:提取metadata.table-name到顶层table-name字段
    3. extractTimestamp:提取metadata.timestamp字符串到顶层event_timestamp_str
    4. convertTimestamp:将字符串时间转换为Timestamp类型,存入event_timestamp
  • partitioner.field:指定用转换后的event_timestamp作为分区时间依据
  • partition.duration.ms:设置为86400000(24小时),确保按天生成分区路径

4. 验证方法

  1. 将上述配置提交到MSK Connect创建连接器。
  2. 向目标MSK主题发送测试消息(如你提供的示例数据)。
  3. 登录AWS S3控制台,检查是否生成路径:bucket-name/customers_dev/customer_con/2024/10/24,并确认路径下的JSON文件包含测试消息内容。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 16:50:21