Adventures in message queues

Salvatore Sanfilippo

訊息佇列的冒險

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

編按:如果你錯過了,Disque 的原始碼現在已經在 http://github.com/antirez/disque 公開。

這幾個月來,我大約花了 15 至 20% 的時間——大多是從夜晚和週末硬擠出來的時數——在打造一個新系統。它是一個訊息中介,叫做 Disque。我已經完成了原始規格中大約 80% 的實作,但還是覺得還沒到可以釋出的階段。既然還無法釋出,至少來寫篇部落格吧……所以這就是它如何開始的故事,以及關於它究竟是什麼的一些細節。

~ 起步 ~

許多開發者把 Redis 當成訊息佇列來用,常常透過某個函式庫把 Redis 底層的原語包裝起來,有時則是直接用 Redis 的原生 API 兜出一個簡單、臨時拼湊的佇列。這個使用情境主要是靠阻斷式 list 操作和 list 的 push 操作來實現。Redis 用在這種場合,顯然同時是最好的也是最糟的系統。說它好,是因為它快速、易於檢視、部署與使用,而且在許多環境中它本來就已經是基礎架構的一部分。然而它也有缺點,因為 Redis 可變的資料結構和不可變的訊息在本質上非常不同。Redis 在高可用性/叢集上的取捨,完全偏向處理大型可變數值,但同樣的取捨用來處理訊息卻不是最佳選擇。

對於一個訊息中介來說,有一件事一定要保證:訊息要不是至少被遞送一次,就是至多被遞送一次。簡言之,既然要保證訊息恰好只被遞送一次(這裡的遞送指的是訊息被 worker 接收「並」處理完成)在實務上幾乎不可能,選擇就剩下訊息中介要能保證 0 次或 1 次的遞送,或是 1 次到無限多次的遞送。這通常被稱為至多一次(at-most-once)語意和至少一次(at-least-once)語意。前者有它的使用場景,但最有趣也最實用的語意是後者,也就是保證訊息至少會被遞送一次,如果發生失敗,就重複遞送。

所以幾個月前,我開始思考一種客戶端協定,讓一組 Redis master(完全不使用複寫或叢集)在提供這些保證的情況下運作。有時候,只要稍微改變 Redis 在某個使用場景下的用法,就有可能得到一個更好的系統。舉例來說,針對分散式鎖,我曾試著整理出一個實作起來極為簡單、卻比單一實例加上容錯移轉的作法更為穩健的演算法(http://redis.io/topics/distlock)。

然而工作了幾天後,我的設計草稿顯示,設計一個專用的系統會是更好的選擇,因為客戶端演算法最後變得太複雜、效率不佳,而且有些我非常想要的功能根本不可能做到,或是非常難以實現。要在 Redis 上再疊加更多東西聽起來不是好主意,它已經做了太多事了,而要把訊息處理好,我需要的東西和 Redis 的運作方式非常不同。但既然世界上已經充滿了訊息中介,為什麼還要設計一個新系統?因為有驚人數量的人用 Redis 來取代那些專為此目的設計的系統,這很奇怪。少數人用錯還說得過去,但這麼多人都這樣做,一定有原因。或許 Redis 低門檻、簡單的 API 和速度,正是大多數人在檢視訊息中介領域時所不習慣的。那個領域似乎充斥著要不是太過陽春、要求應用程式自己做太多事,就是太過複雜、但功能超級齊全的解決方案。也許,中間還有一個「訊息界的 Redis」的空間?

~ 暴力分支 Redis ~

我這輩子第一次沒有馬上開始寫程式。好幾個星期,我不時審視這個設計,把它從一個 Redis 客戶端函式庫轉換成一個全新的系統,並試著以使用者的角度去思考,一個訊息中介要具備什麼才會讓我非常滿意。最初的使用場景始終不變:延遲任務(delayed jobs)。Disque 是一個通用的系統,但在設計時有 90% 的時間,我心中的「參考對象」都是一個需要解決發送訊息——很可能就是待處理的工作——問題的使用者。如果有任何東西和這個使用場景衝突,就把它拿掉。

設計完成後,我終於開始寫程式。但要從哪裡開始?「vi main.c」?幸好 Redis 在某種程度上,是一個用 C 語言撰寫分散式系統的框架。我已經有了協定、網路函式庫、客戶端處理、節點間的訊息匯流排。全部從頭重寫聽起來完全是浪費。與此同時,我希望 Disque 在任何細節上都能徹底脫離 Redis,如果有需要的話;而且我希望它是一個獨立的 side project,不會對 Redis 本身造成影響。所以,與其嘗試把 Redis 拆成一個真正獨立的框架加上 Redis 實作這種艱鉅的工程,我採取了一個更務實的做法:我把程式碼分支出來,然後把所有 Redis 特有的部分從原始碼中移除,最後得到一個骨架。到這個階段,我已經準備好要實作我的規格了。

~ 什麼是 Disque? ~

經過幾個月非常不密集的工作和僅僅 200 次 commit,我終於得到一個不再像玩具的系統:它有好幾個星期看起來都像個玩具,所以我甚至不敢談論它,因為我隨時把整個原始碼目錄刪掉的機率都很大。現在大部分的構想都已經是有測試、能運作的程式碼,我終於確定它未來一定會釋出,也來談談我在設計上所做的取捨。

Disque 預設就是一個分散式系統。既然它是一個 AP 系統,就像 Redis 那樣分成單節點模式和分散式模式就沒有意義了。單一的 Disque 節點只不過是叢集的一個特例,也就是只有一個節點的叢集。所以這是設計上的一個重要重點:容錯、能抵抗分割、而且無論還有多少節點存活都保持可用,也就是所謂的 AP。我也想要一個本質上能在不同場景下擴展的系統,無論是面對大量生產者和消費者以及大量佇列的情況,還是所有這些生產者和消費者都集中在單一佇列、而該佇列可能分散在多個節點上的情況。

我的需求大聲地告訴我一件事……那就是 Disque 將會做出一個重大的設計犧牲。訊息順序。Disque 只提供盡力而為的順序保證。然而,正因為這個犧牲,才有許多收穫……取捨之所以有趣,有時就在於它會徹底打開設計空間。

我可以就這樣繼續跟你講 Disque 是什麼,不過幾個月前我在 Hacker News 上看到一則由 Jacques Chester 寫的留言,見 https://news.ycombinator.com/item?id=8709146 [編按:抱歉,我複製貼上時搞錯了 Adrian 的名字(嗨 Adrian,抱歉誤引了你!)]。Jacques 碰巧和我一樣在 Pivotal 工作,他當時在評論不同的訊息系統有著非常不同的功能、特性,如果沒有細節,幾乎不可能評估不同的選擇,也無法判斷一個系統比另一個快,是因為它有更好的實作,還是只是因為它提供的保證少得多。所以他寫了一組在評估訊息系統時應該要問的問題。我會借用他的問題,再加上幾個,來說明 Disque 是什麼,希望最後不是空談,而是能提供一些實際的資訊。

問:訊息是否至少會被遞送一次?

在 Disque 中,你可以選擇至少一次遞送(預設值)或至多一次遞送。這個屬性可以針對每則訊息個別設定。至多一次遞送其實只是至少一次遞送的特例,只要把訊息的「retry」參數設為 0,並只把訊息複寫到單一節點即可。

問:訊息是否需要由消費者確認?

是的,消費者要告訴系統訊息已正確遞送,唯一的方式就是確認(acknowledge)它。

問:如果未被確認,訊息是否會被重複遞送?

是的,Disque 會在「retry」時間過後自動重新遞送訊息,而且會一直重試(直到訊息達到最大 TTL 時間)。當訊息被確認後,確認訊息會被傳播到所有持有該訊息副本的節點。如果系統認為已經通知到所有人,訊息最終就會被垃圾回收並移除。在記憶體壓力下,已被確認的訊息也會被逐出。

節點會執行一種盡力而為的演算法,避免同一則訊息被多次排入佇列,以便更好地趨近單次遞送。然而在發生故障時,多個節點可能會在同一時間重複遞送同一則訊息多次。

問:佇列是持久的還是暫時的?

持久的。

問:持久性是靠先把每則訊息寫入磁碟,還是靠在伺服器之間複寫訊息來達成?

預設情況下 Disque 僅在記憶體中運作,並使用同步複寫來達成持久性(不過你可以針對每則訊息要求使用非同步複寫)。如果你有需要,也可以像 Redis 一樣開啟 AOF,如果你的環境很可能會遇到大規模重啟之類的情況。當系統要升級時,即使平常沒有使用磁碟持久化,也可以在升級期間把 AOF 寫到磁碟上,以確保重啟後不會遺失狀態。

問:佇列是在一組伺服器之間部分/完全一致,還是為了最大吞吐量而被分散?

為了吞吐量而分散,不過訊息順序仍會以盡力而為的方式保留。每則訊息都有一個不可變的「ctime」,它是由毫秒級的 wall-clock 時間戳記加上同一毫秒內產生的遞增 ID 所組成。節點會使用這個 ctime 來排序並遞送訊息。

問:在壓力下訊息是否可能被完全丟棄?(也就是盡力而為)

不會,不過如果記憶體沒有空間,新的訊息可能會被拒絕。當已使用 75% 的記憶體時,接收訊息的節點會嘗試把訊息向外複寫到其他節點,且不在本地保留副本,但如果其他節點也處於記憶體不足的狀態,就可能無法成功。

問:消費者和生產者能否檢視佇列,還是佇列完全不透明?

有「PEEK」指令可以窺視佇列內容。

問:佇列是無序、FIFO 還是具優先順序的?

如前所述,是盡力而為的類 FIFO。

問:有沒有中介者(broker),還是無中介者?

有中介者,也就是一組 master。客戶端可以和任意節點溝通。

問:中介者是否擁有獨立、具名的佇列(主題、路由等),還是生產者和消費者需要協調彼此的連線?

是具名的佇列。生產者和消費者不需要協調,因為節點會透過聯邦機制在叢集內發現路由,並在消費者需要時把訊息轉送過去。不過,如果客戶端有意願,系統也會提供提示,讓它重新定位到有更多消費者的地方。

問:訊息的發送是否具交易性?

是的,一旦新增訊息的指令回傳,系統就能保證叢集內已經有指定數量的副本。

問:訊息的接收是否具交易性?

我想不算,因為如果訊息未被確認,Disque 會嘗試再次遞送同一則訊息。

問:消費者在接收時會阻斷,還是可以主動檢查是否有新訊息?

兩種行為都支援,預設是阻斷。

問:生產者在發送時會阻斷,還是可以檢查佇列是否已滿?

生產者在新增訊息時,可以要求若本地節點中該佇列的長度已經超過指定值,就回傳錯誤。

此外,如果生產者想盡快離開、讓叢集以盡力而為的方式去複寫訊息,也可以要求以非同步方式複寫訊息。

目前沒有辦法在佇列中訊息過多時阻斷消費者,並在訊息變少時自動解除阻斷。

問:是否支援延遲任務?

支援,精確到秒,最長可達數年。不過它們會占用記憶體。

問:消費者和生產者能否連接到不同的節點?

可以。

我希望透過這篇文章,Disque 已經不那麼像是空中樓閣了。當然,不看程式碼很難斷定,但如果你最關心的功能已經完成了,至少你已經可以開始抱怨了。上面提到的有多少已經實作並且運作良好?除了 AOF 磁碟持久化和幾個我想在 API 上再微調的小地方之外,全部都已經完成了,所以首個釋出版本應該不會太遠,不過因為我做得很零散,要非常快速也很難。

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

留言