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

如何在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客户端测试连接:
    1. 创建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
    
    1. 执行消费命令:
    kafka-console-consumer.sh --bootstrap-server <KAFKA_HOST>:<KAFKA_PORT> --topic <TOPIC_NAME> --consumer.config consumer.properties
    

内容的提问来源于stack exchange,提问作者duskformer

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 13:59:53