【问题标题】:Kinesis Application - Flink 1.11 timeout exceptionKinesis 应用程序 - Flink 1.11 超时异常
【发布时间】:2022-01-27 10:27:04
【问题描述】:

我正在使用 Flink 1.11 处理 Kinesis 应用程序,但启动应用程序时出现以下错误:

java.util.concurrent.TimeoutException: The heartbeat of TaskManager with id 421563c271e57acb4592f9d447d45b42  timed out.
    at org.apache.flink.runtime.resourcemanager.ResourceManager$TaskManagerHeartbeatListener.notifyHeartbeatTimeout(ResourceManager.java:1202)
    at org.apache.flink.runtime.heartbeat.HeartbeatMonitorImpl.run(HeartbeatMonitorImpl.java:109)
    at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515)
    at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)
    at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRunAsync(AkkaRpcActor.java:402)
    at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcMessage(AkkaRpcActor.java:195)
    at org.apache.flink.runtime.rpc.akka.FencedAkkaRpcActor.handleRpcMessage(FencedAkkaRpcActor.java:74)
    at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleMessage(AkkaRpcActor.java:152)
    at akka.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:26)
    at akka.japi.pf.UnitCaseStatement.apply(CaseStatements.scala:21)
    at scala.PartialFunction.applyOrElse(PartialFunction.scala:123)
    at scala.PartialFunction.applyOrElse$(PartialFunction.scala:122)
    at akka.japi.pf.UnitCaseStatement.applyOrElse(CaseStatements.scala:21)
    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:517)
    at akka.actor.Actor.aroundReceive$(Actor.scala:515)
    at akka.actor.AbstractActor.aroundReceive(AbstractActor.scala:225)
    at akka.actor.ActorCell.receiveMessage(ActorCell.scala:592)
    at akka.actor.ActorCell.invoke(ActorCell.scala:561)
    at akka.dispatch.Mailbox.processMailbox(Mailbox.scala:258)
    at akka.dispatch.Mailbox.run(Mailbox.scala:225)
    at akka.dispatch.Mailbox.exec(Mailbox.scala:235)
    at akka.dispatch.forkjoin.ForkJoinTask.doExec(ForkJoinTask.java:260)
    at akka.dispatch.forkjoin.ForkJoinPool$WorkQueue.runTask(ForkJoinPool.java:1339)
    at akka.dispatch.forkjoin.ForkJoinPool.runWorker(ForkJoinPool.java:1979)
    at akka.dispatch.forkjoin.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:107)

我在 Kinesis 和 4 KPU 中使用默认配置。

【问题讨论】:

    标签: scala apache-flink amazon-kinesis amazon-kinesis-analytics


    【解决方案1】:

    从报错来看,TM和JM的心跳已经超时,或许可以给出Job的完整配置信息。

    JM和TM之间超时的原因有很多,比如并发任务大,JM负载高等

    或许你也可以调整heartbeat.timeout参数,参考https://nightlies.apache.org/flink/flink-docs-master/docs/deployment/config/

    【讨论】:

    • 不幸的是,我正在使用 AWS Kinesis 应用程序实施并基于 AWS 文档 Kinesis 使用 Apache Flink 文档中描述的默认配置。我还读到大多数配置参数是不可修改的。
    • 如果你不能调整超时,给集群更多的资源(例如,增加并行度)应该减少超时的可能性。
    【解决方案2】:

    最后我解决了这个问题。就我而言,问题是我正在读取一个对于 Flink 工作人员来说太大的 CSVSource。我使用 dynamodb 作为数据源解决了这个问题。

    【讨论】:

    • 正如目前所写,您的答案尚不清楚。请edit 添加其他详细信息,以帮助其他人了解这如何解决所提出的问题。你可以找到更多关于如何写好答案的信息in the help center
    猜你喜欢
    • 2018-06-13
    • 2019-12-12
    • 1970-01-01
    • 2017-07-08
    • 2013-03-15
    • 2016-07-07
    • 2020-09-23
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多