flink-dev mailing list archives

Site index · List index
Message view « Date » · « Thread »
Top « Date » · « Thread »
From Till Rohrmann <trohrm...@apache.org>
Subject Re: Not able to run beam appliction on Flink cluster environment
Date Tue, 22 Aug 2017 13:13:49 GMT
Hi Ramanji,

do you have the logs of the Flink master running at 192.168.56.1:6123?

Cheers,
Till

On Tue, Aug 22, 2017 at 2:43 PM, P. Ramanjaneya Reddy <ramanjieee@gmail.com>
wrote:

> Hi All,
>
> We followed the steps mentinoned in below link to setup flink cluster
> (Standalone)
> *https://ci.apache.org/projects/flink/flink-docs-
> release-1.2/setup/cluster_setup.html
> <https://ci.apache.org/projects/flink/flink-docs-
> release-1.2/setup/cluster_setup.html>*
>
>
> In the same cluster we are able to run the flink wordcount example, but the
> beam wordcount execution gives below error
>
> *commandline execution:*
> root1@master:~/Projects/beam/examples/java/target$      *mvn package
> exec:java -Dexec.mainClass=org.apache.beam.examples.WordCount \*
> *     -Dexec.args="--runner=FlinkRunner --flinkMaster="192.168.56.1:6123
> <http://192.168.56.1:6123>" --filesToStage=target/word-count-beam-0.1.jar
> \*
> *                  --inputFile=/home/root1/temp/input.txt
> --output=/home/root1/temp/output.txt" -Pflink-runner*
>
> *Logs:*
> NFO: Submitting job with JobID: 9edd3c2e1d318da5d3ffda1cdefa52c7. Waiting
> for job completion.
> Submitting job with JobID: 9edd3c2e1d318da5d3ffda1cdefa52c7. Waiting for
> job completion.
> Aug 22, 2017 2:56:05 PM
> org.apache.flink.client.program.ClusterClient$LazyActorSystemLoader get
> INFO: Starting client actor system.
> Aug 22, 2017 2:56:05 PM akka.event.slf4j.Slf4jLogger$$anonfun$receive$1
> applyOrElse
> INFO: Slf4jLogger started
> Aug 22, 2017 2:56:05 PM
> akka.event.slf4j.Slf4jLogger$$anonfun$receive$1$$anonfun$applyOrElse$3
> apply$mcV$sp
> INFO: Starting remoting
> Aug 22, 2017 2:56:06 PM
> akka.event.slf4j.Slf4jLogger$$anonfun$receive$1$$anonfun$applyOrElse$3
> apply$mcV$sp
> INFO: Remoting started; listening on addresses :[akka.tcp://flink@master
> :44871]
> Aug 22, 2017 2:56:06 PM org.apache.flink.runtime.client.JobClientActor
> disconnectFromJobManager
> *INFO: Disconnect from JobManager null.*
> Aug 22, 2017 2:56:06 PM org.apache.flink.runtime.client.JobClientActor
> handleMessage
> INFO: Received SubmitJobAndWait(JobGraph(jobId:
> 9edd3c2e1d318da5d3ffda1cdefa52c7)) but there is no connection to a
> JobManager yet.
> Aug 22, 2017 2:56:06 PM
> org.apache.flink.runtime.client.JobSubmissionClientActor
> handleCustomMessage
> INFO: Received job wordcount-root1-0822092604-654fbb92
> (9edd3c2e1d318da5d3ffda1cdefa52c7).
> Aug 22, 2017 2:57:06 PM org.apache.flink.runtime.client.JobClientActor
> terminate
> INFO: Terminate JobClientActor.
> Aug 22, 2017 2:57:06 PM org.apache.flink.runtime.client.JobClientActor
> disconnectFromJobManager
> INFO: Disconnect from JobManager null.
> Aug 22, 2017 2:57:06 PM
> akka.event.slf4j.Slf4jLogger$$anonfun$receive$1$$anonfun$applyOrElse$3
> apply$mcV$sp
> INFO: Shutting down remote daemon.
> Aug 22, 2017 2:57:06 PM
> akka.event.slf4j.Slf4jLogger$$anonfun$receive$1$$anonfun$applyOrElse$3
> apply$mcV$sp
> INFO: Remote daemon shut down; proceeding with flushing remote transports.
> Aug 22, 2017 2:57:06 PM
> akka.event.slf4j.Slf4jLogger$$anonfun$receive$1$$anonfun$applyOrElse$3
> apply$mcV$sp
> INFO: Remoting shut down.
> Aug 22, 2017 2:57:06 PM org.apache.beam.runners.flink.FlinkRunner run
> SEVERE: Pipeline execution failed
> org.apache.flink.client.program.ProgramInvocationException: The program
> execution failed: Couldn't retrieve the JobExecutionResult from the
> JobManager.
>     at
> org.apache.flink.client.program.ClusterClient.run(ClusterClient.java:427)
>     at
> org.apache.flink.client.program.StandaloneClusterClient.submitJob(
> StandaloneClusterClient.java:101)
>     at
> org.apache.flink.client.program.ClusterClient.run(ClusterClient.java:400)
>     at
> org.apache.flink.client.program.ClusterClient.run(ClusterClient.java:387)
>     at
> org.apache.flink.client.program.ClusterClient.run(ClusterClient.java:362)
>     at
> org.apache.flink.client.RemoteExecutor.executePlanWithJars(
> RemoteExecutor.java:211)
>     at
> org.apache.flink.client.RemoteExecutor.executePlan(
> RemoteExecutor.java:188)
>     at
> org.apache.flink.api.java.RemoteEnvironment.execute(
> RemoteEnvironment.java:172)
>     at
> org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironm
> ent.executePipeline(FlinkPipelineExecutionEnvironment.java:112)
>     at org.apache.beam.runners.flink.FlinkRunner.run(FlinkRunner.java:119)
>     at org.apache.beam.sdk.Pipeline.run(Pipeline.java:295)
>     at org.apache.beam.sdk.Pipeline.run(Pipeline.java:281)
>     at org.apache.beam.examples.WordCount.main(WordCount.java:184)
>     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.codehaus.mojo.exec.ExecJavaMojo$1.run(ExecJavaMojo.java:293)
>     at java.lang.Thread.run(Thread.java:748)
> Caused by: org.apache.flink.runtime.client.JobExecutionException: Couldn't
> retrieve the JobExecutionResult from the JobManager.
>     at
> org.apache.flink.runtime.client.JobClient.awaitJobResult(JobClient.java:
> 294)
>     at
> org.apache.flink.runtime.client.JobClient.submitJobAndWait(JobClient.
> java:382)
>     at
> org.apache.flink.client.program.ClusterClient.run(ClusterClient.java:423)
>     ... 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:207)
>     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:467)
>     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:238)
>     at akka.dispatch.Mailbox.run(Mailbox.scala:220)
>     at
> akka.dispatch.ForkJoinExecutorConfigurator$AkkaForkJoinTask.exec(
> AbstractDispatcher.scala:397)
>     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)
>
>
> Thanks & Regards,
> Ramanji.
>

Mime
  • Unnamed multipart/alternative (inline, None, 0 bytes)
View raw message