Streams: a new general purpose data structure in Redis.

Salvatore Sanfilippo

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

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

直到幾個月前,對我來說 stream 不過就是訊息傳遞領域中一個有趣且相對單純的概念。Kafka 把這個概念發揚光大之後,我主要是在 Disque 這個情境下研究它的實用性——Disque 是一個訊息佇列,現在正計畫改寫成 Redis 4.2 的模組。後來我認定 Disque 的核心是 AP 訊息,也就是容錯能力與投遞保證,而且不需要客戶端費太多功夫,所以我覺得 stream 的概念在那個情境下並不合適。

然而,與此同時,Redis 本身有個問題一直讓我無法放心,那就是它預設提供的資料結構。在 Redis 的 list、sorted set 和 Pub/Sub 功能之間,存在著某種落差。你當然可以用這些工具來建模一串訊息或事件序列,但各有不同的取捨。Sorted set 很耗記憶體,沒辦法很自然地重複投遞同一則訊息,客戶端也無法阻塞等待新訊息。因為 sorted set 並不是循序的資料結構,它是一個集合,裡面的元素可以透過改變分數被搬來搬去:難怪它不適合拿來處理像時間序列這類資料。List 也有不同的問題,在某些使用情境下造成類似的適用性限制:你無法有效率地探查 list 中間的內容,因為那種存取的時間複雜度是線性的。再者,也無法做到扇出(fan-out),對 list 的阻塞操作一次只會把單一元素交給單一客戶端。List 裡也沒有固定的元素識別碼,讓你說:從那個元素開始把之後的東西給我。對於一對多的 workload,有 Pub/Sub,在很多情況下表現很好,但有些情境你並不想要即發即棄(fire-and-forget):保留歷史紀錄很重要,不只是為了在斷線後重新取得訊息,也因為某些訊息序列本身——例如時間序列——就非常需要用範圍查詢來探索:我在這 10 秒內的溫度讀值是多少?

我曾試圖用一個辦法來解決上述問題,那就是把 sorted set 和 list 泛化、整合成一個更靈活的單一資料結構,但我的設計嘗試幾乎每次都讓結果變得比現有的結構更加生硬、不自然。Redis 有一個優點,就是它所提供的資料結構更像是自然的、電腦科學上的資料結構,而不是「這是 Salvatore 發明的 API」。所以最後,我停止了這些嘗試,心想,好吧,這就是我們目前能提供的東西,也許以後我會幫 Pub/Sub 加上一些歷史紀錄,或是讓 list 的存取模式更有彈性。然而,每當有使用者在研討會上走過來問我「你會怎麼在 Redis 裡建模時間序列?」或類似的問題時,我的臉就一陣發綠。

起源

在 Redis 4.0 引入模組之後,使用者開始嘗試自己動手解決這個問題。其中一位 Timothy Downs 透過 IRC 傳了以下訊息給我:

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

<forkfork> 由訂閱者自己記錄在訊息佇列中的位置,而不是由 Redis 來追蹤每個消費者讀到哪裡、並為每個訂閱者複製訊息

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

當我試圖以 Timothy 的原始想法為起點撰寫規格時,我正好在為 Redis Cluster 實作一個 radix tree,用來優化其內部某些部分。這為實作一個非常節省空間、同時仍能以對數時間取得範圍資料的 log 打下了基礎。同時,我也開始閱讀關於 Kafka streams 的資料,從中尋找能融入設計的其他有趣想法,結果就是借鑒了 Kafka 消費者群組(consumer groups)的概念,並針對 Redis 與記憶體內使用的情境重新加以理想化。不過,這份規格就這樣停留在紙面上好幾個月,到後來我幾乎是從頭重寫了一遍,為的是把這段時間與人們討論這個即將加入 Redis 的新功能時所累積的各種建議都整合進去。我希望 Redis streams 能成為一個非常好的使用案例,特別是針對時間序列,而不只是其他類型的事件與訊息應用。

來寫點程式碼吧

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

把 radix tree 和 listpack 結合起來,就能輕鬆打造出一個同時非常節省空間、又具備索引能力的 log,也就是說,可以透過 ID 和時間進行隨機存取。一旦這部分就緒,我就開始撰寫程式碼來實作 stream 資料結構。我仍在完成最後的實作,不過在這個時間點,Redis 在 GitHub 上的「streams」分支裡已經有足夠的內容可以開始試玩了。我不敢說 API 已經 100% 定案,但有兩個值得注意的地方:其一是到了這個階段,只剩下消費者群組還沒完成,再加上一些較不重要、用來操作 stream 的指令,其他重要的部分都已經實作好了。其二是決定在大約兩個月後、一切看起來穩定時,把所有關於 stream 的成果回溯移植到 4.0 分支。這意味著 Redis 使用者不必等到 Redis 4.2 才能使用 streams,它們會盡快開放給正式環境使用。這之所以可行,是因為作為一個全新的資料結構,幾乎所有的程式碼變動都侷限在新增的程式碼中。唯一的例外是阻塞式的 list 操作:相關程式碼經過重構,讓 streams 和 list 的阻塞操作共用同一份程式碼,也讓 Redis 的內部大幅簡化。

教學:歡迎來到 Redis Streams

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

> XADD mystream * sensor-id 1234 temperature 10.5
1506871964177.0

XADD 指令會把指定的項目以新元素的形式附加到指定的 stream「mystream」中。以上述範例來說,該筆資料有兩個欄位:sensor-id 和 temperature,不過同一個 stream 中的每一筆資料都可以有不同的欄位。使用相同的欄位名稱只會讓記憶體使用更有效率。另一個有趣的地方是,欄位的順序是保證會被保留的。XADD 會回傳剛插入資料的 ID,因為我們在第三個參數使用了星號,要求指令自動產生 ID。這幾乎就是你想要的用法,但也可以強制指定一個特定的 ID,例如為了把指令複製到從節點與 AOF 檔案。

ID 由兩個部分組成:毫秒時間與序號。1506871964177 是毫秒時間,其實就是具備毫秒精度的 Unix 時間。點號後面的數字 0 是序號,用來區分在同一毫秒內加入的資料。這兩個數字都是 64 位元的無號整數。這意味著我們可以在 stream 中加入任意數量的資料,即使是在同一毫秒內也沒問題。ID 的毫秒部分是取產生 ID 的 Redis 伺服器當下本地時間,與 stream 中最後一筆資料的時間兩者之間的較大值。因此,即使例如電腦時鐘往回跳,ID 仍然會保持遞增。在某種程度上,你可以把 stream 資料的 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 則訊息)中的第一則訊息會被當作參考基準,後續具有相同欄位的資料會被壓縮,只用一個旗標表示「與此區塊中第一筆資料的欄位相同」。所以,確實在連續的訊息中使用相同的欄位可以節省大量記憶體,即使欄位的組合隨著時間慢慢改變也一樣。

要從 stream 中取出資料有兩種方式:範圍查詢,由 XRANGE 指令實作;以及串流讀取,由 XREAD 指令實作。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 來逐步迭代整個 stream,每次呼叫都取得指定數量的元素。在 Redis 中有了 *SCAN 系列指令之後,它們讓我們得以迭代那些原本並非為迭代而設計的 Redis 資料結構,我這次避免再犯同樣的錯誤。

透過 XREAD 串流處理:阻塞等待新資料

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

XREAD 指令的設計目的是讓我們能同時從多個 streams 讀取資料,只要指定我們已取得的 stream 中最後一筆資料的 ID 即可。此外,我們還可以要求在沒有資料時進行阻塞,直到有資料到達時再解除阻塞。這跟阻塞式的 list 操作類似,但在這裡資料並不會從 stream 中被消耗掉,多個客戶端可以同時存取相同的資料。

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

> XREAD BLOCK 5000 STREAMS mystream otherstream $ $

它的意思是:從「mystream」和「otherstream」取得資料。如果沒有可用資料,就阻塞客戶端,逾時時間為 5000 毫秒。在 STREAMS 選項之後,我們指定想要監聽的鍵,以及我們手上最後一筆資料的 ID。不過,特殊的 ID「$」代表:假設我已經擁有該 stream 目前所有的元素,所以只要從下一個到達的元素開始給我即可。

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

> 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 選項,以確保客戶端不會被訊息淹沒,伺服器也不必花太多時間只為單一客戶端提供大量訊息。

長度受限的 Stream

到目前為止都很順利……不過 stream 總有需要移除舊訊息的時候。幸好,這可以透過 XADD 指令的 MAXLEN 選項來實現:

> XADD mystream MAXLEN 1000000 * field1 value1 field2 value2

這基本上意味著,如果在加入新元素後發現 stream 擁有超過 100 萬則訊息,就移除舊的訊息,讓長度回到 100 萬筆。這就跟在 list 中使用 RPUSH + LTRIM 一樣,但這次我們有內建的機制來做到這點。不過要注意的是,上述做法意味著每次加入新訊息時,也必須付出從 stream 另一端移除訊息所需的成本。這會消耗一些 CPU,因此可以在 MAXLEN 中的數量前使用「~」符號,來表示我們並非真的要求*恰好* 100 萬則訊息,多一點也沒什麼大礙:

> XADD mystream MAXLEN ~ 1000000 * foo bar

這樣一來,XADD 只會在能夠移除整個節點時才會刪除訊息。相比於一般的 XADD,這會讓使用長度受限的 stream 幾乎不需要額外成本。

消費者群組(開發中)

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

群組功能的一個延伸是,在指定群組時還可以指定「RETRY <milliseconds>」選項:在這種情況下,如果訊息沒有透過 XACK 被確認為已處理,就會在指定的毫秒數後再次投遞。這能在客戶端沒有私有機制來標記訊息已處理的情況下,為訊息的投遞提供某種盡力而為的可靠性。這部分同樣仍在開發中。

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

由於 Redis streams 所採用的設計,其記憶體使用量非常低。這取決於欄位、值及其長度的數量,但對於簡單的訊息來說,每使用 100 MB 記憶體就能存放數百萬則訊息。此外,這種格式被設計成只需要極少的序列化處理:作為 radix tree 節點儲存的 listpack 區塊,在磁碟和記憶體中具有相同的表示方式,因此儲存與讀取都非常輕鬆。舉例來說,Redis 能在 0.3 秒內從 RDB 檔案中讀取 500 萬筆資料。這讓 streams 的複製與持久化變得非常有效率。

未來也計畫允許刪除位於中間的項目。這部分目前只有部分實作,但策略是在項目的旗標中將該筆資料標記為已刪除,而當資料總數與已刪除資料數之間的比例達到一定門檻時,就會重寫該區塊以回收垃圾,必要時還會將其與相鄰的另一個區塊合併,以避免產生碎片化。

結論與預計時程

Redis streams 將在今年年底前成為 Redis 4.0 系列穩定版的一部分。我認為這個通用資料結構將大幅填補 Redis 在許多難以涵蓋的使用情境上的缺口:也就是說,過去你必須發揮創意、濫用現有的資料結構來解決某些問題。其中一個非常重要的使用情境是時間序列,但我的感覺是,透過 TREAD 為其他使用情境進行訊息串流也將會非常有趣,無論是作為需要比即發即棄更高可靠性的 Pub/Sub 應用的替代方案,或是用於全新的使用情境都一樣。目前,如果你想開始在你所面臨的問題中評估這些新功能,只要到 GitHub 上抓取「streams」分支並開始試玩即可。畢竟,我們非常歡迎錯誤回報 :-)

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

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

留言