用的checkpoint后端存储是啥,flink哪个版本的

| |
田向阳
|
|
邮箱:lucas_...@163.com
|

签名由 网易邮箱大师 定制

在2021年05月20日 17:02,cecotw 写道:
2021-05-19 19:04:19
org.apache.flink.runtime.JobException: Recovery is suppressed by
NoRestartBackoffTimeStrategy
     at
org.apache.flink.runtime.executiongraph.failover.flip1.ExecutionFailureHandler.handleFailure(ExecutionFailureHandler.java:116)
     at
org.apache.flink.runtime.executiongraph.failover.flip1.ExecutionFailureHandler.getFailureHandlingResult(ExecutionFailureHandler.java:78)
     at
org.apache.flink.runtime.scheduler.DefaultScheduler.handleTaskFailure(DefaultScheduler.java:192)
     at
org.apache.flink.runtime.scheduler.DefaultScheduler.maybeHandleTaskFailure(DefaultScheduler.java:185)
     at
org.apache.flink.runtime.scheduler.DefaultScheduler.updateTaskExecutionStateInternal(DefaultScheduler.java:179)
     at
org.apache.flink.runtime.scheduler.SchedulerBase.updateTaskExecutionState(SchedulerBase.java:590)
     at
org.apache.flink.runtime.jobmaster.JobMaster.updateTaskExecutionState(JobMaster.java:384)
     at sun.reflect.GeneratedMethodAccessor28.invoke(Unknown Source)
     at
sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
     at java.lang.reflect.Method.invoke(Method.java:498)
     at
org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcInvocation(AkkaRpcActor.java:286)
     at
org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleRpcMessage(AkkaRpcActor.java:201)
     at
org.apache.flink.runtime.rpc.akka.FencedAkkaRpcActor.handleRpcMessage(FencedAkkaRpcActor.java:74)
     at
org.apache.flink.runtime.rpc.akka.AkkaRpcActor.handleMessage(AkkaRpcActor.java:154)
     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)
Caused by: java.lang.IllegalStateException: zip file closed
     at java.util.zip.ZipFile.ensureOpen(ZipFile.java:686)
     at java.util.zip.ZipFile.getEntry(ZipFile.java:315)
     at java.util.jar.JarFile.getEntry(JarFile.java:240)
     at sun.net.www.protocol.jar.URLJarFile.getEntry(URLJarFile.java:128)
     at java.util.jar.JarFile.getJarEntry(JarFile.java:223)
     at sun.misc.URLClassPath$JarLoader.getResource(URLClassPath.java:1054)
     at sun.misc.URLClassPath.getResource(URLClassPath.java:249)
     at java.net.URLClassLoader$1.run(URLClassLoader.java:366)
     at java.net.URLClassLoader$1.run(URLClassLoader.java:363)
     at java.security.AccessController.doPrivileged(Native Method)
     at java.net.URLClassLoader.findClass(URLClassLoader.java:362)
     at java.lang.ClassLoader.loadClass(ClassLoader.java:424)
     at
org.apache.flink.util.FlinkUserCodeClassLoader.loadClassWithoutExceptionHandling(FlinkUserCodeClassLoader.java:61)
     at
org.apache.flink.util.ChildFirstClassLoader.loadClassWithoutExceptionHandling(ChildFirstClassLoader.java:65)
     at
org.apache.flink.util.FlinkUserCodeClassLoader.loadClass(FlinkUserCodeClassLoader.java:48)
     at java.lang.ClassLoader.loadClass(ClassLoader.java:357)
     at
org.apache.flink.kafka.shaded.org.apache.kafka.clients.consumer.internals.ConsumerCoordinator.commitOffsetsAsync(ConsumerCoordinator.java:865)
     at
org.apache.flink.kafka.shaded.org.apache.kafka.clients.consumer.internals.ConsumerCoordinator.doAutoCommitOffsetsAsync(ConsumerCoordinator.java:973)
     at
org.apache.flink.kafka.shaded.org.apache.kafka.clients.consumer.internals.ConsumerCoordinator.maybeAutoCommitOffsetsAsync(ConsumerCoordinator.java:964)
     at
org.apache.flink.kafka.shaded.org.apache.kafka.clients.consumer.internals.ConsumerCoordinator.poll(ConsumerCoordinator.java:488)
     at
org.apache.flink.kafka.shaded.org.apache.kafka.clients.consumer.KafkaConsumer.updateAssignmentMetadataIfNeeded(KafkaConsumer.java:1267)
     at
org.apache.flink.kafka.shaded.org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1235)
     at
org.apache.flink.kafka.shaded.org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1168)
     at
org.apache.flink.streaming.connectors.kafka.internal.KafkaConsumerThread.run(KafkaConsumerThread.java:253)

job重新提交就可以了



--
Sent from: http://apache-flink.147419.n8.nabble.com/

回复