【问题标题】:Call to get EMR results always returns empty state collection even when EMR job done即使 EMR 作业完成,调用以获取 EMR 结果也始终返回空状态集合
【发布时间】:2014-01-08 22:29:03
【问题描述】:

由于 IP 原因,我无法发布完整的源代码。但是,我打电话提交了一个 Amazon Elastic Map Reduce 作业 (EMR),该作业现在运行完成。以前它失败了,基本上是一个找不到文件的错误。

RunJobFlowResult result=emr.runJobFlow(request);

成功,我可以从中获取作业流 ID。

稍后,我有一个循环轮询状态,首先 DescribeJobFlowsRequest 请求=新的 DescribeJobFlowsRequest(jobFlowIdArray); 我通过调用循环检查每个状态 request.getJobFlowStates()

不幸的是,无论作业是在运行、失败还是成功,该调用总是返回一个空集合。我怎样才能至少了解正在发生的事情?

AWSCredentials credentials = new BasicAWSCredentials(accessKey, secretKey);
AmazonElasticMapReduceClient client = new AmazonElasticMapReduceClient(credentials);
client.setEndPoint("elasticmapreduce.us-east-1.amazonaws.com");

StepFactory stepFactory = new StepFactory();
StepConfig enableDebugging = new StepConfig()
.withActionOnFailure("TERMINATE_JOB_FLOW")
.withHadoopjJarStep(stepFactory.newEnableDebuggingStep());


String[] arguments={...} // Custom jar arguments

HadoopJarStepConfig jarConfig = new HadoopJarStepConfig();
jarConfig.setJar(JAR_NAME);
jarConfig.setArgs(Arrays.asList(arguments));

StepConfig runJar = new     StepConfig(JAR_NAME.substring(JAR_NAME.indexOf('/')+1),jarConfig);

RunJobFlowRequest request = new RunJobFlowRequest()
.withName("...")
.withSteps(runJar)
.withLogUri("...")
.withInstances(
     new JobFlowInstancesCOnfig()
       .withHadoopVersion("1.0.3")
        .withInstanceCount(5)
         .withKeepJobFlowAliveWhenNoSteps(false)
         .withMasterInstanceType("m1.small")
          .withSlaveInstanceType("m1.small");

RunJobFlowResult result = client.runJobFlow(request);
String jobFlowID=result.getJobFlowID();
List<String> describeJobFlowIdList=new ArrayList<String>(1);
describeJobFlowIdList.add(jobFlowID);

String lastState="";
boolean jobMonitoringNotDone=true;
while(jobMonitoringNotDone){
    SescribeJobFlowsRequest describeJobFlowsRequest=
        new DescribeJobFlowsRequest(describeJobFlowIdList);
    // Call to describeJobFlowsRequest.getJobFlowStates() always returns
    // empty list even when job succeeds or fails.
    for(String state : describeJobFlowsRequest.getJobFlowStates()){
        if(DONE_STATES.contains(state)){
            jobMonitoringNotDone=false;
        } else if(!lastState.equals(state)){
            lastState = state;
            System.out.println("Job "+state + " at "+ new Date().toString());
        }
    }
    try {
    Thread.sleep(10000);
    } catch (InterruptedException e) {
        e.printStackTrace();
    }
}

【问题讨论】:

  • 你不能创建一个最小功能的例子来重现你的问题吗?
  • 我对其进行了更新,以便对我正在尝试做的事情提供一个极简主义的看法。

标签: java amazon-web-services mapreduce


【解决方案1】:

上面的代码缺少一个类似于

的调用
DescribeJobFlowsResult describeJobFlowsResult =      client.describeJobFlows(describeJobFlowsRequest);

这给了我一个可行的解决方案,但不幸的是,亚马逊弃用了该方法,但没有提供替代方案。我希望我有一个不推荐使用的解决方案,所以这只是部分答案。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2020-11-15
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-10-17
    • 1970-01-01
    • 2017-01-27
    相关资源
    最近更新 更多