最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

Flink實現(xiàn)特定統(tǒng)計的歸約聚合reduce操作

 更新時間:2023年02月08日 11:55:25   作者:響徹天堂丶  
這篇文章主要介紹了Flink實現(xiàn)特定統(tǒng)計的歸約聚合reduce操作,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)吧

如果說簡單聚合是對一些特定統(tǒng)計需求的實現(xiàn),那么 reduce 算子就是一個一般化的聚合統(tǒng)計操作了。從大名鼎鼎的 MapReduce 開始,我們對 reduce 操作就不陌生:它可以對已有的

數(shù)據(jù)進行歸約處理,把每一個新輸入的數(shù)據(jù)和當(dāng)前已經(jīng)歸約出來的值,再做一個聚合計算。與簡單聚合類似,reduce 操作也會將 KeyedStream 轉(zhuǎn)換為 DataStream。它不會改變流的元

素數(shù)據(jù)類型,所以輸出類型和輸入類型是一樣的。調(diào)用 KeyedStream 的 reduce 方法時,需要傳入一個參數(shù),實現(xiàn) ReduceFunction 接口。接口在源碼中的定義如下:

@Public
@FunctionalInterface
public interface ReduceFunction<T> extends Function, Serializable {
    /**
     * The core method of ReduceFunction, combining two values into one value of the same type. The
     * reduce function is consecutively applied to all values of a group until only a single value
     * remains.
     *
     * @param value1 The first value to combine.
     * @param value2 The second value to combine.
     * @return The combined value of both input values.
     * @throws Exception This method may throw exceptions. Throwing an exception will cause the
     *     operation to fail and may trigger recovery.
     */
    T reduce(T value1, T value2) throws Exception;
}

ReduceFunction 接口里需要實現(xiàn) reduce()方法,這個方法接收兩個輸入事件,經(jīng)過轉(zhuǎn)換處理之后輸出一個相同類型的事件;所以,對于一組數(shù)據(jù),我們可以先取兩個進行合并,然后再

將合并的結(jié)果看作一個數(shù)據(jù)、再跟后面的數(shù)據(jù)合并,最終會將它“簡化”成唯一的一個數(shù)據(jù),這也就是 reduce“歸約”的含義。在流處理的底層實現(xiàn)過程中,實際上是將中間“合并的結(jié)果”

作為任務(wù)的一個狀態(tài)保存起來的;之后每來一個新的數(shù)據(jù),就和之前的聚合狀態(tài)進一步做歸約。

其實,reduce 的語義是針對列表進行規(guī)約操作,運算規(guī)則由 ReduceFunction 中的 reduce方法來定義,而在 ReduceFunction 內(nèi)部會維護一個初始值為空的累加器,注意累加器的類型

和輸入元素的類型相同,當(dāng)?shù)谝粭l元素到來時,累加器的值更新為第一條元素的值,當(dāng)新的元素到來時,新元素會和累加器進行累加操作,這里的累加操作就是 reduce 函數(shù)定義的運算規(guī)

則。然后將更新以后的累加器的值向下游輸出。

我們可以單獨定義一個函數(shù)類實現(xiàn) ReduceFunction 接口,也可以直接傳入一個匿名類。當(dāng)然,同樣也可以通過傳入 Lambda 表達式實現(xiàn)類似的功能。與簡單聚合類似,reduce 操作也會將 KeyedStream 轉(zhuǎn)換為 DataStrema。它不會改變流的元素數(shù)據(jù)類型,所以輸出類型和輸入類型是一樣的。下面我們來看一個稍復(fù)雜的例子。

我們將數(shù)據(jù)流按照用戶 id 進行分區(qū),然后用一個 reduce 算子實現(xiàn) sum 的功能,統(tǒng)計每個用戶訪問的頻次;進而將所有統(tǒng)計結(jié)果分到一組,用另一個 reduce 算子實現(xiàn) maxBy 的功能,記錄所有用戶中訪問頻次最高的那個,也就是當(dāng)前訪問量最大的用戶是誰。

package com.rosh.flink.test;
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.api.common.functions.ReduceFunction;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import java.util.ArrayList;
import java.util.List;
import java.util.Random;
/**
 * 我們將數(shù)據(jù)流按照用戶 id 進行分區(qū),然后用一個 reduce 算子實現(xiàn) sum 的功能,統(tǒng)計每個
 * 用戶訪問的頻次;進而將所有統(tǒng)計結(jié)果分到一組,用另一個 reduce 算子實現(xiàn) maxBy 的功能,
 * 記錄所有用戶中訪問頻次最高的那個,也就是當(dāng)前訪問量最大的用戶是誰。
 */
public class TransReduceTest {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);
        //隨機生成數(shù)據(jù)
        Random random = new Random();
        List<Integer> userIds = new ArrayList<>();
        for (int i = 1; i <= 10; i++) {
            userIds.add(random.nextInt(5));
        }
        DataStreamSource<Integer> userIdDS = env.fromCollection(userIds);
        //每個ID訪問記錄一次
        SingleOutputStreamOperator<Tuple2<Integer, Long>> mapDS = userIdDS.map(new MapFunction<Integer, Tuple2<Integer, Long>>() {
            @Override
            public Tuple2<Integer, Long> map(Integer value) throws Exception {
                return new Tuple2<>(value, 1L);
            }
        });
        //統(tǒng)計每個user訪問多少次
        SingleOutputStreamOperator<Tuple2<Integer, Long>> sumDS = mapDS.keyBy(tuple -> tuple.f0).reduce(new ReduceFunction<Tuple2<Integer, Long>>() {
            @Override
            public Tuple2<Integer, Long> reduce(Tuple2<Integer, Long> value1, Tuple2<Integer, Long> value2) throws Exception {
                return new Tuple2<>(value1.f0, value1.f1 + value2.f1);
            }
        });
        sumDS.print("sumDS  ->>>>>>>>>>>>>");
        //把所有分區(qū)合并,求出最大的訪問量
        SingleOutputStreamOperator<Tuple2<Integer, Long>> maxDS = sumDS.keyBy(key -> true).reduce(new ReduceFunction<Tuple2<Integer, Long>>() {
            @Override
            public Tuple2<Integer, Long> reduce(Tuple2<Integer, Long> value1, Tuple2<Integer, Long> value2) throws Exception {
                if (value1.f1 > value2.f1) {
                    return value1;
                } else {
                    return value2;
                }
            }
        });
        maxDS.print("maxDS ->>>>>>>>>>>");
        env.execute("TransReduceTest");
    }
}

到此這篇關(guān)于Flink實現(xiàn)特定統(tǒng)計的歸約聚合reduce操作的文章就介紹到這了,更多相關(guān)Flink歸約聚合內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • springmvc實現(xiàn)自定義類型轉(zhuǎn)換器示例

    springmvc實現(xiàn)自定義類型轉(zhuǎn)換器示例

    本篇文章主要介紹了springmvc實現(xiàn)自定義類型轉(zhuǎn)換器示例,小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2017-02-02
  • Spring Cloud Gateway重試機制的實現(xiàn)

    Spring Cloud Gateway重試機制的實現(xiàn)

    這篇文章主要介紹了Spring Cloud Gateway重試機制的實現(xiàn),小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2019-03-03
  • java計算給定字符串中出現(xiàn)次數(shù)最多的字母和該字母出現(xiàn)次數(shù)的方法

    java計算給定字符串中出現(xiàn)次數(shù)最多的字母和該字母出現(xiàn)次數(shù)的方法

    這篇文章主要介紹了java計算給定字符串中出現(xiàn)次數(shù)最多的字母和該字母出現(xiàn)次數(shù)的方法,涉及java字符串的遍歷、轉(zhuǎn)換及運算相關(guān)操作技巧,需要的朋友可以參考下
    2017-02-02
  • Spring Task定時任務(wù)使用

    Spring Task定時任務(wù)使用

    這篇文章主要介紹了Spring Task定時任務(wù)使用方式,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2024-08-08
  • Java Swing JComboBox下拉列表框的示例代碼

    Java Swing JComboBox下拉列表框的示例代碼

    這篇文章主要介紹了Java Swing JComboBox下拉列表框的示例代碼,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2019-12-12
  • java ThreadPoolExecutor線程池拒絕策略避坑

    java ThreadPoolExecutor線程池拒絕策略避坑

    這篇文章主要為大家介紹了java ThreadPoolExecutor拒絕策略避坑踩坑示例詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪
    2022-07-07
  • SpringData JPA快速上手之關(guān)聯(lián)查詢及JPQL語句書寫詳解

    SpringData JPA快速上手之關(guān)聯(lián)查詢及JPQL語句書寫詳解

    JPA都有SpringBoot的官方直接提供的starter,而Mybatis沒有,直到SpringBoot 3才開始加入到官方模版中,這篇文章主要介紹了SpringData JPA快速上手,關(guān)聯(lián)查詢,JPQL語句書寫的相關(guān)知識,感興趣的朋友一起看看吧
    2023-09-09
  • Spring IOC和DI實現(xiàn)原理及實例解析

    Spring IOC和DI實現(xiàn)原理及實例解析

    這篇文章主要介紹了Spring IOC和DI實現(xiàn)原理及實例解析,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友可以參考下
    2020-06-06
  • JAVA集成本地部署的DeepSeek的圖文教程

    JAVA集成本地部署的DeepSeek的圖文教程

    本文主要介紹了JAVA集成本地部署的DeepSeek的圖文教程,包含配置環(huán)境變量及下載DeepSeek-R1模型并啟動,具有一定的參考價值,感興趣的可以了解一下
    2025-03-03
  • java實現(xiàn)簡易局域網(wǎng)聊天功能

    java實現(xiàn)簡易局域網(wǎng)聊天功能

    這篇文章主要為大家詳細(xì)介紹了java實現(xiàn)簡易局域網(wǎng)聊天功能,使用UDP模式編寫一個聊天程序,具有一定的參考價值,感興趣的小伙伴們可以參考一下
    2018-04-04

最新評論

白水县| 罗城| 桃园县| 灵山县| 大渡口区| 姚安县| 南木林县| 玛纳斯县| 娱乐| 富宁县| 舒城县| 和龙市| 安庆市| 开远市| 安远县| 瑞金市| 沁水县| 江西省| 宁蒗| 沁水县| 景宁| 冕宁县| 阿克苏市| 通江县| 同心县| 松桃| 益阳市| 新化县| 巨鹿县| 云霄县| 汤阴县| 娄烦县| 谢通门县| 漠河县| 中山市| 贺州市| 隆林| 密山市| 祥云县| 巴塘县| 建瓯市|