Package org.apache.flink.streaming.api.environment

Examples of org.apache.flink.streaming.api.environment.StreamExecutionEnvironment.addSource()


  public static void main(String[] args) throws Exception {

    StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironment(1);

    @SuppressWarnings("unused")
    DataStream<String> stream1 = env
      .addSource(new MyKafkaSource("localhost:2181", "group", "test", 1), SOURCE_PARALELISM)
      .addSink(new MyKafkaPrintSink());

    @SuppressWarnings("unused")
    DataStream<String> stream2 = env
View Full Code Here


    DataStream<String> stream1 = env
      .addSource(new MyKafkaSource("localhost:2181", "group", "test", 1), SOURCE_PARALELISM)
      .addSink(new MyKafkaPrintSink());

    @SuppressWarnings("unused")
    DataStream<String> stream2 = env
      .addSource(new MySource())
      .addSink(new MyKafkaSink("test", "localhost:9092"));

    env.execute();
  }
View Full Code Here

    }

    StreamExecutionEnvironment env = StreamExecutionEnvironment
        .createLocalEnvironment(PARALLELISM);

    DataStream<String> streamSource = env.addSource(new TwitterSource(path, NUMBEROFTWEETS),
        SOURCE_PARALLELISM);


    DataStream<Tuple2<String, Integer>> dataStream = streamSource
        .flatMap(new SelectLanguageFlatMap())
View Full Code Here

    }

    StreamExecutionEnvironment env = StreamExecutionEnvironment
        .createLocalEnvironment(PARALLELISM);

    DataStream<String> streamSource = env.addSource(new TwitterSource(path, NUMBEROFTWEETS),
        SOURCE_PARALLELISM);

    DataStream<Tuple5<Long, Integer, String, String, String>> selectedDataStream = streamSource
        .flatMap(new SelectDataFlatMap());
View Full Code Here

    StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironment(
        PARALLELISM).setBufferTimeout(1000);

    // Build new model on every second of new data
    DataStream<Double[]> model = env.addSource(new TrainingDataSource(), SOURCE_PARALLELISM)
        .window(5000).reduceGroup(new PartialModelBuilder());

    // Use partial model for prediction
    DataStream<Integer> prediction = env.addSource(new NewDataSource(), SOURCE_PARALLELISM)
        .connect(model).map(new Predictor());
View Full Code Here

    // Build new model on every second of new data
    DataStream<Double[]> model = env.addSource(new TrainingDataSource(), SOURCE_PARALLELISM)
        .window(5000).reduceGroup(new PartialModelBuilder());

    // Use partial model for prediction
    DataStream<Integer> prediction = env.addSource(new NewDataSource(), SOURCE_PARALLELISM)
        .connect(model).map(new Predictor());

    prediction.print();

    env.execute();
View Full Code Here

  public static void main(String[] args) throws Exception {

    StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironment(
        PARALLELISM).setBufferTimeout(100);

    DataStream<Tuple2<String, Integer>> grades = env.addSource(new GradeSource(),
        SOURCE_PARALLELISM);
   
    DataStream<Tuple2<String, Integer>> salaries = env.addSource(new SalarySource(),
        SOURCE_PARALLELISM);
View Full Code Here

        PARALLELISM).setBufferTimeout(100);

    DataStream<Tuple2<String, Integer>> grades = env.addSource(new GradeSource(),
        SOURCE_PARALLELISM);
   
    DataStream<Tuple2<String, Integer>> salaries = env.addSource(new SalarySource(),
        SOURCE_PARALLELISM);

    DataStream<Tuple3<String, Integer, Integer>> joinedStream = grades.connect(salaries)
        .flatMap(new JoinTask());
   
View Full Code Here

  public static void main(String[] args) throws Exception {

    StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironment(1);

    @SuppressWarnings("unused")
    DataStream<String> dataStream1 = env.addSource(new MyFlumeSource("localhost", 41414))
        .addSink(new MyFlumePrintSink());

    @SuppressWarnings("unused")
    DataStream<String> dataStream2 = env.fromElements("one", "two", "three", "four", "five",
        "q").addSink(new MyFlumeSink("localhost", 42424));
View Full Code Here

  public static void main(String[] args) throws Exception {

    StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironment(1);

    @SuppressWarnings("unused")
    DataStream<String> dataStream1 = env
      .addSource(new MyRMQSource("localhost", "hello"))
      .addSink(new MyRMQPrintSink());

    @SuppressWarnings("unused")
    DataStream<String> dataStream2 = env
View Full Code Here

TOP
Copyright © 2018 www.massapi.com. 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.