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

Flink DataStream基礎(chǔ)框架源碼分析

 更新時(shí)間:2022年12月01日 14:29:26   作者:xiangel  
這篇文章主要為大家介紹了Flink DataStream基礎(chǔ)框架源碼分析,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪

引言

希望通過對Flink底層源碼的學(xué)習(xí)來更深入了解Flink的相關(guān)實(shí)現(xiàn)邏輯。這里新開一個(gè)Flink源碼解析的系列來深入介紹底層源碼邏輯。說明:這里默認(rèn)相關(guān)讀者具備Flink相關(guān)基礎(chǔ)知識和開發(fā)經(jīng)驗(yàn),所以不會過多介紹相關(guān)的基礎(chǔ)概念相關(guān)內(nèi)容,F(xiàn)link使用的版本為1.15.2。初步確定按如下幾個(gè)大的方面來介紹

計(jì)算模型

  • DataStream基礎(chǔ)框架

部署&調(diào)度

存儲體系

底層支撐

概覽

本篇是第一篇,介紹計(jì)算模型的基礎(chǔ)DataStream的相關(guān)內(nèi)容,這一篇只介紹DataStream的基礎(chǔ)內(nèi)容,如如何實(shí)現(xiàn)相關(guān)的操作,數(shù)據(jù)結(jié)構(gòu)等,不會涉及到窗口、事件事件和狀態(tài)等信息

DataStream是對數(shù)據(jù)流的一個(gè)抽象,其提供了豐富的操作算子(例如過濾、map、join、聚合、定義窗口等)來對數(shù)據(jù)流進(jìn)行處理,下圖描述了Flink中源數(shù)據(jù)通過DataStream的轉(zhuǎn)換最后輸出的整個(gè)過程。

通過上圖可以來構(gòu)想下,一般一個(gè)DataStream具有如下主要屬性

屬性說明
上游依賴標(biāo)識上游依賴信息,這樣能把整個(gè)處理流程串聯(lián)起來
并行度處理邏輯的并行度信息,這樣可以提高處理的速度
輸入格式指定輸入數(shù)據(jù)的格式,如InnerType { public int id; public String text; }
輸出格式指定輸出數(shù)據(jù)的格式
處理邏輯上游datastream轉(zhuǎn)換到目前的datastream的具體邏輯操作,如map的具體邏輯信息。

最終整個(gè)數(shù)據(jù)流會生成一個(gè)DAG圖(有向無環(huán)圖),通過這個(gè)DAG圖就可以生成對應(yīng)的任務(wù)來運(yùn)行了。下面來具體分析DataStream的實(shí)現(xiàn)和生成DAG圖(Flink中叫Graph)

深入DataStream

首先我們通過下圖來看看,DataStream中的一些主要的輔助類,DataStream類本身主要邏輯是對各類轉(zhuǎn)換關(guān)系和sink的操作,而前面說到的一些主要屬性信息都是通過輔助類來處理的。

Transformation:本身主要管理了輸出格式、上游依賴、并行度、id編號等信息以及StreamOperator的工廠類(StreamOperatorFactory)

StreamOperator:主要是各類操作的具體處理邏輯

Function:用戶自定義函數(shù)的接口,如DataStream中map處理時(shí)需要傳入的MapFunction就是Function的子接口

DataStream

屬性和方法

DataStream的屬性比較簡單,就2個(gè),1個(gè)是實(shí)行的環(huán)境信息,另一個(gè)是Transformation。
DataStream中的方法主要分為以下幾類

  • 基礎(chǔ)屬性信息:如獲取并行度,id,輸出格式等,大多數(shù)是代理來調(diào)用Transformation中對應(yīng)的方法
  • 轉(zhuǎn)換操作:各類的轉(zhuǎn)換處理,如map、filter、shuffle、join等
  • 輸出處理:各類輸出的sink處理,如保存為文本等,不過大多數(shù)方法都不推薦使用了,這里主要的方法是addSink()
  • 觸發(fā)執(zhí)行: 如executeAndCollect,內(nèi)部是調(diào)用了env.executeAsync來執(zhí)行streaming dataflow

除了轉(zhuǎn)換操作外,其他幾類的邏輯都比較直觀和簡單,這里重點(diǎn)介紹下轉(zhuǎn)換操作的處理,轉(zhuǎn)換操作這里分為3類,1.返回是一個(gè)DataStream。如map、filter、union、keyBy等;2.返回的是一個(gè)Streams,即輸入是多個(gè)DataStream,這類的操作主要是多流關(guān)聯(lián)的操作,如join、coGroup。這些Streams的類中實(shí)現(xiàn)了一些方法,來返回一個(gè)DataStream; 3.window類,返回的是AllWindowedStream類型,同樣這些類中也是有方法,來返回一個(gè)DataStream。

說明:如上的各個(gè)分類都是個(gè)人基于理解上做的各個(gè)分類處理,非官方定義

類體系

DataStream的類圖關(guān)系比較簡單,就如下這幾個(gè)類,具體每個(gè)子類的信息見下表

子類名說明
SingleOutputStreamOperator只有1個(gè)輸出的DataStream
IterativeStream迭代的DataStream,具體使用場景后面分析
DataStreamSource最開始的DataStream,里面有source的信息
KeyedStream有一個(gè)Key信息的DataStream

Transformation

屬性和方法

DataStream類本身主要是提供了給外部的編程接口的支持,而對Streaming flow算子節(jié)點(diǎn)本身的一些屬性和操作則由Transformation來負(fù)責(zé)

從上圖可以看出其主要屬性有節(jié)點(diǎn)id,名稱,并行度,輸出類型還有一些與資源相關(guān)的內(nèi)容,還有一個(gè)是上游的輸入Transformation,由于這個(gè)因不同的Transformation會有不同的數(shù)據(jù)個(gè)數(shù),所以這個(gè)信息是放在各個(gè)子類中的。如ReduceTransformation是有一個(gè)input屬性來記錄上游依賴,而如TwoInputTransformation則是有2個(gè)屬性input1和input2來記錄上游依賴,另外如SourceTransformation,這個(gè)是源頭的Transformation,是沒有上游依賴的Transformation,所以沒有屬性來記錄,但是有個(gè)Source屬性來記錄Source輸入
具體對數(shù)據(jù)的操作處理,在Transformation里面有個(gè)StreamOperatorFactory屬性,其中的StreamOperator實(shí)現(xiàn)了各種的處理算子。注意這里不是所有的Transformation都包含StreamOperatorFactory,如SourceTransformation中就沒有,這個(gè)具體大家可以看看相關(guān)的代碼。

Transformation的方法基本上是對上述屬性的get和set操作,這里重點(diǎn)要說明一下的是PhysicalTransformation(下面類體系來介紹)中的setChainingStrategy方法,這里的ChainingStrategy是一個(gè)枚舉類,主要是控制多個(gè)連續(xù)的算子是否可以進(jìn)行鏈?zhǔn)教幚?,這個(gè)具體的我們在下面介紹StreamGraph時(shí)再介紹

類體系

Transformation的大多數(shù)類均為PhysicalTransformation的子類,PhysicalTransformation為有物理操作的,重點(diǎn)是這類的子類是支持Chaining操作的。我們先來看看其重要子類

子類名說明
SourceTransformation連接Source的Transformation,是整個(gè)streaming flow的最開始的轉(zhuǎn)換處理
SinkTransformation輸出的轉(zhuǎn)換處理,是整個(gè)streaming flow的最后一個(gè)
OneInputTransformation只有1個(gè)輸入的轉(zhuǎn)換處理,如map、filter這類的處理
TwoInputTransformation有2個(gè)輸入的轉(zhuǎn)換處理

StreamOperator

屬性和方法

StreamOperator負(fù)責(zé)對數(shù)據(jù)進(jìn)行處理的具體邏輯,如map處理的StreamMap,由于各個(gè)Operator的處理方式的不同,這里主要以AbstractStreamOperator來介紹一些主要的屬性,如output的數(shù)據(jù),StreamConfig,StreamingRuntimeContext等。

下面我們重點(diǎn)介紹下相關(guān)的方法

StreamOperator接口有定義了重要的3個(gè)方法(這里只介紹與數(shù)據(jù)基礎(chǔ)處理相關(guān)的部分)

方法說明
open()數(shù)據(jù)處理的前處理,如算子的初始化操作等
finish()數(shù)據(jù)處理的后處理,如緩存數(shù)據(jù)的flush操作等
close()該方法在算子生命周期的最后調(diào)用,不管是算子運(yùn)行成功還是失敗或者取消,主要是對算子使用到的資源的各種釋放處理

另外關(guān)注的對數(shù)據(jù)進(jìn)行實(shí)際處理的方法,

接口方法說明
OneInputStreamOperatorprocessElement()對數(shù)據(jù)元素進(jìn)行處理,實(shí)際該接口在OneInputStreamOperator的父接口Input中定義
TwoInputStreamOperatorprocessElement1()對input1的數(shù)據(jù)元素進(jìn)行處理
processElement2()對input2的數(shù)據(jù)元素進(jìn)行處理

類體系

StreamOpterator的子類非常多,包括測試類的一起有287個(gè),這些大致可以歸屬到如下3個(gè)子類中,

類名說明
OneInputStreamOperator只有1個(gè)輸入的源
TwoInputStreamOperator有2個(gè)輸入源
AbstractStreamOperator

Function

Function是針對所有的用戶自定義的函數(shù),各子類主要是實(shí)現(xiàn)對應(yīng)的,這里定義了種類豐富的各類Function的子接口類來適配各種不同的加工場景,具體的就看源碼了,這里就不詳細(xì)介紹了

本節(jié)的最后,我們通過一個(gè)例子來看看這幾個(gè)類是怎么組合的。如下是一個(gè)常見的對DataStream進(jìn)行map處理的操作

      text
        .flatMap(new Tokenizer)

處理后對應(yīng)的DataStream的結(jié)構(gòu)如下圖

DataStream生成提交執(zhí)行的Graph

前面分析了DataStream,是單個(gè)節(jié)點(diǎn)的,接下來看看整個(gè)streaming flow在flink中是怎么轉(zhuǎn)換為可以執(zhí)行的邏輯的。一般整個(gè)數(shù)據(jù)流我們叫做DAG,那在Flink中叫PipeLine,其實(shí)現(xiàn)類是StreamGraph。這里先介紹2個(gè)概念

  • StreamNode:streaming流中的一個(gè)節(jié)點(diǎn),代表對應(yīng)的算子
  • StreamEdge:Graph中的邊,來連接上下游的StreamNode

    如上圖所示,圓形為StreamNode,箭頭為StreamEdge,這樣通過這2者就可以構(gòu)建一個(gè)StreamGraph了。
    StreamGraph是最原始的Graph,而其中會做一些優(yōu)化生成JobGraph,最后會生成待執(zhí)行的ExecutionGraph,這里我們先介紹下基礎(chǔ)概念,后面會深入介紹相關(guān)的內(nèi)容。
  • JobGraph: 優(yōu)化后的StreamGraph,具體做的優(yōu)化就是把相連的算子,如果支持chaining的,合并到一個(gè)StreamNode;
  • ExecutionGraph: 和JobGraph結(jié)構(gòu)一致

StreamGraph

下面我們來看看StreamGraph的主要屬性和方法,以及如何從DataStream轉(zhuǎn)換為StreamGraph的。

屬性和方法

重要屬性如下(這里只介紹與生成圖相關(guān)的屬性,還有一些如狀態(tài),存儲類的后面介紹)

屬性說明
Map<Integer, StreamNode> streamNodesStreamNode數(shù)據(jù),kv格式,key為Transformation的id
Set sourcesStreamGraph的所有source集合,存儲的是Transformation的id
Set sinksStreamGraph的sink集合

說明:StreamGraph只記錄了StreamNode的信息,StreamEdge的信息是記錄在StreamNode中的。如下2個(gè)屬性記錄了StreamNode的輸入Edge和輸出Edge

    //StreamNode.java
    private List<StreamEdge> inEdges = new ArrayList<StreamEdge>();
    private List<StreamEdge> outEdges = new ArrayList<StreamEdge>();

主要方法

方法說明
addSource()添加source節(jié)點(diǎn)
addSink()添加sink節(jié)點(diǎn)
addOperator()添加算子節(jié)點(diǎn)
addVirtualSideOutputNode()添加一個(gè)虛擬的siteOutput節(jié)點(diǎn)

StreamGraph生成

下面我們來看看DataStream是如何生成StreamGraph的。通過前面對DataStream的分析可知,DataStream的前后依賴關(guān)系是通過Transformation來存儲的,這里StreamExecutionEnvironment有個(gè)transformations記錄了所有的Transformation

    //StreamExecutionEnvironment.java
    List<Transformation<?>> transformations

這里的數(shù)據(jù)是在DataStream進(jìn)行轉(zhuǎn)換處理生成了新的Transformation,同時(shí)會把該實(shí)例添加到transformations里面,使用的是如下方法

    getExecutionEnvironment().addOperator(resultTransform);

而具體轉(zhuǎn)換為StreamGraph是通過StreamExecutionEnvironment的 getStreamGraph()方法。最終轉(zhuǎn)換的邏輯是通過StreamGraphGenerator類來實(shí)現(xiàn)。

這里要介紹一個(gè)新的類體系TransformationTranslator,有各種的子類來轉(zhuǎn)換對應(yīng)類型的Transformation。這里有定義了2個(gè)方法分別支持轉(zhuǎn)換Streaming和Batch。

 //TransformationTranslator.java
 Collection<Integer> translateForBatch(final T transformation, final Context context);
 Collection<Integer> translateForStreaming(final T transformation, final Context context);

對應(yīng)的映射關(guān)系存儲在StreamGraphGenerator類的translatorMap中。

//StreamGraphGenerator.java 
private static final Map<
                    Class<? extends Transformation>,
                    TransformationTranslator<?, ? extends Transformation>>
            translatorMap;
static {
        @SuppressWarnings("rawtypes")
        Map<Class<? extends Transformation>, TransformationTranslator<?, ? extends Transformation>>
                tmp = new HashMap<>();
        tmp.put(OneInputTransformation.class, new OneInputTransformationTranslator<>());
        tmp.put(TwoInputTransformation.class, new TwoInputTransformationTranslator<>());
        tmp.put(MultipleInputTransformation.class, new MultiInputTransformationTranslator<>());
        tmp.put(KeyedMultipleInputTransformation.class, new MultiInputTransformationTranslator<>());
        tmp.put(SourceTransformation.class, new SourceTransformationTranslator<>());
        tmp.put(SinkTransformation.class, new SinkTransformationTranslator<>());
    ...

下面我們通過OneInputTransformationTranslator為例來看看是如何進(jìn)行轉(zhuǎn)換的。具體邏輯如下

 //調(diào)用addOperator添加StreamNode
 streamGraph.addOperator(
                transformationId,
                slotSharingGroup,
                transformation.getCoLocationGroupKey(),
                operatorFactory,
                inputType,
                transformation.getOutputType(),
                transformation.getName());
//獲取上游依賴的transformations,然后添加邊
        final List<Transformation<?>> parentTransformations = transformation.getInputs();
        for (Integer inputId : context.getStreamNodeIds(parentTransformations.get(0))) {
            streamGraph.addEdge(inputId, transformationId, 0);
        }

除了添加節(jié)點(diǎn)和邊外,還有一些如設(shè)置節(jié)點(diǎn)的并行度等操作,這塊大家可以去看看具體的代碼。
這樣當(dāng)把所有的Transformation都轉(zhuǎn)換完,這樣StreamGraph就生成好了。

JobGraph

有了StreamGraph,為什么還需要一個(gè)JobGraph呢,這個(gè)和Spark中的Stage類似,如果有多個(gè)算子能夠合并到一起處理,那這樣性能可以提高很多。所以這里 根據(jù)一定的規(guī)則進(jìn)行,先我們介紹相關(guān)的類

  • JobVertex:job的頂點(diǎn),即對應(yīng)的計(jì)算邏輯(這里用的是Vertex, 而前面用的是Node,有點(diǎn)差異),通過inputs記錄了所有來源的Edge,而輸出是ArrayList來記錄
  • JobEdge: job的邊,記錄了源Vertex和咪表Vertex.
  • IntermediateDataSet: 定義了一個(gè)中間數(shù)據(jù)集,但并沒有存儲,只是記錄了一個(gè)Producer(JobVertex)和一個(gè)Consumer(JobEdge)
    主要的概念就這些,下面我們看看JobGraph的結(jié)構(gòu)以及如何從StreamGraph轉(zhuǎn)換為JobGraph

屬性和方法

JobGraph的屬性主要是通過Map<JobVertexID, JobVertex> taskVertices記錄了JobVertex的信息。
另外這個(gè)JobGraph是提交到集群去執(zhí)行的,所以會有一些執(zhí)行相關(guān)的信息,相關(guān)的如下:

    private JobID jobID;
    private final String jobName;
    private SerializedValue<ExecutionConfig> serializedExecutionConfig;
    /** Set of JAR files required to run this job. */
    private final List<Path> userJars = new ArrayList<Path>();
     /** Set of custom files required to run this job. */
    private final Map<String, DistributedCache.DistributedCacheEntry> userArtifacts =
            new HashMap<>();
    /** Set of blob keys identifying the JAR files required to run this job. */
    private final List<PermanentBlobKey> userJarBlobKeys = new ArrayList<>();
    /** List of classpaths required to run this job. */
    private List<URL> classpaths = Collections.emptyList();

而相關(guān)的方法主要是

方法說明
addVertex()添加頂點(diǎn)
getVertices()獲取頂點(diǎn)

而如何從StreamGraph轉(zhuǎn)換到JobGraph這塊的內(nèi)容還是比較多,這塊后續(xù)我們單獨(dú)開一篇來介紹

總結(jié)

本篇從0開始介紹了DataStream的相關(guān)內(nèi)容,并深入介紹了DataStream、Transformation、StreamOperator和Function之間的關(guān)系。另外介紹了streaming flow轉(zhuǎn)換為提交執(zhí)行的StreamGraph的過程及StreamGraph的存儲結(jié)構(gòu)。而從StreamGraph->JobGraph->ExecutionGraph這塊涉及的內(nèi)容也較多,且還涉及到提交部署的內(nèi)容,這塊后面單獨(dú)來介紹。最后本篇介紹的DataStream只是介紹了最基礎(chǔ)的計(jì)算框架,沒有涉及到flink的streaming flow中的時(shí)間、狀態(tài)、window等內(nèi)容,更多關(guān)于Flink DataStream基礎(chǔ)的資料請關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • JDK動(dòng)態(tài)代理之WeakCache緩存的實(shí)現(xiàn)機(jī)制

    JDK動(dòng)態(tài)代理之WeakCache緩存的實(shí)現(xiàn)機(jī)制

    這篇文章主要介紹了JDK動(dòng)態(tài)代理之WeakCache緩存的實(shí)現(xiàn)機(jī)制
    2018-02-02
  • spark之Standalone模式部署配置詳解

    spark之Standalone模式部署配置詳解

    這篇文章主要介紹了spark之Standalone模式部署配置詳解,小編覺得挺不錯(cuò)的,這里分享給大家,供各位參考。
    2017-10-10
  • Java實(shí)戰(zhàn)寵物醫(yī)院預(yù)約掛號系統(tǒng)的實(shí)現(xiàn)流程

    Java實(shí)戰(zhàn)寵物醫(yī)院預(yù)約掛號系統(tǒng)的實(shí)現(xiàn)流程

    只學(xué)書上的理論是遠(yuǎn)遠(yuǎn)不夠的,只有在實(shí)戰(zhàn)中才能獲得能力的提升,本篇文章手把手帶你用java+JSP+Spring+SpringBoot+MyBatis+html+layui+maven+Mysql實(shí)現(xiàn)一個(gè)寵物醫(yī)院預(yù)約掛號系統(tǒng),大家可以在過程中查缺補(bǔ)漏,提升水平
    2022-01-01
  • eclipse 如何創(chuàng)建 user library 方法詳解

    eclipse 如何創(chuàng)建 user library 方法詳解

    這篇文章主要介紹了eclipse 如何創(chuàng)建 user library 方法詳解的相關(guān)資料,需要的朋友可以參考下
    2017-04-04
  • Java8并行流中自定義線程池操作示例

    Java8并行流中自定義線程池操作示例

    這篇文章主要介紹了Java8并行流中自定義線程池操作,結(jié)合實(shí)例形式分析了并行流的相關(guān)概念、定義及自定義線程池的相關(guān)操作技巧,需要的朋友可以參考下
    2019-05-05
  • Spring中SpEL表達(dá)式的使用全解

    Spring中SpEL表達(dá)式的使用全解

    SpEL是Spring框架中用于表達(dá)式語言的一種方式,本文主要介紹了Spring中SpEL表達(dá)式的使用全解,文中通過示例代碼介紹的非常詳細(xì),需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2024-04-04
  • SpringBoot集成Swagger2的方法

    SpringBoot集成Swagger2的方法

    這篇文章主要介紹了SpringBoot集成Swagger2的方法,本文給大家介紹的非常詳細(xì),對大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2020-12-12
  • 使用abstract格式修飾抽象方法

    使用abstract格式修飾抽象方法

    abstract是抽象的意思,用于修飾方法方法和類,修飾的方法是抽象方法,修飾的類是抽象類,這篇文章主要介紹了怎樣使用abstract格式修飾抽象方法,需要的朋友可以參考下
    2023-05-05
  • 解析Spring事件發(fā)布與監(jiān)聽機(jī)制

    解析Spring事件發(fā)布與監(jiān)聽機(jī)制

    本篇文章給大家介紹Spring事件發(fā)布與監(jiān)聽機(jī)制,通過 ApplicationEvent 事件類和 ApplicationListener 監(jiān)聽器接口,可以實(shí)現(xiàn) ApplicationContext 事件發(fā)布與處理,需要的朋友參考下吧
    2021-06-06
  • Java實(shí)現(xiàn)HashMap排序方法的示例詳解

    Java實(shí)現(xiàn)HashMap排序方法的示例詳解

    這篇文章主要通過一些示例為大家介紹了Java對HashMap進(jìn)行排序的方法,幫助大家更好的理解和使用Java,感興趣的朋友可以了解一下
    2022-05-05

最新評論

永嘉县| 河北省| 南丰县| 余姚市| 铜陵市| 集安市| 汝阳县| 民丰县| 邵阳县| 利津县| 陈巴尔虎旗| 林周县| 利川市| 墨江| 渝北区| 鱼台县| 哈密市| 临夏县| 锡林浩特市| 阳新县| 温泉县| 威海市| 潍坊市| 庆安县| 黄冈市| 湖南省| 天门市| 旌德县| 明星| 丹东市| 扶余县| 镇宁| 翁牛特旗| 上饶县| 湖口县| 如东县| 台中县| 介休市| 柳州市| 肇庆市| 新闻|