溫馨提示×

溫馨提示×

您好,登錄后才能下訂單哦!

密碼登錄×
登錄注冊×
其他方式登錄
點擊 登錄注冊 即表示同意《億速云用戶服務(wù)條款》

如何用Flink Apply對窗口內(nèi)的數(shù)據(jù)流進(jìn)行處理

發(fā)布時間:2021-12-31 10:19:42 來源:億速云 閱讀:192 作者:iii 欄目:大數(shù)據(jù)

這篇文章主要講解了“如何用Flink  Apply對窗口內(nèi)的數(shù)據(jù)流進(jìn)行處理”,文中的講解內(nèi)容簡單清晰,易于學(xué)習(xí)與理解,下面請大家跟著小編的思路慢慢深入,一起來研究和學(xué)習(xí)“如何用Flink  Apply對窗口內(nèi)的數(shù)據(jù)流進(jìn)行處理”吧!

Apply算子:對窗口內(nèi)的數(shù)據(jù)流進(jìn)行處理

示例環(huán)境

java.version: 1.8.x
flink.version: 1.11.1

示例數(shù)據(jù)源 (項目碼云下載)

Flink 系例 之 搭建開發(fā)環(huán)境與數(shù)據(jù)

Apply.java

import com.flink.examples.DataSource;
import org.apache.flink.api.java.functions.KeySelector;
import org.apache.flink.api.java.tuple.Tuple3;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.windowing.WindowFunction;
import org.apache.flink.streaming.api.windowing.windows.GlobalWindow;
import org.apache.flink.util.Collector;

import java.util.Iterator;
import java.util.List;

/**
 * @Description Apply方法:對窗口內(nèi)的數(shù)據(jù)流進(jìn)行處理
 */
public class Apply {

    /**
     * 遍歷集合,分別打印不同性別的總?cè)藬?shù)與年齡之和
     * @param args
     * @throws Exception
     */
    public static void main(String[] args) throws Exception {
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        List<Tuple3<String, String, Integer>> tuple3List = DataSource.getTuple3ToList();
        DataStream<String> dataStream = env.fromCollection(tuple3List)
                .keyBy((KeySelector<Tuple3<String, String, Integer>, String>) k -> k.f1)
                //按數(shù)量窗口滾動,每3個輸入窗口數(shù)據(jù)流,計算一次
                .countWindow(3)
                //只能基于Windowed窗口Stream進(jìn)行調(diào)用
                .apply(
                        //WindowFunction<IN, OUT, KEY, W extends Window>
                        new WindowFunction<Tuple3<String, String, Integer>, String, String, GlobalWindow>() {
                            /**
                             * 處理窗口數(shù)據(jù)集合
                             * @param s         從keyBy里返回的key值
                             * @param window    窗口類型
                             * @param input     從窗口獲取的所有分區(qū)數(shù)據(jù)流
                             * @param out       輸出數(shù)據(jù)流對象
                             * @throws Exception
                             */
                            @Override
                            public void apply(String s, GlobalWindow window, Iterable<Tuple3<String, String, Integer>> input, Collector<String> out) throws Exception {
                                Iterator<Tuple3<String, String, Integer>> iterator = input.iterator();
                                int total = 0;
                                int i = 0;
                                while (iterator.hasNext()){
                                    Tuple3<String, String, Integer> tuple3 = iterator.next();
                                    total += tuple3.f2;
                                    i ++ ;
                                }
                                out.collect(s + "共:"+i+"人,累加總年齡:" + total);
                            }
                        });
        dataStream.print();
        env.execute("flink Filter job");
    }
}

打印結(jié)果

4> girl共:3人,累加總年齡:74
2> man共:3人,累加總年齡:79

感謝各位的閱讀,以上就是“如何用Flink  Apply對窗口內(nèi)的數(shù)據(jù)流進(jìn)行處理”的內(nèi)容了,經(jīng)過本文的學(xué)習(xí)后,相信大家對如何用Flink  Apply對窗口內(nèi)的數(shù)據(jù)流進(jìn)行處理這一問題有了更深刻的體會,具體使用情況還需要大家實踐驗證。這里是億速云,小編將為大家推送更多相關(guān)知識點的文章,歡迎關(guān)注!

向AI問一下細(xì)節(jié)

免責(zé)聲明:本站發(fā)布的內(nèi)容(圖片、視頻和文字)以原創(chuàng)、轉(zhuǎn)載和分享為主,文章觀點不代表本網(wǎng)站立場,如果涉及侵權(quán)請聯(lián)系站長郵箱:is@yisu.com進(jìn)行舉報,并提供相關(guān)證據(jù),一經(jīng)查實,將立刻刪除涉嫌侵權(quán)內(nèi)容。

AI