【问题标题】:Flink taskmanager stucks (100% cpu usage) after failing to make a checkpoint/savepointFlink taskmanager 在创建检查点/保存点失败后卡住(100% cpu 使用率)
【发布时间】:2020-08-27 17:12:54
【问题描述】:

-- 通过将状态后端从文件系统更改为rocksdb来解决问题--

在 AWS EMR 上运行 Flink 1.9。 Flink 应用程序使用 kinesis 流作为输入数据,使用另一个 kinesis 流作为输出。 最近检查点大小已增长到 1 GB(由于更多数据)。 有时,在尝试获取检查点的过程中 - 应用程序开始利用整个处理器资源(一天发生几次)

指标:

LA (emr ec2 core node with job/task managers)

Run Loop Time - kinesis consumer

Records Per Fetch - kinesis consumer

Task manager GC

jobmanager 日志

{"level":"INFO","timestamp":"2020-08-25 04:55:27,399","thread":"Checkpoint Timer","file":"CheckpointCoordinator.java","line":"617","message":"Triggering checkpoint 1232 @ 1598331327244 for job 0039825bafae26bc34db88e037a1dae3."}

{"level":"INFO","timestamp":"2020-08-25 04:58:24,509","thread":"flink-akka.actor.default-dispatcher-7010","file":"ResourceManager.java","line":"1144","message":"The heartbeat of TaskManager with id container_1597960565773_0003_01_000002 timed out."}

{"level":"INFO","timestamp":"2020-08-25 04:58:24,510","thread":"flink-akka.actor.default-dispatcher-7010","file":"ResourceManager.java","line":"805","message":"Closing TaskExecutor connection container_1597960565773_0003_01_000002 because: The heartbeat of TaskManager with id container_1597960565773_0003_01_000002  timed out."}
{"level":"INFO","timestamp":"2020-08-25 04:58:24,514","thread":"flink-akka.actor.default-dispatcher-7015","file":"Execution.java","line":"1493","message":"Sink: kinesis-events-sink (1/1) (573401e241fe0a0ac0a8a54c81c4eefd) switched from RUNNING to FAILED."}
java.util.concurrent.TimeoutException: The heartbeat of TaskManager with id container_1597960565773_0003_01_000002  timed out.
        at org.apache.flink.runtime.resourcemanager.ResourceManager$TaskManagerHeartbeatListener.notifyHeartbeatTimeout(ResourceManager.java:1146)
        at org.apache.flink.runtime.heartbeat.HeartbeatMonitorImpl.run(HeartbeatMonitorImpl.java:109)
        at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
        at java.util.concurrent.FutureTask.run(FutureTask.java:266)
        at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRunAsync(AkkaRpcActor.java:397)
        at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcMessage(AkkaRpcActor.java:190)
        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$class.applyOrElse(PartialFunction.scala:123)
        at akka.japi.pf.UnitCaseStatement.applyOrElse(CaseStatements.scala:21)
        at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:170)
        at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:171)
        at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:171)
        at akka.actor.Actor$class.aroundReceive(Actor.scala:517)
        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)
{"level":"INFO","timestamp":"2020-08-25 04:58:24,515","thread":"flink-akka.actor.default-dispatcher-7015","file":"ExecutionGraph.java","line":"1324","message":"Job JOB_NAME_HIDDEN (0039825bafae26bc34db88e037a1dae3) switched from state RUNNING to FAILING."}
java.util.concurrent.TimeoutException: The heartbeat of TaskManager with id container_1597960565773_0003_01_000002  timed out.
        at org.apache.flink.runtime.resourcemanager.ResourceManager$TaskManagerHeartbeatListener.notifyHeartbeatTimeout(ResourceManager.java:1146)
        at org.apache.flink.runtime.heartbeat.HeartbeatMonitorImpl.run(HeartbeatMonitorImpl.java:109)
        at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
        at java.util.concurrent.FutureTask.run(FutureTask.java:266)
        at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRunAsync(AkkaRpcActor.java:397)
        at org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcMessage(AkkaRpcActor.java:190)
        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$class.applyOrElse(PartialFunction.scala:123)
        at akka.japi.pf.UnitCaseStatement.applyOrElse(CaseStatements.scala:21)
        at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:170)
        at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:171)
        at scala.PartialFunction$OrElse.applyOrElse(PartialFunction.scala:171)
        at akka.actor.Actor$class.aroundReceive(Actor.scala:517)
        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)

任务管理器日志

{"level":"INFO","timestamp":"2020-08-25 04:58:17,504","thread":"kpl-daemon-0003","file":"LogInputStreamReader.java","line":"59","message":"[2020-08-25 04:58:02.469832] [0x00007394][0x00007f2673d70700] [info] [processing_statistics_logger.cc:127] (sensing-events-test) Aver
age Processing Time: -nan ms"}
{"level":"INFO","timestamp":"2020-08-25 04:58:17,504","thread":"kpl-daemon-0003","file":"LogInputStreamReader.java","line":"59","message":"[2020-08-25 04:58:17.469922] [0x00007394][0x00007f2673d70700] [info] [processing_statistics_logger.cc:109] Stage 1 Triggers: { stream
: 'sensing-events-test', manual: 0, count: 0, size: 0, matches: 0, timed: 0, UserRecords: 0, KinesisRecords: 0 }"}
{"level":"INFO","timestamp":"2020-08-25 04:58:17,504","thread":"kpl-daemon-0003","file":"LogInputStreamReader.java","line":"59","message":"[2020-08-25 04:58:17.469977] [0x00007394][0x00007f2673d70700] [info] [processing_statistics_logger.cc:112] Stage 2 Triggers: { stream
: 'sensing-events-test', manual: 0, count: 0, size: 0, matches: 0, timed: 0, KinesisRecords: 0, PutRecords: 0 }"}
{"level":"INFO","timestamp":"2020-08-25 04:58:17,504","thread":"kpl-daemon-0003","file":"LogInputStreamReader.java","line":"59","message":"[2020-08-25 04:58:17.469992] [0x00007394][0x00007f2673d70700] [info] [processing_statistics_logger.cc:127] (sensing-events-test) Aver
age Processing Time: -nan ms"}
{"level":"ERROR","timestamp":"2020-08-25 04:58:28,535","thread":"flink-akka.actor.default-dispatcher-628","file":"FatalExitExceptionHandler.java","line":"40","message":"FATAL: Thread 'flink-akka.actor.default-dispatcher-628' produced an uncaught exception. Stopping the pr
ocess..."}
java.lang.OutOfMemoryError: GC overhead limit exceeded
        at sun.reflect.UTF8.encode(UTF8.java:36)
        at sun.reflect.ClassFileAssembler.emitConstantPoolUTF8(ClassFileAssembler.java:103)
        at sun.reflect.MethodAccessorGenerator.generate(MethodAccessorGenerator.java:331)
        at sun.reflect.MethodAccessorGenerator.generateSerializationConstructor(MethodAccessorGenerator.java:112)
        at sun.reflect.ReflectionFactory.generateConstructor(ReflectionFactory.java:398)
        at sun.reflect.ReflectionFactory.newConstructorForSerialization(ReflectionFactory.java:360)
        at java.io.ObjectStreamClass.getSerializableConstructor(ObjectStreamClass.java:1588)
        at java.io.ObjectStreamClass.access$1500(ObjectStreamClass.java:79)
        at java.io.ObjectStreamClass$3.run(ObjectStreamClass.java:519)
        at java.io.ObjectStreamClass$3.run(ObjectStreamClass.java:494)
        at java.security.AccessController.doPrivileged(Native Method)
        at java.io.ObjectStreamClass.<init>(ObjectStreamClass.java:494)
        at java.io.ObjectStreamClass.lookup(ObjectStreamClass.java:391)
        at java.io.ObjectStreamClass.initNonProxy(ObjectStreamClass.java:681)
        at java.io.ObjectInputStream.readNonProxyDesc(ObjectInputStream.java:1941)
        at java.io.ObjectInputStream.readClassDesc(ObjectInputStream.java:1807)
        at java.io.ObjectInputStream.readClass(ObjectInputStream.java:1770)
        at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1595)
        at java.io.ObjectInputStream.readObject(ObjectInputStream.java:464)
        at java.io.ObjectInputStream.readObject(ObjectInputStream.java:422)
        at org.apache.flink.runtime.rpc.messages.RemoteRpcInvocation$MethodInvocation.readObject(RemoteRpcInvocation.java:204)
        at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
        at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
        at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
        at java.lang.reflect.Method.invoke(Method.java:498)
        at java.io.ObjectStreamClass.invokeReadObject(ObjectStreamClass.java:1184)
        at java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2234)
        at java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2125)
        at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1624)
        at java.io.ObjectInputStream.readObject(ObjectInputStream.java:464)
        at java.io.ObjectInputStream.readObject(ObjectInputStream.java:422)
        at org.apache.flink.util.InstantiationUtil.deserializeObject(InstantiationUtil.java:576)

更新。 flink-conf.yaml

state.backend.fs.checkpointdir: s3a://s3-bucket-with-checkpoints/flink-checkpoints
taskmanager.numberOfTaskSlots: 1
state.backend: filesystem
taskmanager.heap.size: 3057m
state.checkpoints.dir: s3a://s3-bucket-with-checkpoints/external-checkpoints

Flink Checkpoints

【问题讨论】:

    标签: apache-flink flink-streaming amazon-kinesis


    【解决方案1】:

    通过将状态后端从文件系统更改为rocksdb解决了问题

    【讨论】:

      【解决方案2】:

      我认为,这可能与SlidingEventTimeWindow 有关,据我从检查点屏幕截图中了解,它是一个大小为 2 分钟的窗口,带有 2 秒的窗口幻灯片。 Flink 为每个元素所属的每个窗口创建一个副本。因此,对于滑动窗口,它会创建大约 60 个元素副本,因此状态大小是滚动窗口的 60 倍。

      我猜,在检查点 flink 尝试序列化状态并且没有足够的内存,因此 GC 启动,最后你的内存用完了。

      【讨论】:

      • 从图表来看,看起来当 flink 开始创建保存点时,kinesis 客户端停止读取,然后在检查点之后,flink 得到“冲击”负载。应该是这样吗?我需要调整什么内存参数才能有足够的内存?在图表上,我没有看到内存耗尽。这在相同的负载下偶然发生。
      • 同样实现 CheckpointedFunction 接口的源必须确保状态检查点、内部状态更新和元素发射不会同时进行。这是通过使用检查点锁对象来保护同步块中的状态更新和元素发射来实现的。因此,新记录仅在同步块内收集,因此当 snaphost 启动时,消费者等待快照完成。从内存的角度来看,每个 TM 上的 flink 托管内存量是多少?您能否提供您的flink-conf.yaml 设置?谢谢。
      • 在顶帖中附上flink-conf.yaml
      • 很抱歉问了太多问题,但请您澄清一下您是在独立模式还是在任何容器化环境中启动 flink 集群?另外,你使用什么版本的flink?自 1.10 以来,内存管理发生了变化。通常,在文件系统状态后端的情况下,建议将托管内存设置为零。这将确保为 JVM 上的用户代码分配最大的内存量。
      • 通过将状态后端从文件系统更改为rocksdb解决了问题
      猜你喜欢
      • 2020-09-16
      • 1970-01-01
      • 1970-01-01
      • 2023-01-18
      • 2021-01-27
      • 1970-01-01
      • 1970-01-01
      • 2019-07-25
      • 1970-01-01
      相关资源
      最近更新 更多