EKS上Strimzi Kafka Connect部署Confluent S3 Sink连接器异常排查
问题:Strimzi Kafka Connect部署Confluent S3 Sink连接器报错NoSuchMethodException
报错信息
(io.confluent.connect.storage.partitioner.PartitionerConfig) [task-thread-abc-dp-test-sink-connector-0] 2023-05-31 08:43:28,563 ERROR [abc-dp-test-sink-connector|task-0] WorkerSinkTask{id=abc-dp-test-sink-connector-0} Task threw an uncaught and unrecoverable exception. Task is being killed and will not recover until manually restarted (org.apache.kafka.connect.runtime.WorkerTask) [task-thread-abc-dp-test-sink-connector-0] org.apache.kafka.connect.errors.ConnectException: java.lang.NoSuchMethodException: io.confluent.connect.s3.S3SinkConnector.<init>(io.confluent.connect.s3.S3SinkConnectorConfig, java.lang.String) at io.confluent.connect.storage.StorageFactory.createStorage(StorageFactory.java:55) at io.confluent.connect.s3.S3SinkTask.start(S3SinkTask.java:108) at org.apache.kafka.connect.runtime.WorkerSinkTask.initializeAndStart(WorkerSinkTask.java:312) at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:186) at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:243) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515) at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) at java.base/java.lang.Thread.run(Thread.java:829) Caused by: java.lang.NoSuchMethodException: io.confluent.connect.s3.S3SinkConnector.<init>(io.confluent.connect.s3.S3SinkConnectorConfig, java.lang.String) at java.base/java.lang.Class.getConstructor0(Class.java:3349) at java.base/java.lang.Class.getConstructor(Class.java:2151) at io.confluent.connect.storage.StorageFactory.createStorage(StorageFactory.java:49) ... 9 more
相关配置
插件镜像Dockerfile
FROM quay.io/strimzi/kafka:0.31.1-kafka-3.1.0-amd64 USER root:root RUN mkdir -p /opt/kafka/plugins/confluentinc-connect-transforms-1.4.3 RUN mkdir -p /opt/kafka/plugins/confluentinc-kafka-connect-avro-converter-7.2.2 RUN mkdir -p /opt/kafka/plugins/confluentinc-kafka-connect-protobuf-converter-7.2.2 RUN mkdir -p /opt/kafka/plugins/confluentinc-kafka-connect-s3-10.5.0-3 COPY ./confluentinc-connect-transforms-1.4.3/* /opt/kafka/plugins/confluentinc-connect-transforms-1.4.3 COPY ./confluentinc-kafka-connect-avro-converter-7.2.2/* /opt/kafka/plugins/confluentinc-kafka-connect-avro-converter-7.2.2 COPY ./confluentinc-kafka-connect-protobuf-converter-7.2.2/* /opt/kafka/plugins/confluentinc-kafka-connect-protobuf-converter-7.2.2 COPY ./confluentinc-kafka-connect-s3-10.5.0-3/* /opt/kafka/plugins/confluentinc-kafka-connect-s3-10.5.0-3 USER 1001
Strimzi Connect集群配置片段
... spec: version: 3.1.0 replicas: 2 image: s3-sink-image-repo:latest bootstrapServers: bootstrap_server1:9094,bootstrap_server2:9094 config: group.id: connect-cluster offset.storage.topic: connect-cluster-offsets-v1 config.storage.topic: connect-cluster-configs-v1 status.storage.topic: connect-cluster-status-v1 key.converter: org.apache.kafka.connect.json.JsonConverter value.converter: org.apache.kafka.connect.json.JsonConverter key.converter.schemas.enable: true value.converter.schemas.enable: true config.storage.replication.factor: 2 offset.storage.replication.factor: 2 status.storage.replication.factor: 2 ...
S3 Sink连接器配置
... spec: class: io.confluent.connect.s3.S3SinkConnector tasksMax: 1 config: topics: test_strimzi_kafka behavior.on.null.values: ignore s3.region: ap-southeast-1 s3.bucket.name: my-bucket-for-test flush.size: 1 consumer.request.timeout.ms: 30000 errors.retry.delay.max.ms: 30000 storage.class: io.confluent.connect.s3.S3SinkConnector format.class: io.confluent.connect.s3.format.json.JsonFormat key.converter: org.apache.kafka.connect.json.JsonConverter value.converter: org.apache.kafka.connect.json.JsonConverter key.converter.schemas.enable: false value.converter.schemas.enable: false errors.log.enable: true errors.tolerance: all errors.retry.timeout: 600000 errors.log.include.messages: false s3.part.size: 5242880 schema.compatibility: NONE consumer.fetch.max.bytes: 1000000 consumer.fetch.max.wait.ms: 500
问题分析与解决方案
核心原因
报错根源是连接器配置中的storage.class参数设置错误:当前将storage.class指定为连接器主类io.confluent.connect.s3.S3SinkConnector,但该类并没有(S3SinkConnectorConfig, String)的构造方法,而StorageFactory会尝试用这个构造方法实例化存储类,因此抛出NoSuchMethodException。
storage.class应该指定的是处理S3存储交互的具体实现类,而非连接器本身。
修复步骤
- 修改S3 Sink连接器配置中的
storage.class为正确的存储实现类:storage.class: io.confluent.connect.s3.storage.S3Storage - 重启连接器任务,让配置生效。
补充说明
- 版本兼容性:使用的Confluent S3 Sink 10.5.0(对应Confluent Platform 7.2.x)与Strimzi的Kafka 3.1.0版本兼容,无需调整版本。
- MSK Connect正常运行的原因:推测在MSK Connect的配置中,
storage.class参数被正确设置,或MSK Connect对该参数有默认处理,因此未出现此错误。
内容的提问来源于stack exchange,提问作者Surat Suntong
相关产品推荐
相关产品推荐

