【问题标题】:Flink slot removed exceptionFlink 插槽移除异常
【发布时间】:2019-01-08 15:37:01
【问题描述】:

我收到以下异常

org.apache.flink.util.FlinkException: The assigned slot container_1546939492951_0001_01_003659_0 was removed.
at org.apache.flink.runtime.resourcemanager.slotmanager.SlotManager.removeSlot(SlotManager.java:789)
at org.apache.flink.runtime.resourcemanager.slotmanager.SlotManager.removeSlots(SlotManager.java:759)
at org.apache.flink.runtime.resourcemanager.slotmanager.SlotManager.internalUnregisterTaskManager(SlotManager.java:951)
at org.apache.flink.runtime.resourcemanager.slotmanager.SlotManager.unregisterTaskManager(SlotManager.java:372)
at org.apache.flink.runtime.resourcemanager.ResourceManager.closeTaskManagerConnection(ResourceManager.java:823)
at org.apache.flink.yarn.YarnResourceManager.lambda$onContainersCompleted$0(YarnResourceManager.java:346)
at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRunAsync(AkkaRpcActor.java:332)
at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcMessage(AkkaRpcActor.java:158)
at org.apache.flink.runtime.rpc.akka.FencedAkkaRpcActor.handleRpcMessage(FencedAkkaRpcActor.java:70)
at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.onReceive(AkkaRpcActor.java:142)
at org.apache.flink.runtime.rpc.akka.FencedAkkaRpcActor.onReceive(FencedAkkaRpcActor.java:40)
at akka.actor.UntypedActor$$anonfun$receive$1.applyOrElse(UntypedActor.scala:165)
at akka.actor.Actor$class.aroundReceive(Actor.scala:502)
at akka.actor.UntypedActor.aroundReceive(UntypedActor.scala:95)
at akka.actor.ActorCell.receiveMessage(ActorCell.scala:526)
at akka.actor.ActorCell.invoke(ActorCell.scala:495)
at akka.dispatch.Mailbox.processMailbox(Mailbox.scala:257)
at akka.dispatch.Mailbox.run(Mailbox.scala:224)
at akka.dispatch.Mailbox.exec(Mailbox.scala:234)
at scala.concurrent.forkjoin.ForkJoinTask.doExec(ForkJoinTask.java:260)
at scala.concurrent.forkjoin.ForkJoinPool$WorkQueue.runTask(ForkJoinPool.java:1339)
at scala.concurrent.forkjoin.ForkJoinPool.runWorker(ForkJoinPool.java:1979)
at scala.concurrent.forkjoin.ForkJoinWorkerThread.run(ForkJoinWorkerThread.java:107)

在运行涉及连接两个非常大的数据集的批处理时。

这是我在概览中看到的内容。失败发生在没有得到任何输入的任务管理器上。奇怪的是,尽管前面进行了重新平衡,但前一组(分区 -> 平面地图 -> 地图)并未向该任务管理器发送任何内容。

我在 EMR 上运行它。我看到有一个 slot.idle.timeout,这会产生影响吗?如果有,我该如何为该工作指定它?可以在命令行上完成吗?

【问题讨论】:

    标签: apache-flink


    【解决方案1】:

    这可能是一个超时问题,但通常当这种情况发生在我身上时,是因为出现了故障(例如,YARN 杀死了容器,因为它的运行超出了 pmem 或 vmem 限制)。我建议仔细检查 JobManager 和所有 TaskManager 日志文件。

    【讨论】:

    • 它在 YARN 之上的 EMR 上运行。作业完成后,我找不到访问任务管理器日志的方法,除非激活了到 S3 的日志记录。我使用了一个具有更多磁盘空间的集群,并且从那以后就能够完成这项工作。
    • 嗨 Julien - 很高兴这只是磁盘空间问题。仅供参考,我们通常在调试 Flink 作业时保持 EMR 集群处于活动状态,因为有时我们必须在盒子上四处寻找工作失败时出现的问题。是的,我们还将日志保存到 S3。
    • 您是否考虑过使用独立集群直接在 EC2 上运行它?除了节省 EMR 的额外成本之外,是否更容易找到任务的日志?
    • 我们已经完成了这两种方法,但理论上 EMR 简化了正确配置 YARN 集群的设置和运行。
    【解决方案2】:

    您可以在java代码中添加以下行。

    env.getCheckpointConfig().enableExternalizedCheckpoints(ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);

    那么您的工作将在取消时自动开始。

    【讨论】:

    • 能否添加一个指向 ApiDoc 的链接?这将使它成为一个更好的答案。
    【解决方案3】:

    我有一个类似的问题,结果证明是我们的 Flink 作业中的过度登录。我猜这会导致任务管理器超时。删除或减少日志记录的数量解决了这个问题

    【讨论】:

      【解决方案4】:

      我刚刚在 Kubernetes 上运行 Flink 时遇到了类似的问题,原来是 TaskManager 被 OOMKilled 并重新启动。如果你还在 Kubernetes 上运行 Flink,你可以检查 TaskManager pod 的状态:

      kubectl describe pods <pod>
      

      如果您看到容器之前被 OOKilled,这可能是原因:

          Last State:     Terminated
            Reason:       OOMKilled
            Exit Code:    137
      

      【讨论】:

        【解决方案5】:

        这个问题并不总是由oom引起并被纱线杀死,如果你有这样的日志:“Closing TaskExecutor connection container_e86_1590402668190_3503_01_000015 because: Container release on a lost node”,它在你之前错误日志。我猜这个问题是Nodemanager宕机引起的。大约10分钟,flink ResourceManager无法与NodeManager通信,ResourceManager将开始删除slot,并重新启动(如果你有重启策略)。

        【讨论】:

          猜你喜欢
          • 2018-01-28
          • 2014-08-20
          • 2014-07-18
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          相关资源
          最近更新 更多