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

Flink部署集群整體架構(gòu)源碼分析

 更新時(shí)間:2022年12月01日 11:35:00   作者:xiangel  
這篇文章主要為大家介紹了Flink部署集群及整體架構(gòu)示例詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪

概覽

本篇我們來(lái)了解Flink的部署模式和Flink集群的整體架構(gòu)

部署模式

Flink支持如下三種運(yùn)行模式

運(yùn)行模式描述
Application ModeFlink Cluster只執(zhí)行提交的整個(gè)job,然后退出;main方法在cluster中執(zhí)行;支持yarn和k8s;官方建議yarn和k8s上的運(yùn)行方式
pre-job modeFlink Cluster只執(zhí)行提交的整個(gè)job,然后退出;main方法在client中執(zhí)行;支持yarn;官方建議yarn上運(yùn)行方式, 該模式在Flink 1.15中被廢棄了,建議用application mode
session mode支持在一個(gè)Flink Cluster中提交多個(gè)任務(wù);main方法在client中執(zhí)行;支持yarn和k8s

Flink的部署步驟分為如下2步:

  • 部署啟動(dòng)一個(gè)Flink Cluster,負(fù)責(zé)接收job提交請(qǐng)求和管理job信息;
  • 向Flink Cluster提交job; 根據(jù)Flink Cluster可以運(yùn)行的任務(wù)的數(shù)量(1個(gè)或多個(gè))和提交job請(qǐng)求的地點(diǎn)(遠(yuǎn)端或Cluster端)的不同,從而有了不同的運(yùn)行模式。由于pre-job模式已經(jīng)被廢棄了,下面我們主要來(lái)學(xué)習(xí)下Application mode和session mode

Application mode

Application mode是Flink Cluster運(yùn)行1個(gè)job,提交任務(wù)的地點(diǎn)為Cluster端。其提交方式如下

./bin/flink run-application -t yarn-application ./examples/streaming/TopSpeedWindowing.jar

其處理流程為,客戶端提交部署請(qǐng)求,服務(wù)端啟動(dòng)Flink Cluster, 服務(wù)端運(yùn)行Flink Application提交Job到Cluster。下面我們分析下具體實(shí)現(xiàn)細(xì)節(jié)。

客戶端提交請(qǐng)求

通過(guò)flink命令提交請(qǐng)求,其運(yùn)行的類為CliFrontend。為支持部署到不同的資源管理平臺(tái),所以有和對(duì)應(yīng)資源管理系統(tǒng)交互的類,具體如下:

  • CliFrontend:flink命令對(duì)應(yīng)的類,發(fā)起提交請(qǐng)求,后面session mode的提交Flink Application也是由該類負(fù)責(zé)
  • ClusterClientFactory:集群客戶端工廠類,負(fù)責(zé)生成不同資源管理平臺(tái)的客戶端
  • ClusterDescriptor:負(fù)責(zé)和對(duì)應(yīng)的資源管理平臺(tái)交互,申請(qǐng)資源和提交請(qǐng)求
  • ClusterEntrypoint:在資源管理平臺(tái)運(yùn)行的類,啟動(dòng)Flink Cluster。 針對(duì)不同資源管理平臺(tái)的對(duì)應(yīng)實(shí)現(xiàn)類如下:
接口類yarnkubernates
ClusterClientFactoryYarnClusterClientFactoryKubernetesClusterClientFactory
ClusterDescriptorYarnClusterDescriptorKubernetesClusterDescriptor
ClusterEntrypointYarnApplicationClusterEntryPointKubernetesApplicationClusterEntrypoint

服務(wù)端啟動(dòng)&提交Application

服務(wù)端啟動(dòng)對(duì)應(yīng)的ClusterEntrypoint,其中會(huì)啟動(dòng)一個(gè)REST Server來(lái)接受提交Flink Application,另外有個(gè)Dispatcher負(fù)責(zé)作業(yè)的調(diào)度,其他部分后面我們分析運(yùn)行流程時(shí)再展開介紹。作業(yè)的提交請(qǐng)求是在Dispatcher中的DispatcherBootstrap屬性實(shí)例化的時(shí)候觸發(fā)。 Flink Application運(yùn)行時(shí),是在StreamExecutionEnvironment.execute()方法來(lái)觸發(fā)實(shí)際提交,提交相關(guān)的調(diào)用鏈如下:

這幾個(gè)都是接口類,在Application模式下對(duì)應(yīng)的實(shí)現(xiàn)類如下

接口類實(shí)現(xiàn)類
PipelineExecutorServiceLoaderEmbeddedExecutorServiceLoader
PipelineExecutorFactoryEmbeddedExecutorFactory
PipelineExecutorEmbeddedExecutor

session mode

session mode是一個(gè)Flink Cluster可以來(lái)運(yùn)行多個(gè)Flink job。那這里的提交會(huì)分為2個(gè)步驟

// 提交啟動(dòng)session cluster
// yarn session
./bin/yarn-session.sh --detached
// kubernates session
./bin/kubernetes-session.sh -Dkubernetes.cluster-id=my-first-flink-cluster
// 提交job
./bin/flink run ./examples/streaming/TopSpeedWindowing.jar
  • 通過(guò)yarn-session.sh (或kubernates-session.sh) 來(lái)提交部署Flink Cluster,這塊和前面application mode類似,以yarn模式為例,底層也是調(diào)用了YarnClusterDescriptor來(lái)提交相應(yīng)的請(qǐng)求,提交到服務(wù)器的是YarnSessionClusterEntrypoint類。
  • 提交Job,這塊是在client端來(lái)單獨(dú)提交的,直接提交信息到服務(wù)器的REST Server,根據(jù)提交的目標(biāo)資源管理系統(tǒng)的不同,使用了不同的實(shí)現(xiàn)類
接口類實(shí)現(xiàn)類yarn實(shí)現(xiàn)類kubernates
PipelineExecutorServiceLoaderDefaultExecutorServiceLoaderDefaultExecutorServiceLoader
PipelineExecutorFactoryYarnSessionClusterExecutorFactoryYarnSessionClusterExecutorFactory
PipelineExecutorYarnSessionClusterExecutorKubernetesSessionClusterExecutor

Cluster架構(gòu)

Flink是一個(gè)Master/Worker的架構(gòu),Master節(jié)點(diǎn)負(fù)責(zé)整個(gè)任務(wù)的管理,Worker節(jié)點(diǎn)負(fù)責(zé)執(zhí)行對(duì)應(yīng)的任務(wù)。其整體結(jié)構(gòu)如下:

* JobManager: Master節(jié)點(diǎn)的統(tǒng)稱,目前版本沒有該類,其中有幾個(gè)重點(diǎn)的服務(wù),如上圖所示,目前的代碼中對(duì)應(yīng)的組合了這些服務(wù)的類為:

Dispatcher

ResourceManager

Component。

* Dispatcher: Job調(diào)度器,負(fù)責(zé)接收J(rèn)ob的提交,保存Job和管理JobMaster來(lái)執(zhí)行作業(yè)。前面我們提到的提交作業(yè)到Cluster,實(shí)際上是提交給了Dispatcher的。

* ResourceManager: 負(fù)責(zé)和不同的資源調(diào)度系統(tǒng)交互,管理資源申請(qǐng)。

* WebMonitorEndpoint: 負(fù)責(zé)web界面的Rest請(qǐng)求處理

* JobMaster: 負(fù)責(zé)運(yùn)行單個(gè)JobGraph,包括TaskManager的管理,任務(wù)的調(diào)度等。

* TaskManager: 負(fù)責(zé)任務(wù)的執(zhí)行,也沒有TaskManager的類,對(duì)應(yīng)的類為TaskExecutor,來(lái)執(zhí)行多個(gè)Task

說(shuō)明:JobManager可能是原來(lái)的JobMaster,具體通過(guò)Dispatcher.java的如下代碼可以看出,重點(diǎn)在對(duì)其具體結(jié)構(gòu)的理解,這個(gè)變化的邏輯我們就不考究了。

 private JobManagerRunner createJobMasterRunner(JobGraph jobGraph) throws Exception

Cluster的啟動(dòng)流程

上面介紹了Cluster的整體架構(gòu),接下來(lái)我們看看Cluster的啟動(dòng)流程。以Application mode部署到Y(jié)arn為例(其他模式的啟動(dòng)類似,只是啟動(dòng)的主類不同)。該方式下的主類為:YarnApplicationClusterEntryPoint,其內(nèi)部調(diào)用了ClusterEntrypoint的方法,最終是通過(guò)ClusterEntrypoint類的runCluster()方法來(lái)創(chuàng)建DispatcherResourceManagerComponent對(duì)象。

DispatcherResourceManagerComponent

接下來(lái)我們看看DispatcherResourceManagerComponent中的具體屬性信息

    @Nonnull private final DispatcherRunner dispatcherRunner;
    @Nonnull private final ResourceManagerService resourceManagerService;
    @Nonnull private final RestService webMonitorEndpoint;
    @Nonnull private final LeaderRetrievalService dispatcherLeaderRetrievalService;
    @Nonnull private final LeaderRetrievalService resourceManagerRetrievalService;

Runner代碼

這里我們并沒有看到Dispatcher,而是一個(gè)類似名字的DispatcherRunner.DispatcherRunner是來(lái)管理Dispatcher如何運(yùn)行的。類似ResourceManagerService是來(lái)管理ResourceManager的生命周期的。

HA代碼框架

另外由于這些服務(wù)都有雙機(jī)容錯(cuò)機(jī)制(HA), 所以這里在看相關(guān)代碼的時(shí)候會(huì)產(chǎn)生一定的干擾,本篇的最后我們來(lái)介紹下這塊HA的相關(guān)的機(jī)制,這樣對(duì)大家來(lái)梳理相關(guān)的流程會(huì)更清晰。 Leader的選舉,是通過(guò)LeaderElectionService(選舉服務(wù),實(shí)現(xiàn)類為DefaultLeaderElectionService)和LeaderContender(競(jìng)選者)共同來(lái)完成的。具體過(guò)程為L(zhǎng)eaderElectionService.start(LeaderContender),啟動(dòng)選舉服務(wù),傳入LeaderContender信息,等選舉成功后,會(huì)回調(diào)LeaderContender的grantLeadership()方法,F(xiàn)link中相關(guān)的服務(wù)都實(shí)現(xiàn)了LeaderContender接口。所以理清這個(gè)邏輯后,我們?cè)诳吹较嚓P(guān)服務(wù)的start()方法中只調(diào)用了leaderElectionService.start方法時(shí)也不用懵了,直接看該服務(wù)的grantLeadership方法來(lái)梳理邏輯。 LeaderElectionDriver:進(jìn)行Leader的選舉和保存Leader的信息,具體的實(shí)現(xiàn)有ZooKeeperLeaderElectionDriver和KubernetesLeaderElectionDriver

那如何獲取Leader的地址呢,也提供了相應(yīng)的接口LeaderRetrievalService和LeaderRetrievalLister,啟動(dòng)一個(gè)對(duì)Leader地址的監(jiān)聽,leader有變化時(shí)會(huì)得到通知。

總結(jié)

本篇我們了解了Flink的部署模式,按Job提交方式和一個(gè)集群可同時(shí)運(yùn)行任務(wù)的數(shù)量的不同,分為ApplicationMode和SessionMode2種模式。接著介紹了Cluster的整體架構(gòu)和啟動(dòng)流程,主要包括Dispatcher、ResourceManager和WebMonitorEndpoint。最后介紹了HA處理的整體框架,便于大家更好的梳理核心流程。

以上就是Flink部署集群整體架構(gòu)源碼分析的詳細(xì)內(nèi)容,更多關(guān)于Flink部署集群架構(gòu)的資料請(qǐng)關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • SpringMVC通過(guò)攔截器實(shí)現(xiàn)IP黑名單

    SpringMVC通過(guò)攔截器實(shí)現(xiàn)IP黑名單

    這篇文章主要為大家詳細(xì)介紹了SpringMVC通過(guò)攔截器實(shí)現(xiàn)IP黑名單,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2019-08-08
  • Spring Boot 2.0.0 終于正式發(fā)布-重大修訂版本

    Spring Boot 2.0.0 終于正式發(fā)布-重大修訂版本

    北京時(shí)間 2018 年 3 月 1 日早上,如約發(fā)布的 Spring Boot 2.0 在同步至 Maven 倉(cāng)庫(kù)時(shí)出現(xiàn)問(wèn)題,導(dǎo)致在 GitHub 上發(fā)布的 v2.0.0.RELEASE 被撤回
    2018-03-03
  • Java實(shí)現(xiàn)FTP上傳與下載功能

    Java實(shí)現(xiàn)FTP上傳與下載功能

    這篇文章主要為大家詳細(xì)介紹了Java實(shí)現(xiàn)FTP上傳與下載功能,文中示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2022-02-02
  • SpringBoot-JWT生成Token和攔截器的使用(訪問(wèn)受限資源)

    SpringBoot-JWT生成Token和攔截器的使用(訪問(wèn)受限資源)

    本文主要介紹了SpringBoot-JWT生成Token和攔截器的使用(訪問(wèn)受限資源),文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2022-05-05
  • Spring注解之Service用法及示例詳解

    Spring注解之Service用法及示例詳解

    使用 @Service 注解可以將一個(gè)類聲明為業(yè)務(wù)邏輯組件,并將其對(duì)象存入 Spring 容器中,在控制器類中,通過(guò)注入該組件的實(shí)例,即可調(diào)用其中的方法,這篇文章主要介紹了Spring注解之Service用法及示例詳解,需要的朋友可以參考下
    2024-04-04
  • Java并發(fā)(Runnable+Thread)實(shí)現(xiàn)硬盤文件搜索功能

    Java并發(fā)(Runnable+Thread)實(shí)現(xiàn)硬盤文件搜索功能

    這篇文章主要介紹了Java并發(fā)(Runnable+Thread)實(shí)現(xiàn)硬盤文件搜索,本文給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2021-01-01
  • Spring Boot快速搭建Spring框架教程

    Spring Boot快速搭建Spring框架教程

    這篇文章主要為大家詳細(xì)介紹了Spring Boot快速搭建Spring框架教程,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2017-10-10
  • SpringBoot整合RedisTemplate實(shí)現(xiàn)緩存信息監(jiān)控的基本操作

    SpringBoot整合RedisTemplate實(shí)現(xiàn)緩存信息監(jiān)控的基本操作

    SpringBoot中的 redistemplate 是一個(gè)用于操作 Redis 數(shù)據(jù)庫(kù)的高級(jí)模板類,它提供了一組方法,可以方便地執(zhí)行常見的 Redis 操作,如存儲(chǔ)、檢索和刪除數(shù)據(jù),本文給大家介紹了SpringBoot整合RedisTemplate實(shí)現(xiàn)緩存信息監(jiān)控的基本操作,需要的朋友可以參考下
    2025-02-02
  • SpringBoot項(xiàng)目中出現(xiàn)不同端口跨域問(wèn)題的解決方法

    SpringBoot項(xiàng)目中出現(xiàn)不同端口跨域問(wèn)題的解決方法

    這篇文章主要介紹了SpringBoot項(xiàng)目中出現(xiàn)不同端口跨域問(wèn)題的解決方法,文中介紹了兩種解決方法,并給出了詳細(xì)的代碼供大家參考,具有一定的參考價(jià)值,需要的朋友可以參考下
    2024-03-03
  • 基于java swing實(shí)現(xiàn)答題系統(tǒng)

    基于java swing實(shí)現(xiàn)答題系統(tǒng)

    這篇文章主要為大家詳細(xì)介紹了基于java swing實(shí)現(xiàn)答題系統(tǒng),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2018-01-01

最新評(píng)論

如东县| 项城市| 瓮安县| 巴中市| 尼勒克县| 梁平县| 乌兰浩特市| 岗巴县| 东乌珠穆沁旗| 墨竹工卡县| 高清| 固原市| 九龙坡区| 库车县| 平山县| 宜春市| 河南省| 宁南县| 育儿| 加查县| 阜平县| 华阴市| 安龙县| 独山县| 玉林市| 雅安市| 奉贤区| 进贤县| 卢龙县| 喀喇| 保亭| 蓬莱市| 澎湖县| 名山县| 阿巴嘎旗| 延庆县| 洛南县| 蓬安县| 金阳县| 微博| 靖安县|