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

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存储交互的具体实现类,而非连接器本身。

修复步骤

  1. 修改S3 Sink连接器配置中的storage.class为正确的存储实现类:
    storage.class: io.confluent.connect.s3.storage.S3Storage
    
  2. 重启连接器任务,让配置生效。

补充说明

  • 版本兼容性:使用的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 00:22:03