【问题标题】:How can I submit a Cascading job to a remote YARN cluster from Java?如何从 Java 向远程 YARN 集群提交级联作业?
【发布时间】:2014-12-22 16:42:11
【问题描述】:

我知道我可以通过将 Cascading 作业打包到 JAR 中来提交,如 Cascading 用户指南中所述。如果我使用hadoop jar CLI 命令手动提交该作业,那么该作业将在我的集群上运行。

但是,在最初的 Hadoop 1 Cascading 版本中,可以通过在 Hadoop JobConf 上设置某些属性来向集群提交作业。设置fs.defaultFSmapred.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.address
  • mapreduce.framework.name
  • yarn.resourcemanager.address
  • yarn.resourcemanager.host
  • yarn.resourcemanager.hostname
  • yarn.resourcemanager.resourcetracker.address

这些都不起作用,它们只会导致作业在本地模式下运行(除非还设置了mapred.job.tracker)。

【问题讨论】:

    标签: java hadoop mapreduce hadoop-yarn cascading


    【解决方案1】:

    我现在已经解决了这个问题。它来自于尝试使用 Cloudera 分发的旧 Hadoop 类,尤其是 JobClient。如果您将hadoop-core 与提供的2.5.0-mr1-cdh5.2.1 版本一起使用,或者将hadoop-client 依赖项与相同的版本号一起使用,则会发生这种情况。虽然这个号称是MR1版本,而且我们使用的是MR1 API提交,但是这个版本实际上只支持提交到Hadoop1 JobTracker,不支持YARN。

    为了允许提交到 YARN,您必须将 hadoop-client 依赖项与非 MR1 2.5.0-cdh5.2.1 版本一起使用,该版本仍然支持将 MR1 作业提交到 YARN。

    【讨论】:

      猜你喜欢
      • 2016-12-20
      • 2018-08-22
      • 2015-08-20
      • 2015-09-29
      • 1970-01-01
      • 2019-05-30
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多