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

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


  @Test
  public void test() throws Exception {
    LocalStreamEnvironment env = StreamExecutionEnvironment.createLocalEnvironment(1);
   
    @SuppressWarnings("unused")
    DataStream<Tuple1<Integer>> dataStream1 = env.addSource(new MySource1(), 1).writeAsText(PREFIX + "test1.txt");

    fillExpected1();

    @SuppressWarnings("unused")
    DataStream<Tuple1<Integer>> dataStream2 = env.addSource(new MySource1(), 1).writeAsText(PREFIX + "test2.txt", 5);
View Full Code Here


    DataStream<Tuple1<Integer>> dataStream1 = env.addSource(new MySource1(), 1).writeAsText(PREFIX + "test1.txt");

    fillExpected1();

    @SuppressWarnings("unused")
    DataStream<Tuple1<Integer>> dataStream2 = env.addSource(new MySource1(), 1).writeAsText(PREFIX + "test2.txt", 5);

    fillExpected2();

    @SuppressWarnings("unused")
    DataStream<Tuple1<Integer>> dataStream3 = env.addSource(new MySource1(), 1).writeAsText(PREFIX + "test3.txt", 10);
View Full Code Here

    DataStream<Tuple1<Integer>> dataStream2 = env.addSource(new MySource1(), 1).writeAsText(PREFIX + "test2.txt", 5);

    fillExpected2();

    @SuppressWarnings("unused")
    DataStream<Tuple1<Integer>> dataStream3 = env.addSource(new MySource1(), 1).writeAsText(PREFIX + "test3.txt", 10);

    fillExpected3();

    @SuppressWarnings("unused")
    DataStream<Tuple1<Integer>> dataStream4 = env.addSource(new MySource1(), 1).writeAsText(PREFIX + "test4.txt", 10, new Tuple1<Integer>(26));
View Full Code Here

    DataStream<Tuple1<Integer>> dataStream3 = env.addSource(new MySource1(), 1).writeAsText(PREFIX + "test3.txt", 10);

    fillExpected3();

    @SuppressWarnings("unused")
    DataStream<Tuple1<Integer>> dataStream4 = env.addSource(new MySource1(), 1).writeAsText(PREFIX + "test4.txt", 10, new Tuple1<Integer>(26));

    fillExpected4();

    @SuppressWarnings("unused")
    DataStream<Tuple1<Integer>> dataStream5 = env.addSource(new MySource1(), 1).writeAsText(PREFIX + "test5.txt", 10, new Tuple1<Integer>(14));
View Full Code Here

    DataStream<Tuple1<Integer>> dataStream4 = env.addSource(new MySource1(), 1).writeAsText(PREFIX + "test4.txt", 10, new Tuple1<Integer>(26));

    fillExpected4();

    @SuppressWarnings("unused")
    DataStream<Tuple1<Integer>> dataStream5 = env.addSource(new MySource1(), 1).writeAsText(PREFIX + "test5.txt", 10, new Tuple1<Integer>(14));

    fillExpected5();

    env.executeTest(MEMORYSIZE);
View Full Code Here

  @Test
  public void runStream() throws Exception {
    LocalStreamEnvironment env = StreamExecutionEnvironment
        .createLocalEnvironment(SOURCE_PARALELISM);

    env.addSource(new MySource(), SOURCE_PARALELISM).map(new MyTask()).addSink(new MySink());

    env.executeTest(MEMORYSIZE);
    assertEquals(10, data.keySet().size());

    for (Integer k : data.keySet()) {
View Full Code Here

  @Test
  public void test() throws Exception {
    LocalStreamEnvironment env = StreamExecutionEnvironment.createLocalEnvironment(1);

    @SuppressWarnings("unused")
    DataStream<Tuple1<Integer>> dataStream1 = env.addSource(new MySource1(), 1).writeAsCsv(PREFIX + "test1.txt");

    fillExpected1();

    @SuppressWarnings("unused")
    DataStream<Tuple1<Integer>> dataStream2 = env.addSource(new MySource1(), 1).writeAsCsv(PREFIX + "test2.txt", 5);
View Full Code Here

    DataStream<Tuple1<Integer>> dataStream1 = env.addSource(new MySource1(), 1).writeAsCsv(PREFIX + "test1.txt");

    fillExpected1();

    @SuppressWarnings("unused")
    DataStream<Tuple1<Integer>> dataStream2 = env.addSource(new MySource1(), 1).writeAsCsv(PREFIX + "test2.txt", 5);

    fillExpected2();

    @SuppressWarnings("unused")
    DataStream<Tuple1<Integer>> dataStream3 = env.addSource(new MySource1(), 1).writeAsCsv(PREFIX + "test3.txt", 10);
View Full Code Here

    DataStream<Tuple1<Integer>> dataStream2 = env.addSource(new MySource1(), 1).writeAsCsv(PREFIX + "test2.txt", 5);

    fillExpected2();

    @SuppressWarnings("unused")
    DataStream<Tuple1<Integer>> dataStream3 = env.addSource(new MySource1(), 1).writeAsCsv(PREFIX + "test3.txt", 10);

    fillExpected3();

    @SuppressWarnings("unused")
    DataStream<Tuple1<Integer>> dataStream4 = env.addSource(new MySource1(), 1).writeAsCsv(PREFIX + "test4.txt", 10, new Tuple1<Integer>(26));
View Full Code Here

    DataStream<Tuple1<Integer>> dataStream3 = env.addSource(new MySource1(), 1).writeAsCsv(PREFIX + "test3.txt", 10);

    fillExpected3();

    @SuppressWarnings("unused")
    DataStream<Tuple1<Integer>> dataStream4 = env.addSource(new MySource1(), 1).writeAsCsv(PREFIX + "test4.txt", 10, new Tuple1<Integer>(26));

    fillExpected4();

    @SuppressWarnings("unused")
    DataStream<Tuple1<Integer>> dataStream5 = env.addSource(new MySource1(), 1).writeAsCsv(PREFIX + "test5.txt", 10, new Tuple1<Integer>(14));
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.