What is the way operators are specified in Flink
Today, I will talk to you about the way of specifying operators in Flink. Many people may not know much about it. In order to make you understand better, the editor has summarized the following content for you. I hope you can get something according to this article.
When we were using flatMap, we passed an new FlatMapFunction anonymous inner class. And this is just one of them.
Method 1: implement the MapFunction interface
The easiest way is to implement a MapFunction interface, such as:
Text.flatMap (new MyFlatMapFunction ()) .keyby (new KeySelector () {@ Override public Object getKey (WC value) throws Exception {return value.word;}}) .timewindow (Time.seconds (5)) .sum ("count") .keyby () .setParallelism (1) Public static class MyFlatMapFunction implements FlatMapFunction {@ Override public void flatMap (String value, Collector out) throws Exception {String [] tokens = value.toLowerCase () .split (","); for (String token: tokens) {if (token.length () > 0) {out.collect (new WC (token, 1)) Method 2: anonymous inner class
This is the way we've been using it before.
Java8 Lambdas mode 4: Rich functionstext.flatMap (new RichFlatMapFunction () {@ Override public void flatMap (String value, Collector out) throws Exception {String [] tokens = value.toLowerCase () .split (",") For (String token: tokens) {if (token.length () > 0) {out.collect (new WC (token, 1));})
Inherit a RichFlatMapFunction class
After reading the above, do you have any further understanding of the way operators are specified in Flink? If you want to know more knowledge or related content, please follow the industry information channel, thank you for your support.