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

AWS MSK集群下S3 Source Connector元数据获取错误排查求助

问题:AWS MSK S3 Source Connector元数据获取错误

我在使用AWS MSK集群,通过S3 Source Connector将S3数据导入自动创建的Kafka Topic时,遇到元数据获取错误,相关日志如下:

[Worker-001d22b042f681d7a] [2023-10-21 17:49:44,127] WARN [source-connector|task-0] [Producer clientId=connector-producer-source-connector-0] Error while fetching metadata with correlation id 1 : {source-topic=UNKNOWN_TOPIC_OR_PARTITION} (org.apache.kafka.clients.NetworkClient:1119)
[Worker-001d22b042f681d7a] [2023-10-21 17:49:44,127] INFO [source-connector|task-0] [Producer clientId=connector-producer-source-connector-0] Cluster ID: J4cTke2TRmOSsoYBO5uZdA (org.apache.kafka.clients.Metadata:279)
[Worker-001d22b042f681d7a] [2023-10-21 17:49:44,478] WARN [source-connector|task-0] [Producer clientId=connector-producer-source-connector-0] Error while fetching metadata with correlation id 3 : {source-topic=UNKNOWN_TOPIC_OR_PARTITION} (org.apache.kafka.clients.NetworkClient:1119)
[Worker-001d22b042f681d7a] [2023-10-21 17:49:44,581] WARN [source-connector|task-0] [Producer clientId=connector-producer-source-connector-0] Error while fetching metadata with correlation id 4 : {source-topic=UNKNOWN_TOPIC_OR_PARTITION} (org.apache.kafka.clients.NetworkClient:1119)

我的S3 Source Connector配置:

connector.class=io.confluent.connect.s3.source.S3SourceConnector
useAccelerateMode=true
s3.region=us-east-1
confluent.topic.bootstrap.servers=b-4.sink.2uql3y.c4.kafka.us-east-1.amazonaws.com:90921.amazonaws.com:9092
auto.create.topics.enable=true
flush.size=7
schema.compatibility=NONE
tasks.max=2
topics=target-topic
pathStyleAccess=true
schema.enable=false
key.converter.schemas.enable=false
format.class=io.confluent.connect.s3.format.json.JsonFormat
aws.region=us-east-1
partitioner.class=io.confluent.connect.storage.partitioner.DefaultPartitioner
value.converter=org.apache.kafka.connect.storage.StringConverter
storage.class=io.confluent.connect.s3.storage.S3Storage
errors.log.enable=true
s3.bucket.name=bucket-name
key.converter=org.apache.kafka.connect.storage.StringConverter

创建_confluent-command主题后,我用以下命令尝试接收S3消息到自动创建的Topic:

bin/kafka-console-consumer.sh --topic _confluent-command --consumer.config /home/ec2-user/kafka_2.12-3.5.1/config/consumer.properties --from-beginning --bootstrap-server <bootstrap-server>

收到的消息内容:

�
�eyJhbGciOiJub25lIn0.eyJpc3MiOiJDb25mbHVlbnQiLCJhdWQiOiJ0cmlhbCIsImV4cCI6MTcwMDUwNTg3OCwianRpIjoiamoySWFWQTV6ckVjSG94ZUw5X1dsUSIsImlhdCI6MTY5NzkxMzg3NywibmJmIjoxNjk3OTEzNzU3LCJzdWIiOiJDb25mbHVlbnQgRW50ZXJwcmlzZSIsIm1vbml0b3JpbmciOnRydWUsImxpY2Vuc2VUeXBlIjoidHJpYWwifQ.

解决方案

1. 修复Bootstrap Servers配置错误

你的confluent.topic.bootstrap.servers配置存在明显的地址拼接错误:

confluent.topic.bootstrap.servers=b-4.sink.2uql3y.c4.kafka.us-east-1.amazonaws.com:90921.amazonaws.com:9092

正确格式应为MSK集群的有效bootstrap服务器地址,多节点用逗号分隔,示例:

confluent.topic.bootstrap.servers=b-4.sink.2uql3y.c4.kafka.us-east-1.amazonaws.com:9092,b-3.sink.2uql3y.c4.kafka.us-east-1.amazonaws.com:9092

该错误会导致连接器无法正常连接MSK集群,进而无法识别或创建目标Topic。

2. 确认自动创建Topic权限

检查MSK集群的IAM策略或ACL规则,确保连接器使用的服务账户拥有CreateTopics权限。即使开启auto.create.topics.enable=true,若无对应权限,连接器仍无法自动创建Topic。

3. 手动验证目标Topic

  • 连接器配置中topics=target-topic是数据写入的目标Topic,可先手动创建该Topic测试:
bin/kafka-topics.sh --create --topic target-topic --bootstrap-server <你的MSK bootstrap地址> --partitions 2 --replication-factor 2

创建完成后重启连接器,观察错误是否消失。

4. 解析JWT消息

你收到的是JWT格式消息,可通过本地工具或命令解码查看详细内容,了解连接器的指令执行状态,辅助排查问题。

5. 验证S3访问权限

确认连接器使用的IAM角色具备目标S3桶的s3:GetObject、s3:ListBucket权限,确保能正常读取S3中的数据文件,权限不足也可能间接引发Topic相关错误。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 06:16:07