Thank you Aljoscha. With above steps working wordcount beam using quick start program.
When running on actual beam source tree getting following error. root1@master:~/Projects/*beam*/examples/java$ *git branch * master * release-2.0.0 * ==> beam source code* root1@master:~/Projects/beam/examples/java$ root1@master:~/Projects/beam/examples/java$* mvn dependency:tree -Pflink-runner |grep flink* [INFO] \- org.apache.beam:beam-runners-flink_2.10:jar:2.2.0-SNAPSHOT:runtime [INFO] +- org.apache.flink:flink-clients_2.10:jar:1.3.0:runtime [INFO] | +- org.apache.flink:flink-optimizer_2.10:jar:1.3.0:runtime [INFO] | \- org.apache.flink:force-shading:jar:1.3.0:runtime [INFO] +- org.apache.flink:flink-core:jar:1.3.0:runtime [INFO] | +- org.apache.flink:flink-annotations:jar:1.3.0:runtime [INFO] +- org.apache.flink:flink-metrics-core:jar:1.3.0:runtime [INFO] +- org.apache.flink:flink-java:jar:1.3.0:runtime [INFO] | +- org.apache.flink:flink-shaded-hadoop2:jar:1.3.0:runtime [INFO] +- org.apache.flink:flink-runtime_2.10:jar:1.3.0:runtime [INFO] +- org.apache.flink:flink-streaming-java_2.10:jar:1.3.0:runtime root1@master:~/Projects/beam/examples/java$ root1@master:~/Projects/*beam*/examples/java$ *mvn package exec:java -Dexec.mainClass=org.apache.be <http://org.apache.be>am.examples.WordCount -Dexec.args="--runner=FlinkRunner --flinkMaster=192.168.56.1:6123 <http://192.168.56.1:6123/> --filesToStage=/home/root1/Projects/beam/examples/java/target/beam-examples-java-2.0.0.jar --inputFile=hdfs://master:9000/test/wordcount_input.txt --output=hdfs://master:9000/test/wordcount_output919" -Pflink-runner -Dcheckstyle.skip=true -DskipTests* *Error Log:* INFO: Received job wordcount-root1-0830134254-67bc7d88 (02066e0dc345cdd6f34f20258a4c807e). Aug 30, 2017 7:12:56 PM org.apache.flink.runtime.client.JobClientActor disconnectFromJobManager INFO: Disconnect from JobManager null. Aug 30, 2017 7:12:56 PM org.apache.flink.runtime.client.JobClientActor connectToJobManager INFO: Connect to JobManager Actor[akka.tcp://flink@master: 6123/user/jobmanager#-1763674796]. Aug 30, 2017 7:12:56 PM org.apache.flink.runtime.client.JobClientActor logAndPrintMessage INFO: Connected to JobManager at Actor[akka.tcp://flink@master: 6123/user/jobmanager#-1763674796] with leader session id 00000000-0000-0000-0000-000000000000. Connected to JobManager at Actor[akka.tcp://flink@master: 6123/user/jobmanager#-1763674796] with leader session id 00000000-0000-0000-0000-000000000000. Aug 30, 2017 7:12:56 PM org.apache.flink.runtime.client.JobSubmissionClientActor tryToSubmitJob INFO: Sending message to JobManager akka.tcp://flink@master:6123/user/jobmanager to submit job wordcount-root1-0830134254-67bc7d88 (02066e0dc345cdd6f34f20258a4c807e) and wait for progress Aug 30, 2017 7:12:56 PM org.apache.flink.runtime.client.JobSubmissionClientActor$1 call INFO: Upload jar files to job manager akka.tcp://flink@master:6123/u ser/jobmanager. Aug 30, 2017 7:12:56 PM org.apache.flink.runtime.blob.BlobClient uploadJarFiles INFO: Blob client connecting to akka.tcp://flink@master:6123/user/jobmanager Aug 30, 2017 7:12:56 PM org.apache.flink.runtime.client.JobSubmissionClientActor$1 call INFO: Submit job to the job manager akka.tcp://flink@master:6123/u ser/jobmanager. Aug 30, 2017 7:12:56 PM org.apache.flink.runtime.client.JobClientActor terminate INFO: Terminate JobClientActor. Aug 30, 2017 7:12:56 PM org.apache.flink.runtime.client.JobClientActor disconnectFromJobManager INFO: Disconnect from JobManager Actor[akka.tcp://flink@master: 6123/user/jobmanager#-1763674796]. Aug 30, 2017 7:12:56 PM org.apache.flink.runtime.client.JobClient awaitJobResult INFO: Job execution failed Aug 30, 2017 7:12:56 PM akka.event.slf4j.Slf4jLogger$$ anonfun$receive$1$$anonfun$applyOrElse$3 apply$mcV$sp INFO: Shutting down remote daemon. Aug 30, 2017 7:12:56 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 30, 2017 7:12:56 PM akka.event.slf4j.Slf4jLogger$$ anonfun$receive$1$$anonfun$applyOrElse$3 apply$mcV$sp INFO: Remoting shut down. Aug 30, 2017 7:12:56 PM org.apache.beam.runners.flink.FlinkRunner run SEVERE: Pipeline execution failed org.apache.flink.client.program.ProgramInvocationException: The program execution failed: Cannot initialize task 'DataSource (at Read(CreateSource) (org.apache.beam.runners.flink.translation.wrappers.SourceInputFormat))': Deserializing the InputFormat (org.apache.beam.runners.flink .translation.wrappers.SourceInputFormat@7ef64f) failed: Could not read the user code wrapper: org.apache.beam.runners.flink. translation.wrappers.SourceInputFormat at org.apache.flink.client.program.ClusterClient.run(ClusterClient.java:478) at org.apache.flink.client.program.StandaloneClusterClient.subm itJob(StandaloneClusterClient.java:105) at org.apache.flink.client.program.ClusterClient.run(ClusterClient.java:442) at org.apache.flink.client.program.ClusterClient.run(ClusterClient.java:429) at org.apache.flink.client.program.ClusterClient.run(ClusterClient.java:404) at org.apache.flink.client.RemoteExecutor.executePlanWithJars(R emoteExecutor.java:211) at org.apache.flink.client.RemoteExecutor.executePlan(RemoteExe cutor.java:188) at org.apache.flink.api.java.RemoteEnvironment.execute(RemoteEn vironment.java:172) at org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironm ent.executePipeline(FlinkPipelineExecutionEnvironment.java:114) at org.apache.beam.runners.flink.FlinkRunner.run(FlinkRunner.java:118) at org.apache.beam.sdk.Pipeline.run(Pipeline.java:303) at org.apache.beam.sdk.Pipeline.run(Pipeline.java:289) at org.apache.beam.examples.WordCount.main(WordCount.java:188) at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAcce ssorImpl.java:62) at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMe thodAccessorImpl.java:43) at java.lang.reflect.Method.invoke(Method.java:498) at org.codehaus.mojo.exec.ExecJavaMojo$1.run(ExecJavaMojo.java:294) at java.lang.Thread.run(Thread.java:748) Caused by: org.apache.flink.runtime.client.JobExecutionException: Cannot initialize task 'DataSource (at Read(CreateSource) (org.apache.beam.runners.flink.translation.wrappers.SourceInputFormat))': Deserializing the InputFormat (org.apache.beam.runners.flink .translation.wrappers.SourceInputFormat@7ef64f) failed: Could not read the user code wrapper: org.apache.beam.runners.flink. translation.wrappers.SourceInputFormat at org.apache.flink.runtime.executiongraph.ExecutionGraphBuilde r.buildGraph(ExecutionGraphBuilder.java:153) at org.apache.flink.runtime.jobmanager.JobManager.org <http://org.apache.flink.runtime.jobmanager.jobmanager.org/>$apache $flink$runtime$jobmanager$JobManager$$submitJob(JobManager.scala:1315) at org.apache.flink.runtime.jobmanager.JobManager$$anonfun$ handleMessage$1.applyOrElse(JobManager.scala:495) at scala.runtime.AbstractPartialFunction$mcVL$sp.apply$mcVL$sp( AbstractPartialFunction.scala:33) at scala.runtime.AbstractPartialFunction$mcVL$sp.apply(Abstract PartialFunction.scala:33) at scala.runtime.AbstractPartialFunction$mcVL$sp.apply(Abstract PartialFunction.scala:25) at org.apache.flink.runtime.LeaderSessionMessageFilter$$anonfun $receive$1.applyOrElse(LeaderSessionMessageFilter.scala:38) at scala.runtime.AbstractPartialFunction$mcVL$sp.apply$mcVL$sp( AbstractPartialFunction.scala:33) at scala.runtime.AbstractPartialFunction$mcVL$sp.apply(Abstract PartialFunction.scala:33) at scala.runtime.AbstractPartialFunction$mcVL$sp.apply(Abstract PartialFunction.scala:25) at org.apache.flink.runtime.LogMessages$$anon$1.apply(LogMessages.scala:33) at org.apache.flink.runtime.LogMessages$$anon$1.apply(LogMessages.scala:28) at scala.PartialFunction$class.applyOrElse(PartialFunction.scala:118) at org.apache.flink.runtime.LogMessages$$anon$1.applyOrElse( LogMessages.scala:28) at akka.actor.Actor$class.aroundReceive(Actor.scala:467) at org.apache.flink.runtime.jobmanager.JobManager.aroundReceive (JobManager.scala:125) 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(ForkJoinW orkerThread.java:107) Caused by: java.lang.Exception: Deserializing the InputFormat (org.apache.beam.runners.flink.translation.wrappers.SourceInputFormat@7ef64f) failed: Could not read the user code wrapper: org.apache.beam.runners.flink. translation.wrappers.SourceInputFormat at org.apache.flink.runtime.jobgraph.InputFormatVertex.initiali zeOnMaster(InputFormatVertex.java:66) at org.apache.flink.runtime.executiongraph.ExecutionGraphBuilde r.buildGraph(ExecutionGraphBuilder.java:150) ... 24 more Caused by: org.apache.flink.runtime.operators.util.CorruptConfigurationException: Could not read the user code wrapper: org.apache.beam.runners.flink. translation.wrappers.SourceInputFormat at org.apache.flink.runtime.operators.util.TaskConfig.getStubWr apper(TaskConfig.java:290) at org.apache.flink.runtime.jobgraph.InputFormatVertex.initiali zeOnMaster(InputFormatVertex.java:63) ... 25 more Caused by: java.lang.ClassNotFoundException: org.apache.beam.runners.flink. translation.wrappers.SourceInputFormat at java.net.URLClassLoader.findClass(URLClassLoader.java:381) at java.lang.ClassLoader.loadClass(ClassLoader.java:424) at java.lang.ClassLoader.loadClass(ClassLoader.java:357) at java.lang.Class.forName0(Native Method) at java.lang.Class.forName(Class.java:348) at org.apache.flink.util.InstantiationUtil$ClassLoaderObjectInp utStream.resolveClass(InstantiationUtil.java:64) at java.io.ObjectInputStream.readNonProxyDesc(ObjectInputStream.java:1826) at java.io.ObjectInputStream.readClassDesc(ObjectInputStream.java:1713) at java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2000) at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1535) at java.io.ObjectInputStream.defaultReadFields(ObjectInputStream.java:2245) at java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2169) at java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2027) at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1535) at java.io.ObjectInputStream.readObject(ObjectInputStream.java:422) at org.apache.flink.util.InstantiationUtil.deserializeObject(In stantiationUtil.java:290) at org.apache.flink.util.InstantiationUtil.readObjectFromConfig (InstantiationUtil.java:248) at org.apache.flink.runtime.operators.util.TaskConfig.getStubWr apper(TaskConfig.java:288) ... 26 more [WARNING] java.lang.reflect.InvocationTargetException at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAcce ssorImpl.java:62) at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMe thodAccessorImpl.java:43) at java.lang.reflect.Method.invoke(Method.java:498) at org.codehaus.mojo.exec.ExecJavaMojo$1.run(ExecJavaMojo.java:294) 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:121) at org.apache.beam.sdk.Pipeline.run(Pipeline.java:303) at org.apache.beam.sdk.Pipeline.run(Pipeline.java:289) at org.apache.beam.examples.WordCount.main(WordCount.java:188) ... 6 more Caused by: org.apache.flink.client.program.ProgramInvocationException: The program execution failed: Cannot initialize task 'DataSource (at Read(CreateSource) (org.apache.beam.runners.flink .translation.wrappers.SourceInputFormat))': Deserializing the InputFormat (org.apache.beam.runners.flink.translation.wrappers.SourceInputFormat@7ef64f) failed: Could not read the user code wrapper: org.apache.beam.runners.flink. translation.wrappers.SourceInputFormat at org.apache.flink.client.program.ClusterClient.run(ClusterClient.java:478) at org.apache.flink.client.program.StandaloneClusterClient.subm itJob(StandaloneClusterClient.java:105) at org.apache.flink.client.program.ClusterClient.run(ClusterClient.java:442) at org.apache.flink.client.program.ClusterClient.run(ClusterClient.java:429) at org.apache.flink.client.program.ClusterClient.run(ClusterClient.java:404) at org.apache.flink.client.RemoteExecutor.executePlanWithJars(R emoteExecutor.java:211) at org.apache.flink.client.RemoteExecutor.executePlan(RemoteExe cutor.java:188) at org.apache.flink.api.java.RemoteEnvironment.execute(RemoteEn vironment.java:172) at org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironm ent.executePipeline(FlinkPipelineExecutionEnvironment.java:114) at org.apache.beam.runners.flink.FlinkRunner.run(FlinkRunner.java:118) ... 9 more Caused by: org.apache.flink.runtime.client.JobExecutionException: Cannot initialize task 'DataSource (at Read(CreateSource) (org.apache.beam.runners.flink.translation.wrappers.SourceInputFormat))': Deserializing the InputFormat (org.apache.beam.runners.flink .translation.wrappers.SourceInputFormat@7ef64f) failed: Could not read the user code wrapper: org.apache.beam.runners.flink. translation.wrappers.SourceInputFormat at org.apache.flink.runtime.executiongraph.ExecutionGraphBuilde r.buildGraph(ExecutionGraphBuilder.java:153) at org.apache.flink.runtime.jobmanager.JobManager.org <http://org.apache.flink.runtime.jobmanager.jobmanager.org/>$apache $flink$runtime$jobmanager$JobManager$$submitJob(JobManager.scala:1315) at org.apache.flink.runtime.jobmanager.JobManager$$anonfun$ handleMessage$1.applyOrElse(JobManager.scala:495) at scala.runtime.AbstractPartialFunction$mcVL$sp.apply$mcVL$sp( AbstractPartialFunction.scala:33) at scala.runtime.AbstractPartialFunction$mcVL$sp.apply(Abstract PartialFunction.scala:33) at scala.runtime.AbstractPartialFunction$mcVL$sp.apply(Abstract PartialFunction.scala:25) at org.apache.flink.runtime.LeaderSessionMessageFilter$$anonfun $receive$1.applyOrElse(LeaderSessionMessageFilter.scala:38) at scala.runtime.AbstractPartialFunction$mcVL$sp.apply$mcVL$sp( AbstractPartialFunction.scala:33) at scala.runtime.AbstractPartialFunction$mcVL$sp.apply(Abstract PartialFunction.scala:33) at scala.runtime.AbstractPartialFunction$mcVL$sp.apply(Abstract PartialFunction.scala:25) at org.apache.flink.runtime.LogMessages$$anon$1.apply(LogMessages.scala:33) at org.apache.flink.runtime.LogMessages$$anon$1.apply(LogMessages.scala:28) at scala.PartialFunction$class.applyOrElse(PartialFunction.scala:118) at org.apache.flink.runtime.LogMessages$$anon$1.applyOrElse( LogMessages.scala:28) at akka.actor.Actor$class.aroundReceive(Actor.scala:467) at org.apache.flink.runtime.jobmanager.JobManager.aroundReceive (JobManager.scala:125) 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(ForkJoinW orkerThread.java:107) Caused by: java.lang.Exception: Deserializing the InputFormat (org.apache.beam.runners.flink.translation.wrappers.SourceInputFormat@7ef64f) failed: Could not read the user code wrapper: org.apache.beam.runners.flink. translation.wrappers.SourceInputFormat at org.apache.flink.runtime.jobgraph.InputFormatVertex.initiali zeOnMaster(InputFormatVertex.java:66) at org.apache.flink.runtime.executiongraph.ExecutionGraphBuilde r.buildGraph(ExecutionGraphBuilder.java:150) ... 24 more Caused by: org.apache.flink.runtime.operators.util.CorruptConfigurationException: Could not read the user code wrapper: org.apache.beam.runners.flink. translation.wrappers.SourceInputFormat at org.apache.flink.runtime.operators.util.TaskConfig.getStubWr apper(TaskConfig.java:290) at org.apache.flink.runtime.jobgraph.InputFormatVertex.initiali zeOnMaster(InputFormatVertex.java:63) ... 25 more Caused by: java.lang.ClassNotFoundException: org.apache.beam.runners.flink. translation.wrappers.SourceInputFormat at java.net.URLClassLoader.findClass(URLClassLoader.java:381) at java.lang.ClassLoader.loadClass(ClassLoader.java:424) at java.lang.ClassLoader.loadClass(ClassLoader.java:357) at java.lang.Class.forName0(Native Method) at java.lang.Class.forName(Class.java:348) at org.apache.flink.util.InstantiationUtil$ClassLoaderObjectInp utStream.resolveClass(InstantiationUtil.java:64) at java.io.ObjectInputStream.readNonProxyDesc(ObjectInputStream.java:1826) at java.io.ObjectInputStream.readClassDesc(ObjectInputStream.java:1713) at java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2000) at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1535) at java.io.ObjectInputStream.defaultReadFields(ObjectInputStream.java:2245) at java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2169) at java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2027) at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1535) at java.io.ObjectInputStream.readObject(ObjectInputStream.java:422) at org.apache.flink.util.InstantiationUtil.deserializeObject(In stantiationUtil.java:290) at org.apache.flink.util.InstantiationUtil.readObjectFromConfig (InstantiationUtil.java:248) at org.apache.flink.runtime.operators.util.TaskConfig.getStubWr apper(TaskConfig.java:288) ... 26 more [INFO] ------------------------------------------------------------ ------------ [INFO] BUILD FAILURE