【问题标题】:in aws emr job flow, does each step receive the output from the previous step?在 aws emr 作业流程中,每个步骤是否都收到上一步的输出?
【发布时间】:2023-01-08 23:25:28
【问题描述】:

我正在用 Java 制作一个 map reduce 程序,它有 4 个步骤。 每一步都对上一步的输出进行操作。

到目前为止,我在本地手动运行了这些步骤,我想开始使用作业流程在 AWS EMR 上运行。
我的教授给了我们一些代码来配置作业流程的步骤,但现在我们面临一个问题:

我的每一步都期望输入和输出路径作为其主要功能的参数。 JobFlow 允许我将 args 传输到每个步骤,但据我了解,作业流程中的每个步骤都应该自动接收上一步的输出

有没有人知道这是不是真的? 步骤中的 map-reduce 应用程序如何实现其输入位置?路径是否作为 JobFlow 的参数隐式传递给它?

我正在使用适用于 Java 的 AWS SDK 2。

我的代码:

 public static void main(String args[]) throws IOException, ClassNotFoundException, InterruptedException {
                // AwsCredentialsProvider credentialsProvider = StaticCredentialsProvider
                // .create(ProfileCredentialsProvider.create().resolveCredentials());

                EmrClient mapReduce = EmrClient.builder().credentialsProvider(ProfileCredentialsProvider.create())
                                .build();
                List<StepConfig> steps = new LinkedList<StepConfig>();

                HadoopJarStepConfig hadoopJarStepConfig = HadoopJarStepConfig.builder()
                                .jar("s3n://" + myBucketName + "/" + NCount + jarPostfix)
                                .mainClass(packageName + NCount)
                                .args(??????????????????????)
                                .build();
                steps.add(StepConfig.builder().name(NCount).hadoopJarStep(hadoopJarStepConfig)
                                .actionOnFailure("TERMINATE_JOB_FLOW").build());

                HadoopJarStepConfig hadoopJarStepConfig2 = HadoopJarStepConfig.builder()
                                .jar("s3n://" + myBucketName + "/" + CountNrTr + jarPostfix)
                                .mainClass(packageName + CountNrTr)
                                .args(??????????????????????)
                                .build();
                steps.add(StepConfig.builder().name(CountNrTr).hadoopJarStep(hadoopJarStepConfig2)
                                .actionOnFailure("TERMINATE_JOB_FLOW").build());

                HadoopJarStepConfig hadoopJarStepConfig3 = HadoopJarStepConfig.builder()
                                .jar("s3n://" + myBucketName + "/" + JoinAndCalculate + jarPostfix)
                                .mainClass(packageName + JoinAndCalculate)
                                .args(??????????????????????)
                                .build();
                steps.add(StepConfig.builder().name(JoinAndCalculate).hadoopJarStep(hadoopJarStepConfig3)
                                .actionOnFailure("TERMINATE_JOB_FLOW").build());

                HadoopJarStepConfig hadoopJarStepConfig4 = HadoopJarStepConfig.builder()
                                .jar("s3n://" + myBucketName + "/" + ValueToKeySort + jarPostfix)
                                .mainClass(packageName + ValueToKeySort)
                                .args(??????????????????????)
                                .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();

【问题讨论】:

    标签: java amazon-web-services hadoop mapreduce amazon-emr


    【解决方案1】:

    EMR 与问题无关。不,它不是自动的。

    我们需要查看您执行的 JAR 的代码,但我只假设它是您使用 FileInputFormat 的传统 mapreduce 代码,并且可能有类似 Path(args[0]) 的代码,如果是这样,那很可能是您的输入。然后 Path(args[1]) 可能是输出。

    因此,您只需在每个步骤中将这些参数链接在一起......

    step1 = ...
       .args(new String[] {"/in", "/stage1" })
    ...
    final = ...
       .args(new String[] {"/stageN", "/out" }) 
    

    或者,将您的代码转换为 Spark/Flink 或 Hive 查询,其中多个 mapreduce 阶段自动处理

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-11-15
      • 1970-01-01
      • 2023-02-07
      • 1970-01-01
      • 1970-01-01
      • 2018-10-10
      相关资源
      最近更新 更多