【发布时间】:2014-12-22 16:42:11
【问题描述】:
我知道我可以通过将 Cascading 作业打包到 JAR 中来提交,如 Cascading 用户指南中所述。如果我使用hadoop jar CLI 命令手动提交该作业,那么该作业将在我的集群上运行。
但是,在最初的 Hadoop 1 Cascading 版本中,可以通过在 Hadoop JobConf 上设置某些属性来向集群提交作业。设置fs.defaultFS 和mapred.job.tracker 会导致本地Hadoop 库自动尝试将作业提交到Hadoop1 JobTracker。但是,设置这些属性似乎不适用于较新的版本。使用 Cascading 版本 2.5.3(将 CDH5 列为支持的平台)提交到 CDH5 5.2.1 Hadoop 集群会在与服务器协商时导致 IPC 异常,如下所述。
我相信这个平台组合——Cascading 2.5.6、Hadoop 2、CDH 5、YARN 和用于提交的 MR1 API——是基于compatibility table 的受支持组合(请参阅“先前版本”标题下)。并且使用hadoop jar 提交作业在同一个集群上工作正常。提交主机和 ResourceManager 之间的 8031 端口是开放的。在服务器端的 ResourceManager 日志中发现具有相同消息的错误。
我正在使用cascading-hadoop2-mr1 库。
Exception in thread "main" cascading.flow.FlowException: unhandled exception
at cascading.flow.BaseFlow.complete(BaseFlow.java:894)
at WordCount.main(WordCount.java:91)
at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:57)
at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
at java.lang.reflect.Method.invoke(Method.java:606)
at com.intellij.rt.execution.application.AppMain.main(AppMain.java:134)
Caused by: org.apache.hadoop.ipc.RemoteException(org.apache.hadoop.ipc.RpcServerException): Unknown rpc kind in rpc headerRPC_WRITABLE
at org.apache.hadoop.ipc.Client.call(Client.java:1411)
at org.apache.hadoop.ipc.Client.call(Client.java:1364)
at org.apache.hadoop.ipc.WritableRpcEngine$Invoker.invoke(WritableRpcEngine.java:231)
at org.apache.hadoop.mapred.$Proxy11.getStagingAreaDir(Unknown Source)
at org.apache.hadoop.mapred.JobClient.getStagingAreaDir(JobClient.java:1368)
at org.apache.hadoop.mapreduce.JobSubmissionFiles.getStagingDir(JobSubmissionFiles.java:102)
at org.apache.hadoop.mapred.JobClient$2.run(JobClient.java:982)
at org.apache.hadoop.mapred.JobClient$2.run(JobClient.java:976)
at java.security.AccessController.doPrivileged(Native Method)
at javax.security.auth.Subject.doAs(Subject.java:415)
at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1614)
at org.apache.hadoop.mapred.JobClient.submitJobInternal(JobClient.java:976)
at org.apache.hadoop.mapred.JobClient.submitJob(JobClient.java:950)
at cascading.flow.hadoop.planner.HadoopFlowStepJob.internalNonBlockingStart(HadoopFlowStepJob.java:105)
at cascading.flow.planner.FlowStepJob.blockOnJob(FlowStepJob.java:196)
at cascading.flow.planner.FlowStepJob.start(FlowStepJob.java:149)
at cascading.flow.planner.FlowStepJob.call(FlowStepJob.java:124)
at cascading.flow.planner.FlowStepJob.call(FlowStepJob.java:43)
at java.util.concurrent.FutureTask.run(FutureTask.java:262)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1145)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:615)
at java.lang.Thread.run(Thread.java:745)
演示代码如下,与 Cascading 用户指南中的 WordCount 示例基本相同。
public class WordCount {
public static void main(String[] args) {
String inputPath = "/user/vagrant/wordcount/input";
String outputPath = "/user/vagrant/wordcount/output";
Scheme sourceScheme = new TextLine( new Fields( "line" ) );
Tap source = new Hfs( sourceScheme, inputPath );
Scheme sinkScheme = new TextDelimited( new Fields( "word", "count" ) );
Tap sink = new Hfs( sinkScheme, outputPath, SinkMode.REPLACE );
Pipe assembly = new Pipe( "wordcount" );
String regex = "(?<!\\pL)(?=\\pL)[^ ]*(?<=\\pL)(?!\\pL)";
Function function = new RegexGenerator( new Fields( "word" ), regex );
assembly = new Each( assembly, new Fields( "line" ), function );
assembly = new GroupBy( assembly, new Fields( "word" ) );
Aggregator count = new Count( new Fields( "count" ) );
assembly = new Every( assembly, count );
Properties properties = AppProps.appProps()
.setName( "word-count-application" )
.setJarClass( WordCount.class )
.buildProperties();
properties.put("fs.defaultFS", "hdfs://192.168.30.101");
properties.put("mapred.job.tracker", "192.168.30.101:8032");
FlowConnector flowConnector = new HadoopFlowConnector( properties );
Flow flow = flowConnector.connect( "word-count", source, sink, assembly );
flow.complete();
}
}
我还尝试设置一些其他属性以使其正常工作:
mapreduce.jobtracker.addressmapreduce.framework.nameyarn.resourcemanager.addressyarn.resourcemanager.hostyarn.resourcemanager.hostnameyarn.resourcemanager.resourcetracker.address
这些都不起作用,它们只会导致作业在本地模式下运行(除非还设置了mapred.job.tracker)。
【问题讨论】:
标签: java hadoop mapreduce hadoop-yarn cascading