This is an automated email from the ASF dual-hosted git repository. dwysakowicz pushed a commit to branch master in repository https://gitbox.apache.org/repos/asf/flink.git
commit 7e7a9e9fa302fa7aa8d9c1d61463810cc07f71d1 Author: [email protected] <sanshi@WWDZ1234> AuthorDate: Thu Aug 27 16:28:05 2020 +0800 [FLINK-18797][examples] Update deprecated forms of keyBy in examples --- .../java/org/apache/flink/streaming/examples/async/AsyncIOExample.java | 2 +- .../apache/flink/streaming/examples/sideoutput/SideOutputExample.java | 2 +- .../apache/flink/streaming/examples/socket/SocketWindowWordCount.java | 2 +- .../org/apache/flink/streaming/examples/twitter/TwitterExample.java | 2 +- .../examples/windowing/GroupedProcessingTimeWindowExample.java | 2 +- .../org/apache/flink/streaming/examples/windowing/SessionWindowing.java | 2 +- .../apache/flink/streaming/examples/windowing/TopSpeedWindowing.java | 2 +- .../org/apache/flink/streaming/examples/windowing/WindowWordCount.java | 2 +- .../java/org/apache/flink/streaming/examples/wordcount/WordCount.java | 2 +- .../flink/streaming/scala/examples/socket/SocketWindowWordCount.scala | 2 +- .../apache/flink/streaming/scala/examples/twitter/TwitterExample.scala | 2 +- .../scala/examples/windowing/GroupedProcessingTimeWindowExample.scala | 2 +- .../flink/streaming/scala/examples/windowing/SessionWindowing.scala | 2 +- .../flink/streaming/scala/examples/windowing/TopSpeedWindowing.scala | 2 +- .../flink/streaming/scala/examples/windowing/WindowWordCount.scala | 2 +- .../org/apache/flink/streaming/scala/examples/wordcount/WordCount.scala | 2 +- 16 files changed, 16 insertions(+), 16 deletions(-) diff --git a/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/async/AsyncIOExample.java b/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/async/AsyncIOExample.java index 1e2260611..5f9aabc 100644 --- a/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/async/AsyncIOExample.java +++ b/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/async/AsyncIOExample.java @@ -304,7 +304,7 @@ public class AsyncIOExample { public void flatMap(String value, Collector<Tuple2<String, Integer>> out) throws Exception { out.collect(new Tuple2<>(value, 1)); } - }).keyBy(0).sum(1).print(); + }).keyBy(value -> value.f0).sum(1).print(); // execute the program env.execute("Async IO Example"); diff --git a/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/sideoutput/SideOutputExample.java b/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/sideoutput/SideOutputExample.java index f7fcf55..a3f37cd 100644 --- a/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/sideoutput/SideOutputExample.java +++ b/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/sideoutput/SideOutputExample.java @@ -95,7 +95,7 @@ public class SideOutputExample { }); DataStream<Tuple2<String, Integer>> counts = tokenized - .keyBy(0) + .keyBy(value -> value.f0) .window(TumblingEventTimeWindows.of(Time.seconds(5))) // group by the tuple field "0" and sum up tuple field "1" .sum(1); diff --git a/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/socket/SocketWindowWordCount.java b/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/socket/SocketWindowWordCount.java index 921ab3a..660b74a 100644 --- a/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/socket/SocketWindowWordCount.java +++ b/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/socket/SocketWindowWordCount.java @@ -76,7 +76,7 @@ public class SocketWindowWordCount { } }) - .keyBy("word") + .keyBy(value -> value.word) .timeWindow(Time.seconds(5)) .reduce(new ReduceFunction<WordWithCount>() { diff --git a/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/twitter/TwitterExample.java b/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/twitter/TwitterExample.java index 427832d..b2234a1 100644 --- a/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/twitter/TwitterExample.java +++ b/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/twitter/TwitterExample.java @@ -91,7 +91,7 @@ public class TwitterExample { // selecting English tweets and splitting to (word, 1) .flatMap(new SelectEnglishAndTokenizeFlatMap()) // group by words and sum their occurrences - .keyBy(0).sum(1); + .keyBy(value -> value.f0).sum(1); // emit result if (params.has("output")) { diff --git a/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/windowing/GroupedProcessingTimeWindowExample.java b/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/windowing/GroupedProcessingTimeWindowExample.java index 1837314..056ba5e 100644 --- a/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/windowing/GroupedProcessingTimeWindowExample.java +++ b/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/windowing/GroupedProcessingTimeWindowExample.java @@ -47,7 +47,7 @@ public class GroupedProcessingTimeWindowExample { DataStream<Tuple2<Long, Long>> stream = env.addSource(new DataSource()); stream - .keyBy(0) + .keyBy(value -> value.f0) .timeWindow(Time.of(2500, MILLISECONDS), Time.of(500, MILLISECONDS)) .reduce(new SummingReducer()) diff --git a/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/windowing/SessionWindowing.java b/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/windowing/SessionWindowing.java index 0c02d5b..38a85cf 100644 --- a/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/windowing/SessionWindowing.java +++ b/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/windowing/SessionWindowing.java @@ -81,7 +81,7 @@ public class SessionWindowing { // We create sessions for each id with max timeout of 3 time units DataStream<Tuple3<String, Long, Integer>> aggregated = source - .keyBy(0) + .keyBy(value -> value.f0) .window(EventTimeSessionWindows.withGap(Time.milliseconds(3L))) .sum(2); diff --git a/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/windowing/TopSpeedWindowing.java b/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/windowing/TopSpeedWindowing.java index ee06cd4..5276845 100644 --- a/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/windowing/TopSpeedWindowing.java +++ b/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/windowing/TopSpeedWindowing.java @@ -70,7 +70,7 @@ public class TopSpeedWindowing { double triggerMeters = 50; DataStream<Tuple4<Integer, Integer, Double, Long>> topSpeeds = carData .assignTimestampsAndWatermarks(new CarTimestamp()) - .keyBy(0) + .keyBy(value -> value.f0) .window(GlobalWindows.create()) .evictor(TimeEvictor.of(Time.of(evictionSec, TimeUnit.SECONDS))) .trigger(DeltaTrigger.of(triggerMeters, diff --git a/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/windowing/WindowWordCount.java b/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/windowing/WindowWordCount.java index b454f48..e1fc0d6 100644 --- a/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/windowing/WindowWordCount.java +++ b/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/windowing/WindowWordCount.java @@ -75,7 +75,7 @@ public class WindowWordCount { // split up the lines in pairs (2-tuples) containing: (word,1) text.flatMap(new WordCount.Tokenizer()) // create windows of windowSize records slided every slideSize records - .keyBy(0) + .keyBy(value -> value.f0) .countWindow(windowSize, slideSize) // group by the tuple field "0" and sum up tuple field "1" .sum(1); diff --git a/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/wordcount/WordCount.java b/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/wordcount/WordCount.java index 4fc8fea..7ab9f67 100644 --- a/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/wordcount/WordCount.java +++ b/flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/wordcount/WordCount.java @@ -83,7 +83,7 @@ public class WordCount { // split up the lines in pairs (2-tuples) containing: (word,1) text.flatMap(new Tokenizer()) // group by the tuple field "0" and sum up tuple field "1" - .keyBy(0).sum(1); + .keyBy(value -> value.f0).sum(1); // emit result if (params.has("output")) { diff --git a/flink-examples/flink-examples-streaming/src/main/scala/org/apache/flink/streaming/scala/examples/socket/SocketWindowWordCount.scala b/flink-examples/flink-examples-streaming/src/main/scala/org/apache/flink/streaming/scala/examples/socket/SocketWindowWordCount.scala index bdb1561..4b49e4c 100644 --- a/flink-examples/flink-examples-streaming/src/main/scala/org/apache/flink/streaming/scala/examples/socket/SocketWindowWordCount.scala +++ b/flink-examples/flink-examples-streaming/src/main/scala/org/apache/flink/streaming/scala/examples/socket/SocketWindowWordCount.scala @@ -67,7 +67,7 @@ object SocketWindowWordCount { val windowCounts = text .flatMap { w => w.split("\\s") } .map { w => WordWithCount(w, 1) } - .keyBy("word") + .keyBy(_.word) .timeWindow(Time.seconds(5)) .sum("count") diff --git a/flink-examples/flink-examples-streaming/src/main/scala/org/apache/flink/streaming/scala/examples/twitter/TwitterExample.scala b/flink-examples/flink-examples-streaming/src/main/scala/org/apache/flink/streaming/scala/examples/twitter/TwitterExample.scala index 048e7ac..de10d93 100644 --- a/flink-examples/flink-examples-streaming/src/main/scala/org/apache/flink/streaming/scala/examples/twitter/TwitterExample.scala +++ b/flink-examples/flink-examples-streaming/src/main/scala/org/apache/flink/streaming/scala/examples/twitter/TwitterExample.scala @@ -98,7 +98,7 @@ object TwitterExample { // selecting English tweets and splitting to (word, 1) .flatMap(new SelectEnglishAndTokenizeFlatMap) // group by words and sum their occurrences - .keyBy(0).sum(1) + .keyBy(_._1).sum(1) // emit result if (params.has("output")) { diff --git a/flink-examples/flink-examples-streaming/src/main/scala/org/apache/flink/streaming/scala/examples/windowing/GroupedProcessingTimeWindowExample.scala b/flink-examples/flink-examples-streaming/src/main/scala/org/apache/flink/streaming/scala/examples/windowing/GroupedProcessingTimeWindowExample.scala index 090b3f5..6ef66df 100644 --- a/flink-examples/flink-examples-streaming/src/main/scala/org/apache/flink/streaming/scala/examples/windowing/GroupedProcessingTimeWindowExample.scala +++ b/flink-examples/flink-examples-streaming/src/main/scala/org/apache/flink/streaming/scala/examples/windowing/GroupedProcessingTimeWindowExample.scala @@ -41,7 +41,7 @@ object GroupedProcessingTimeWindowExample { val stream: DataStream[(Long, Long)] = env.addSource(new DataSource) stream - .keyBy(0) + .keyBy(_._1) .timeWindow(Time.of(2500, MILLISECONDS), Time.of(500, MILLISECONDS)) .reduce((value1, value2) => (value1._1, value1._2 + value2._2)) .addSink(new SinkFunction[(Long, Long)]() { diff --git a/flink-examples/flink-examples-streaming/src/main/scala/org/apache/flink/streaming/scala/examples/windowing/SessionWindowing.scala b/flink-examples/flink-examples-streaming/src/main/scala/org/apache/flink/streaming/scala/examples/windowing/SessionWindowing.scala index 3723674..7fe483c 100644 --- a/flink-examples/flink-examples-streaming/src/main/scala/org/apache/flink/streaming/scala/examples/windowing/SessionWindowing.scala +++ b/flink-examples/flink-examples-streaming/src/main/scala/org/apache/flink/streaming/scala/examples/windowing/SessionWindowing.scala @@ -75,7 +75,7 @@ object SessionWindowing { // We create sessions for each id with max timeout of 3 time units val aggregated: DataStream[(String, Long, Int)] = source - .keyBy(0) + .keyBy(_._1) .window(EventTimeSessionWindows.withGap(Time.milliseconds(3L))) .sum(2) diff --git a/flink-examples/flink-examples-streaming/src/main/scala/org/apache/flink/streaming/scala/examples/windowing/TopSpeedWindowing.scala b/flink-examples/flink-examples-streaming/src/main/scala/org/apache/flink/streaming/scala/examples/windowing/TopSpeedWindowing.scala index bb66ead..33f9076 100644 --- a/flink-examples/flink-examples-streaming/src/main/scala/org/apache/flink/streaming/scala/examples/windowing/TopSpeedWindowing.scala +++ b/flink-examples/flink-examples-streaming/src/main/scala/org/apache/flink/streaming/scala/examples/windowing/TopSpeedWindowing.scala @@ -103,7 +103,7 @@ object TopSpeedWindowing { val topSpeeds = cars .assignAscendingTimestamps( _.time ) - .keyBy("carId") + .keyBy(_.carId) .window(GlobalWindows.create) .evictor(TimeEvictor.of(Time.of(evictionSec * 1000, TimeUnit.MILLISECONDS))) .trigger(DeltaTrigger.of(triggerMeters, new DeltaFunction[CarEvent] { diff --git a/flink-examples/flink-examples-streaming/src/main/scala/org/apache/flink/streaming/scala/examples/windowing/WindowWordCount.scala b/flink-examples/flink-examples-streaming/src/main/scala/org/apache/flink/streaming/scala/examples/windowing/WindowWordCount.scala index 349f253..07efd90 100644 --- a/flink-examples/flink-examples-streaming/src/main/scala/org/apache/flink/streaming/scala/examples/windowing/WindowWordCount.scala +++ b/flink-examples/flink-examples-streaming/src/main/scala/org/apache/flink/streaming/scala/examples/windowing/WindowWordCount.scala @@ -79,7 +79,7 @@ object WindowWordCount { .flatMap(_.toLowerCase.split("\\W+")) .filter(_.nonEmpty) .map((_, 1)) - .keyBy(0) + .keyBy(_._1) // create windows of windowSize records slided every slideSize records .countWindow(windowSize, slideSize) // group by the tuple field "0" and sum up tuple field "1" diff --git a/flink-examples/flink-examples-streaming/src/main/scala/org/apache/flink/streaming/scala/examples/wordcount/WordCount.scala b/flink-examples/flink-examples-streaming/src/main/scala/org/apache/flink/streaming/scala/examples/wordcount/WordCount.scala index 74b726f..271d737 100644 --- a/flink-examples/flink-examples-streaming/src/main/scala/org/apache/flink/streaming/scala/examples/wordcount/WordCount.scala +++ b/flink-examples/flink-examples-streaming/src/main/scala/org/apache/flink/streaming/scala/examples/wordcount/WordCount.scala @@ -74,7 +74,7 @@ object WordCount { .filter(_.nonEmpty) .map((_, 1)) // group by the tuple field "0" and sum up tuple field "1" - .keyBy(0) + .keyBy(_._1) .sum(1) // emit result
