背景:?
最近flink發(fā)布新版本1.11, 除了優(yōu)化舊版本已有的缺陷, 還增加了一些新功能,? 其中我發(fā)現(xiàn)有一些改變適合用于現(xiàn)在負(fù)責(zé)的flink項(xiàng)目?
我們當(dāng)前的flink項(xiàng)目的問(wèn)題是生成checkpoint失敗較多,造成checkpoint失敗的原因是某幾個(gè)subtask的快照超時(shí)導(dǎo)致整體的checkpoint生成失敗,隨著每天的處理的任務(wù)越多, 這個(gè)問(wèn)題越發(fā)突顯出來(lái), 而后果是:
引用的答案:
目前的 Checkpoint 算法在大多數(shù)情況下運(yùn)行良好,然而當(dāng)作業(yè)出現(xiàn)反壓時(shí)谋右,阻塞式的 Barrier 對(duì)齊反而會(huì)加劇作業(yè)的反壓惶看,甚至導(dǎo)致作業(yè)的不穩(wěn)定幼苛。首先, Chandy-Lamport 分布式快照的結(jié)束依賴(lài)于 Marker 的流動(dòng)飒硅,而反壓則會(huì)限制 Marker 的流動(dòng)运沦,導(dǎo)致快照的完成時(shí)間變長(zhǎng)甚至超時(shí)。無(wú)論是哪種情況掘而,都會(huì)導(dǎo)致 Checkpoint 的時(shí)間點(diǎn)落后于實(shí)際數(shù)據(jù)流較多。這時(shí)作業(yè)的計(jì)算進(jìn)度是沒(méi)有被持久化的于购,處于一個(gè)比較脆弱的狀態(tài)袍睡,如果作業(yè)出于異常被動(dòng)重啟或者被用戶(hù)主動(dòng)重啟,作業(yè)會(huì)回滾丟失一定的進(jìn)度肋僧。如果 Checkpoint 連續(xù)超時(shí)且沒(méi)有很好的監(jiān)控女蜈,回滾丟失的進(jìn)度可能高達(dá)一天以上,對(duì)于實(shí)時(shí)業(yè)務(wù)這通常是不可接受的色瘩。更糟糕的是伪窖,回滾后的作業(yè)落后的 Lag 更大,通常帶來(lái)更大的反壓居兆,形成一個(gè)惡性循環(huán)覆山。
優(yōu)化1: rebalance分區(qū)改為rescale分區(qū)
rebalance使用Round-ribon思想將數(shù)據(jù)均勻分配到各實(shí)例上。Round-ribon是負(fù)載均衡領(lǐng)域經(jīng)常使用的均勻分配的方法泥栖,上游的數(shù)據(jù)會(huì)輪詢(xún)式地分配到下游的所有的實(shí)例上簇宽。如下圖所示,上游的算子會(huì)將數(shù)據(jù)依次發(fā)送給下游所有算子實(shí)例吧享。
dataStream.rebalance()
rescale與rebalance很像魏割,也是將數(shù)據(jù)均勻分布到各下游各實(shí)例上,但它的傳輸開(kāi)銷(xiāo)更小钢颂,因?yàn)閞escale并不是將每個(gè)數(shù)據(jù)輪詢(xún)地發(fā)送給下游每個(gè)實(shí)例钞它,而是就近發(fā)送給下游實(shí)例。
dataStream.rescale()
如上圖所示殊鞭,當(dāng)上游有兩個(gè)實(shí)例時(shí)遭垛,上游第一個(gè)實(shí)例將數(shù)據(jù)發(fā)送給下游第一個(gè)和第二個(gè)實(shí)例,上游第二個(gè)實(shí)例將數(shù)據(jù)發(fā)送給下游第三個(gè)和第四個(gè)實(shí)例操灿,相比rebalance將數(shù)據(jù)發(fā)送給下游每個(gè)實(shí)例锯仪,rescale的傳輸開(kāi)銷(xiāo)更小。下圖則展示了當(dāng)上游有四個(gè)實(shí)例趾盐,上游前兩個(gè)實(shí)例將數(shù)據(jù)發(fā)送給下游第一個(gè)實(shí)例庶喜,上游后兩個(gè)實(shí)例將數(shù)據(jù)發(fā)送給下游第二個(gè)實(shí)例。
優(yōu)化2: 升級(jí)flink1.11, 使用Unaligned Checkpoint + rocksdb生成Checkpoint
flink1.11新特性相關(guān)介紹:?https://www.h5w3.com/33867.html
Rocksdb state ssd:? 使用rocksdb緩存checkpoint, 并且從原來(lái)的全量生成改為增量生成的方式, 速度更快
Unaligned Checkpoint
Flink 現(xiàn)有的 Checkpoint 機(jī)制下救鲤,每個(gè)算子需要等到收到所有上游發(fā)送的 Barrier 對(duì)齊后才可以進(jìn)行 Snapshot 并繼續(xù)向后發(fā)送 barrier久窟。在反壓的情況下,Barrier 從上游算子傳送到下游可能需要很長(zhǎng)的時(shí)間蜒简,從而導(dǎo)致 Checkpoint 超時(shí)的問(wèn)題瘸羡。
針對(duì)這一問(wèn)題漩仙,F(xiàn)link 1.11 增加了 Unaligned Checkpoint 機(jī)制搓茬。開(kāi)啟 Unaligned Checkpoint 后當(dāng)收到第一個(gè) barrier 時(shí)就可以執(zhí)行 checkpoint犹赖,并把上下游之間正在傳輸?shù)臄?shù)據(jù)也作為狀態(tài)保存到快照中,這樣 checkpoint 的完成時(shí)間大大縮短卷仑,不再依賴(lài)于算子的處理能力峻村,解決了反壓場(chǎng)景下 checkpoint 長(zhǎng)期做不出來(lái)的問(wèn)題。
可以通過(guò) env.getCheckpointConfig().enableUnalignedCheckpoints();開(kāi)啟unaligned Checkpoint 機(jī)制锡凝。
總的來(lái)說(shuō), 新特性一定程度解決了Checkpoint與反壓的耦合
分析過(guò)程:?
首先測(cè)查算子間是否存在反壓, 在flink web ui后臺(tái)可以查看:
我的flink作業(yè)沒(méi)有反壓的問(wèn)題
定位問(wèn)題的原因是:?部分幾個(gè)subtask處理速度跟不上, 導(dǎo)致barrier流向慢, input緩沖區(qū)占滿(mǎn), barrier對(duì)齊不了, 導(dǎo)致整體的checkpoint生成失敗
flink作業(yè)operator處理數(shù)據(jù)的效率不均的原因主要是:
數(shù)據(jù)的多樣性: 不同數(shù)據(jù)的類(lèi)型或大小不一致, 導(dǎo)致處理的時(shí)間不一致,
如果使用了rebalance分區(qū)策略, 還是會(huì)負(fù)載均衡地分配到每個(gè)subtask上, 本來(lái)負(fù)載高的subtask還是會(huì)發(fā)配到任務(wù)處理, 導(dǎo)致了惡性循環(huán)
Flink 現(xiàn)有的物理分區(qū)策略全是靜態(tài)的負(fù)載均衡策略粘昨,沒(méi)有動(dòng)態(tài)根據(jù)負(fù)載能力進(jìn)行負(fù)載均衡的策略
未升級(jí)之前:?
網(wǎng)上看到一篇分析得很好的文章, 恰好就是現(xiàn)在內(nèi)容引入出現(xiàn)的問(wèn)題:?Flink 中的木桶效應(yīng):?jiǎn)蝹€(gè) subtask 卡死導(dǎo)致整個(gè)任務(wù)卡死? 建議大家看一看`~
引用如下:?
代碼實(shí)現(xiàn):? 略過(guò), 可以參考官方文檔
優(yōu)化后的效果:
參考文獻(xiàn):
Flink 1.11 Unaligned Checkpoint 解析
Flink 中的木桶效應(yīng):?jiǎn)蝹€(gè) subtask 卡死導(dǎo)致整個(gè)任務(wù)卡死
flink消費(fèi)kafka時(shí)出現(xiàn)數(shù)據(jù)傾斜的原因和處理方式