Flink的Job啟動TaskManager端(源碼分析)

来源:https://www.cnblogs.com/ljygz/archive/2019/09/03/11453579.html
-Advertisement-
Play Games

前面說到了 Flink的TaskManager啟動(源碼分析) 啟動了TaskManager 然後 Flink的Job啟動JobManager端(源碼分析) 說到JobManager會將轉化得到的TDD發送到TaskManager的RPC 這篇主要就講一下,Job在TaskManager端是如何啟動 ...


前面說到了  Flink的TaskManager啟動(源碼分析)  啟動了TaskManager

然後  Flink的Job啟動JobManager端(源碼分析)  說到JobManager會將轉化得到的TDD發送到TaskManager的RPC

這篇主要就講一下,Job在TaskManager端是如何啟動的

先來看一下,TaskManager端用來接收JobManager發送過來的TDD對象的RPC介面

在TaskExecutor.java中

 這個方法用於接收了一個TaskDeploymentDescriptor對象用於啟動任務(上一篇知道這裡executionGraph的每一個並行度都會調用deploy方法生成一個TDD)

來看一下具體接收到以後做了什麼

 

創建了一個Task並且將其內部的一個線程啟動起來了

註意這裡從TDD中得到了InputGate,Partition的信息,用於創建InputGate,ResultPartition

InputGate用於對接上游產生的數據(消費)

ResultPartition用於往下游發送自己產生的數據(生產)

來看一下Task創建,在Task的構造方法中

 

這裡看到創建了對應往下游發送數據的ResultPartition

ResultPartition中創建的SubPartition具體分為

可以看到就是說三個參數分別對應

  PIPELINED    可以邊消費邊生產,是有背壓的,這個partition沒有buffer數量的限制

      (因為背壓的控制是通過,接收數據端公用同一個指定大小的bufferPool,以後背壓的時候講)

  其他同理

這裡看一下不同類型的ResultPartitionType是創建的什麼subpartitions

BLOCKING  這種創建了一個SpillableSubpartition並且傳進去了一個ioManager(這個ioManager以後io管理細講)

大致看了一下就是說這種Subpartition是會落盤的

PIPELINED  而這種方式是完全基於記憶體的

根據上游的信息創建好ResultPartition以後

 接著創建了InputGate用於接收上游的數據,並且在create方法中

會根據partition的位置創建對應的channel,這裡可以分為

Local      就是說下游和自己是在同一臺機器

Remote  下游是需要通過網路發送的

 並且在這裡將inputGate和它所有的inputChannels關聯了起來

創建完inputGate以後Task就初始化完了,然後會被start()起來,來看下Task的run方法

在run方法中

 這個地方會為初始化inputGate與ResultPartition的bufferPool(以後講到反壓在講)

繼續

這裡通過反射創建了一個StreamTask的實例

並且

調用了他的invoke()方法,這裡也是Job開始的邏輯,來看一下invoke方法

在invoke方法中

只要知道這裡會初始化OperatorChain這裡包含了我們用戶運算元的邏輯(這裡不細講,隨緣講到Task操作責任鏈的時候講)

然後得到了operatorChain的頭headoperator其實這裡的頭就包含了用戶的第一個運算元邏輯在裡面

然後init()方法中用上面的headoperator初始化了一個inputProcess對象並且關聯上了上面創建的inputGate(也是留到責任鏈講)

接著

 

 這裡就是上面在init方法中創建的inputProcess,並且調用了他的processInput方法

 重頭戲來了,來看一下processInput方法

 

這裡有個while(true)也就是說這裡會一直迴圈下去

來看一下他迴圈做什麼

這裡!!!!這個streamOperator就是上面構造inputProcess時傳入的headOperator

這個processElement方法裡面就是調用用戶的方法啦

也就是不停的從上游接收到數據以後,調用用戶具體的處理邏輯

這裡job就啟動完成了

 

註意這個while迴圈內既然開始走我們用戶的邏輯,那肯定會先從inputGate關聯到的上游獲取數據

這裡就非常重要了,因為接收數據就包含了很多的機制的實現

包含了watermark處理的邏輯,水印對齊的邏輯,水印更新的邏輯,如下

 以及idle停滯流邏輯,流狀態更新邏輯

以及如何接收數據邏輯,接收端反壓的邏輯,barriers對齊的邏輯,checkpoint觸發的邏輯

 所以這個StreamInputProcessor.processInput()方法是一個非常重要的方法,以後隨緣更新各種機制的時候也會經常看到

 


您的分享是我們最大的動力!

-Advertisement-
Play Games
更多相關文章
  • 條件判斷: [ condition ],condition前後都有空格 常用的判斷條件: 1)兩個整數的比較 = 字元串比較 -lt 小於 -le 小於等於 -eq 等於 -gt 大於 -ge 大於等於 -ne 不等於 2)按照文件許可權進行判斷 -r有讀的許可權 -w有寫的許可權 -x有執行的許可權 3) ...
  • 一、數據挖掘 中文分詞 • 一段文字不僅僅在於字面上是什麼,還在於怎麼切分和理解。• 例如: – 阿三炒飯店: – 阿三 / 炒飯 / 店 阿三 / 炒 / 飯店• 和英文不同,中文詞之間沒有空格,所以實現中文搜索引擎,比英文多了一項分詞的任務。• 如果沒有中文分詞會出現: – 搜索“達內”,會出現 ...
  • 公司一SQL Server鏡像發生了故障轉移(主備切換),檢查SQL Server鏡像發生主備切換的原因,在錯誤日誌中發現下麵錯誤: Date 2019/8/31 14:09:17 Log SQL Server (Archive #4 - 2019/9/1 0:00:00) Source spid3... ...
  • redis是key-value的數據,所以每個數據都是一個鍵值對。 數據操作的全部命令,可以查看中文網站。 鍵的類型是字元串 值的類型分為五種: 字元串string 哈希hash 列表list 集合set 有序集合zset 字元串string 哈希hash 列表list 集合set 有序集合zset ...
  • 1.獨立模式(standalone|local) nothing! 本地文件系統。 不需要啟用單獨進程。 2.pesudo(偽分佈模式) 等同於完全分散式,只有一個節點。 SSH: //(Socket), //public + private /server : sshd ps -Af | grep ...
  • SELECT select的完整語法: 上述如果都有:執行順序from->where->group by->having->order by->limit->select 列的結果顯示 1、去掉重覆的數據:distinct(針對於記錄而言,不是針對於列的數據而言) 2、運算符:+、-、*、/、%(只 ...
  • 1.一個問題 InnoDB一棵B+樹可以存放多少行數據?這個問題的簡單回答是:約2千萬。為什麼是這麼多呢?因為這是可以算出來的,要搞清楚這個問題,我們先從InnoDB索引數據結構、數據組織方式說起。 我們都知道電腦在存儲數據的時候,有最小存儲單元,這就好比我們今天進行現金的流通最小單位是一毛。在計 ...
  • 定義 各類別的出現概率不均衡的情況 如信用風險中正常用戶遠多於逾期、違約用戶;流失風險中留存用戶多於流失用戶 隱患 降低對少類樣本的靈敏性。但我們建模就是要找到這少類樣本,所以必須對數據加以處理,來提高靈敏性。 解決方案 1. 過採樣 對壞的人群提高權重,即複製壞樣本,提高壞樣本的占比。 優點: 簡 ...
一周排行
    -Advertisement-
    Play Games
  • 移動開發(一):使用.NET MAUI開發第一個安卓APP 對於工作多年的C#程式員來說,近來想嘗試開發一款安卓APP,考慮了很久最終選擇使用.NET MAUI這個微軟官方的框架來嘗試體驗開發安卓APP,畢竟是使用Visual Studio開發工具,使用起來也比較的順手,結合微軟官方的教程進行了安卓 ...
  • 前言 QuestPDF 是一個開源 .NET 庫,用於生成 PDF 文檔。使用了C# Fluent API方式可簡化開發、減少錯誤並提高工作效率。利用它可以輕鬆生成 PDF 報告、發票、導出文件等。 項目介紹 QuestPDF 是一個革命性的開源 .NET 庫,它徹底改變了我們生成 PDF 文檔的方 ...
  • 項目地址 項目後端地址: https://github.com/ZyPLJ/ZYTteeHole 項目前端頁面地址: ZyPLJ/TreeHoleVue (github.com) https://github.com/ZyPLJ/TreeHoleVue 目前項目測試訪問地址: http://tree ...
  • 話不多說,直接開乾 一.下載 1.官方鏈接下載: https://www.microsoft.com/zh-cn/sql-server/sql-server-downloads 2.在下載目錄中找到下麵這個小的安裝包 SQL2022-SSEI-Dev.exe,運行開始下載SQL server; 二. ...
  • 前言 隨著物聯網(IoT)技術的迅猛發展,MQTT(消息隊列遙測傳輸)協議憑藉其輕量級和高效性,已成為眾多物聯網應用的首選通信標準。 MQTTnet 作為一個高性能的 .NET 開源庫,為 .NET 平臺上的 MQTT 客戶端與伺服器開發提供了強大的支持。 本文將全面介紹 MQTTnet 的核心功能 ...
  • Serilog支持多種接收器用於日誌存儲,增強器用於添加屬性,LogContext管理動態屬性,支持多種輸出格式包括純文本、JSON及ExpressionTemplate。還提供了自定義格式化選項,適用於不同需求。 ...
  • 目錄簡介獲取 HTML 文檔解析 HTML 文檔測試參考文章 簡介 動態內容網站使用 JavaScript 腳本動態檢索和渲染數據,爬取信息時需要模擬瀏覽器行為,否則獲取到的源碼基本是空的。 本文使用的爬取步驟如下: 使用 Selenium 獲取渲染後的 HTML 文檔 使用 HtmlAgility ...
  • 1.前言 什麼是熱更新 游戲或者軟體更新時,無需重新下載客戶端進行安裝,而是在應用程式啟動的情況下,在內部進行資源或者代碼更新 Unity目前常用熱更新解決方案 HybridCLR,Xlua,ILRuntime等 Unity目前常用資源管理解決方案 AssetBundles,Addressable, ...
  • 本文章主要是在C# ASP.NET Core Web API框架實現向手機發送驗證碼簡訊功能。這裡我選擇是一個互億無線簡訊驗證碼平臺,其實像阿裡雲,騰訊雲上面也可以。 首先我們先去 互億無線 https://www.ihuyi.com/api/sms.html 去註冊一個賬號 註冊完成賬號後,它會送 ...
  • 通過以下方式可以高效,並保證數據同步的可靠性 1.API設計 使用RESTful設計,確保API端點明確,並使用適當的HTTP方法(如POST用於創建,PUT用於更新)。 設計清晰的請求和響應模型,以確保客戶端能夠理解預期格式。 2.數據驗證 在伺服器端進行嚴格的數據驗證,確保接收到的數據符合預期格 ...