MSK Connector创建失败求助:部署S3 Sink连接器遇未知错误
问题描述
尝试使用kafka-connect-aws-s3-kafka-2-8-4.0.0创建MSK连接器,将MSK数据下沉至S3,但连接器一直处于**Creating...**状态,最终因未知错误失败。连接器配置如下:
connector.class=io.lenses.streamreactor.connect.aws.s3.sink.S3SinkConnector behavior.on.null.values=ignore s3.region=eu-central-1 flush.size=5 tasks.max=8 topics=MyTopic s3.part.size=5242880 connect.s3.vhost.bucket=true schema.enable=false format.bytearray.separator=\n--==SEPARATOR==--\n key.converter.schemas.enable=false connect.s3.kcql=INSERT INTO kafka-connect-aws2:evelin-msk SELECT * FROM MyTopic WITH_FLUSH_INTERVAL = 300 format.class=io.confluent.connect.s3.format.bytearray.ByteArrayFormat partitioner.class=io.confluent.connect.storage.partitioner.DefaultPartitioner value.converter.schemas.enable=false connect.s3.aws.region=eu-central-1 value.converter=org.apache.kafka.connect.converters.ByteArrayConverter storage.class=io.confluent.connect.s3.storage.S3Storage errors.log.enable=true s3.bucket.name=kafka-connect-aws
可能的故障原因
- 配置项冲突:同时设置了
s3.region和connect.s3.aws.region,Stream Reactor S3连接器仅需connect.s3.aws.region,多余的s3.region会导致配置解析混乱。 - KCQL与Bucket配置不一致:KCQL语句中指定的目标bucket是
kafka-connect-aws2,但配置里的s3.bucket.name=kafka-connect-aws,两者不匹配会导致连接器无法找到正确的存储bucket,引发权限或资源不存在错误。 - 存储类/格式类不兼容:使用了Stream Reactor的连接器类,但搭配了Confluent的
storage.class和format.class,应该替换为Stream Reactor对应的实现类:storage.class改为io.lenses.streamreactor.connect.aws.s3.storage.S3Storageformat.class改为io.lenses.streamreactor.connect.aws.s3.format.bytearray.ByteArrayFormat
- 权限不足:MSK Connect执行角色未获得S3所需权限(如
s3:PutObject、s3:ListBucket、s3:GetBucketLocation等),或目标S3 bucket的访问策略未允许该角色操作。 - Vhost Bucket配置问题:
connect.s3.vhost.bucket=true开启了虚拟主机模式,若bucket名称包含特殊字符(如点号)导致SSL证书验证失败,或当前区域不支持该模式,会引发S3连接超时。 - 任务数超出集群资源:
tasks.max=8设置过高,若MSK Connect集群的CPU、内存资源不足,任务无法正常启动,最终导致连接器创建超时失败。 - 分区器类不兼容:当前使用Confluent的
DefaultPartitioner,建议替换为Stream Reactor对应的分区器类io.lenses.streamreactor.connect.aws.s3.partitioner.DefaultPartitioner。
内容的提问来源于stack exchange,提问作者MBA
相关产品推荐
相关产品推荐

