LocalFlinkMiniCluster啟動DataStream任務(wù)的流程

LocalFlinkMiniCluster 集群的actor 模型


  • 相關(guān)的主要類圖如下:
image-20190415180352502.png
  • AkkaRpcActor持有一個rpcEndpoint成員,接收到消息后進行基礎(chǔ)解析后調(diào)用rpcEndpoint的的對應(yīng)方法來進行處理。

  • 其中RpcGateway及RpcEndPoint的類圖


    image-20190415175424225.png
  • 支持的消息類型

    • 其中主要使用RpcInvocation基于反射調(diào)用RPCEndpoint的對應(yīng)函數(shù)
    • FencedMessage 將message進行封裝成payload,通過fencingToken進行校驗,保證請求的合法性


      image-20190415175913490.png

LocalFlinkMiniCluster集群的角色


  • ResouceManager

    • 負責(zé)容器的分配
    • 使用FencedAkkaRpcActor實現(xiàn),其rpcEndpoint為 org.apache.flink.runtime.resourcemanager.ResourceManager
  • JobMaster

    • 負責(zé)任務(wù)執(zhí)行計劃的調(diào)度和執(zhí)行,

    • 使用FencedAkkaRpcActor實現(xiàn),其rpcEndpoint為 org.apache.flink.runtime.jobmaster.JobMaster

      • JobMaster持有一個SlotPool的Actor,用來暫存TaskExecutor提供給JobMaster并被接受的slot。JobMaster的Scheduler組件從這個SlotPool中獲取資源以調(diào)度job的task
  • Dispatcher

    • 主要職責(zé)是接收從Client端提交過來的job并生成一個JobMaster去負責(zé)這個job在集群資源管理器上執(zhí)行。

      • 不是所有部署方式都需要用到dispatcher,比如yarn-cluster 的部署方式可能就不需要
    • 使用FencedAkkaRpcActor實現(xiàn),其rpcEndpoint為 org.apache.flink.runtime.dispatcher.StandaloneDispatcher

  • TaskExecutor

    • TaskExecutor會與ResouceManager和 JobMaster兩者進行通信。

      • 會向ResourceManager報告自身的可用資源;并維護本身slot的狀態(tài)
      • 根據(jù)slot的分配結(jié)果,接收JobMaster的命令在對應(yīng)的slot上執(zhí)行指定的task。
      • TaskExecutor還需要向以上兩者定時上報心跳信息。
    • 使用AkkaRpcActor實現(xiàn),其rpcEndpoint為org.apache.flink.runtime.taskexecutor.TaskExecutor

啟動DataStream任務(wù)的主體流程


image-20190417172051347.png
image-20190417174333612.png

參考資料


最后編輯于
?著作權(quán)歸作者所有,轉(zhuǎn)載或內(nèi)容合作請聯(lián)系作者
【社區(qū)內(nèi)容提示】社區(qū)部分內(nèi)容疑似由AI輔助生成,瀏覽時請結(jié)合常識與多方信息審慎甄別。
平臺聲明:文章內(nèi)容(如有圖片或視頻亦包括在內(nèi))由作者上傳并發(fā)布,文章內(nèi)容僅代表作者本人觀點,簡書系信息發(fā)布平臺,僅提供信息存儲服務(wù)。

相關(guān)閱讀更多精彩內(nèi)容

友情鏈接更多精彩內(nèi)容