Package org.apache.flink.api.common

Examples of org.apache.flink.api.common.Plan


      .input2(source2)
      .name("Sort Join")
      .build();
    GenericDataSink sink = new GenericDataSink(new DiscardingOutputFormat(), matcher, "Sink");
   
    Plan p = new Plan(sink);
    p.setDefaultParallelism(64);
   
    runAndCancelJob(p, 3000, 20*1000);
  }
 
View Full Code Here


    CsvOutputFormat.configureRecordFormat(finalResult)
      .recordDelimiter('\n')
      .fieldDelimiter(' ')
      .field(StringValue.class, 0);

    Plan plan = new Plan(finalResult, "Iteration with AllReducer (keyless Reducer)");
    plan.setDefaultParallelism(numSubTasks);
    Assert.assertTrue(plan.getDefaultParallelism() > 1);
    return plan;
  }
View Full Code Here

   
    CsvOutputFormat.configureRecordFormat(out).recordDelimiter('\n')
        .fieldDelimiter(' ').field(StringValue.class, 0)
        .field(IntValue.class, 1);

    Plan plan = new Plan(out, "WordCount Example");
    plan.setDefaultParallelism(numSubTasks);
    return plan;
  }
View Full Code Here

    if (args.length < 3) {
      System.err.println(wc.getDescription());
      System.exit(1);
    }

    Plan plan = wc.getPlan(args);

    JobExecutionResult result = LocalExecutor.execute(plan);

    // Accumulators can be accessed by their name.
    System.out.println("Number of lines counter: "+ result.getAccumulatorResult(TokenizeLine.ACCUM_NUM_LINES));
View Full Code Here

      .recordDelimiter('\n')
      .fieldDelimiter(',')
      .field(IntValue.class, 0)
      .field(IntValue.class, 1);
   
    Plan p = new Plan(sink);
    p.setDefaultParallelism(dop);
    return p;
  }
View Full Code Here

    compareResultsByLinesInMemory(EXPECTED, resultPath);
  }

  @Override
  protected Plan getTestJob() {
    Plan plan = getTestPlanPlan(DOP, dataPath, resultPath);
    return plan;
  }
View Full Code Here

    CsvOutputFormat.configureRecordFormat(finalResult)
      .recordDelimiter('\n')
      .fieldDelimiter(' ')
      .field(StringValue.class, 0);

    Plan plan = new Plan(finalResult, "Iteration with AllReducer (keyless Reducer)");
   
    plan.setDefaultParallelism(numSubTasks);
    Assert.assertTrue(plan.getDefaultParallelism() > 1);
   
    return plan;
  }
View Full Code Here

    output.setDegreeOfParallelism(1);

    output.setInput(testMapper);
    testMapper.setInput(input);

    return new Plan(output);
  }
View Full Code Here

      .name("Identity Mapper")
      .build();
    GenericDataSink sink = new GenericDataSink(new DiscardingOutputFormat(), mapper, "Sink");
   
   
    Plan p = new Plan(sink);
    p.setDefaultParallelism(DOP);
   
    runAndCancelJob(p, 5 * 1000, 10 * 1000);
  }
 
View Full Code Here

    List<FileDataSink> sinks = new ArrayList<FileDataSink>();
    sinks.add(finalClusters);
    sinks.add(clusterAssignments);
   
    // return the PACT plan
    Plan plan = new Plan(sinks, "Iterative KMeans");
    plan.setDefaultParallelism(numSubTasks);
    return plan;
  }
View Full Code Here

TOP

Related Classes of org.apache.flink.api.common.Plan

Copyright © 2018 www.massapicom. All rights reserved.
All source code are property of their respective owners. Java is a trademark of Sun Microsystems, Inc and owned by ORACLE Inc. Contact coftware#gmail.com.