Flink Kafka Connector获取元数据失败问题求助
问题描述
作为Flink新手,编写Java作业实现从Kafka读取数据写入ClickHouse,Kafka为数据源。编译后通过Flink UI提交作业时触发如下错误:
org.apache.flink.runtime.JobException: Recovery is suppressed by NoRestartBackoffTimeStrategy at org.apache.flink.runtime.executiongraph.failover.ExecutionFailureHandler.handleFailure(ExecutionFailureHandler.java:219) at org.apache.flink.runtime.executiongraph.failover.ExecutionFailureHandler.handleFailureAndReport(ExecutionFailureHandler.java:166) at org.apache.flink.runtime.executiongraph.failover.ExecutionFailureHandler.getGlobalFailureHandlingResult(ExecutionFailureHandler.java:140) at org.apache.flink.runtime.scheduler.DefaultScheduler.handleGlobalFailure(DefaultScheduler.java:324) at org.apache.flink.runtime.operators.coordination.OperatorCoordinatorHolder$LazyInitializedCoordinatorContext.lambda$failJob$0(OperatorCoordinatorHolder.java:669) at org.apache.flink.runtime.rpc.pekko.PekkoRpcActor.lambda$handleRunAsync$4(PekkoRpcActor.java:460) at org.apache.flink.runtime.concurrent.ClassLoadingUtils.runWithContextClassLoader(ClassLoadingUtils.java:68) at org.apache.flink.runtime.rpc.pekko.PekkoRpcActor.handleRunAsync(PekkoRpcActor.java:460) at org.apache.flink.runtime.rpc.pekko.PekkoRpcActor.handleRpcMessage(PekkoRpcActor.java:225) at org.apache.flink.runtime.rpc.pekko.FencedPekkoRpcActor.handleRpcMessage(FencedPekkoRpcActor.java:88) at org.apache.flink.runtime.rpc.pekko.PekkoRpcActor.handleMessage(PekkoRpcActor.java:174) at org.apache.pekko.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:33) at org.apache.pekko.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:29) at scala.PartialFunction.applyOrElse(PartialFunction.scala:127) at scala.PartialFunction.applyOrElse$(PartialFunction.scala:126) at org.apache.pekko.japi.pf.UnitCaseStatement.applyOrElse(CaseStatements.scala:29) at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:175) at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:176) at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:176) at org.apache.pekko.actor.Actor.aroundReceive(Actor.scala:547) at org.apache.pekko.actor.Actor.aroundReceive$(Actor.scala:545) at org.apache.pekko.actor.AbstractActor.aroundReceive(AbstractActor.scala:229) at org.apache.pekko.actor.ActorCell.receiveMessage(ActorCell.scala:590) at org.apache.pekko.actor.ActorCell.invoke(ActorCell.scala:557) at org.apache.pekko.dispatch.Mailbox.processMailbox(Mailbox.scala:280) at org.apache.pekko.dispatch.Mailbox.run(Mailbox.scala:241) at org.apache.pekko.dispatch.Mailbox.exec(Mailbox.scala:253) at java.base/java.util.concurrent.ForkJoinTask.doExec(Unknown Source) at java.base/java.util.concurrent.ForkJoinPool$WorkQueue.topLevelExec(Unknown Source) at java.base/java.util.concurrent.ForkJoinPool.scan(Unknown Source) at java.base/java.util.concurrent.ForkJoinPool.runWorker(Unknown Source) at java.base/java.util.concurrent.ForkJoinWorkerThread.run(Unknown Source) Caused by: org.apache.flink.util.FlinkException: Global failure triggered by OperatorCoordinator for 'Source: Kafka Source -> Sink: Unnamed' (operator cbc357ccb763df2852fee8c4fc7d55f2). at org.apache.flink.runtime.operators.coordination.OperatorCoordinatorHolder$LazyInitializedCoordinatorContext.failJob(OperatorCoordinatorHolder.java:651) at org.apache.flink.runtime.operators.coordination.RecreateOnResetOperatorCoordinator$QuiesceableContext.failJob(RecreateOnResetOperatorCoordinator.java:259) at org.apache.flink.runtime.source.coordinator.SourceCoordinatorContext.failJob(SourceCoordinatorContext.java:432) at org.apache.flink.runtime.source.coordinator.SourceCoordinatorContext.handleUncaughtExceptionFromAsyncCall(SourceCoordinatorContext.java:445) at org.apache.flink.util.ThrowableCatchingRunnable.run(ThrowableCatchingRunnable.java:42) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Unknown Source) at java.base/java.util.concurrent.FutureTask.run(Unknown Source) at java.base/java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(Unknown Source) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(Unknown Source) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(Unknown Source) at java.base/java.lang.Thread.run(Unknown Source) Caused by: org.apache.flink.util.FlinkRuntimeException: Failed to list subscribed topic partitions due to at org.apache.flink.connector.kafka.source.enumerator.KafkaSourceEnumerator.checkPartitionChanges(KafkaSourceEnumerator.java:248) at org.apache.flink.runtime.source.coordinator.ExecutorNotifier.lambda$null$1(ExecutorNotifier.java:83) at org.apache.flink.util.ThrowableCatchingRunnable.run(ThrowableCatchingRunnable.java:40) ... 6 more Caused by: java.lang.RuntimeException: Failed to get metadata for topics [mytopic]. at org.apache.flink.connector.kafka.source.enumerator.subscriber.KafkaSubscriberUtils.getTopicMetadata(KafkaSubscriberUtils.java:47) at org.apache.flink.connector.kafka.source.enumerator.subscriber.TopicListSubscriber.getSubscribedTopicPartitions(TopicListSubscriber.java:52) at org.apache.flink.connector.kafka.source.enumerator.KafkaSourceEnumerator.getSubscribedTopicPartitions(KafkaSourceEnumerator.java:233) at org.apache.flink.runtime.source.coordinator.ExecutorNotifier.lambda$notifyReadyAsync$2(ExecutorNotifier.java:80) ... 6 more Caused by: java.util.concurrent.ExecutionException: org.apache.kafka.common.errors.TimeoutException: The AdminClient thread has exited. Call: describeTopics at java.base/java.util.concurrent.CompletableFuture.reportGet(Unknown Source) at java.base/java.util.concurrent.CompletableFuture.get(Unknown Source) at org.apache.kafka.common.internals.KafkaFutureImpl.get(KafkaFutureImpl.java:165) at org.apache.flink.connector.kafka.source.enumerator.subscriber.KafkaSubscriberUtils.getTopicMetadata(KafkaSubscriberUtils.java:44) ... 9 more Caused by: org.apache.kafka.common.errors.TimeoutException: The AdminClient thread has exited. Call: describeTopics
连接Kafka的代码片段:
KafkaSource<String> source = KafkaSource.<String>builder() .setBootstrapServers("boostrap.server:32100") .setTopics(parameter.get("topics", "newtopic")) .setGroupId(parameter.get("group", "mygroup")) .setStartingOffsets(OffsetsInitializer.earliest()) .setValueOnlyDeserializer(new SimpleStringSchema()) .setProperty("security.protocol", "SASL_SSL") .setProperty("ssl.truststore.location", "/opt/flink/kafkatruststore.jks") .setProperty("ssl.truststore.password", "changeit") .setProperty("ssl.keystore.location", "/opt/flink/kafkatruststore.jks") .setProperty("ssl.keystore.password", "changeit") .setProperty("sasl.mechanism", "SCRAM-SHA-256") .setProperty("sasl.jaas.config", "org.apache.kafka.common.security.scram.ScramLoginModule required username=\"username\" password=\"password\";)") .build(); //read data from kafka DataStream<String> text = env.fromSource(source, WatermarkStrategy.noWatermarks(), "Kafka Source");
环境说明:Kafka和Flink均部署在Kubernetes集群,通过UI提交作业,Kafka启用SSL加密。
解决方案
1. 修复JAAS配置语法错误
代码中sasl.jaas.config的配置末尾多了一个多余的括号),这会导致SASL认证失败,进而触发AdminClient超时。修正后的配置应为:
.setProperty("sasl.jaas.config", "org.apache.kafka.common.security.scram.ScramLoginModule required username=\"username\" password=\"password\";")
2. 验证证书文件的可用性
- 确认Flink Pod中
/opt/flink/kafkatruststore.jks路径存在:如果是通过K8s Secret/ConfigMap挂载的证书,检查挂载配置是否正确,路径是否匹配。 - 确保Flink进程对该文件有读取权限:可以在Pod中执行
ls -l /opt/flink/kafkatruststore.jks和cat /opt/flink/kafkatruststore.jks验证权限。
3. 延长Kafka客户端超时参数
添加以下配置,避免因K8s网络延迟导致AdminClient超时:
.setProperty("admin.request.timeout.ms", "30000") // 延长AdminClient请求超时 .setProperty("socket.connection.setup.timeout.ms", "30000") // 延长连接建立超时 .setProperty("metadata.max.age.ms", "30000") // 缩短元数据刷新间隔,避免过期
4. 验证网络连通性
在Flink Pod中执行命令,确认能连通Kafka的bootstrap服务器:
nc -zv boostrap.server 32100
如果是K8s内部服务,确认service名称正确,DNS解析正常(可以用nslookup boostrap.server验证)。
5. 检查Kafka ACL权限
确保配置的SASL用户拥有目标topic的Describe权限,否则AdminClient无法获取topic元数据。可以通过Kafka命令行工具验证:
kafka-acls.sh --bootstrap-server boostrap.server:32100 --command-config client.properties --list --topic mytopic
(client.properties需包含对应的SASL_SSL认证配置)
6. 确认版本兼容性
检查Flink和Kafka的版本是否兼容:
- Flink 1.17+ 推荐搭配 Kafka 2.8+
- Flink 1.15-1.16 推荐搭配 Kafka 2.6+
版本不兼容可能导致客户端与集群的协议交互异常,触发超时。
内容的提问来源于stack exchange,提问作者Ayush Goyal
相关产品推荐
相关产品推荐

