Streams: a new general purpose data structure in Redis.

Salvatore Sanfilippo

Streams:Redis 全新的通用資料結構

直到幾個月前,對我來說 Streams(串流) 不過是在訊息傳遞領域中一個有趣且相對單純的概念。在 Kafka 讓這個概念普及之後,我主要研究它在 Disque(一個訊息佇列,如今正準備轉化為 Redis 4.2 模組)中的實用性。後來我認為 Disque 的核心在於 AP 訊息傳遞,也就是在客戶端無需付出太多努力的情況下提供容錯能力與投遞保證,因此我認定 Streams 的概念在該情境下並不合適。

然而,與此同時,Redis 預設提供的資料結構卻讓我始終無法安心。Redis 的列表、有序集合與 Pub/Sub 功能之間存在著某種落差。你當然可以運用這些工具來建模一連串的訊息或事件,但各有不同的取捨。有序集合非常耗費記憶體,無法自然地表達同一則訊息被反覆投遞的情況,客戶端也無法阻塞等待新訊息。因為有序集合並非循序資料結構,而是一個可透過改變分數來移動元素的集合:難怪它不適合用來處理像 Time Series(時間序列) 這類需求。列表在某些使用情境下則有不同的問題,造成類似的適用性限制:你無法有效檢視列表中間的內容,因為在該情況下的存取時間是線性的。此外,也無法做到扇出,列表的阻塞操作一次只會將單一元素提供給單一客戶端。列表中也沒有固定的元素識別碼,讓你無法說:從該元素開始給我之後的資料。對於一對多的負載,有 Pub/Sub,在許多情況下表現出色,但對於某些需求,你並不想要發送即忘的模式:保留歷史紀錄很重要,不僅是為了在斷線後重新取得訊息,也因為某些訊息序列,例如 Time Series,非常需要透過範圍查詢來探索:在這 10 秒區間內我的溫度讀數是多少?

我當時嘗試解決上述問題的方式,是計畫將有序集合與列表整合成一個更具彈性的單一資料結構,然而我的設計嘗試幾乎總是讓最終產生的資料結構比現有的結構更加生硬、不自然。Redis 的一個優點在於,它所提供的資料結構更貼近自然的電腦科學資料結構,而非「這是 Salvatore(薩爾瓦多)發明的 API」。所以最後,我停止了嘗試,心想:好吧,這就是目前我們能提供的,或許未來我會為 Pub/Sub 加上歷史紀錄,或讓列表的存取模式更具彈性。然而,每當在研討會上有使用者走過來問我「你會如何在 Redis 中建模 Time Series?」或類似的問題時,我的臉色就會發綠。

起源

在 Redis 4.0 引入模組之後,使用者開始嘗試自行解決這個問題。其中一位 Timothy Downs(提摩西·唐斯)透過 IRC 寫了以下訊息給我:

<forkfork> 我打算做的模組是新增一種類似交易日誌的資料型別——意思是讓非常大量的訂閱者能夠做到類似 Pub/Sub 的事情,而不會造成 Redis 記憶體大幅增長

<forkfork> 由訂閱者在訊息佇列中自行記錄位置,而不是由 Redis 來維護每個消費者的進度並為每個訂閱者複製訊息

這段話激發了我的想像。我思考了幾天,意識到這或許就是能一次解決上述所有問題的時刻。我需要的是重新想像「日誌」這個概念。它是程式設計中一個基本的元素,大家對它都很熟悉,因為它就跟以附加模式開啟一個檔案、並以某種格式寫入資料一樣簡單。然而,Redis 的資料結構必須是抽象的。它們位於記憶體中,而我們使用 RAM 不只是因為偷懶,而是因為運用幾個指標,我們就能將資料結構概念化並加以抽象,使其擺脫顯而易見的限制。舉例來說,一般的日誌有幾個問題:位移不是邏輯上的,而是實際的位元組位移,如果我們想要與項目插入時間相關的邏輯位移呢?這樣我們就能免費獲得範圍查詢。同樣地,日誌通常難以進行垃圾回收:如何在僅供附加的資料結構中移除舊元素?嗯,在我們理想化的日誌中,我們只要說最多想要保留這個數量的項目,舊的項目就會自動消失,依此類推。

當我嘗試以提摩西·唐斯的初始想法為基礎撰寫規格時,我正為 Redis Cluster 開發一個 Radix Tree(基數樹) 的實作,用來最佳化其內部的某些部分。這為實作一個非常節省空間、同時仍能以對數時間取得範圍資料的日誌奠定了基礎。與此同時,我開始研讀 Kafka Streams,以尋找其他能融入設計的有趣想法,這讓我吸收了 Kafka Consumer Groups(消費者群組) 的概念,並為 Redis 與記憶體內使用情境重新加以理想化。然而,這份規格在好幾個月內就只是一份規格,以至於過了一段時間後,我幾乎從頭重寫了一遍,用我在與人們討論這項即將加入 Redis 的新功能時所累積的許多提示來加以改進。我希望 Redis Streams 能成為特別適合 Time Series 的絕佳使用案例,而不僅僅適用於其他類型的事件與訊息傳遞應用。

來寫些程式碼吧

從 Redis Conf 回來後的夏天期間,我正在實作一個名為「listpack」的函式庫。這個函式庫其實就是 ziplist.c 的後繼者,也就是一種能在單一記憶體配置中表示字串元素列表的資料結構。它只是一種非常特殊的序列化格式,其特別之處在於也能以反向、從右至左的方式解析:這是在所有使用情境中取代 ziplists 所需要的特性。

結合 Radix Tree 與 listpack,就能輕鬆打造出一個同時非常節省空間且具備索引的日誌,也就是說,能夠透過 ID 與時間進行隨機存取。完成這部分後,我便開始撰寫程式碼來實作 Streams 資料結構。我仍在完成實作的最後階段,不過在此刻,Redis 在 Github 上的「streams」分支中已有足夠的功能可以開始把玩了。我不敢說 API 已 100% 定案,但有兩個值得注意的事實:其一是到了這個階段,只剩下 Consumer Groups 以及若干較不重要的操作 Streams 的指令尚未完成,但所有重要的部分皆已實作。其二是決定在約兩個月後、一切看來穩定時,將所有 Streams 相關的工作回移植到 4.0 分支。這意味著 Redis 使用者不必等到 Redis 4.2 就能使用 Streams,它們將盡快可供正式環境使用。這之所以可行,是因為作為一個全新的資料結構,幾乎所有的程式碼變更都侷限在新的程式碼中。唯一的例外是列表的阻塞操作:程式碼經過重構,讓 Streams 與列表的阻塞操作共用同一份程式碼,大幅簡化了 Redis 的內部實作。

教學:歡迎使用 Redis Streams

在某種程度上,你可以把 Streams 想成是 Redis 列表的強化版。Streams 的元素不只是一個單一字串,而是由欄位與值所組成的物件。範圍查詢是可行且快速的。Streams 中的每一筆項目都有一個 ID,也就是邏輯位移。不同的客戶端可以阻塞等待 ID 大於指定值的新元素。Redis Streams 的一個基礎指令是 XADD。沒錯,所有 Redis Streams 指令都是以「X」為前綴。

> XADD mystream * sensor-id 1234 temperature 10.5
1506871964177.0

XADD 指令會將指定的項目作為新元素附加到指定的 Streams「mystream」中。以上例而言,該項目有兩個欄位:sensor-id 與 temperature,然而同一個 Streams 中的每個項目都可以擁有不同的欄位。使用相同的欄位名稱只會帶來更好的記憶體使用率。另一個有趣的地方是,欄位的順序保證會被保留。XADD 會回傳剛插入項目的 ID,因為我們在第三個參數中使用星號,要求指令自動產生 ID。這幾乎總是你想要的行為,但也可以強制指定特定的 ID,例如為了將指令複製到備庫與 AOF 檔案。

ID 由兩個部分組成:毫秒時間與序號。1506871964177 是毫秒時間,也就是具有毫秒精度的 Unix 時間。小數點後的數字 0 是序號,用來區分在同一毫秒內加入的項目。兩個數字皆為 64 位元無號整數。這意味著我們可以在同一個 Streams 中加入任意數量的項目,即使是在同一毫秒內也沒問題。ID 的毫秒部分是取產生該 ID 的 Redis 伺服器當前本地時間與 Streams 中最後一筆項目兩者中的最大值而得。因此,即使電腦時鐘往回跳,ID 仍會持續遞增。在某種程度上,你可以把 Streams 項目的 ID 想成一個完整的 128 位元數字。然而,正因為它們與加入時所在執行個體的本地時間有所關聯,意味著我們能免費獲得毫秒精度的範圍查詢。

如你所料,以非常快速的方式加入兩筆項目,結果只會讓序號遞增。我們可以簡單地用一個 MULTI/EXEC 區塊來模擬「快速插入」:

> MULTI
OK
> XADD mystream * foo 10
QUEUED
> XADD mystream * bar 20
QUEUED
> EXEC
1) 1506872463535.0
2) 1506872463535.1

上述範例也顯示了我們如何在不需要事先指定任何結構描述的情況下,為不同項目使用不同的欄位。然而實際運作方式是,每個區塊(通常包含約 50 至 150 則訊息)的第一則訊息會被作為參考,而後續具有相同欄位的項目會以單一旗標壓縮,該旗標表示「與此區塊中第一筆項目的欄位相同」。因此,確實在連續訊息中使用相同欄位可節省大量記憶體,即使欄位集合會隨著時間緩慢變動也是如此。

要從 Streams 中擷取資料有兩種方式:由 XRANGE 指令實作的範圍查詢,以及由 XREAD 指令實作的 Streams 讀取。XRANGE 只是從起始到結束(含頭尾)擷取一個範圍內的項目。舉例來說,如果我知道某個項目的 ID,就可以用以下方式擷取單一項目:

> XRANGE mystream 1506871964177.0 1506871964177.0
1) 1) 1506871964177.0
   2) 1) "sensor-id"
      2) "1234"
      3) "temperature"
      4) "10.5"

不過,你也可以使用特殊起始符號「-」與特殊結束符號「+」來表示最小與最大的可能 ID。也可以使用 COUNT 選項來限制回傳的項目數量。以下是一個較複雜的 XRANGE 範例:

> XRANGE mystream - + COUNT 2
1) 1) 1506871964177.0
   2) 1) "sensor-id"
      2) "1234"
      3) "temperature"
      4) "10.5"
2) 1) 1506872463535.0
   2) 1) "foo"
      2) "10"

在這裡我們是以 ID 範圍的角度來思考,不過你也可以使用 XRANGE 來取得在特定時間範圍內的元素,因為你可以省略 ID 的「序號」部分。所以你可以直接指定以毫秒為單位的時間。以下指令的意思是:「從 Unix 時間 1506872463 開始,給我 10 筆項目」:

127.0.0.1:6379> XRANGE mystream 1506872463000 + COUNT 10
1) 1) 1506872463535.0
   2) 1) "foo"
      2) "10"
2) 1) 1506872463535.1
   2) 1) "bar"
      2) "20"

關於 XRANGE 最後一個重要的注意事項是,由於我們會在回覆中收到 ID,而緊接的下一個 ID 只要將 ID 的序號部分遞增即可輕易取得,因此可以使用 XRANGE 以增量方式遍歷整個 Streams,每次呼叫取得指定數量的元素。在 Redis 中 *SCAN 系列指令讓我們得以遍歷原本並非為遍歷而設計的資料結構之後,我避免再次犯同樣的錯誤。

使用 XREAD 進行 Streams:阻塞等待新資料

當我們想要透過 ID 或時間取得範圍資料,或透過 ID 取得單一元素時,XRANGE 非常適合。然而,對於不同客戶端必須在資料到達時即時取用的 Streams 來說,這樣還不夠好,且會需要某種形式的輪詢(對於某些只是偶爾連線來取得資料的應用程式而言,這或許是個不錯的主意)。

XREAD 指令的設計目的是讓你同時從多個 Streams 讀取資料,只需指定我們已取得的 Streams 中最後一筆項目的 ID。此外,我們還可以要求在沒有資料時進行阻塞,直到有資料到達時再解除阻塞。這與列表的阻塞操作類似,但在此資料不會從 Streams 中被消耗,且多個客戶端可以同時存取相同的資料。

以下是一個典型的 XREAD 呼叫範例:

> XREAD BLOCK 5000 STREAMS mystream otherstream $ $

它的意思是:從「mystream」與「otherstream」取得資料。如果沒有可用資料,就將客戶端阻塞,逾時時間為 5000 毫秒。在 STREAMS 選項之後,我們指定想要監聽的鍵以及我們已擁有的最後一個 ID。然而,特殊的 ID「$」表示:假設我已擁有 Streams 中目前所有的元素,因此只從下一個到達的元素開始給我資料。

如果我從另一個客戶端發送以下指令:

> XADD otherstream * message "Hi There"

這是在 XREAD 那一端會發生的情況:

1) 1) "otherstream"
   2) 1) 1) 1506935385635.0
         2) 1) "message"
            2) "Hi There"

我們會取得收到資料的鍵以及所收到的資料。在下一次呼叫中,我們很可能會使用最後收到的訊息 ID:

> XREAD BLOCK 5000 STREAMS mystream otherstream $ 1506935385635.0

依此類推。然而請注意,採用這種使用模式時,客戶端有可能在經過很長的延遲後才再次連線(因為處理訊息花了時間,或基於其他任何原因)。在這種情況下,期間可能會堆積大量訊息,因此明智的做法是始終在 XREAD 中使用 COUNT 選項,以確保客戶端不會被訊息淹沒,且伺服器也不必花費太多時間僅為單一客戶端提供大量訊息。

容量受限的 Streams

到目前為止都很順利……不過 Streams 在某個時間點必須移除舊訊息。幸好這可以透過 XADD 指令的 MAXLEN 選項來達成:

> XADD mystream MAXLEN 1000000 * field1 value1 field2 value2

這基本上意味著,如果在加入新元素後發現 Streams 擁有超過 100 萬則訊息,就移除舊訊息,使長度回到 100 萬個元素。這就像在列表中使用 RPUSH + LTRIM 一樣,但這次我們有內建機制來達成。不過請注意,上述做法意味著每次加入新訊息時,也必須負擔從 Streams 另一端移除訊息所需的工作。這會消耗一些 CPU,因此可以在 MAXLEN 的數量前使用「~」符號,以表示我們並非真的要求*恰好* 100 萬則訊息,多一些也無妨:

> XADD mystream MAXLEN ~ 1000000 * foo bar

透過這種方式,XADD 只有在能夠移除整個節點時才會移除訊息。這使得維護容量受限的 Streams 相較於普通的 XADD 幾乎不需額外成本。

Consumer Groups(開發中)

這是 Redis 中尚未實作、仍在開發中的第一項功能。它也是明顯受到 Kafka 啟發的想法,儘管在此處的實作方式相當不同。重點在於,使用 XREAD 時,客戶端還可以加上「GROUP <name>」選項。同一個群組中的所有客戶端將自動取得*不同*的訊息。當然,可能會有多個群組同時從同一個 Streams 讀取,在這種情況下,所有群組都會收到 Streams 中新到訊息的重複副本,但在每個群組內部,訊息不會重複。

對群組的延伸是,當指定群組時,將可以指定「RETRY <milliseconds>」選項:在此情況下,若訊息未透過 XACK 確認已處理,它們將在指定的毫秒數後再次投遞。這為訊息的投遞提供了某種盡力而為的可靠性,以防客戶端沒有私有的方式來標記訊息已處理。這部分同樣仍在開發中。

記憶體使用量與儲存/載入時間

由於用來建模 Redis Streams 的設計,其記憶體使用量非常低。這取決於欄位、值及其長度的數量,但對於簡單的訊息,每 100 MB 的已用記憶體可容納數百萬則訊息。此外,該格式被設計為只需極少量的序列化:作為 Radix Tree 節點儲存的 listpack 區塊在磁碟與記憶體中具有相同的表示方式,因此儲存與讀取都非常簡單。舉例來說,Redis 能在 0.3 秒內從 RDB 檔案讀取 500 萬筆項目。這使得 Streams 的複寫與持久化非常有效率。

未來也計畫允許刪除中間的項目。此功能目前僅部分實作,但策略是在項目的旗標中將項目標記為已刪除,當項目與已刪除項目之間的比例達到特定門檻時,就會重寫該區塊以回收垃圾,並在需要時將其與相鄰的另一個區塊黏合,以避免碎片化。

結論與預計時程

Redis Streams 將在年底前成為 Redis 4.0 穩定版的一部分。我認為這個通用型資料結構將大幅填補缺口,讓 Redis 能夠涵蓋許多過去難以涵蓋的使用情境:也就是說,過去你必須發揮創意、濫用現有資料結構來解決某些問題。其中一個非常重要的使用情境是 Time Series,但我的感覺是,透過 TREAD 為其他使用情境進行的訊息 Streams 也將非常有趣,無論是作為需要比發送即忘更高可靠性的 Pub/Sub 應用程式的替代方案,或是用於全新的使用情境皆然。就目前而言,如果你想開始針對自身問題評估這些新功能,只要在 Github 上取得「streams」分支並開始試用即可。畢竟,我們非常歡迎錯誤回報 :-)

如果你喜歡影片,這裡有一場即時展示 Streams 的課程:https://www.youtube.com/watch?v=ELDzy9lCFHQ

原文由 Salvatore Sanfilippo 發布

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