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

Flink?側流輸出源碼示例解析

 更新時間:2022年09月15日 11:24:20   作者:JasonLee實時計算  
這篇文章主要為大家介紹了Flink?側流輸出源碼示例解析,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪

Flink 側流輸出源碼解析

Flink 的 side output 為我們提供了側流(分流)輸出的功能,根據(jù)條件可以把一條流分為多個不同的流,之后做不同的處理邏輯,下面就來看下側流輸出相關的源碼。

先來看下面的一個 Demo,一個流被分成了 3 個流,一個主流,兩個側流輸出。

SingleOutputStreamOperator<JasonLeePOJO> process =
        kafka_source1.process(
                new ProcessFunction<JasonLeePOJO, JasonLeePOJO>() {
                    @Override
                    public void processElement(
                            JasonLeePOJO value,
                            ProcessFunction<JasonLeePOJO, JasonLeePOJO>.Context ctx,
                            Collector<JasonLeePOJO> out)
                            throws Exception {
                        // 這個是主流輸出
                        if (value.getName().equals("flink")) {
                            out.collect(value);
                        // 下面兩個是測流輸出
                        } else if (value.getName().equals("spark")) {
                            ctx.output(test, value);
                        // 測流
                        } else if (value.getName().equals("hadoop")) {
                            ctx.output(test1, value);
                        }
                    }
                });

為了更加清楚的查看每一個算子,我禁用了 operator chain,任務的 DAG 圖如下所示:

這樣就比較清晰了,很明顯從 process 算子開始,1 個數(shù)據(jù)流分為了 3 個數(shù)據(jù)流,當然,在默認情況下沒有禁止

operator chain 所有的算子都是 chain 在一起的。

源碼解析

我們先來看第一個主流輸出也就是 out.collect(value) 的源碼,這里的 out 實際上是 TimestampedCollector 對象。

TimestampedCollector#collect

@Override
public void collect(T record) {
    output.collect(reuse.replace(record));
}

在 collect 方法中持有一個 output 對象,用來輸出數(shù)據(jù),在這里實際上是一個 CountingOutput 它是一個包裝了 Output 的對象,主要用于更新發(fā)送數(shù)據(jù)的 metric,并輸出數(shù)據(jù)。

CountingOutput#collect

@Override
public void collect(StreamRecord<OUT> record) {
    numRecordsOut.inc();
    output.collect(record);
}

在 CountingOutput 中也持有一個 output 對象,但是這里的 output 是 BroadcastingOutputCollector 對象,從名字就可以看出它是往下游廣播數(shù)據(jù)的,這里就有一個疑問?把數(shù)據(jù)廣播到下游,那豈不是下游的每個數(shù)據(jù)流都有這條數(shù)據(jù)嗎?這樣的話是怎么實現(xiàn)分流的呢?帶著這個疑問,我們來看 BroadcastingOutputCollector 的 collect 方法是怎么實現(xiàn)的。

BroadcastingOutputCollector#collect

@Override
public void collect(StreamRecord<T> record) {
    // 這里的 outputs 數(shù)組有三個 output 分別對應上面的三個輸出流
    for (Output<StreamRecord<T>> output : outputs) {
        output.collect(record);
    }
}

在 BroadcastingOutputCollector 對象里也持有一個 output 對象,其實他們都實現(xiàn)了 Output 接口,用來往下游發(fā)送數(shù)據(jù),這里的 outputs 是一個 Output 數(shù)組,代表了下游的所有 Output,因為上面有三個輸出流,所以數(shù)組里面就包含了 3 個 Output 對象。

循環(huán)的調用 output 的 collect 方法往下游發(fā)送數(shù)據(jù),因為我打斷了 operator chain,所以 process 算子和下游的 Print 算子不在同一個 operatorChain 內,那么上下游算子之間數(shù)據(jù)傳輸用的就是 RecordWriterOutput,否則用的是 CopyingChainingOutput 或者 ChainingOutput,具體使用的是哪個 Output 這里就不多介紹了,后面有時間的話會單獨介紹。

RecordWriterOutput#collect

@Override
public void collect(StreamRecord<OUT> record) {
    // 主流是沒有 outputTag 的,只有測流有 outputTag
    if (this.outputTag != null) {
        // we are not responsible for emitting to the main output.
        return;
    }

    pushToRecordWriter(record);
}

接著來看 RecordWriterOutput 的 collect 方法,在 collect 方法里面會先判斷 outputTag 是否為空,如果不為空不做任何處理,直接返回,否則就把數(shù)據(jù)推送到下游算子,只有側流輸出才需要定義 outputTag,主流(正常流)是沒有 outputTag 的,所以這里會走 pushToRecordWriter 方法把數(shù)據(jù)寫入到下游,也就是說雖然會以廣播的形式把數(shù)據(jù)廣播到所有下游,但其實另外兩個側流是直接返回的,只有主流才會把數(shù)據(jù)推送到下游,這也就解釋了上面的疑問。

然后再來看第二個側流輸出 ctx.output(test, value) 的源碼,這里的 ctx 實際上是 ProcessOperator#ContextImpl 對象。

ProcessOperator#ContextImpl#output

@Override
public <X> void output(OutputTag<X> outputTag, X value) {
    if (outputTag == null) {
        throw new IllegalArgumentException("OutputTag must not be null.");
    }
    output.collect(outputTag, new StreamRecord<>(value, element.getTimestamp()));
}

如果 outputTag 是空,直接拋出異常,因為這個是側流,所以必須要定義 OutputTag。這里的 output 實際上是父類 AbstractStreamOperator 所持有的變量,如果 outputTag 不為空,就調用 output 的 collect 方法把數(shù)據(jù)發(fā)送到下游,這里的 output 和上面的一樣是 CountingOutput 但是 collect 方法是另外一個重載的方法。

CountingOutput#collect

@Override
public <X> void collect(OutputTag<X> outputTag, StreamRecord<X> record) {
    numRecordsOut.inc();
    output.collect(outputTag, record);
}

可以發(fā)現(xiàn),這個 collect 方法比上面那個多了一個 OutputTag 參數(shù),也就是使用側流輸出的時候定義的 OutputTag 對象,然后調用 output 的 collect 方法發(fā)送數(shù)據(jù),這個也和上面的一樣,同樣是 BroadcastingOutputCollector 對象的另外一個重載方法,多了一個 OutputTag 參數(shù)。

BroadcastingOutputCollector#collect

@Override
public <X> void collect(OutputTag<X> outputTag, StreamRecord<X> record) {
    for (Output<StreamRecord<T>> output : outputs) {
        output.collect(outputTag, record);
    }
}

這里的邏輯和上面是一樣的,同樣的循環(huán)調用 collect 方法發(fā)送數(shù)據(jù)。

RecordWriterOutput#collect

@Override
public <X> void collect(OutputTag<X> outputTag, StreamRecord<X> record) {
    // 先要判斷兩個 OutputTag 是否一樣
    if (OutputTag.isResponsibleFor(this.outputTag, outputTag)) {
        pushToRecordWriter(record);
    }
}

在這個 collect 方法中會先判斷傳入的 OutputTag 對象和成員變量 this.outputTag 是不是相等,如果是的話,就發(fā)送數(shù)據(jù),否則不做任何處理,所以這里每次只會選擇一個下游側流輸出數(shù)據(jù),這樣就實現(xiàn)了所謂的分流。

OutputTag#isResponsibleFor

public static boolean isResponsibleFor(
        @Nullable OutputTag<?> owner, @Nonnull OutputTag<?> other) {
    return other.equals(owner);
}

可以看到在 isResponsibleFor 方法內是直接調用 OutputTag 的 equals 方法判斷兩個對象是否相等的。

第三個側流 test1 ctx.output(test1, value) 和第二個側流 test 是完全一樣的情況,這里就不在看代碼了。

上面是完成了分流操作,那怎么獲取到分流后結果呢(數(shù)據(jù)流)?我們可以通過 getSideOutput 方法獲取。

DataStream<JasonLeePOJO> sideOutput = process.getSideOutput(test);
DataStream<JasonLeePOJO> sideOutput1 = process.getSideOutput(test1);

getSideOutput 源碼

public <X> DataStream<X> getSideOutput(OutputTag<X> sideOutputTag) {
    sideOutputTag = clean(requireNonNull(sideOutputTag));

    // make a defensive copy
    sideOutputTag = new OutputTag<X>(sideOutputTag.getId(), sideOutputTag.getTypeInfo());

    TypeInformation<?> type = requestedSideOutputs.get(sideOutputTag);
    if (type != null && !type.equals(sideOutputTag.getTypeInfo())) {
        throw new UnsupportedOperationException(
                "A side output with a matching id was "
                        + "already requested with a different type. This is not allowed, side output "
                        + "ids need to be unique.");
    }

    requestedSideOutputs.put(sideOutputTag, sideOutputTag.getTypeInfo());

    SideOutputTransformation<X> sideOutputTransformation =
            new SideOutputTransformation<>(this.getTransformation(), sideOutputTag);
    return new DataStream<>(this.getExecutionEnvironment(), sideOutputTransformation);
}

getSideOutput 方法里先是構建了一個 SideOutputTransformation 對象,然后又構建了 DataStream 對象,這樣我們就可以基于分流后的 DataStream 做不同的處理邏輯了,從而實現(xiàn)了把一個 DataStream 分流成多個 DataStream 功能。

總結

通過對側流輸出的源碼進行解析,在分流的時候,數(shù)據(jù)是通過廣播的方式發(fā)送到下游算子的,對于主流的數(shù)據(jù)來說,只有 OutputTag 為空的才會處理,側流因為 OutputTag 不為空,所以直接返回,不做任何處理,那對于側流的數(shù)據(jù)來說,是通過判斷兩個 OutputTag 是否相等,所以每次只會把數(shù)據(jù)發(fā)送到下游對應的那一個側流上去,這樣即可實現(xiàn)分流邏輯。

以上就是Flink 側流輸出源碼示例解析的詳細內容,更多關于Flink 側流輸出的資料請關注腳本之家其它相關文章!

相關文章

  • RSync實現(xiàn)文件備份同步詳解

    RSync實現(xiàn)文件備份同步詳解

    rsync,remote synchronize顧名思意就知道它是一款實現(xiàn)遠程同步功能的軟件,它在同步文件的同時,可以保持原來文件的權限、時間、軟硬鏈接等附加信息
    2016-03-03
  • 批量修改所有服務器的dbmail配置(推薦)

    批量修改所有服務器的dbmail配置(推薦)

    這篇文章主要介紹了批量修改所有服務器的dbmail配置的相關資料,需要的朋友可以參考下
    2017-08-08
  • Ubuntu20.04美化桌面dock欄居中方式

    Ubuntu20.04美化桌面dock欄居中方式

    這篇文章主要介紹了Ubuntu20.04美化桌面dock欄居中方式,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2024-09-09
  • 阿里云服務器?jdk1.8?安裝配置教程

    阿里云服務器?jdk1.8?安裝配置教程

    這篇文章主要介紹了阿里云服務器?jdk1.8?安裝配置,本文給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2022-12-12
  • 高命中率的varnish緩存配置分享

    高命中率的varnish緩存配置分享

    這篇文章主要介紹了高命中率的varnish緩存配置分享,本文直接給出配置代碼,需要的朋友可以參考下
    2014-12-12
  • 服務器Apache與Tomcat和Nginx的理解和對比分析詳解

    服務器Apache與Tomcat和Nginx的理解和對比分析詳解

    今天小編就為大家分享一篇關于服務器Apache與Tomcat和Nginx的理解和對比分析詳解,小編覺得內容挺不錯的,現(xiàn)在分享給大家,具有很好的參考價值,需要的朋友一起跟隨小編來看看吧
    2019-04-04
  • k8s查看各組件日志的方法圖文詳解

    k8s查看各組件日志的方法圖文詳解

    這篇文章主要給大家介紹了關于k8s查看各組件日志的方法,Kubernetes(簡稱K8s)已成為現(xiàn)代容器化應用程序管理的主要平臺之一,文中通過代碼介紹的非常詳細,需要的朋友可以參考下
    2023-09-09
  • Mac OSX下使用MAMP安裝配置PHP開發(fā)環(huán)境

    Mac OSX下使用MAMP安裝配置PHP開發(fā)環(huán)境

    本部分描述如何在 Mac 上安裝 MAMP。將通過一個操作安裝 Apache Web 服務器、MySQL 和phpMyAdmin,需要的朋友可以參考下
    2017-09-09
  • vscode內網訪問服務器的方法

    vscode內網訪問服務器的方法

    這篇文章主要介紹了vscode內網訪問服務器的相關知識,本文給大家介紹的非常詳細,對大家的學習或工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2022-06-06
  • 獨立IP與共享IP的區(qū)別

    獨立IP與共享IP的區(qū)別

    做網站選擇獨立IP還是共享IP?相信很多站長都在此糾結過,自己不使用服務器的時候從來沒有關心過獨立IP和共享IP的究竟有什么具體的差別。但當自己真正用到的時候,才發(fā)現(xiàn):同樣都是IP,差別不是一般的大,獨立IP的強悍,不用的人是沒有辦法體會的
    2015-12-12

最新評論

桑日县| 金山区| 舒城县| 彰武县| 鲁甸县| 宕昌县| 岳阳市| 南投市| 全南县| 延安市| 平塘县| 富民县| 泾阳县| 黄陵县| 台江县| 阿勒泰市| 田阳县| 恭城| 乾安县| 运城市| 平凉市| 和静县| 临沧市| 新乡市| 溧水县| 吉木萨尔县| 南昌市| 武夷山市| 盘山县| 留坝县| 青岛市| 闽侯县| 建始县| 松溪县| 北宁市| 凤庆县| 建宁县| 延安市| 宜丰县| 霍州市| 黔江区|