The binary search of distributed programming

Salvatore Sanfilippo

分散式程式設計的二分搜尋

昨晚我重新閱讀了 Martin Kleppmann(馬丁·克萊普曼)撰寫的 Redlock 分析(http://martin.kleppmann.com/2016/02/08/how-to-do-distributed-locking.html)。在某個段落中,馬丁·克萊普曼思考是否有什麼好方法可以用 Redis 來產生單調遞增的 ID。

這個看似簡單的問題,初看之下可能比想像中複雜,因為它必須確保在所有情況下,都有一個始終受到保證的 safety property(安全性質):所產生的 ID 永遠大於過去所有已產生的 ID,且同一個 ID 不會被產生多次。這必須在發生網路分割及其他故障時仍然成立。系統在可連線的節點少於過半數時,可以直接變為不可用,但絕不能提供錯誤的答案(注意:如後文所述,這個演算法在高負載請求期間還有另一個 liveness issue(活性問題))。

所以,為了多玩一下分散式系統演算法,並在過程中多學一點東西,我試著尋找一個解法。其實我本來就知道一個可以解決這個問題的演算法。它效率不高,不適合每秒產生大量的 ID。許多複雜的分散式演算法,像是 Raft 和 Paxos,都會將它作為其中一個步驟來取得單調遞增的 ID,作為建構其所需完整特性集的基礎。這個演算法非常迷人,因為它極易理解與實作,而且能非常直觀地理解為什麼它會有效。我可以說,它就是分散式演算法中的二分搜尋,夠簡單卻又夠聰明,能讓剛接觸分散式程式設計的新手產生恍然大悟的時刻。

不過,我必須修改這個演算法,使其能在客戶端實作。希望它仍然是正確的(歡迎提供回饋)。雖然我不打算用這個演算法來改進 Redlock(請參閱我先前的部落格文章),但我認為嘗試解決這類問題既是很好的練習,對於初次接觸分散式系統、想在真實系統中找些簡單問題來實作看看的人來說,也可能是有趣的閱讀材料。

運作原理是什麼?

這個演算法的需求有以下兩點:

  1. 支援 set_if_less_than() 操作的資料儲存。
  2. 在回覆客戶端之前,能在寫入時將資料 fsync() 至磁碟的資料儲存。

以上包含了幾乎所有的 *SQL 伺服器、Redis 以及許多其他儲存系統。

我們有一組 N 個節點,為了簡化說明,假設 N = 5。我們透過將一個名為「current」的鍵設為 0 來初始化系統,以 Redis 的語法來說,就是執行:

SET current 0

在全部 5 個實例上都要這麼做。這是初始化的一部分,且只有在初始化一個新的「叢集」時才需要執行。這個步驟可以省略,但會讓說明更為簡單。

為了產生一個新的 ID,我們這樣做:

  1. 從過半數的實例取得「current」的值(若 N=5,則為 3 個或以上)。
  2. 如果未能連上 3 個實例,則回到步驟 1。
  3. 在取得的值中取最大值,並將其加 1。我們稱之為 $NEXTID
  4. 將以下寫入操作傳送給所有能夠連上的節點。
    IF current < $NEXTID THEN
        SET current $NEXTID
        return $NEXTID
    ELSE
        return NULL
    END
  5. 如果有 3 個或以上的實例回覆 $NEXTID,則演算法成功,我們已成功產生一個新的單調遞增 ID。
  6. 否則,若未達到過半數,則回到步驟 1。

我們在步驟 4 傳送的內容,可以輕鬆地轉譯為一個簡單的 Redis Lua 腳本:

local val = tonumber(redis.call('get',KEYS[1]))
local nextid = tonumber(ARGV[1])

if val < nextid then
    redis.call('set',KEYS[1],nextid)
    return nextid
else
    return nil
end

這樣安全嗎?

除了它以修改後的形式被用作經過深入分析的演算法中的一個步驟之外,我直觀地相信它可行,原因在於:

如果我們能夠取得過半數的「票數」,根據定義,任何其他客戶端都不可能為大於或等於我們所產生之 ID 的值取得過半數。否則,就表示已經有 3 個或以上的實例其 current 的值 >= $NEXTID,那我們就不可能取得過半數。因此,產生的 ID 永遠大於過去的 ID,且基於同樣的條件,兩個實例也不可能產生相同的 ID。

或許有熱心的讀者能夠指出這個演算法中的錯誤,或是提供此演算法在其他已分析系統中被使用的相關分析,不過鑑於上述內容是為了在客戶端執行而改寫的,實際上涉及了更多的處理程序,應該重新進行分析以證明其等價性。

為什麼這是個緩慢的演算法?

這個演算法的問題在於並行存取。如果許多客戶端同時嘗試產生新的 ID,可能沒有任何一方能取得過半數,他們將需要用更大的數字再次重試。請注意,這也意味著產生的 ID 序列中可能會出現「空洞」,因此客戶端可能會產生如 1、2、6、10、11、21、… 的序列。因為許多數字可能會在並行存取所造成的 split brain(腦裂) 狀態下被「燒掉」。

(請注意,上述句子中的「split brain」並不表示節點之間存在不一致的狀態,而是指無法取得過半數來就某個特定的 ID 達成共識。通常 split brain 指的是設定上的衝突,例如,多個節點同時宣稱自己是 master。不過,在 Raft 論文中,split brain 一詞的使用方式與我在這裡的用法相同)。

在不因並行存取而頻繁失敗的情況下,每秒能產生多少個 ID,取決於網路 RTT 與並行客戶端的數量。然而有趣的是,你可以透過建立一個與叢集溝通的「IDs server」,並以序列化方式逐一調解客戶端對新 ID 產生的存取,來讓演算法更具可擴展性。這並不會造成單點故障,因為你不需要只有單一的 ID 伺服器,你可以為了備援而執行數個,並讓數百個客戶端連線到它們。

使用這種架構來每秒產生 5k 個 ID 應該是可行的,特別是如果客戶端是以聰明的方式實作,嘗試同時使用多工或多執行緒的方式將請求發送給 5 個節點的話。

另一種在有許多客戶端且沒有節點居中調解存取時的做法,是在某一輪演算法失敗後,使用隨機與指數退避的延遲來再次聯繫節點。

為什麼需要 fsync?

在這裡,每次寫入時進行 fsync 是強制性的,因為如果節點當機並重新啟動,它們必須擁有「current」鍵的最新值。如果 current 的值回溯,就有可能違反我們那個「新產生的 ID 永遠大於過去產生的任何其他 ID」的 safety property。然而,如果你使用完整複寫的 FSM(有限狀態機) 來達成相同的目標,你無論如何都需要 fsync(但你就不會有並行存取的問題。例如,在正常情況下使用 Raft 時,你會有一個可以發送請求的單一 leader)。

因此,以 Redis 來說,必須將其設定為啟用 AOF,並將 AOF 的 fsync 策略設為 always,以確保寫入在回覆客戶端之前一定會被持久化。

這些 ID 能用來做什麼

像這樣的一組 ID,具有一種稱為「total ordering(全序)」的特性,因此在不同的情境中都非常有用。通常在分散式運算中,很難說清楚什麼先發生、什麼後發生。使用這些 ID,你就能始終知道某些事件的順序。

舉一個簡單的例子:不同的處理程序使用這個系統,可以計算出一份項目清單,並將各自的子清單保存在本地儲存中。最後,它們可以合併多個清單並得到一份順序正確的最終清單,就好像從一開始就有一個單一的共享清單,讓每個處理程序都能在上面附加項目一樣。

這個演算法的起源

這裡所描述的內容非常類似於 Paxos 的第一階段以及 Raft 的 leader election。然而,它似乎只是 Lamport timestamp(Lamport 時間戳記) 的一個特例,其中利用過半數來建立 total ordering。

非常感謝 Max Neunhoeffer(馬克斯·諾伊霍弗)與馬丁·克萊普曼對本文初稿提供的回饋。請注意,任何錯誤皆為我個人所致。

原文由 Salvatore Sanfilippo 發布

本文章由 muse-spark-1.2-contributor 進行翻譯