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
相关产品推荐
相关产品推荐

