【问题标题】:Apache flink - PartitionNotFoundExceptionApache flink - PartitionNotFoundException
【发布时间】:2019-01-23 18:32:25
【问题描述】:

我们在 kubernetes 和 azure 上运行一个 5 节点的 flink 集群(每个 8 gb 内存,总共 40 个插槽)。我们正在运行四个作业,所有作业都使用来自 kafka 的数据(每个都在不同的消费者组中)。 几天前,随着我们数据负载的增加,我们将生产者转移到 5 个 kafka 分区上生成数据,并将作业并行度增加到 5 个。 从那时起,我们时不时地(平均每小时)在我们的一个任务管理器上遇到以下异常:

NFO|N||-|||Flink-4jc| 2019-01-22 16:00:32,032 Task:917 - org=[] - Map (2/5) (949a8349e7bdcf3fe3b8f992f52d249c) switched from RUNNING to FAILED.
org.apache.flink.runtime.io.network.partition.PartitionNotFoundException: Partition 86656e59799eb529f24bac704ea06790@b1955e1a072e3b2f9e1f969fea509841 not found.
        at org.apache.flink.runtime.io.network.partition.consumer.RemoteInputChannel.failPartitionRequest(RemoteInputChannel.java:273)
        at org.apache.flink.runtime.io.network.partition.consumer.RemoteInputChannel.retriggerSubpartitionRequest(RemoteInputChannel.java:182)
        at org.apache.flink.runtime.io.network.partition.consumer.SingleInputGate.retriggerPartitionRequest(SingleInputGate.java:400)
        at org.apache.flink.runtime.taskmanager.Task.onPartitionStateUpdate(Task.java:1293)
        at org.apache.flink.runtime.taskmanager.Task.lambda$triggerPartitionProducerStateCheck$1(Task.java:1150)
        at java.util.concurrent.CompletableFuture.uniWhenComplete(CompletableFuture.java:760)
        at java.util.concurrent.CompletableFuture$UniWhenComplete.tryFire(CompletableFuture.java:736)
        at java.util.concurrent.CompletableFuture$Completion.run(CompletableFuture.java:442)
        at akka.dispatch.TaskInvocation.run(AbstractDispatcher.scala:39)
        at akka.dispatch.ForkJoinExecutorConfigurator$AkkaForkJoinTask.exec(AbstractDispatcher.scala:415)
        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)

异常发生在不同的任务和不同的工作上。 我已阅读以下主题: http://apache-flink-user-mailing-list-archive.2336050.n4.nabble.com/PartitionNotFoundException-when-running-in-yarn-session-td16081.html 这给了我一些可能导致异常的提示,但我仍然无法弄清楚是什么导致了我的情况(增加超时和网络缓冲区大小没有帮助,我不明白为什么 jar 文件大小很重要)

谁能引导我了解如何调查正在发生的事情、我应该打开哪些日志、要更改哪些配置等? 如果需要任何其他细节,我很乐意提供。

谢谢!

【问题讨论】:

  • 您能与我们分享集群入口点/作业管理器日志吗?理想情况下在 DEBUG 日志级别。

标签: apache-flink


【解决方案1】:

可能是你节点的网卡满了

【讨论】:

  • 这应该被视为评论,而不是答案
猜你喜欢
  • 1970-01-01
  • 2016-11-13
  • 2016-09-24
  • 1970-01-01
  • 2020-05-17
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多