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

Apache Flink的網(wǎng)絡(luò)協(xié)議棧詳細(xì)介紹

  發(fā)布時間:2019-06-28 17:26:24   作者:佚名   我要評論
Flink 的網(wǎng)絡(luò)協(xié)議棧是組成 flink-runtime 模塊的核心組件之一,本文中介紹了Apache Flink網(wǎng)絡(luò)協(xié)議棧,感興趣的朋友可以閱讀本文參考一下

▼ 造成反壓(2)

與沒有流量控制的接收端反壓機制不同,Credit 提供了更直接的控制:如果接收端的處理速度跟不上,最終它的 Credit 會減少成 0,此時發(fā)送端就不會在向網(wǎng)絡(luò)中發(fā)送數(shù)據(jù)(數(shù)據(jù)會被序列化到 Buffer 中并緩存在發(fā)送端)。由于反壓只發(fā)生在邏輯鏈路上,因此沒必要阻斷從多路復(fù)用的 TCP 連接中讀取數(shù)據(jù),也就不會影響其他的接收者接收和處理數(shù)據(jù)。

▼ Credit-based 的優(yōu)勢與問題

由于通過 Credit-based 流控機制,多路復(fù)用中的一個信道不會由于反壓阻塞其他邏輯信道,因此整體資源利用率會增加。此外,通過完全控制正在發(fā)送的數(shù)據(jù)量,我們還能夠加快 Checkpoint alignment:如果沒有流量控制,通道需要一段時間才能填滿網(wǎng)絡(luò)協(xié)議棧的內(nèi)部緩沖區(qū)并表明接收端不再讀取數(shù)據(jù)了。在這段時間里,大量的 Buffer 不會被處理。任何 Checkpoint barrier(觸發(fā) Checkpoint 的消息)都必須在這些數(shù)據(jù) Buffer 后排隊,因此必須等到所有這些數(shù)據(jù)都被處理后才能夠觸發(fā) Checkpoint(“Barrier 不會在數(shù)據(jù)之前被處理!”)。

但是,來自接收方的附加通告消息(向發(fā)送端通知 Credit)可能會產(chǎn)生一些額外的開銷,尤其是在使用 SSL 加密信道的場景中。此外,單個輸入通道( Input channel)不能使用緩沖池中的所有 Buffer,因為存在無法共享的 Exclusive buffer。新的流控協(xié)議也有可能無法做到立即發(fā)送盡可能多的數(shù)據(jù)(如果生成數(shù)據(jù)的速度快于接收端反饋 Credit 的速度),這時則可能增長發(fā)送數(shù)據(jù)的時間。雖然這可能會影響作業(yè)的性能,但由于其所有優(yōu)點,通常新的流量控制會表現(xiàn)得更好??赡軙ㄟ^增加單個通道的獨占 Buffer 數(shù)量,這會增大內(nèi)存開銷。然而,與先前實現(xiàn)相比,總體內(nèi)存使用可能仍然會降低,因為底層的網(wǎng)絡(luò)協(xié)議棧不再需要緩存大量數(shù)據(jù),因為我們總是可以立即將其傳輸?shù)?Flink(一定會有相應(yīng)的 Buffer 接收數(shù)據(jù))。

在使用新的 Credit-based 流量控制時,可能還會注意到另一件事:由于我們在發(fā)送方和接收方之間緩沖較少的數(shù)據(jù),反壓可能會更早的到來。然而,這是我們所期望的,因為緩存更多數(shù)據(jù)并沒有真正獲得任何好處。如果要緩存更多的數(shù)據(jù)并且保留 Credit-based 流量控制,可以考慮通過增加單個輸入共享 Buffer 的數(shù)量。

注意:如果需要關(guān)閉 Credit-based 流量控制,可以將這個配置添加到 flink-conf.yaml 中:taskmanager.network.credit-model:false。但是,此參數(shù)已過時,最終將與非 Credit-based 流控制代碼一起刪除。

4.序列號與反序列化

下圖從上面的擴展了更高級別的視圖,其中包含網(wǎng)絡(luò)協(xié)議棧及其周圍組件的更多詳細(xì)信息,從發(fā)送算子發(fā)送記錄(Record)到接收算子獲取它:

在生成 Record 并將其傳遞出去之后,例如通過 Collector#collect(),它被傳遞給 RecordWriter,RecordWriter 會將 Java 對象序列化為字節(jié)序列,最終存儲在 Buffer 中按照上面所描述的在網(wǎng)絡(luò)協(xié)議棧中進(jìn)行處理。RecordWriter 首先使用 SpanningRecordSerializer 將 Record 序列化為靈活的堆上字節(jié)數(shù)組。然后,它嘗試將這些字節(jié)寫入目標(biāo)網(wǎng)絡(luò) Channel 的 Buffer 中。我們將在下面的章節(jié)回到這一部分。

在接收方,底層網(wǎng)絡(luò)協(xié)議棧(Netty)將接收到的 Buffer 寫入相應(yīng)的輸入通道(Channel)。流任務(wù)的線程最終從這些隊列中讀取并嘗試在 RecordReader 的幫助下通過 SpillingAdaptiveSpanningRecordDeserializer 將累積的字節(jié)反序列化為 Java 對象。與序列化器類似,這個反序列化器還必須處理特殊情況,例如跨越多個網(wǎng)絡(luò) Buffer 的 Record,或者因為記錄本身比網(wǎng)絡(luò)緩沖區(qū)大(默認(rèn)情況下為32KB,通過 taskmanager.memory.segment-size 設(shè)置)或者因為序列化 Record 時,目標(biāo) Buffer 中已經(jīng)沒有足夠的剩余空間保存序列化后的字節(jié)數(shù)據(jù),在這種情況下,F(xiàn)link 將使用這些字節(jié)空間并繼續(xù)將其余字節(jié)寫入新的網(wǎng)絡(luò) Buffer 中。

4.1 將網(wǎng)絡(luò) Buffer 寫入 Netty

在上圖中,Credit-based 流控制機制實際上位于“Netty Server”(和“Netty Client”)組件內(nèi)部,RecordWriter 寫入的 Buffer 始終以空狀態(tài)(無數(shù)據(jù))添加到 Subpartition 中,然后逐漸向其中填寫序列化后的記錄。但是 Netty 在什么時候真正的獲取并發(fā)送這些 Buffer 呢?顯然,不能是 Buffer 中只要有數(shù)據(jù)就發(fā)送,因為跨線程(寫線程與發(fā)送線程)的數(shù)據(jù)交換與同步會造成大量的額外開銷,并且會造成緩存本身失去意義(如果是這樣的話,不如直接將將序列化后的字節(jié)發(fā)到網(wǎng)絡(luò)上而不必引入中間的 Buffer)。

在 Flink 中,有三種情況可以使 Netty 服務(wù)端使用(發(fā)送)網(wǎng)絡(luò) Buffer:

寫入 Record 時 Buffer 變滿,或者 Buffer 超時未被發(fā)送,或 發(fā)送特殊消息,例如 Checkpoint barrier。

▼ 在 Buffer 滿后發(fā)送

RecordWriter 將 Record 序列化到本地的序列化緩沖區(qū)中,并將這些序列化后的字節(jié)逐漸寫入位于相應(yīng) Result subpartition 隊列中的一個或多個網(wǎng)絡(luò) Buffer中。雖然單個 RecordWriter 可以處理多個 Subpartition,但每個 Subpartition 只會有一個 RecordWriter 向其寫入數(shù)據(jù)。另一方面,Netty 服務(wù)端線程會從多個 Result subpartition 中讀取并像上面所說的那樣將數(shù)據(jù)寫入適當(dāng)?shù)亩嗦窂?fù)用信道。這是一個典型的生產(chǎn)者 - 消費者模式,網(wǎng)絡(luò)緩沖區(qū)位于生產(chǎn)者與消費者之間,如下圖所示。在(1)序列化和(2)將數(shù)據(jù)寫入 Buffer 之后,RecordWriter 會相應(yīng)地更新緩沖區(qū)的寫入索引。一旦 Buffer 完全填滿,RecordWriter 會(3)為當(dāng)前 Record 剩余的字節(jié)或者下一個 Record 從其本地緩沖池中獲取新的 Buffer,并將新的 Buffer 添加到相應(yīng) Subpartition 的隊列中。這將(4)通知 Netty服務(wù)端線程有新的數(shù)據(jù)可發(fā)送(如果 Netty 還不知道有可用的數(shù)據(jù)的話4)。每當(dāng) Netty 有能力處理這些通知時,它將(5)從隊列中獲取可用 Buffer 并通過適當(dāng)?shù)?TCP 通道發(fā)送它。

注釋4:如果隊列中有更多已完成的 Buffer,我們可以假設(shè) Netty 已經(jīng)收到通知。

▼ 在 Buffer 超時后發(fā)送

為了支持低延遲應(yīng)用,我們不能只等到 Buffer 滿了才向下游發(fā)送數(shù)據(jù)。因為可能存在這種情況,某種通信信道沒有太多數(shù)據(jù),等到 Buffer 滿了在發(fā)送會不必要地增加這些少量 Record 的處理延遲。因此,F(xiàn)link 提供了一個定期 Flush 線程(the output flusher)每隔一段時間會將任何緩存的數(shù)據(jù)全部寫出??梢酝ㄟ^ StreamExecutionEnvironment#setBufferTimeout 配置 Flush 的間隔,并作為延遲5的上限(對于低吞吐量通道)。下圖顯示了它與其他組件的交互方式:RecordWriter 如前所述序列化數(shù)據(jù)并寫入網(wǎng)絡(luò) Buffer,但同時,如果 Netty 還不知道有數(shù)據(jù)可以發(fā)送,Output flusher 會(3,4)通知 Netty 服務(wù)端線程數(shù)據(jù)可讀(類似與上面的“buffer已滿”的場景)。當(dāng) Netty 處理此通知(5)時,它將消費(獲取并發(fā)送)Buffer 中的可用數(shù)據(jù)并更新 Buffer 的讀取索引。Buffer 會保留在隊列中——從 Netty 服務(wù)端對此 Buffer 的任何進(jìn)一步操作將在下次從讀取索引繼續(xù)讀取。

注釋5:嚴(yán)格來說,Output flusher 不提供任何保證——它只向 Netty 發(fā)送通知,而 Netty 線程會按照能力與意愿進(jìn)行處理。這也意味著如果存在反壓,則 Output flusher 是無效的。

▼ 特殊消息后發(fā)送

一些特殊的消息如果通過 RecordWriter 發(fā)送,也會觸發(fā)立即 Flush 緩存的數(shù)據(jù)。其中最重要的消息包括 Checkpoint barrier 以及 end-of-partition 事件,這些事件應(yīng)該盡快被發(fā)送,而不應(yīng)該等待 Buffer 被填滿或者 Output flusher 的下一次 Flush。

▼ 進(jìn)一步的討論

與小于 1.5 版本的 Flink 不同,請注意(a)網(wǎng)絡(luò) Buffer 現(xiàn)在會被直接放在 Subpartition 的隊列中,(b)網(wǎng)絡(luò) Buffer 不會在 Flush 之后被關(guān)閉。這給我們帶來了一些好處:

同步開銷較少(Output flusher 和 RecordWriter 是相互獨立的) 在高負(fù)荷情況下,Netty 是瓶頸(直接的網(wǎng)絡(luò)瓶頸或反壓),我們?nèi)匀豢梢栽谖赐瓿傻?Buffer 中填充數(shù)據(jù) Netty 通知顯著減少

但是,在低負(fù)載情況下,可能會出現(xiàn) CPU 使用率和 TCP 數(shù)據(jù)包速率的增加。這是因為,F(xiàn)link 將使用任何可用的 CPU 計算能力來嘗試維持所需的延遲。一旦負(fù)載增加,F(xiàn)link 將通過填充更多的 Buffer 進(jìn)行自我調(diào)整。由于同步開銷減少,高負(fù)載場景不會受到影響,甚至可以實現(xiàn)更高的吞吐。

4.2 BufferBuilder 和 BufferConsumer

更深入地了解 Flink 中是如何實現(xiàn)生產(chǎn)者 - 消費者機制,需要仔細(xì)查看 Flink 1.5 中引入的 BufferBuilder 和 BufferConsumer 類。雖然讀取是以 Buffer 為粒度,但寫入它是按 Record 進(jìn)行的,因此是 Flink 中所有網(wǎng)絡(luò)通信的核心路徑。因此,我們需要在任務(wù)線程(Task thread)和 Netty 線程之間實現(xiàn)輕量級連接,這意味著盡量小的同步開銷。你可以通過查看源代碼獲取更加詳細(xì)的信息。

5. 延遲與吞吐

引入網(wǎng)絡(luò) Buffer 的目是獲得更高的資源利用率和更高的吞吐,代價是讓 Record 在 Buffer 中等待一段時間。雖然可以通過 Buffer 超時給出此等待時間的上限,但可能很想知道有關(guān)這兩個維度(延遲和吞吐)之間權(quán)衡的更多信息,顯然,無法兩者同時兼得。下圖顯示了不同的 Buffer 超時時間下的吞吐,超時時間從 0 開始(每個 Record 直接 Flush)到 100 毫秒(默認(rèn)值),測試在具有 100 個節(jié)點每個節(jié)點 8 個 Slot 的群集上運行,每個節(jié)點運行沒有業(yè)務(wù)邏輯的 Task 因此只用于測試網(wǎng)絡(luò)協(xié)議棧的能力。為了進(jìn)行比較,我們還測試了低延遲改進(jìn)(如上所述)之前的 Flink 1.4 版本。

如圖,使用 Flink 1.5+,即使是非常低的 Buffer 超時(例如1ms)(對于低延遲場景)也提供高達(dá)超時默認(rèn)參數(shù)(100ms)75% 的最大吞吐,但會緩存更少的數(shù)據(jù)。

6.結(jié)論

了解 Result partition,批處理和流式計算的不同網(wǎng)絡(luò)連接以及調(diào)度類型,Credit-Based 流量控制以及 Flink 網(wǎng)絡(luò)協(xié)議棧內(nèi)部的工作機理,有助于更好的理解網(wǎng)絡(luò)協(xié)議棧相關(guān)的參數(shù)以及作業(yè)的行為。后續(xù)我們會推出更多 Flink 網(wǎng)絡(luò)棧的相關(guān)內(nèi)容,并深入更多細(xì)節(jié),包括運維相關(guān)的監(jiān)控指標(biāo)(Metrics),進(jìn)一步的網(wǎng)絡(luò)調(diào)優(yōu)策略以及需要避免的常見錯誤等。

以上就是小編為大家?guī)淼腁pache Flink的網(wǎng)絡(luò)協(xié)議棧詳細(xì)介紹的全部內(nèi)容,希望能對您有所幫助,小伙伴們有空可以來腳本之家網(wǎng)站,我們的網(wǎng)站上還有許多其它的資料等著小伙伴來挖掘哦!

相關(guān)文章

最新評論

溧阳市| 久治县| 怀集县| 清水河县| 石嘴山市| 芮城县| 东丰县| 霍林郭勒市| 阿克苏市| 九龙坡区| 同江市| 娱乐| 罗城| 旅游| 庆云县| 榆树市| 缙云县| 普格县| 定襄县| 清流县| 潮州市| 郸城县| 罗山县| 利川市| 苍山县| 海南省| 肃北| 汝城县| 宜兴市| 灵山县| 仲巴县| 右玉县| 正镶白旗| 同仁县| 饶河县| 昌邑市| 蓬溪县| 枝江市| 万荣县| 丽江市| 古浪县|