前提条件

  1. テストプログラムの Jar パッケージを準備します。 パッケージの名前を "mapreduce-examples.jar" 、ローカルストレージパスを "data\resources" とします。
  2. MultiJobs をテストするためのテーブルとリソースを準備します。
    • テーブルを作成します。
      create table mr_empty (key string, value string);
      create table mr_multijobs_out (value bigint);
    • リソースを追加します。
      add table mr_multijobs_out as multijobs_res_table -f;
      Add jar data \ resources \ mapreduce-examples.jar-f;

手順

odpscmd で MultiJobs を実行します。
jar -resources mapreduce-examples.jar,multijobs_res_table -classpath data\resources\mapreduce-examples.jar 
 com.aliyun.odps.mapred.open.example.MultiJobs mr_multijobs_out;

予想される出力

出力テーブル 'mr_multijobs_out' は次のとおりです。
+------------+
| value |
+------------+
| 0 |
+------------+

サンプルコード

    package com.aliyun.odps.mapred.open.example;
    import java.io.IOException;
    import java.util.Iterator;
    import com.aliyun.odps.data.Record;
    import com.aliyun.odps.data.TableInfo;
    import com.aliyun.odps.mapred.JobClient;
    import com.aliyun.odps.mapred.MapperBase;
    import com.aliyun.odps.mapred.RunningJob;
    import com.aliyun.odps.mapred.TaskContext;
    import com.aliyun.odps.mapred.conf.JobConf;
    import com.aliyun.odps.mapred.utils.InputUtils;
    import com.aliyun.odps.mapred.utils.OutputUtils;
    import com.aliyun.odps.mapred.utils.SchemaUtils;
    /**
     * MultiJobs
     *
     * 複数のジョブを実行
     *
     **/
    public class MultiJobs {
      public static class InitMapper extends MapperBase {
        @Override
        public void setup(TaskContext context) throws IOException{
          Record record = context.createOutputRecord();
          long v = context.getJobConf().getLong("multijobs.value", 2);
          record.set(0, v);
          context.write(record);
        }
      }
      public static class DecreaseMapper extends MapperBase {
        @Override
        public void cleanup(TaskContext context) throws IOException {
          // main 関数で定義された可変値を JobConf から取得
          long expect = context.getJobConf().getLong("multijobs.expect.value", -1);
          long v = -1;
          int count = 0;
          // リソーステーブルのデータ (直前のジョブの出力テーブル) の読み取り
          Iterator<Record> iter = context.readResourceTable("multijobs_res_table");
          while (iter.hasNext()) {
            Record r = iter.next();
            v = (Long) r.get(0);
            if (expect ! = v) {
              throw new IOException("expect: " + expect + ", but: " + v);
            }
            count++;
          }
          if (count ! = 1) {
            throw new IOException("res_table should have 1 record, but: " + count);
          }
          Record record = context.createOutputRecord();
          v--;
          record.set(0, v);
          context.write(record);
          // カウンターの設定。ジョブが正常終了した後、main 関数で取得できます。
          context.getCounter("multijobs", "value").setValue(v);
        }
      }
      public static void main(String[] args) throws Exception {
        if (args.length ! = 1) {
          System.err.println("Usage: TestMultiJobs <table>");
          System.exit(1);
        }
        String tbl = args[0];
        long iterCount = 2;
        System.err.println("Start to run init job.") ;
        JobConf initJob = new JobConf();
        initJob.setLong("multijobs.value", iterCount);
        initJob.setMapperClass(InitMapper.class);
        InputUtils.addTable(TableInfo.builder().tableName("mr_empty").build(), initJob);
        OutputUtils.addTable(TableInfo.builder().tableName(tbl).build(), initJob);
        initJob.setMapOutputKeySchema(SchemaUtils.fromString("key:string"));
        initJob.setMapOutputValueSchema(SchemaUtils.fromString("value:string"));
        // Maponly ジョブは、reducer 数を明示的に 0 に設定する必要があります
        initJob.setNumReduceTasks(0);
        JobClient.runJob(initJob);
        while (true) {
          System.err.println("Start to run iter job, count: " + iterCount);
          JobConf decJob = new JobConf();
          decJob.setLong("multijobs.expect.value", iterCount);
          decJob.setMapperClass(DecreaseMapper.class);
          InputUtils.addTable(TableInfo.builder().tableName("mr_empty").build(), decJob);
          OutputUtils.addTable(TableInfo.builder().tableName(tbl).build(), decJob);
          // Maponly job needs to explicitly set reducer number to 0
          decJob.setNumReduceTasks(0);
          RunningJob rJob = JobClient.runJob(decJob);
          iterCount--;
          // Exit the loop if the number of iterations has been reached
          if (rJob.getCounters().findCounter("multijobs", "value").getValue() == 0) {
            break;
          }
        }
        if (iterCount ! = 0) {
          throw new IOException("Job failed.") ;
        }
      }
    }