【发布时间】: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