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)
>

Reply via email to