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

AWS MSK Connect用Lenses插件同步数据至S3报错求助

问题:AWS MSK Connect搭配Lenses S3 Sink插件连接第三方Kafka集群报错

报错信息:

[Worker-001b25e1c610b1241] org.apache.kafka.connect.errors.ConnectException: Could not look up partition metadata for offset backing store topic in allotted period. This could indicate a connectivity issue, unavailable topic partitions, or if this is your first use of the topic it may have taken too long to create.

已通过EC2服务器使用kafka-consumer成功从目标Kafka集群获取数据,但配置MSK Connect使用Lenses S3 Sink插件导出数据至S3时出现上述错误,连接器配置如下:

{
    "connectorConfiguration": {
        "connector.class":"io.lenses.streamreactor.connect.aws.s3.sink.S3SinkConnector",
        "key.converter.schemas.enable":"false",
        "connect.s3.kcql":"INSERT INTO bigdata-XXXX:output SELECT * FROM topic_name `JSON` WITH_FLUSH_INTERVAL = 5",
        "aws.region":"eu-central-1",
        "tasks.max":"1",
        "topics":"topic_name",
        "schema.enable":"false",
        "value.converter":"org.apache.kafka.connect.storage.StringConverter",
        "errors.log.enable":"true",
        "key.converter":"org.apache.kafka.connect.storage.StringConverter",
        "allow.auto.create.topics " : "false",
        "connect.s3.aws.region": "eu-central-1",
        "connect.s3.vhost.bucket": "true",
        "aws.custom.endpoint":"https://s3.eu-central-1.amazonaws.com/"

    },
    "connectorName": "bigdata-transactions-connector",
    "kafkaCluster": {
        "apacheKafkaCluster": {
            "bootstrapServers": "kafka.XXXXXX:9092",
            "vpc": {
                "subnets": [
                    "subnet-XXXX",
                    "subnet-XXXX",
                    "subnet-XXXX"
                ],
                "securityGroups": ["sg-XXXXX"]
            }
        }
    },
    "capacity": {
        "provisionedCapacity": {
            "mcuCount": 1,
            "workerCount": 1
        }
    },
    "kafkaConnectVersion": "2.7.1",
    "serviceExecutionRoleArn": "arn:aws:iam::XXXXX",
    "plugins": [
        {
            "customPlugin": {
                "customPluginArn": "arn:aws:XXXXX",
                "revision": 1
            }
        }
    ],
    "logDelivery": { 
      "workerLogDelivery": { 
         "cloudWatchLogs": { 
            "enabled": true,
            "logGroup": "big_XXXXX"
         }
      }
   },
   "workerConfiguration": { 
      "revision": 1,
      "workerConfigurationArn": "arn:XXXXX"
   },
    "kafkaClusterEncryptionInTransit": {"encryptionType": "TLS"},
    "kafkaClusterClientAuthentication": {"authenticationType": "NONE"}
}

排查与解决步骤

1. 确认偏移量存储主题的存在性与权限

MSK Connect默认依赖connect-offsets主题存储任务偏移量,因你设置了allow.auto.create.topics = false,需手动在第三方Kafka集群创建该主题:

  • 主题需配置至少3个分区、2个副本,保证高可用
  • 验证MSK Connect使用的IAM角色对应的Worker身份,拥有该主题的Read和Write权限

2. 修正配置语法错误

连接器配置中allow.auto.create.topics 存在多余空格,需修改为:

"allow.auto.create.topics": "false"

空格会导致配置解析异常,干扰Connect初始化时的主题检查逻辑。

3. 验证MSK Connect与第三方Kafka的网络连通性

即便EC2能正常消费,MSK Connect Worker所在VPC的子网/安全组仍可能存在限制:

  • 确认安全组sg-XXXXX已放行出站9092端口到第三方Kafka集群的IP/网段
  • 检查子网subnet-XXXX的路由表,确保存在指向第三方Kafka网络的路由(如公网网关、VPN或专线)
  • 若第三方Kafka使用自定义TLS证书,需在Worker配置中添加证书信任配置

4. 调整偏移量存储的超时参数

在Worker配置中增加以下参数,延长元数据查询的超时时间,适配第三方集群的网络延迟:

offset.storage.retry.backoff.ms=5000
offset.storage.request.timeout.ms=30000

5. 检查插件与MSK Connect版本兼容性

确认使用的Lenses S3 Sink插件版本与MSK Connect 2.7.1兼容:

  • 插件需支持Kafka Connect 2.7.x版本,避免API不匹配导致的初始化失败
  • 尝试更新插件至最新兼容版本

内容的提问来源于stack exchange,提问作者uno-2017

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 17:40:47