如何在Amazon EMR上解决Flink连接Kafka的AdminClient线程退出错误
问题
我尝试通过PyFlink Table API连接部署在DigitalOcean上的Kafka数据源,本地运行正常,但提交到Amazon EMR上的Flink集群时,作业立即被取消并抛出超时错误。
我的代码如下:
logging.basicConfig(stream=sys.stdout, level=logging.INFO, format="%(message)s") env = StreamExecutionEnvironment.get_execution_environment() env.set_parallelism(1) env.add_jars("file:///opt/flink/lib/flink-sql-connector-kafka-1.16.0.jar", "file:///opt/flink/lib/flink-sql-connector-elasticsearch7-1.16.0.jar", "file:///opt/flink/lib/flink-sql-avro-1.16.0.jar", "file:///opt/flink/lib/flink-sql-avro-confluent-registry-1.16.0.jar") t_env = StreamTableEnvironment.create(stream_execution_environment=env) t_env.get_config().get_configuration().set_boolean("python.fn-execution.memory.managed", True) create_kafka_source_ddl = """ CREATE TABLE KafkaTable ( **TABLE DETAILS** WITH ( 'connector' = 'kafka', 'scan.startup.mode' = 'latest-offset', 'topic' = '**TOPIC NAME**', 'properties.bootstrap.servers' = '**HOST:PORT**', 'properties.group.id' = 'GROUP ID', 'properties.security.protocol' = 'SASL_SSL', 'properties.sasl.mechanism' = 'SCRAM-SHA-512', 'properties.sasl.jaas.config' = 'org.apache.flink.kafka.shaded.org.apache.kafka.common.security.scram.ScramLoginModule required username=\"**USERNAME**\" password=\"**PASSOWORD**\";', 'properties.ssl.truststore.type' = 'PEM', 'properties.ssl.truststore.location' = '**/PATH/TO/PEM**', **REST ARE FORMATTING DETAILS** """ create_print_sink_ddl = """ CREATE TABLE print_sink( id STRING, name STRING, country_code STRING ) with ( 'connector' = 'print' ); """ t_env.execute_sql(create_kafka_source_ddl) t_env.execute_sql(create_print_sink_ddl) t_env.from_path("KafkaTable") \ .select(col('id'), col('name'), col('country_code')) \ .execute_insert("print_sink")
抛出的错误信息:
org.apache.flink.runtime.JobException: Recovery is suppressed by NoRestartBackoffTimeStrategy at org.apache.flink.runtime.executiongraph.failover.flip1.ExecutionFailureHandler.handleFailure(ExecutionFailureHandler.java:139) at org.apache.flink.runtime.executiongraph.failover.flip1.ExecutionFailureHandler.getGlobalFailureHandlingResult(ExecutionFailureHandler.java:102) at org.apache.flink.runtime.scheduler.DefaultScheduler.handleGlobalFailure(DefaultScheduler.java:299) at org.apache.flink.runtime.operators.coordination.OperatorCoordinatorHolder$LazyInitializedCoordinatorContext.lambda$failJob$0(OperatorCoordinatorHolder.java:635) at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.lambda$handleRunAsync$4(AkkaRpcActor.java:453) at org.apache.flink.runtime.concurrent.akka.ClassLoadingUtils.runWithContextClassLoader(ClassLoadingUtils.java:68) at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRunAsync(AkkaRpcActor.java:453) at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcMessage(AkkaRpcActor.java:218) at org.apache.flink.runtime.rpc.akka.FencedAkkaRpcActor.handleRpcMessage(FencedAkkaRpcActor.java:84) at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleMessage(AkkaRpcActor.java:168) at akka.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:24) at akka.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:20) at scala.PartialFunction.applyOrElse(PartialFunction.scala:123) at scala.PartialFunction.applyOrElse$(PartialFunction.scala:122) at akka.japi.pf.UnitCaseStatement.applyOrElse(CaseStatements.scala:20) at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:171) at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:172) at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:172) at akka.actor.Actor.aroundReceive(Actor.scala:537) at akka.actor.Actor.aroundReceive$(Actor.scala:535) at akka.actor.AbstractActor.aroundReceive(AbstractActor.scala:220) at akka.actor.ActorCell.receiveMessage(ActorCell.scala:580) at akka.actor.ActorCell.invoke(ActorCell.scala:548) at akka.dispatch.Mailbox.processMailbox(Mailbox.scala:270) at akka.dispatch.Mailbox.run(Mailbox.scala:231) at akka.dispatch.Mailbox.exec(Mailbox.scala:243) at java.util.concurrent.ForkJoinTask.doExec(ForkJoinTask.java:289) at java.util.concurrent.ForkJoinPool$WorkQueue.runTask(ForkJoinPool.java:1056) at java.util.concurrent.ForkJoinPool.runWorker(ForkJoinPool.java:1692) at java.util.concurrent.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:175) Caused by: org.apache.flink.util.FlinkException: Global failure triggered by OperatorCoordinator for 'Source: KafkaTable[1] -> Calc[2] -> Sink: print_sink[3]' (operator cbc357ccb763df2852fee8c4fc7d55f2). at org.apache.flink.runtime.operators.coordination.OperatorCoordinatorHolder$LazyInitializedCoordinatorContext.failJob(OperatorCoordinatorHolder.java:617) at org.apache.flink.runtime.operators.coordination.RecreateOnResetOperatorCoordinator$QuiesceableContext.failJob(RecreateOnResetOperatorCoordinator.java:237) at org.apache.flink.runtime.source.coordinator.SourceCoordinatorContext.failJob(SourceCoordinatorContext.java:360) at org.apache.flink.runtime.source.coordinator.SourceCoordinatorContext.handleUncaughtExceptionFromAsyncCall(SourceCoordinatorContext.java:373) at org.apache.flink.util.ThrowableCatchingRunnable.run(ThrowableCatchingRunnable.java:42) at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511) at java.util.concurrent.FutureTask.run(FutureTask.java:266) at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.access$201(ScheduledThreadPoolExecutor.java:180) at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:293) 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:750) 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:234) at org.apache.flink.runtime.source.coordinator.ExecutorNotifier.lambda$null$1(ExecutorNotifier.java:83) at org.apache.flink.util.ThrowableCatchingRunnable.run(ThrowableCatchingRunnable.java:40) ... 7 more Caused by: java.lang.RuntimeException: Failed to get metadata for topics [**TOPIC NAME**]. 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:219) at org.apache.flink.runtime.source.coordinator.ExecutorNotifier.lambda$notifyReadyAsync$2(ExecutorNotifier.java:80) ... 7 more Caused by: java.util.concurrent.ExecutionException: org.apache.flink.kafka.shaded.org.apache.kafka.common.errors.TimeoutException: The AdminClient thread has exited. Call: describeTopics at java.util.concurrent.CompletableFuture.reportGet(CompletableFuture.java:357) at java.util.concurrent.CompletableFuture.get(CompletableFuture.java:1908) at org.apache.flink.kafka.shaded.org.apache.kafka.common.internals.KafkaFutureImpl.get(KafkaFutureImpl.java:165) at org.apache.flink.connector.kafka.source.enumerator.subscriber.KafkaSubscriberUtils.getTopicMetadata(KafkaSubscriberUtils.java:44) ... 10 more Caused by: org.apache.flink.kafka.shaded.org.apache.kafka.common.errors.TimeoutException: The AdminClient thread has exited. Call: describeTopics
解决方案
核心问题是EMR集群无法在超时时间内连接到DigitalOcean的Kafka集群,无法获取Topic元数据,可按以下步骤排查修复:
1. 验证网络连通性
- 检查EMR集群安全组,确保允许出站访问DigitalOcean Kafka的端口(通常为9093)。
- 在EMR节点上执行命令测试连通性:
nc -zv <KAFKA_HOST> <KAFKA_PORT>,若不通,调整EMR安全组或DigitalOcean防火墙规则。 - 确认DigitalOcean Kafka集群已将EMR集群的公网IP加入访问白名单。
2. 延长Kafka客户端超时时间
在Kafka Source的WITH参数中添加以下配置,增加连接和元数据获取的等待时长:
'properties.request.timeout.ms' = '30000', 'properties.admin.timeout.ms' = '30000', 'properties.metadata.max.age.ms' = '30000'
3. 确认SSL证书配置有效性
- 确保PEM证书已上传至EMR集群的指定路径,且Flink进程拥有读取权限。
- 在EMR节点上执行命令验证证书:
openssl s_client -connect <KAFKA_HOST>:<KAFKA_PORT> -CAfile /PATH/TO/PEM。
4. 清理重复依赖
移除代码中env.add_jars这一行,因为指定的jar包已存在于/opt/flink/lib/目录下,重复加载可能引发冲突。
5. 调整作业资源配置
- 尝试提高作业并行度(如设置为2或更高),避免单任务压力过大。
- 提交作业时指定更多资源,例如:
-ys 4 -ym 4096(分配4个TaskManager,每个4GB内存)。
6. 验证Kafka权限与Topic存在性
- 确认配置的Topic名称正确,且Kafka账号拥有该Topic的读取权限。
- 在EMR节点上用Kafka客户端测试连接:
- 创建
consumer.properties文件:
security.protocol=SASL_SSL sasl.mechanism=SCRAM-SHA-512 sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required username="<USERNAME>" password="<PASSWORD>"; ssl.truststore.type=PEM ssl.truststore.location=/PATH/TO/PEM- 执行消费命令:
kafka-console-consumer.sh --bootstrap-server <KAFKA_HOST>:<KAFKA_PORT> --topic <TOPIC_NAME> --consumer.config consumer.properties - 创建
内容的提问来源于stack exchange,提问作者duskformer
相关产品推荐
相关产品推荐

