借助一個文件寫入的例子來說明兩階段提交,在Flink中使用兩階段提交,需要實現TwoPhaseCommitSinkFunction這個抽象類的四個方法,我們下面來說明。
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 - 在事務開始前,我們在目標文件系統上面的臨時目錄上創建一個臨時文件。隨后,我們在程序處理的時候可以將數據寫入到這個文件。
2. preCommit - 在預提交階段,我們刷新文件到磁盤,關閉文件。
3. commit - 在提交階段,我們原子性的將預提交階段的文件移動到真正的目標目錄。需要注意的是,這增加了輸出數據的可見性的延遲,因為不mv是看不到數據的,延遲時間就是設定的checkpoint的時間。
4. abort - 在終止階段,我們刪除臨時文件 *如果步驟中有任何錯誤,Flink會通過最新的checkpoint來恢復程序狀態。
比如預提交成功了,在通知到達operator之前失敗了。
這時候,Flink將operator的狀態恢復到預提交階段,即還未真正提交的時候。
為了能在重啟的時候能夠正確的終止或者提交事務,我們需要在預提交階段將足夠的信息保存到checkpoint中。
在這個例子中,這些信息是臨時文件以及目標目錄的地址, 當從checpoint恢復時,Flink會先執行一個Commit操作。