MSK Connect连接Confluent Cloud Kafka主题超时问题求助
问题描述
尝试通过MSK Connect将iceberg-kafka-connect连接器连接到Confluent Cloud中的Kafka主题。此前已通过同VPC、子网和安全组的AWS EC2实例成功访问该主题,说明VPC与Confluent Cloud的Kafka集群连通性正常,但MSK Connect配置相同网络参数后无法获取数据,怀疑是sasl.jaas.config格式错误或超时问题。
连接器配置
iceberg.kafka.sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule required username="<CONFLUENT_API_KEY>" password="<CONFLUENT_API_SECRET>"; iceberg.kafka.security.protocol=SASL_SSL iceberg.kafka.sasl.mechanism=PLAIN connector.class=io.tabular.iceberg.connect.IcebergSinkConnector table.write-format=parquet iceberg.tables.evolve-schema-enabled=true iceberg.fs.s3a.path.style.access=true table.namespace=pocckafkaflink2 iceberg.catalog.catalog-impl=org.apache.iceberg.aws.glue.GlueCatalog tasks.max=1 topics=confluent_test_topic2 iceberg.catalog.io-impl=org.apache.iceberg.aws.s3.S3FileIO iceberg.catalog.client.region=us-east-1 iceberg.fs.s3a.aws.credentials.provider=com.amazonaws.auth.DefaultAWSCredentialsProviderChain iceberg.fs.s3a.impl=org.apache.hadoop.fs.s3a.S3AFileSystem iceberg.tables=pocckafkaflink2.mskmytopicmaib value.converter.schemas.enable=false iceberg.catalog.warehouse=s3://bucket/data4 value.converter=org.apache.kafka.connect.json.JsonConverter table.auto-create=true
错误日志
[Worker-06cf2e74aeb2c59e6] org.apache.kafka.common.errors.TimeoutException: Call(callName=fetchMetadata, deadlineMs=1722837422758, tries=1, nextAllowedTryMs=1722837422859) timed out at 1722837422759 after 1 attempt(s) [Worker-06cf2e74aeb2c59e6] Caused by: org.apache.kafka.common.errors.TimeoutException: Timed out waiting to send the call. Call: fetchMetadata [Worker-06cf2e74aeb2c59e6] [2024-08-05 05:57:32,759] INFO App info kafka.admin.client for adminclient-1 unregistered (org.apache.kafka.common.utils.AppInfoParser:83) [Worker-06cf2e74aeb2c59e6] [2024-08-05 05:57:32,759] INFO [AdminClient clientId=adminclient-1] Metadata update failed (org.apache.kafka.clients.admin.internals.AdminMetadataManager:235) [Worker-06cf2e74aeb2c59e6] org.apache.kafka.common.errors.TimeoutException: Call(callName=fetchMetadata, deadlineMs=1722837452759, tries=1, nextAllowedTryMs=-9223372036854775709) timed out at 9223372036854775807 after 1 attempt(s) [Worker-06cf2e74aeb2c59e6] Caused by: org.apache.kafka.common.errors.TimeoutException: Timed out waiting to send the call. Call: fetchMetadata [Worker-06cf2e74aeb2c59e6] [2024-08-05 05:57:32,766] INFO Metrics scheduler closed (org.apache.kafka.common.metrics.Metrics:668) [Worker-06cf2e74aeb2c59e6] [2024-08-05 05:57:32,766] INFO Closing reporter org.apache.kafka.common.metrics.JmxReporter (org.apache.kafka.common.metrics.Metrics:672) [Worker-06cf2e74aeb2c59e6] [2024-08-05 05:57:32,766] INFO Metrics reporters closed (org.apache.kafka.common.metrics.Metrics:678) [Worker-06cf2e74aeb2c59e6] [2024-08-05 05:57:32,767] ERROR Stopping due to error (org.apache.kafka.connect.cli.ConnectDistributed:86) [Worker-06cf2e74aeb2c59e6] org.apache.kafka.connect.errors.ConnectException: Failed to connect to and describe Kafka cluster. Check worker's broker connection and security properties. [Worker-06cf2e74aeb2c59e6] at org.apache.kafka.connect.util.ConnectUtils.lookupKafkaClusterId(ConnectUtils.java:70) [Worker-06cf2e74aeb2c59e6] at org.apache.kafka.connect.util.ConnectUtils.lookupKafkaClusterId(ConnectUtils.java:51) [Worker-06cf2e74aeb2c59e6] at org.apache.kafka.connect.cli.ConnectDistributed.startConnect(ConnectDistributed.java:97) [Worker-06cf2e74aeb2c59e6] at org.apache.kafka.connect.cli.ConnectDistributed.main(ConnectDistributed.java:80) [Worker-06cf2e74aeb2c59e6] Caused by: java.util.concurrent.ExecutionException: org.apache.kafka.common.errors.TimeoutException: Call(callName=listNodes, deadlineMs=1722837452757, tries=1, nextAllowedTryMs=1722837452858) timed out at 1722837452758 after 1 attempt(s) [Worker-06cf2e74aeb2c59e6] at org.apache.kafka.common.internals.KafkaFutureImpl.wrapAndThrow(KafkaFutureImpl.java:45) [Worker-06cf2e74aeb2c59e6] at org.apache.kafka.common.internals.KafkaFutureImpl.access$000(KafkaFutureImpl.java:32) [Worker-06cf2e74aeb2c59e6] at org.apache.kafka.common.internals.KafkaFutureImpl$SingleWaiter.await(KafkaFutureImpl.java:89) [Worker-06cf2e74aeb2c59e6] at org.apache.kafka.common.internals.KafkaFutureImpl.get(KafkaFutureImpl.java:260) [Worker-06cf2e74aeb2c59e6] at org.apache.kafka.connect.util.ConnectUtils.lookupKafkaClusterId(ConnectUtils.java:64) [Worker-06cf2e74aeb2c59e6] ... 3 more [Worker-06cf2e74aeb2c59e6] Caused by: org.apache.kafka.common.errors.TimeoutException: Call(callName=listNodes, deadlineMs=1722837452757, tries=1, nextAllowedTryMs=1722837452858) timed out at 1722837452758 after 1 attempt(s) [Worker-06cf2e74aeb2c59e6] Caused by: org.apache.kafka.common.errors.TimeoutException: Timed out waiting for a node assignment. Call: listNodes [Worker-06cf2e74aeb2c59e6] MSK Connect encountered errors and failed.
排查与解决建议
1. 修正JAAS配置格式
当前sasl.jaas.config中的引号转义可能存在问题,在MSK Connect环境中无需转义双引号,直接使用以下格式:
iceberg.kafka.sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule required username="<CONFLUENT_API_KEY>" password="<CONFLUENT_API_SECRET>";
若通过AWS控制台配置,需确保JAAS配置中的引号未被额外转义,避免解析错误。
2. 补全Kafka Bootstrap地址
连接器配置中遗漏了核心的iceberg.kafka.bootstrap.servers参数,需添加Confluent Cloud提供的集群bootstrap地址:
iceberg.kafka.bootstrap.servers=<CONFLUENT_BOOTSTRAP_SERVERS>
3. 延长Kafka客户端超时时间
超时错误表明客户端无法在默认时间内获取集群元数据,添加以下参数延长超时:
iceberg.kafka.request.timeout.ms=30000 iceberg.kafka.retry.backoff.ms=1000 iceberg.kafka.metadata.max.age.ms=300000
4. 验证MSK Connect网络连通性
尽管EC2实例能正常访问,MSK Connect执行环境仍可能存在额外限制:
- 检查MSK Connect安全组的出站规则,确保允许访问Confluent Cloud Kafka集群的9092端口及对应IP范围(可从Confluent Cloud控制台获取集群IP列表)。
- 确认MSK Connect使用的子网已配置NAT网关,确保能访问公网(Confluent Cloud Kafka集群为公网服务时)。
5. 验证认证参数有效性
- 确认
<CONFLUENT_API_KEY>和<CONFLUENT_API_SECRET>正确,可通过EC2实例上的kafka-console-consumer.sh工具再次验证认证信息。 - 确保
iceberg.kafka.security.protocol和sasl.mechanism与Confluent Cloud配置一致(PLAIN认证对应SASL_SSL)。
内容的提问来源于stack exchange,提问作者Malina Dale
相关产品推荐
相关产品推荐

