如何用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:依次执行字段提取与格式转换:extractSchema:提取metadata.schema-name到顶层schema-name字段extractTable:提取metadata.table-name到顶层table-name字段extractTimestamp:提取metadata.timestamp字符串到顶层event_timestamp_strconvertTimestamp:将字符串时间转换为Timestamp类型,存入event_timestamp
partitioner.field:指定用转换后的event_timestamp作为分区时间依据partition.duration.ms:设置为86400000(24小时),确保按天生成分区路径
4. 验证方法
- 将上述配置提交到MSK Connect创建连接器。
- 向目标MSK主题发送测试消息(如你提供的示例数据)。
- 登录AWS S3控制台,检查是否生成路径:
bucket-name/customers_dev/customer_con/2024/10/24,并确认路径下的JSON文件包含测试消息内容。
内容的提问来源于stack exchange,提问作者Shashwat Awasthi
相关产品推荐
相关产品推荐

