Hi Amir, what does the JM logs say?
Cheers, Till On Tue, Nov 8, 2016 at 9:33 PM, amir bahmanyari <amirto...@yahoo.com> wrote: > Hi colleagues, > I started the cluster all fine. Started the Beam app running in the Flink > Cluster fine. > Dashboard showed all tasks being consumed and open for business. > I started sending data to the Beam app, and all of the sudden the Flink JM > crashed. > Exceptions below. > Thanks+regards > Amir > > java.lang.RuntimeException: Pipeline execution failed > at org.apache.beam.runners.flink.FlinkRunner.run(FlinkRunner. > java:113) > at org.apache.beam.runners.flink.FlinkRunner.run(FlinkRunner. > java:48) > at org.apache.beam.sdk.Pipeline.run(Pipeline.java:183) > at > benchmark.flinkspark.flink.BenchBeamRunners.main(BenchBeamRunners.java:622) > //p.run(); > 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 org.apache.flink.client.program.PackagedProgram.callMainMethod( > PackagedProgram.java:505) > at org.apache.flink.client.program.PackagedProgram. > invokeInteractiveModeForExecution(PackagedProgram.java:403) > at org.apache.flink.client.program.Client.runBlocking( > Client.java:248) > at org.apache.flink.client.CliFrontend.executeProgramBlocking( > CliFrontend.java:866) > at org.apache.flink.client.CliFrontend.run(CliFrontend.java:333) > at org.apache.flink.client.CliFrontend.parseParameters( > CliFrontend.java:1189) > at org.apache.flink.client.CliFrontend.main(CliFrontend.java:1239) > Caused by: org.apache.flink.client.program.ProgramInvocationException: > The program execution failed: Communication with JobManager failed: Lost > connection to the JobManager. > at org.apache.flink.client.program.Client.runBlocking( > Client.java:381) > at org.apache.flink.client.program.Client.runBlocking( > Client.java:355) > at org.apache.flink.streaming.api.environment. > StreamContextEnvironment.execute(StreamContextEnvironment.java:65) > at org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironm > ent.executePipeline(FlinkPipelineExecutionEnvironment.java:118) > at org.apache.beam.runners.flink.FlinkRunner.run(FlinkRunner. > java:110) > ... 14 more > Caused by: org.apache.flink.runtime.client.JobExecutionException: > Communication with JobManager failed: Lost connection to the JobManager. > at org.apache.flink.runtime.client.JobClient. > submitJobAndWait(JobClient.java:140) > at org.apache.flink.client.program.Client.runBlocking( > Client.java:379) > ... 18 more > Caused by: > org.apache.flink.runtime.client.JobClientActorConnectionTimeoutException: > Lost connection to the JobManager. > at org.apache.flink.runtime.client.JobClientActor. > handleMessage(JobClientActor.java:244) > at org.apache.flink.runtime.akka.FlinkUntypedActor. > handleLeaderSessionID(FlinkUntypedActor.java:88) > at org.apache.flink.runtime.akka.FlinkUntypedActor.onReceive( > FlinkUntypedActor.java:68) > at akka.actor.UntypedActor$$anonfun$receive$1.applyOrElse( > UntypedActor.scala:167) > at akka.actor.Actor$class.aroundReceive(Actor.scala:465) > at akka.actor.UntypedActor.aroundReceive(UntypedActor.scala:97) > at akka.actor.ActorCell.receiveMessage(ActorCell.scala:516) > at akka.actor.ActorCell.invoke(ActorCell.scala:487) > at akka.dispatch.Mailbox.processMailbox(Mailbox.scala:254) > at akka.dispatch.Mailbox.run(Mailbox.scala:221) > at akka.dispatch.Mailbox.exec(Mailbox.scala:231) > at scala.concurrent.forkjoin.ForkJoinTask.doExec( > ForkJoinTask.java:260) > at scala.concurrent.forkjoin.ForkJoinPool$WorkQueue. > pollAndExecAll(ForkJoinPool.java:1253) > at scala.concurrent.forkjoin.ForkJoinPool$WorkQueue. > runTask(ForkJoinPool.java:1346) > at scala.concurrent.forkjoin.ForkJoinPool.runWorker( > ForkJoinPool.java:1979) > at scala.concurrent.forkjoin.ForkJoinWorkerThread.run( > ForkJoinWorkerThread.java:107) >