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
隨機一篇部落格
留言
登入後參與討論