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

Flink DataStream基礎框架源碼分析

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

引言

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

計算模型

  • DataStream基礎框架

部署&調(diào)度

存儲體系

底層支撐

概覽

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

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

通過上圖可以來構想下,一般一個DataStream具有如下主要屬性

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

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

深入DataStream

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

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

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

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

DataStream

屬性和方法

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

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

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

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

類體系

DataStream的類圖關系比較簡單,就如下這幾個類,具體每個子類的信息見下表

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

Transformation

屬性和方法

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

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

Transformation的方法基本上是對上述屬性的get和set操作,這里重點要說明一下的是PhysicalTransformation(下面類體系來介紹)中的setChainingStrategy方法,這里的ChainingStrategy是一個枚舉類,主要是控制多個連續(xù)的算子是否可以進行鏈式處理,這個具體的我們在下面介紹StreamGraph時再介紹

類體系

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

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

StreamOperator

屬性和方法

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

下面我們重點介紹下相關的方法

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

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

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

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

類體系

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

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

Function

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

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

      text
        .flatMap(new Tokenizer)

處理后對應的DataStream的結構如下圖

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

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

  • StreamNode:streaming流中的一個節(jié)點,代表對應的算子
  • StreamEdge:Graph中的邊,來連接上下游的StreamNode

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

StreamGraph

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

屬性和方法

重要屬性如下(這里只介紹與生成圖相關的屬性,還有一些如狀態(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個屬性記錄了StreamNode的輸入Edge和輸出Edge

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

主要方法

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

StreamGraph生成

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

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

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

    getExecutionEnvironment().addOperator(resultTransform);

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

這里要介紹一個新的類體系TransformationTranslator,有各種的子類來轉換對應類型的Transformation。這里有定義了2個方法分別支持轉換Streaming和Batch。

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

對應的映射關系存儲在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為例來看看是如何進行轉換的。具體邏輯如下

 //調(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é)點和邊外,還有一些如設置節(jié)點的并行度等操作,這塊大家可以去看看具體的代碼。
這樣當把所有的Transformation都轉換完,這樣StreamGraph就生成好了。

JobGraph

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

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

屬性和方法

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

    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();

而相關的方法主要是

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

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

總結

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

相關文章

  • JDK動態(tài)代理之WeakCache緩存的實現(xiàn)機制

    JDK動態(tài)代理之WeakCache緩存的實現(xiàn)機制

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

    spark之Standalone模式部署配置詳解

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

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

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

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

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

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

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

    Spring中SpEL表達式的使用全解

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

    SpringBoot集成Swagger2的方法

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

    使用abstract格式修飾抽象方法

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

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

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

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

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

最新評論

获嘉县| 萨嘎县| 茶陵县| 额济纳旗| 昌江| 许昌县| 同心县| 勃利县| 江陵县| 阳高县| 沂水县| 开封市| 石嘴山市| 奉贤区| 小金县| 明星| 沭阳县| 长乐市| 嘉荫县| 五河县| 双辽市| 读书| 临桂县| 贵溪市| 江安县| 安丘市| 隆安县| 淳安县| 海林市| 阳泉市| 炎陵县| 杂多县| 屏南县| 金坛市| 东丽区| 潢川县| 廊坊市| 青海省| 蒙自县| 溧水县| 新干县|