、能回溯的消息隊(duì)列)
還記得 List 做隊(duì)列嗎BRPOP一把把消息拿走消費(fèi)者這時(shí)候掛了這條就沒了。Streams 不這么干。本次導(dǎo)航Stream 長(zhǎng)什么樣每條消息一個(gè) ID后面跟一堆字段讀寫XADD、XRANGE、XREAD消費(fèi)者組XREADGROUPXACK掛了還能撈回來Redis 8.6 的冪等投遞IDMPAUTO/IDMP和 Kafka 比什么時(shí)候用 Redis 就夠了發(fā)車前提醒Redis 核心數(shù)據(jù)結(jié)構(gòu)二——List 與消息隊(duì)列 這篇文章已經(jīng)把 List 當(dāng)隊(duì)列的坑寫過了這篇接著往下做。先用 docker 啟動(dòng)容器還是之前的容器dockerexec-itredis-demo redis-cli一、Stream 是一條只往后追加的日志別把它想成 ListList 的元素被POP之后就從 Key 里消失了Stream 更像一本流水賬新消息往后面貼舊的還在除非你主動(dòng)裁。每條消息兩個(gè)部分ID默認(rèn)是毫秒時(shí)間戳-序號(hào)比如1789086136869-0。你那邊數(shù)字肯定不一樣規(guī)則一樣。字段就是一堆 field-value跟 Hash 挺像。一條訂單可以帶user_id、amount、sku。127.0.0.1:6379XADD orders * user_id1001amount299sku sku_10011789086136869-0127.0.0.1:6379XADD orders * user_id1002amount88sku sku_20021789086136938-0127.0.0.1:6379XLEN orders(integer)2*的意思是ID 讓 Redis 自己生成。一般別手寫 ID時(shí)鐘回?fù)?、多?shí)例搶號(hào)都是麻煩。二、先當(dāng)普通日志讀按范圍翻127.0.0.1:6379XRANGE orders - 1)1)1789086136869-02)1)user_id2)10013)amount4)2995)sku6)sku_10012)1)1789086136938-02)1)user_id2)10023)amount4)885)sku6)sku_2002-是最小是最大。只看某一段就寫成XRANGE orders 1789086136869-0 1789086136869-0。從某個(gè) ID 之后接著讀用XREAD。0-0表示從最早開始$表示只看比現(xiàn)在更新的127.0.0.1:6379XREAD COUNT1STREAMS orders0-01)1)orders2)1)1)1789086136869-02)1)user_id2)1001...想卡著等新消息加上BLOCK單位毫秒0表示一直等。空流的時(shí)候它會(huì)掛起有人XADD才醒這點(diǎn)和BRPOP類似XREAD BLOCK10000STREAMS orders $Stream 可以在寫入時(shí)裁一刀XADD orders MAXLEN1000* user_id1003amount50sku sku_3003大概只留最近 1000 條。精確裁用MAXLEN量大時(shí)可以寫成MAXLEN ~ 1000允許 Redis 少裁幾條換更快的速度。老數(shù)據(jù)要?dú)w檔就XTRIM或者定時(shí)把XRANGE掃出來丟到別的地方。三、消費(fèi)者組同一條消息別搶兩遍XREAD是廣播味道的——每個(gè)客戶端自己記游標(biāo)兩個(gè)進(jìn)程都能讀到同一條。做任務(wù)隊(duì)列通常不是這個(gè)需求。你要的是一組工人搶活一條訂單只被一個(gè)人處理。這就是消費(fèi)者組。# 從頭開始讀歷史。只想收新消息把 0-0 換成 $127.0.0.1:6379XGROUP CREATE orders packer0-0 OK$和0-0搞反是新手常踩的坑。組已經(jīng)建了、流還不存在時(shí)后面加MKSTREAM讓 Redis 順手建一個(gè)空流。兩個(gè)工人分別叫node-a、node-b127.0.0.1:6379XREADGROUP GROUP packer node-a COUNT1STREAMS orders1)1)orders2)1)1)1789086136869-02)... user_id1001...127.0.0.1:6379XREADGROUP GROUP packer node-b COUNT1STREAMS orders1)1)orders2)1)1)1789086136938-02)... user_id1002...表示給我一條組里還沒分過的新消息。node-a拿走第一單node-b拿走第二單不會(huì)撞車。消息被讀走之后還在 Stream 里。Redis 另外記了一本賬誰(shuí)領(lǐng)了、還沒簽字。這本賬叫 PELPending Entries List。127.0.0.1:6379XPENDING orders packer1)(integer)22)1789086136869-03)1789086136938-04)1)1)node-a2)12)1)node-b2)1處理完了簽個(gè)字127.0.0.1:6379XACK orders packer1789086136869-0(integer)1XACK只是從 PEL 里劃掉不是刪消息。別的組還能再讀同一條你自己以后用XRANGE也能翻到。這就是 List 做不到的回溯。再建一個(gè)組試試127.0.0.1:6379XGROUP CREATE orders stock0-0 OK127.0.0.1:6379XREADGROUP GROUP stock warehouse-1 COUNT1STREAMS orders# 又拿到 1001 那單packer負(fù)責(zé)打包stock負(fù)責(zé)扣庫(kù)存各記各的進(jìn)度。組與組之間互不打擾。四、工人掛了活怎么轉(zhuǎn)出去node-a領(lǐng)了單還沒XACK就進(jìn)程沒了。消息不會(huì)丟它還掛在 PEL 里。別人可以用XPENDING看到再用XAUTOCLAIM認(rèn)領(lǐng)超時(shí)的單# 把 packer 組里 idle 超過 60000 毫秒的未確認(rèn)消息轉(zhuǎn)給 node-bXAUTOCLAIM orders packer node-b600000-0 COUNT10XCLAIM是指定 ID 硬搶XAUTOCLAIM6.2更省事按空閑時(shí)間批量撈。線上常見寫法工人循環(huán)XREADGROUP拿新活隔一會(huì)兒再XAUTOCLAIM掃一遍別人留下的殘骸。還沒 ACK 的時(shí)候同一個(gè)工人把 ID 寫成0再讀一次讀到的是自己 PEL 里的舊單不是新單。重啟之后先用這個(gè)把上次沒做完的活做掉再去搶。五、生產(chǎn)者重試別寫出兩條一樣的Redis 8.6 給XADD加了冪等。網(wǎng)絡(luò)抖一下客戶端重發(fā)以前會(huì)在流里留下兩條一模一樣的訂單。完整寫法要帶生產(chǎn)者 IDID 仍然用*127.0.0.1:6379XADD pay_log IDMPAUTO pay-svc-1 * order_id9001amount2991789086143272-0# 內(nèi)容一樣再發(fā)一次ID 還是這條流的長(zhǎng)度不變127.0.0.1:6379XADD pay_log IDMPAUTO pay-svc-1 * order_id9001amount2991789086143272-0127.0.0.1:6379XLEN pay_log(integer)1IDMPAUTO按消息內(nèi)容算一個(gè)內(nèi)部指紋內(nèi)容相同就當(dāng)重試服務(wù)重啟后同一個(gè)進(jìn)程要繼續(xù)用同一個(gè)生產(chǎn)者 ID上面的pay-svc-1換一個(gè)就失效了。內(nèi)容可能真的重復(fù)——比如用戶連點(diǎn)兩下兩筆金額一樣——就別用IDMPAUTO自己帶業(yè)務(wù)單號(hào)XADD pay_log IDMP pay-svc-1 txn-9001 * order_id9001amount299XADD pay_log IDMP pay-svc-1 txn-9001 * order_id9001amount299# 重試還是同一條XADD pay_log IDMP pay-svc-1 txn-9002 * order_id9001amount299# 另一筆會(huì)新增冪等記錄默認(rèn)只記大約 100 秒、每個(gè)生產(chǎn)者最多 100 個(gè) ID。隔太久再重試Redis 可能已經(jīng)忘了又會(huì)寫入一條。業(yè)務(wù)上重試窗口比較長(zhǎng)的話用XCFGSET調(diào)大XCFGSET pay_log IDMP-DURATION600IDMP-MAXSIZE1000注意XCFGSET會(huì)清掉當(dāng)前這條流上已經(jīng)記下的冪等指紋改完配置再跑你的重試邏輯。六、和 Kafka 比Kafka 分區(qū)、副本、消費(fèi)積壓、跨機(jī)房這些不是 Redis Streams 要做的事。Streams 吃的是 Redis 已經(jīng)在那、消息量不大、丟了會(huì)心疼但還沒到要單獨(dú)養(yǎng)一套 Kafka 的場(chǎng)景下單后發(fā)短信、扣積分、刷新庫(kù)存、服務(wù)內(nèi)部的任務(wù)隊(duì)列。幾個(gè)硬差別數(shù)據(jù)在 Redis 內(nèi)存里堆積能力看你給 Redis 的內(nèi)存和MAXLEN不是磁盤日志高可用靠 Redis 自己的復(fù)制 / 集群不是 Kafka 那套 ISR消費(fèi)者組夠用管理命令也少?zèng)]有 Broker、Topic、再均衡那一長(zhǎng)串運(yùn)維如果允許偶爾丟、邏輯又簡(jiǎn)單List 就夠了。要確認(rèn)、要重放、不要消息重復(fù)用 Streams。真要日千萬(wàn)級(jí)、還得按分區(qū)擴(kuò)去看 Kafka 或 RocketMQ不要逼 Redis 干這活。七、練手XADD三條訂單進(jìn)ordersXRANGE orders - 能看到。XGROUP CREATE orders packer 0-0兩個(gè)窗口分別用node-a、node-b做XREADGROUP確認(rèn)兩人拿到的 ID 不同。只 ACK 其中一條XPENDING里應(yīng)該還剩一條。Redis 8.6 的話用IDMPAUTO把同一條支付記錄發(fā)兩遍XLEN仍是 1。能走完這四步List 當(dāng)年那幾個(gè)坑基本都有著落了消息還在、能簽字、能換人接著干、重試不容易寫重。下期講講緩存的穿透、雪崩、擊穿。本次導(dǎo)航結(jié)束歡迎關(guān)注、點(diǎn)贊、轉(zhuǎn)發(fā)。