借助一個(gè)文件寫入的例子來說明兩階段提交,在Flink中使用兩階段提交,需要實(shí)現(xiàn)TwoPhaseCommitSinkFunction這個(gè)抽象類的四個(gè)方法,我們下面來說明。
protected abstract TXN beginTransaction() throws Exception; protected abstract void preCommit(TXN transaction) throws Exception; protected abstract void commit(TXN transaction); protected abstract void abort(TXN transaction);
1. beginTransaction - 在事務(wù)開始前,我們?cè)谀繕?biāo)文件系統(tǒng)上面的臨時(shí)目錄上創(chuàng)建一個(gè)臨時(shí)文件。隨后,我們?cè)诔绦蛱幚淼臅r(shí)候可以將數(shù)據(jù)寫入到這個(gè)文件。
2. preCommit - 在預(yù)提交階段,我們刷新文件到磁盤,關(guān)閉文件。
3. commit - 在提交階段,我們?cè)有缘膶㈩A(yù)提交階段的文件移動(dòng)到真正的目標(biāo)目錄。需要注意的是,這增加了輸出數(shù)據(jù)的可見性的延遲,因?yàn)椴籱v是看不到數(shù)據(jù)的,延遲時(shí)間就是設(shè)定的checkpoint的時(shí)間。
4. abort - 在終止階段,我們刪除臨時(shí)文件 *如果步驟中有任何錯(cuò)誤,F(xiàn)link會(huì)通過最新的checkpoint來恢復(fù)程序狀態(tài)。
比如預(yù)提交成功了,在通知到達(dá)operator之前失敗了。
這時(shí)候,F(xiàn)link將operator的狀態(tài)恢復(fù)到預(yù)提交階段,即還未真正提交的時(shí)候。
為了能在重啟的時(shí)候能夠正確的終止或者提交事務(wù),我們需要在預(yù)提交階段將足夠的信息保存到checkpoint中。
在這個(gè)例子中,這些信息是臨時(shí)文件以及目標(biāo)目錄的地址, 當(dāng)從checpoint恢復(fù)時(shí),F(xiàn)link會(huì)先執(zhí)行一個(gè)Commit操作。