分布式模式下Kafka Connect启动超时问题求助
Kafka Connect分布式启动超时问题
问题背景
启动分布式模式的Kafka Connect时出现超时错误,已调整connect-distributed.properties中的部分参数,且必须以分布式模式启动Connect。
connect-distributed.properties 配置内容
offset.flush.interval.ms=5000 offset.flush.timeout.ms=300000 producer.buffer.memory=300097152 status.storage.topic=connect-status status.storage.replication.factor=1 config.storage.topic=connect-configs config.storage.replication.factor=1 offset.storage.topic=connect-offsets offset.storage.replication.factor=1 key.converter.schemas.enable=false value.converter.schemas.enable=false
错误日志
INFO [AdminClient clientId=adminclient-8] Disconnecting from node 0 due to request timeout. (org.apache.kafka.clients.NetworkClient:797) INFO [AdminClient clientId=adminclient-8] Cancelled in-flight METADATA request with correlation id 530256 due to node 0 being disconnected (elapsed time since creation: 1ms, elapsed time since send: 1ms, request timeout: 0ms) (org.apache.kafka.clients.NetworkClient:341) WARN Attempt 2 to list offsets for topic partitions resulted in RetriableException; retrying automatically. Reason: Timed out while waiting to get end offsets for topic 'connect-offsets' on brokers at *.*.*.*:9092 (org.apache.kafka.connect.util.RetryUtil:84) org.apache.kafka.common.errors.TimeoutException: Timed out while waiting to get end offsets for topic 'connect-offsets' on brokers at *.*.*.*:9092 Caused by: java.util.concurrent.ExecutionException: org.apache.kafka.common.errors.TimeoutException: Call(callName=metadata, deadlineMs=1680677446792, tries=1, nextAllowedTryMs=1680677446893) timed out at 1680677446793 after 1 attempt(s) at java.util.concurrent.CompletableFuture.reportGet(CompletableFuture.java:357) at java.util.concurrent.CompletableFuture.get(CompletableFuture.java:1908) at org.apache.kafka.common.internals.KafkaFutureImpl.get(KafkaFutureImpl.java:165) at org.apache.kafka.connect.util.TopicAdmin.endOffsets(TopicAdmin.java:672) at org.apache.kafka.connect.util.TopicAdmin.lambda$retryEndOffsets$6(TopicAdmin.java:725) at org.apache.kafka.connect.util.RetryUtil.retryUntilTimeout(RetryUtil.java:82) at org.apache.kafka.connect.util.TopicAdmin.retryEndOffsets(TopicAdmin.java:724) at org.apache.kafka.connect.util.KafkaBasedLog.readEndOffsets(KafkaBasedLog.java:375) at org.apache.kafka.connect.util.KafkaBasedLog.readToLogEnd(KafkaBasedLog.java:335) at org.apache.kafka.connect.util.KafkaBasedLog.start(KafkaBasedLog.java:201) at org.apache.kafka.connect.storage.KafkaOffsetBackingStore.start(KafkaOffsetBackingStore.java:151) at org.apache.kafka.connect.runtime.Worker.start(Worker.java:186) at org.apache.kafka.connect.runtime.AbstractHerder.startServices(AbstractHerder.java:135) at org.apache.kafka.connect.runtime.distributed.DistributedHerder.run(DistributedHerder.java:320) at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511) at java.util.concurrent.FutureTask.run(FutureTask.java:266) at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) at java.lang.Thread.run(Thread.java:748) Caused by: org.apache.kafka.common.errors.TimeoutException: Call(callName=metadata, deadlineMs=1680677446792, tries=1, nextAllowedTryMs=1680677446893) timed out at 1680677446793 after 1 attempt(s) Caused by: org.apache.kafka.common.errors.DisconnectException: Cancelled metadata request with correlation id 530256 due to node 0 being disconnected
connect-offsets 主题详情
kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 -describe --topic connect-offsets Topic: connect-offsets TopicId: naW7oTihQwK PartitionCount: 25 ReplicationFactor: 1 Configs: cleanup.policy=compact,segment.bytes=1073741824 Topic: connect-offsets Partition: 1 Leader: 0 Replicas: 0 Isr: 0 Topic: connect-offsets Partition: 2 Leader: 0 Replicas: 0 Isr: 0 Topic: connect-offsets Partition: 3 Leader: 0 Replicas: 0 Isr: 0 Topic: connect-offsets Partition: 0 Leader: 0 Replicas: 0 Isr: 0 Topic: connect-offsets Partition: 4 Leader: 0 Replicas: 0 Isr: 0 Topic: connect-offsets Partition: 5 Leader: 0 Replicas: 0 Isr: 0 Topic: connect-offsets Partition: 6 Leader: 0 Replicas: 0 Isr: 0 Topic: connect-offsets Partition: 7 Leader: 0 Replicas: 0 Isr: 0 Topic: connect-offsets Partition: 8 Leader: 0 Replicas: 0 Isr: 0 Topic: connect-offsets Partition: 9 Leader: 0 Replicas: 0 Isr: 0 Topic: connect-offsets Partition: 10 Leader: 0 Replicas: 0 Isr: 0 Topic: connect-offsets Partition: 11 Leader: 0 Replicas: 0 Isr: 0 Topic: connect-offsets Partition: 12 Leader: 0 Replicas: 0 Isr: 0 Topic: connect-offsets Partition: 13 Leader: 0 Replicas: 0 Isr: 0 Topic: connect-offsets Partition: 14 Leader: 0 Replicas: 0 Isr: 0 Topic: connect-offsets Partition: 15 Leader: 0 Replicas: 0 Isr: 0 Topic: connect-offsets Partition: 16 Leader: 0 Replicas: 0 Isr: 0 Topic: connect-offsets Partition: 17 Leader: 0 Replicas: 0 Isr: 0 Topic: connect-offsets Partition: 18 Leader: 0 Replicas: 0 Isr: 0 Topic: connect-offsets Partition: 19 Leader: 0 Replicas: 0 Isr: 0 Topic: connect-offsets Partition: 20 Leader: 0 Replicas: 0 Isr: 0 Topic: connect-offsets Partition: 21 Leader: 0 Replicas: 0 Isr: 0 Topic: connect-offsets Partition: 22 Leader: 0 Replicas: 0 Isr: 0 Topic: connect-offsets Partition: 23 Leader: 0 Replicas: 0 Isr: 0 Topic: connect-offsets Partition: 24 Leader: 0 Replicas: 0 Isr: 0
排查与解决思路
验证Broker连通性
- 检查Connect节点能否通过
telnet *.*.*.* 9092或nc -zv *.*.*.* 9092测试端口连通性,排除网络防火墙或路由问题。 - 确认Broker的
listeners配置未仅绑定localhost,Connect的bootstrap.servers需与Broker对外监听地址完全一致。
- 检查Connect节点能否通过
调整超时参数
- 日志中出现
request timeout: 0ms,属于异常默认值,在connect-distributed.properties中添加以下参数:admin.client.connection.timeout.ms=30000 admin.client.request.timeout.ms=60000 request.timeout.ms=30000 metadata.max.age.ms=300000 - 可将
offset.flush.timeout.ms从300000(5分钟)调至600000(10分钟),给偏移量提交足够时间。
- 日志中出现
检查存储主题可用性
- 用
kafka-console-consumer.sh尝试消费connect-offsets主题,验证主题是否能正常读写:kafka/bin/kafka-console-consumer.sh --bootstrap-server *.*.*.*:9092 --topic connect-offsets --from-beginning - 若主题异常,可先停掉Connect服务,删除该主题后重启Connect,让其自动重建存储主题。
- 用
检查Broker资源状态
- 查看Broker的CPU、内存、磁盘使用率,若资源过载会导致请求超时;同时检查Broker日志,确认是否存在GC超时、磁盘IO过高的记录。
内容的提问来源于stack exchange,提问作者Emrahall
相关产品推荐
相关产品推荐

