Giter Site home page Giter Site logo

beamexample's People

Contributors

dhalperi avatar eljefe6a avatar francesperry avatar kennknowles avatar takidau avatar

Stargazers

 avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar

Watchers

 avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar  avatar

beamexample's Issues

Remote Flink Job Execution Error: Couldn't retrieve the JobExecutionResult from the JobManager. Lost connection to the JobManager

Hi,

Thanks for the sample code, very helpful in learning Apache Beam.

I am trying to run job on Remote Apache Flink Cluster. I notice that after compilation, it takes 60 sec to timeout.

Command:

mvn compile exec:java -Dexec.mainClass=org.apache.beam.examples.tutorial.game.solution.Exercise1 -Dexec.args='--runner=FlinkRunner --flinkMaster=192.168.8.10:6123 --filesToStage=./target/BeamTutorial-bundled-flink.jar' -Pflink-runner

Output

[INFO] Scanning for projects...
[INFO]
[INFO] ------------------------------------------------------------------------
[INFO] Building BeamTutorial 0.0.1-SNAPSHOT
[INFO] ------------------------------------------------------------------------
[INFO]
[INFO] --- maven-resources-plugin:2.6:resources (default-resources) @ BeamTutorial ---
[WARNING] Using platform encoding (ANSI_X3.4-1968 actually) to copy filtered resources, i.e. build is platform dependent!
[INFO] skip non existing resourceDirectory /beamexample/BeamTutorial/src/main/resources
[INFO]
[INFO] --- maven-compiler-plugin:3.5.1:compile (default-compile) @ BeamTutorial ---
[INFO] Nothing to compile - all classes are up to date
[INFO]
[INFO] --- exec-maven-plugin:1.4.0:java (default-cli) @ BeamTutorial ---
SLF4J: Class path contains multiple SLF4J bindings.
SLF4J: Found binding in [jar:file:/root/.m2/repository/org/slf4j/slf4j-jdk14/1.7.14/slf4j-jdk14-1.7.14.jar!/org/slf4j/impl/StaticLoggerBinder.class]
SLF4J: Found binding in [jar:file:/root/.m2/repository/org/slf4j/slf4j-log4j12/1.7.7/slf4j-log4j12-1.7.7.jar!/org/slf4j/impl/StaticLoggerBinder.class]
SLF4J: See http://www.slf4j.org/codes.html#multiple_bindings for an explanation.
SLF4J: Actual binding is of type [org.slf4j.impl.JDK14LoggerFactory]
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkRunner run
INFO: Executing pipeline using FlinkRunner.
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkRunner run
INFO: Translating pipeline to Flink program.
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironment createBatchExecutionEnvironment
INFO: Creating the required Batch Execution Environment.
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator enterCompositeTransform
INFO: enterCompositeTransform-
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator enterCompositeTransform
INFO: | enterCompositeTransform- Input.BoundedGenerator
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator visitPrimitiveTransform
INFO: | | visitPrimitiveTransform- Input.BoundedGenerator/Read(InjectorBoundedSource)
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator leaveCompositeTransform
INFO: | leaveCompositeTransform- Input.BoundedGenerator
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator enterCompositeTransform
INFO: | enterCompositeTransform- ExtractUserScore
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator enterCompositeTransform
INFO: | | enterCompositeTransform- ExtractUserScore/MapElements
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator enterCompositeTransform
INFO: | | | enterCompositeTransform- ExtractUserScore/MapElements/Map
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator visitPrimitiveTransform
INFO: | | | | visitPrimitiveTransform- ExtractUserScore/MapElements/Map/ParMultiDo(Anonymous)
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator leaveCompositeTransform
INFO: | | | leaveCompositeTransform- ExtractUserScore/MapElements/Map
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator leaveCompositeTransform
INFO: | | leaveCompositeTransform- ExtractUserScore/MapElements
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator enterCompositeTransform
INFO: | | enterCompositeTransform- ExtractUserScore/Combine.perKey(SumInteger)
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator enterCompositeTransform
INFO: | | | translated- ExtractUserScore/Combine.perKey(SumInteger)
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator leaveCompositeTransform
INFO: | | leaveCompositeTransform- ExtractUserScore/Combine.perKey(SumInteger)
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator leaveCompositeTransform
INFO: | leaveCompositeTransform- ExtractUserScore
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator enterCompositeTransform
INFO: | enterCompositeTransform- Output.WriteUserScoreSums
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator enterCompositeTransform
INFO: | | enterCompositeTransform- Output.WriteUserScoreSums/MapContextElements
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator enterCompositeTransform
INFO: | | | enterCompositeTransform- Output.WriteUserScoreSums/MapContextElements/Map
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator visitPrimitiveTransform
INFO: | | | | visitPrimitiveTransform- Output.WriteUserScoreSums/MapContextElements/Map/ParMultiDo(Anonymous)
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator leaveCompositeTransform
INFO: | | | leaveCompositeTransform- Output.WriteUserScoreSums/MapContextElements/Map
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator leaveCompositeTransform
INFO: | | leaveCompositeTransform- Output.WriteUserScoreSums/MapContextElements
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator enterCompositeTransform
INFO: | | enterCompositeTransform- Output.WriteUserScoreSums/TextIO.Write
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator enterCompositeTransform
INFO: | | | enterCompositeTransform- Output.WriteUserScoreSums/TextIO.Write/WriteFiles
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator enterCompositeTransform
INFO: | | | | enterCompositeTransform- Output.WriteUserScoreSums/TextIO.Write/WriteFiles/Window.Into()
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator visitPrimitiveTransform
INFO: | | | | | visitPrimitiveTransform- Output.WriteUserScoreSums/TextIO.Write/WriteFiles/Window.Into()/Window.Assign
Aug 29, 2017 12:40:02 PM org.apache.flink.api.java.typeutils.TypeExtractor analyzePojo
INFO: No fields detected for class org.apache.beam.sdk.util.WindowedValue. Cannot be used as a PojoType. Will be handled as GenericType
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator leaveCompositeTransform
INFO: | | | | leaveCompositeTransform- Output.WriteUserScoreSums/TextIO.Write/WriteFiles/Window.Into()
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator enterCompositeTransform
INFO: | | | | enterCompositeTransform- Output.WriteUserScoreSums/TextIO.Write/WriteFiles/WriteBundles
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator visitPrimitiveTransform
INFO: | | | | | visitPrimitiveTransform- Output.WriteUserScoreSums/TextIO.Write/WriteFiles/WriteBundles/ParMultiDo(WriteUnwindowedBundles)
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator leaveCompositeTransform
INFO: | | | | leaveCompositeTransform- Output.WriteUserScoreSums/TextIO.Write/WriteFiles/WriteBundles
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator enterCompositeTransform
INFO: | | | | enterCompositeTransform- Output.WriteUserScoreSums/TextIO.Write/WriteFiles/View.AsIterable
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator visitPrimitiveTransform
INFO: | | | | | visitPrimitiveTransform- Output.WriteUserScoreSums/TextIO.Write/WriteFiles/View.AsIterable/View.CreatePCollectionView
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator leaveCompositeTransform
INFO: | | | | leaveCompositeTransform- Output.WriteUserScoreSums/TextIO.Write/WriteFiles/View.AsIterable
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator enterCompositeTransform
INFO: | | | | enterCompositeTransform- Output.WriteUserScoreSums/TextIO.Write/WriteFiles/Create.Values
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator visitPrimitiveTransform
INFO: | | | | | visitPrimitiveTransform- Output.WriteUserScoreSums/TextIO.Write/WriteFiles/Create.Values/Read(CreateSource)
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator leaveCompositeTransform
INFO: | | | | leaveCompositeTransform- Output.WriteUserScoreSums/TextIO.Write/WriteFiles/Create.Values
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator enterCompositeTransform
INFO: | | | | enterCompositeTransform- Output.WriteUserScoreSums/TextIO.Write/WriteFiles/Finalize
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator visitPrimitiveTransform
INFO: | | | | | visitPrimitiveTransform- Output.WriteUserScoreSums/TextIO.Write/WriteFiles/Finalize/ParMultiDo(Anonymous)
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator leaveCompositeTransform
INFO: | | | | leaveCompositeTransform- Output.WriteUserScoreSums/TextIO.Write/WriteFiles/Finalize
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator leaveCompositeTransform
INFO: | | | leaveCompositeTransform- Output.WriteUserScoreSums/TextIO.Write/WriteFiles
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator leaveCompositeTransform
INFO: | | leaveCompositeTransform- Output.WriteUserScoreSums/TextIO.Write
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator leaveCompositeTransform
INFO: | leaveCompositeTransform- Output.WriteUserScoreSums
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkBatchPipelineTranslator leaveCompositeTransform
INFO: leaveCompositeTransform-
Aug 29, 2017 12:40:02 PM org.apache.beam.runners.flink.FlinkRunner run
INFO: Starting execution of Flink program.
Aug 29, 2017 12:40:03 PM org.apache.flink.api.java.ExecutionEnvironment createProgramPlan
INFO: The job has 0 registered types and 0 default Kryo serializers
Aug 29, 2017 12:40:04 PM org.apache.flink.client.program.ClusterClient logAndSysout
INFO: Submitting job with JobID: 4e46320ec2a0eae86718e4b4637c07d8. Waiting for job completion.
Submitting job with JobID: 4e46320ec2a0eae86718e4b4637c07d8. Waiting for job completion.
Aug 29, 2017 12:40:04 PM org.apache.flink.client.program.ClusterClient$LazyActorSystemLoader get
INFO: Starting client actor system.
Aug 29, 2017 12:40:05 PM akka.event.slf4j.Slf4jLogger$$anonfun$receive$1 applyOrElse
INFO: Slf4jLogger started
Aug 29, 2017 12:40:06 PM akka.event.slf4j.Slf4jLogger$$anonfun$receive$1$$anonfun$applyOrElse$3 apply$mcV$sp
INFO: Starting remoting
Aug 29, 2017 12:40:06 PM akka.event.slf4j.Slf4jLogger$$anonfun$receive$1$$anonfun$applyOrElse$3 apply$mcV$sp
INFO: Remoting started; listening on addresses :[akka.tcp://flink@2e56b291bd16:44755]
Aug 29, 2017 12:40:06 PM org.apache.flink.runtime.client.JobClientActor disconnectFromJobManager
INFO: Disconnect from JobManager null.
Aug 29, 2017 12:40:07 PM org.apache.flink.runtime.client.JobClientActor handleMessage
INFO: Received SubmitJobAndWait(JobGraph(jobId: 4e46320ec2a0eae86718e4b4637c07d8)) but there is no connection to a JobManager yet.
Aug 29, 2017 12:40:07 PM org.apache.flink.runtime.client.JobSubmissionClientActor handleCustomMessage
INFO: Received job exercise1-root-0829124002-a4d765bd (4e46320ec2a0eae86718e4b4637c07d8).
Aug 29, 2017 12:41:07 PM org.apache.flink.runtime.client.JobClientActor terminate
INFO: Terminate JobClientActor.
Aug 29, 2017 12:41:07 PM org.apache.flink.runtime.client.JobClientActor disconnectFromJobManager
INFO: Disconnect from JobManager null.
Aug 29, 2017 12:41:07 PM akka.event.slf4j.Slf4jLogger$$anonfun$receive$1$$anonfun$applyOrElse$3 apply$mcV$sp
INFO: Shutting down remote daemon.
Aug 29, 2017 12:41:07 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 29, 2017 12:41:07 PM akka.event.slf4j.Slf4jLogger$$anonfun$receive$1$$anonfun$applyOrElse$3 apply$mcV$sp
INFO: Remoting shut down.
Aug 29, 2017 12:41:07 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.FlinkPipelineExecutionEnvironment.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.tutorial.game.solution.Exercise1.main(Exercise1.java:136)
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)

[WARNING]
java.lang.reflect.InvocationTargetException
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: java.lang.RuntimeException: Pipeline execution failed
at org.apache.beam.runners.flink.FlinkRunner.run(FlinkRunner.java:122)
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.tutorial.game.solution.Exercise1.main(Exercise1.java:136)
... 6 more
Caused by: 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.FlinkPipelineExecutionEnvironment.executePipeline(FlinkPipelineExecutionEnvironment.java:112)
at org.apache.beam.runners.flink.FlinkRunner.run(FlinkRunner.java:119)
... 9 more
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)
[INFO] ------------------------------------------------------------------------
[INFO] BUILD FAILURE
[INFO] ------------------------------------------------------------------------
[INFO] Total time: 01:15 min
[INFO] Finished at: 2017-08-29T12:41:07+00:00
[INFO] Final Memory: 25M/193M
[INFO] ------------------------------------------------------------------------
[ERROR] Failed to execute goal org.codehaus.mojo:exec-maven-plugin:1.4.0:java (default-cli) on project BeamTutorial: An exception occured while executing the Java class. null: InvocationTargetException: Pipeline execution failed: The program execution failed: Couldn't retrieve the JobExecutionResult from the JobManager. Lost connection to the JobManager. -> [Help 1]
[ERROR]
[ERROR] To see the full stack trace of the errors, re-run Maven with the -e switch.
[ERROR] Re-run Maven using the -X switch to enable full debug logging.
[ERROR]
[ERROR] For more information about the errors and possible solutions, please read the following articles:
[ERROR] [Help 1] http://cwiki.apache.org/confluence/display/MAVEN/MojoExecutionException

How can I correct the error? Is there any timeout issue?

Thanks

Spark dependencies are provided by design, not a bug.

In https://github.com/eljefe6a/beamexample/blob/master/BeamTutorial/pom.xml#L93 it says that the dependencies need to be explicitly added because of a bug. That is wrong. They need to be explicitly added because they are scope provided in the runner pom.

Spark is usually deployed in production clusters and applications depend on it with scope provided or runtime.
One good reason is the fact that the Spark distribution used needs to match the backend: YARN, Mesos, Stand Alone.
Another reason is the fact that if using Spark with Hadoop, you have to use a Hadoop-specific build - see here.

HDFS example

Now that we've established that there is no bug in the SparkRunner over HDFS as described here.
I was wondering if you have plans to try this on GS ?

Running Exercise1~3 on Windows with local runner will fail with InvalidPathException exception

When I running Excercise1 to Excercise 3 with local runner in Intellij it will fail with below error stack:
Exception in thread "main" org.apache.beam.sdk.Pipeline$PipelineExecutionException: java.nio.file.InvalidPathException: Illegal char <> at index 23: output/user_score-temp-
at org.apache.beam.runners.direct.DirectRunner.run(DirectRunner.java:279)
at org.apache.beam.runners.direct.DirectRunner.run(DirectRunner.java:72)
at org.apache.beam.sdk.Pipeline.run(Pipeline.java:182)
at org.apache.beam.examples.tutorial.game.solution.Exercise1.main(Exercise1.java:140)
at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
at sun.reflect.NativeMethodAccessorImpl.invoke(Unknown Source)
at sun.reflect.DelegatingMethodAccessorImpl.invoke(Unknown Source)
at java.lang.reflect.Method.invoke(Unknown Source)
at com.intellij.rt.execution.application.AppMain.main(AppMain.java:147)
Caused by: java.nio.file.InvalidPathException: Illegal char <*> at index 23: output/user_score-temp-*
at sun.nio.fs.WindowsPathParser.normalize(Unknown Source)
at sun.nio.fs.WindowsPathParser.parse(Unknown Source)
at sun.nio.fs.WindowsPathParser.parse(Unknown Source)
at sun.nio.fs.WindowsPath.parse(Unknown Source)
at sun.nio.fs.WindowsFileSystem.getPath(Unknown Source)
at java.nio.file.Paths.get(Unknown Source)
at org.apache.beam.sdk.util.FileIOChannelFactory.specToFile(FileIOChannelFactory.java:64)
at org.apache.beam.sdk.util.FileIOChannelFactory.match(FileIOChannelFactory.java:72)
at org.apache.beam.sdk.io.FileBasedSink$FileBasedWriteOperation.removeTemporaryFiles(FileBasedSink.java:483)
at org.apache.beam.sdk.io.FileBasedSink$FileBasedWriteOperation.finalize(FileBasedSink.java:411)
at org.apache.beam.sdk.io.Write$Bound$2.processElement(Write.java:417)
Looks program generated a file name which contains "*" and which is not acceptable for windows file system.
And Excerise4, Excerise5 run smoothly with expected result.

Recommend Projects

  • React photo React

    A declarative, efficient, and flexible JavaScript library for building user interfaces.

  • Vue.js photo Vue.js

    ๐Ÿ–– Vue.js is a progressive, incrementally-adoptable JavaScript framework for building UI on the web.

  • Typescript photo Typescript

    TypeScript is a superset of JavaScript that compiles to clean JavaScript output.

  • TensorFlow photo TensorFlow

    An Open Source Machine Learning Framework for Everyone

  • Django photo Django

    The Web framework for perfectionists with deadlines.

  • D3 photo D3

    Bring data to life with SVG, Canvas and HTML. ๐Ÿ“Š๐Ÿ“ˆ๐ŸŽ‰

Recommend Topics

  • javascript

    JavaScript (JS) is a lightweight interpreted programming language with first-class functions.

  • web

    Some thing interesting about web. New door for the world.

  • server

    A server is a program made to process requests and deliver data to clients.

  • Machine learning

    Machine learning is a way of modeling and interpreting data that allows a piece of software to respond intelligently.

  • Game

    Some thing interesting about game, make everyone happy.

Recommend Org

  • Facebook photo Facebook

    We are working to build community through open source technology. NB: members must have two-factor auth.

  • Microsoft photo Microsoft

    Open source projects and samples from Microsoft.

  • Google photo Google

    Google โค๏ธ Open Source for everyone.

  • D3 photo D3

    Data-Driven Documents codes.