Adventures in message queues

Salvatore Sanfilippo

訊息佇列探險記

編按:如果你還沒看到,Disque 的原始碼現已開放,位於 http://github.com/antirez/disque

這幾個月來,我大約花了 15 至 20% 的時間,大多是從夜晚與週末擠出來的零碎時間,投入開發一個新系統。它是一個 message broker(訊息代理),名叫 Disque。我已經實作了原始規格中約 80% 的功能,但仍覺得還沒準備好釋出。既然還無法發布,至少來寫篇網誌……所以這就是它的起源故事,以及一些關於它是什麼的細節。

~ 起步 ~

許多開發者把 Redis 當作 message queue(訊息佇列) 來使用,通常是透過某個函式庫來封裝 Redis 底層的原語,有時則是直接用 Redis 原生的 API 打造一個簡單、即席的佇列。這個使用情境主要靠 blocking list operations(阻塞式串列操作) 與 list push 操作來實現。表面上看來,Redis 同時是拿來這樣用的最好與最差的系統。說它好,是因為它快速、易於檢視、部署與使用,而且在許多環境中它本來就是基礎架構的一環。然而它也有缺點,因為 Redis 可變的資料結構與不可變的訊息本質上截然不同。Redis 在 HA / Cluster 上的取捨完全偏向大型可變數值,但同樣的取捨並不是處理訊息的最佳選擇。

對於 message broker 來說,有一件重要的事必須保證:訊息要不是至少被遞送一次,就是至多被遞送一次。簡言之,要保證訊息的精確單次遞送(此處的遞送指的是訊息已被工作者接收「且」處理完成)在實務上幾乎是不可能的,因此 message broker 能保證的選擇只有兩種:0 或 1 次遞送,或是 1 到無限次遞送。這通常被稱為 at-most-once semantics(最多一次語意) 與 at-least-once semantics(至少一次語意)。前者也有其使用場景,但最有趣也最實用的語意是後者,也就是保證訊息至少被遞送一次,若發生失敗則會多次遞送。

所以幾個月前,我開始思考某種 client-side protocol(客戶端協定),用一組 Redis master(完全不使用複寫或叢集)來提供這些保證。有時只要在 Redis 的使用方式上做些微小的改變,就可能得到更好的系統。舉例來說,針對 distributed locks(分散式鎖),我曾試著記錄一個實作上極為簡單、卻比單一實例加容錯移轉的實作更為強健的演算法(http://redis.io/topics/distlock)。

然而工作幾天後,我的設計草稿顯示,打造一個專門的系統會是更好的選擇,因為客戶端演算法最終變得過於複雜、效率不佳,而且我非常想要的某些功能根本不可能或極難實現。想在 Redis 上再疊加更多功能聽起來不是好主意,它已經做了太多事,而要把訊息處理好,我需要的東西與 Redis 的運作方式截然不同。但既然世界上已經充滿了 message broker,為什麼還要設計一個新系統?因為有令人驚訝數量的人把 Redis 拿來當作這些專為此目的設計的系統的替代品,這很奇怪。少數人犯錯或許正常,但這麼多人這樣做一定有其原因。也許是 Redis 的低進入門檻、簡易的 API 與速度,是多數人在檢視 message broker 的版圖時所不習慣見到的。那個領域似乎充斥著不是過於簡單、要求應用程式做太多事,就是過於複雜、卻功能滿載的解決方案。也許「訊息領域的 Redis」還有一些空間?

~ 暴力分岔 Redis ~

這是我生平第一次沒有馬上開始寫程式。好幾個星期以來,我不時審視這個設計,把它從一個 Redis 客戶端函式庫轉換為一個全新的系統,並試著站在使用者的角度去理解,什麼樣的 message broker 會讓我非常滿意。最初的使用情境始終不變:延遲任務。Disque 是一個通用系統,但在設計時有 90% 的時間,「參考情境」都是一個需要解決發送訊息(很可能就是待處理的工作)的使用者。如果某個設計與這個使用情境衝突,就把它拿掉。

當設計就緒後,我終於開始寫程式。但要從哪裡開始?「vi main.c」?幸好 Redis 在某種程度上是一個用 C 撰寫 distributed systems(分散式系統) 的框架。我已經有了協定、網路函式庫、客戶端處理、節點間的 message bus(訊息匯流排)。要從頭重寫這一切聽起來是巨大的浪費。同時,我希望 Disque 在任何細節上都能完全脫離 Redis 獨立演進,並且希望它是一個不影響 Redis 本身的支線專案。因此,與其嘗試把 Redis 拆成一個真正獨立的框架與 Redis 實作這項艱鉅任務,我採取了更務實的做法:我把程式碼分岔(fork),並從原始碼中移除所有 Redis 特有的部分,最終得到一個骨架。到了這個階段,我已經準備好實作我的規格了。

~ 什麼是 Disque? ~

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

Disque 預設就是一個 distributed system。既然它是 AP system(AP 系統),像 Redis 那樣同時提供單節點模式與分散式模式就沒有意義了。單一的 Disque 節點只是叢集的一個特例,也就是只有一個節點的叢集。因此這是設計上的重要原則之一:容錯、可抵禦網路分割,且無論還有多少節點存活都能保持可用,也就是所謂的 AP。我也想要一個本質上能在不同情境下擴展的系統,無論是面對擁有許多佇列的大量生產者與消費者,或是所有這些生產者與消費者都集中在單一佇列上、而該佇列可能分散到多個節點的情況。

我的需求大聲地告訴我一件事……Disque 將會做出一個重大的設計犧牲:訊息順序。Disque 只提供 best-effort ordering(盡力而為的順序)。然而正因為這個犧牲,才有許多收穫……取捨之所以有趣,有時正是因為它們完全打開了設計空間。

我可以就這樣繼續敘述 Disque 是什麼,不過幾個月前我在 Hacker News 上看到一則由 Jacques Chester(賈克·切斯特)撰寫的留言,見 https://news.ycombinator.com/item?id=8709146 [編按:抱歉,我剪貼時弄錯了 Adrian(艾卓恩)的名字(嗨 Adrian,抱歉誤引了你!)]。賈克·切斯特碰巧和我一樣在 Pivotal 工作,他評論道,不同的訊息系統擁有非常不同的功能與特性,若沒有細節,幾乎不可能評估不同的選擇,也難以判斷某個系統比較快,究竟是因為它有更好的實作,還是單純因為它提供的保證少得多。因此他寫了一組在評估訊息系統時應該提出的問題。我將借用他的問題,並補上幾個,來說明 Disque 是什麼,希望不會只是空泛之談,而是提供一些實際的資訊。

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

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

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

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

問:若未被確認,訊息是否會被多次遞送?

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

節點會執行 best-effort algorithm(盡力而為演算法) 來避免同一訊息被多次排入佇列,以更接近單次遞送的理想情況。然而在故障期間,多個節點可能會同時多次重新遞送同一訊息。

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

持久的。

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

預設情況下 Disque 僅在記憶體中運作,並使用 synchronous replication(同步複寫) 來達成持久性(不過你可以針對每則訊息要求使用 asynchronous replication(非同步複寫))。如果部署環境可能會遇到大規模重啟等情況,也可以選擇開啟 AOF(僅附加檔案)(類似 Redis 的做法)。當系統升級時,也可以僅為了升級而將 AOF 寫入磁碟,以便即使平時不使用磁碟持久化,重啟後也不會遺失狀態。

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

為了吞吐量而分割,不過訊息順序仍以盡力而為的方式保留。每則訊息都有一個不可變的「ctime」,它是實體時鐘的毫秒時間戳加上同一毫秒內產生的遞增 ID。節點會使用這個 ctime 來為遞送進行排序。

問:訊息在壓力下是否會被完全丟棄?(也就是 best effort)

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

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

有可用來「PEEK」檢視佇列的指令。

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

如前所述,盡力而為的近似 FIFO。

問:有 broker 還是沒有 broker?

以一組 master 作為 broker。客戶端可以與任意節點通訊。

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

具名佇列。生產者與消費者不需要協調,因為節點會透過聯邦機制在叢集內發現路由,並在消費者需要時傳遞訊息。不過,若客戶端願意遷移到有更多消費者的節點,系統也會提供提示。

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

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

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

我想不是,因為若訊息未被確認,Disque 會嘗試再次遞送同一訊息。

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

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

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

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

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

沒有辦法在佇列中訊息過多時阻塞消費者,並在訊息變少時再喚醒它。

問:是否支援延遲任務?

支援,粒度為秒,最長可達數年。不過它們會占用記憶體。

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

可以。

希望透過這篇文章,Disque 不再那麼像是空中樓閣。當然,不看程式碼很難評斷,但既然最重要的功能已經公開,你至少已經可以開始抱怨了。以上有多少已經實作且運作良好?除了 AOF 磁碟持久化,以及我想在 API 上再精煉的幾個小地方之外,全部都已完成,所以首個版本應該不會太遠,只是由於我投入的時間如此零碎,要非常快速推進確實很難。

原文由 Salvatore Sanfilippo 發布

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