AWS EMR Job Flow中各步骤是否自动接收前序步骤输出?
EMR Job Flow中MapReduce多步骤路径传递问题解答
核心结论
AWS EMR Job Flow不会自动将前序步骤的输出路径传递给后续步骤,也不会隐式向MapReduce主类注入路径参数。你必须手动为每个步骤显式指定输入和输出S3路径——因为你的MapReduce主类依赖命令行参数来获取这些路径,EMR没有内置的隐式传递机制。
具体实现方案
- 为每个步骤定义独立的S3输出路径,前序步骤的输出路径直接作为后续步骤的输入路径
- 在每个
HadoopJarStepConfig的args方法中,传入对应的输入路径、输出路径(以及其他业务所需参数)
修改后的代码示例
public static void main(String args[]) throws IOException, ClassNotFoundException, InterruptedException { EmrClient mapReduce = EmrClient.builder() .credentialsProvider(ProfileCredentialsProvider.create()) .build(); List<StepConfig> steps = new LinkedList<>(); // 定义各步骤的S3路径(根据实际业务调整前缀) String rawInputPath = "s3n://" + myBucketName + "/raw-input-data"; String nCountOutputPath = "s3n://" + myBucketName + "/step1-ncount-output"; String countNrTrOutputPath = "s3n://" + myBucketName + "/step2-countnrtr-output"; String joinCalculateOutputPath = "s3n://" + myBucketName + "/step3-join-calculate-output"; String finalSortOutputPath = "s3n://" + myBucketName + "/step4-final-sort-output"; // 步骤1: NCount HadoopJarStepConfig hadoopJarStepConfig = HadoopJarStepConfig.builder() .jar("s3n://" + myBucketName + "/" + NCount + jarPostfix) .mainClass(packageName + NCount) .args(rawInputPath, nCountOutputPath) // 传入原始输入路径、当前步骤输出路径 .build(); steps.add(StepConfig.builder() .name(NCount) .hadoopJarStep(hadoopJarStepConfig) .actionOnFailure("TERMINATE_JOB_FLOW") .build()); // 步骤2: CountNrTr HadoopJarStepConfig hadoopJarStepConfig2 = HadoopJarStepConfig.builder() .jar("s3n://" + myBucketName + "/" + CountNrTr + jarPostfix) .mainClass(packageName + CountNrTr) .args(nCountOutputPath, countNrTrOutputPath) // 前序步骤输出作为当前输入 .build(); steps.add(StepConfig.builder() .name(CountNrTr) .hadoopJarStep(hadoopJarStepConfig2) .actionOnFailure("TERMINATE_JOB_FLOW") .build()); // 步骤3: JoinAndCalculate HadoopJarStepConfig hadoopJarStepConfig3 = HadoopJarStepConfig.builder() .jar("s3n://" + myBucketName + "/" + JoinAndCalculate + jarPostfix) .mainClass(packageName + JoinAndCalculate) .args(countNrTrOutputPath, joinCalculateOutputPath) .build(); steps.add(StepConfig.builder() .name(JoinAndCalculate) .hadoopJarStep(hadoopJarStepConfig3) .actionOnFailure("TERMINATE_JOB_FLOW") .build()); // 步骤4: ValueToKeySort HadoopJarStepConfig hadoopJarStepConfig4 = HadoopJarStepConfig.builder() .jar("s3n://" + myBucketName + "/" + ValueToKeySort + jarPostfix) .mainClass(packageName + ValueToKeySort) .args(joinCalculateOutputPath, finalSortOutputPath) .build(); steps.add(StepConfig.builder() .name(ValueToKeySort) .hadoopJarStep(hadoopJarStepConfig4) .actionOnFailure("TERMINATE_JOB_FLOW") .build()); JobFlowInstancesConfig instances = JobFlowInstancesConfig.builder() .instanceCount(2) .masterInstanceType("m4.large") .slaveInstanceType("m4.large") .hadoopVersion("3.3.4") .ec2KeyName(myKeyPair) .keepJobFlowAliveWhenNoSteps(false) .placement(PlacementType.builder().availabilityZone("us-east-1a").build()) .build(); }
关键注意事项
- 权限校验:确保EMR集群的IAM角色拥有读写这些S3路径的权限,否则会出现访问拒绝错误
- 路径唯一性:每个步骤的输出路径必须唯一且未提前存在(Hadoop默认要求输出路径不存在,若需覆盖可在主类中添加路径删除逻辑)
内容的提问来源于stack exchange,提问作者אורי זיו
相关产品推荐
相关产品推荐

