Re: "Unable to find registrar for hdfs" on Flink cluster

2017-08-30 Thread P. Ramanjaneya Reddy
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 am.examples.WordCount
-Dexec.args="--runner=FlinkRunner --flinkMaster=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
----.
Connected to JobManager at Actor[akka.tcp://flink@master:
6123/user/jobmanager#-1763674796] with leader session id
----.
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 

"Unable to find registrar for hdfs" on Flink cluster

2017-08-29 Thread P. Ramanjaneya Reddy
Hi All,

build jar file from the beam quickstart. while run the jar on Flinkcluster
got below error.?

anybody got this error?
Could you please help how to resolve this?

root1@master:~/NAI/Tools/flink-1.3.0$ *bin/flink run -c
org.apache.beam.examples.WordCount
/home/root1/NAI/Tools/word-count-beam/target/word-count-beam-bundled-0.1.jar
--runner=FlinkRunner
--filesToStage=/home/root1/NAI/Tools/word-count-beam/target/word-count-beam-bundled-0.1.jar
--inputFile=hdfs://master:9000/test/wordcount_input.txt
 --output=hdfs://master:9000/test/wordcount_output919*


This is the output I get:

Caused by: java.lang.IllegalStateException: Unable to find registrar for
hdfs
at
org.apache.beam.sdk.io.FileSystems.getFileSystemInternal(FileSystems.java:447)
at org.apache.beam.sdk.io.FileSystems.matchNewResource(FileSystems.java:517)
at
org.apache.beam.sdk.io.FileBasedSink.convertToFileResourceIfPossible(FileBasedSink.java:204)
at org.apache.beam.sdk.io.TextIO$Write.to(TextIO.java:296)
at org.apache.beam.examples.WordCount.main(WordCount.java:182)
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:528)
... 13 more


Thanks & Regards,
Ramanji.