Package org.apache.flink.configuration

Examples of org.apache.flink.configuration.Configuration


  public static Collection<Object[]> getConfigurations() throws FileNotFoundException, IOException {

    LinkedList<Configuration> tConfigs = new LinkedList<Configuration>();

    for(int i=1; i <= NUM_PROGRAMS; i++) {
      Configuration config = new Configuration();
      config.setInteger("ProgramId", i);
      tConfigs.add(config);
    }
   
    return toParameterList(tConfigs);
  }
View Full Code Here


   */
  public static <T, F extends FileInputFormat<T>> F openInput(
      Class<F> inputFormatClass, String path, Configuration configuration)
    throws IOException
  {
    configuration = configuration == null ? new Configuration() : configuration;

    Path normalizedPath = normalizePath(new Path(path));
    final F inputFormat = ReflectionUtil.newInstance(inputFormatClass);

    inputFormat.setFilePath(normalizedPath);
View Full Code Here

   * @throws IOException
   *         if an I/O error occurred while accessing the file or initializing the InputFormat.
   */
  public static <T, IS extends InputSplit, F extends InputFormat<T, IS>> F openInput(
      Class<F> inputFormatClass, Configuration configuration) throws IOException {
    configuration = configuration == null ? new Configuration() : configuration;

    final F inputFormat = ReflectionUtil.newInstance(inputFormatClass);
    inputFormat.configure(configuration);
    final IS[] splits = inputFormat.createInputSplits(1);
    inputFormat.open(splits[0]);
View Full Code Here

  {
    final F outputFormat = ReflectionUtil.newInstance(outputFormatClass);
    outputFormat.setOutputFilePath(new Path(path));
    outputFormat.setWriteMode(WriteMode.OVERWRITE);
 
    configuration = configuration == null ? new Configuration() : configuration;
   
    outputFormat.configure(configuration);
    outputFormat.open(0, 1);
    return outputFormat;
  }
View Full Code Here

    {
      TaskConfig taskConfig = new TaskConfig(pointsInput.getConfiguration());
      taskConfig.addOutputShipStrategy(ShipStrategyType.FORWARD);
      taskConfig.setOutputSerializer(serializer);
     
      TaskConfig chainedMapper = new TaskConfig(new Configuration());
      chainedMapper.setDriverStrategy(DriverStrategy.COLLECTOR_MAP);
      chainedMapper.setStubWrapper(new UserCodeObjectWrapper<PointBuilder>(new PointBuilder()));
      chainedMapper.addOutputShipStrategy(ShipStrategyType.FORWARD);
      chainedMapper.setOutputSerializer(serializer);
     
View Full Code Here

    {
      TaskConfig taskConfig = new TaskConfig(modelsInput.getConfiguration());
      taskConfig.addOutputShipStrategy(ShipStrategyType.FORWARD);
      taskConfig.setOutputSerializer(serializer);

      TaskConfig chainedMapper = new TaskConfig(new Configuration());
      chainedMapper.setDriverStrategy(DriverStrategy.COLLECTOR_MAP);
      chainedMapper.setStubWrapper(new UserCodeObjectWrapper<PointBuilder>(new PointBuilder()));
      chainedMapper.addOutputShipStrategy(ShipStrategyType.FORWARD);
      chainedMapper.setOutputSerializer(serializer);
     
View Full Code Here

    // output into iteration. broadcasting the centers
    headConfig.setOutputSerializer(serializer);
    headConfig.addOutputShipStrategy(ShipStrategyType.BROADCAST);
   
    // final output
    TaskConfig headFinalOutConfig = new TaskConfig(new Configuration());
    headFinalOutConfig.setOutputSerializer(serializer);
    headFinalOutConfig.addOutputShipStrategy(ShipStrategyType.FORWARD);
    headConfig.setIterationHeadFinalOutputConfig(headFinalOutConfig);
   
    // the sync
View Full Code Here

  public static Collection<Object[]> getConfigurations() throws FileNotFoundException, IOException {

    LinkedList<Configuration> tConfigs = new LinkedList<Configuration>();

    for(int i=1; i <= NUM_PROGRAMS; i++) {
      Configuration config = new Configuration();
      config.setInteger("ProgramId", i);
      tConfigs.add(config);
    }
   
    return toParameterList(tConfigs);
  }
View Full Code Here

          final ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();

          DataSet<Tuple3<Integer, Long, String>> ds = CollectionDataSets.get3TupleDataSet(env)
              .map(new IdentityMapper<Tuple3<Integer, Long, String>>()).setParallelism(4);

          Configuration cfg = new Configuration();
          cfg.setString(PactCompiler.HINT_SHIP_STRATEGY, PactCompiler.HINT_SHIP_STRATEGY_REPARTITION);
          DataSet<Tuple2<Integer, String>> reduceDs = ds.reduceGroup(new Tuple3AllGroupReduceWithCombine())
              .withParameters(cfg);

          reduceDs.writeAsCsv(resultPath);
          env.execute();
View Full Code Here

    }
    else if (operation instanceof CrossOperatorBase.CrossWithLarge) {
      allowBCsecond = false;
    }
   
    Configuration conf = operation.getParameters();
    String localStrategy = conf.getString(PactCompiler.HINT_LOCAL_STRATEGY, null);
 
    if (localStrategy != null) {
      final OperatorDescriptorDual fixedDriverStrat;
      if (PactCompiler.HINT_LOCAL_STRATEGY_NESTEDLOOP_BLOCKED_OUTER_FIRST.equals(localStrategy)) {
        fixedDriverStrat = new CrossBlockOuterFirstDescriptor(allowBCfirst, allowBCsecond);
View Full Code Here

TOP

Related Classes of org.apache.flink.configuration.Configuration

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.