今天就跟大家聊聊有關Flink中指定算子的方式是什么,可能很多人都不太了解,為了讓大家更加了解,小編給大家總結了以下內容,希望大家根據這篇文章可以有所收獲。
我們之前在使用flatMap時,傳了一個new FlatMapFunction匿名內部類。而這僅僅是其中的一種方式。
最簡單的方式就是實現一個MapFunction接口,例如:
text.flatMap(new MyFlatMapFunction()).keyBy(new KeySelector<WC, Object>() { @Override public Object getKey(WC value) throws Exception { return value.word; } }).timeWindow(Time.seconds(5)).sum("count").print().setParallelism(1); public static class MyFlatMapFunction implements FlatMapFunction<String, WC> { @Override public void flatMap(String value, Collector<WC> out) throws Exception { String[] tokens = value.toLowerCase().split(","); for (String token : tokens) { if (token.length() > 0) { out.collect(new WC(token, 1)); } } } }
這種方式就是我們之前一直使用的方式。
text.flatMap(new RichFlatMapFunction<String, WC>() { @Override public void flatMap(String value, Collector<WC> out) throws Exception { String[] tokens = value.toLowerCase().split(","); for (String token : tokens) { if (token.length() > 0) { out.collect(new WC(token, 1)); } } } })
繼承一個RichFlatMapFunction類
看完上述內容,你們對Flink中指定算子的方式是什么有進一步的了解嗎?如果還想了解更多知識或者相關內容,請關注億速云行業資訊頻道,感謝大家的支持。
免責聲明:本站發布的內容(圖片、視頻和文字)以原創、轉載和分享為主,文章觀點不代表本網站立場,如果涉及侵權請聯系站長郵箱:is@yisu.com進行舉報,并提供相關證據,一經查實,將立刻刪除涉嫌侵權內容。