The binary search of distributed programming

Salvatore Sanfilippo

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

原文由 Salvatore Sanfilippo 發布,訂閱此部落格

昨晚我在重讀 Martin Kleppmann 寫的 Redlock 分析(http://martin.kleppmann.com/2016/02/08/how-to-do-distributed-locking.html)。文中有個地方,Martin 在思考是否有什麼好方法可以用 Redis 來產生單調遞增的 ID。

這個看似簡單的問題,乍看之下可能比想像中複雜,因為它必須在所有情況下都保證一個安全性質(safety property)永遠成立:產生的 ID 永遠大於過去產生過的所有 ID,而且同一個 ID 不能被產生多次。這一點在發生網路分割和其他故障時也必須成立。如果可連上的節點少於多數,系統可以就此變得無法使用,但絕不能給出錯誤的答案(註:如我們稍後會看到的,這個演算法在請求負載很高時,還有另一個關於活性的問題)。

所以,為了多玩一點分散式系統演算法,並在過程中多學一點東西,我試著找一個解法。其實我本來就知道一個可以解決這個問題的演算法。它沒有效率,不適合每秒產生大量的 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. 從多數實例(若 N=5,則為 3 個以上)取得「current」的值。
  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)情況下,許多數字可能會被「燒掉」。

(請注意,上一句中的「腦裂」並不是指節點之間存在不一致的狀態,而是指無法就某個特定 ID 達成多數共識。通常腦裂是指設定上的衝突,例如多個節點同時宣稱自己是主節點。然而在 Raft 論文中,腦裂一詞的用法與我在這裡的用法相同)。

在不因並行存取而頻繁失敗的前提下,每秒能產生多少 ID,取決於網路的 RTT 以及並行客戶端的數量。不過有趣的是,你可以透過建立一個與叢集溝通的「ID 伺服器」來讓這個演算法更具可擴展性,由它來居中協調客戶端的存取,將新 ID 的產生一個接一個地序列化。這並不會造成單點故障,因為你不需要只有單一的 ID 伺服器,你可以為了備援而運行數個,讓數百個客戶端連接到它們。

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

另一種做法是,當有許多客戶端且沒有節點居中協調時,如果某一輪演算法失敗,就使用隨機且指數退避的延遲再去聯繫節點。

為什麼需要 fsync?

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

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

這些 ID 能拿來做什麼

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

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

這個演算法的起源

這裡所描述的內容與 Paxos 的第一階段以及 Raft 的領導者選舉都非常相似。不過,它似乎只是 Lamport 時間戳的一種特例,其中利用多數來建立全序。

非常感謝 Max Neunhoeffer 和 Martin Kleppmann 對這篇部落格文章初稿提供的回饋。請注意,任何錯誤皆由我本人負責。

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

留言