階:管道、事務(wù)與發(fā)布訂閱)
Redis進(jìn)階:管道、事務(wù)與發(fā)布訂閱摘要: 本篇深入Redis高級(jí)操作講解Pipeline管道批量執(zhí)行原理與使用、MULTI/EXEC事務(wù)配合WATCH實(shí)現(xiàn)樂(lè)觀鎖、Pub/Sub發(fā)布訂閱模式的消息收發(fā)實(shí)現(xiàn)分享Pipeline中混入耗時(shí)操作導(dǎo)致后續(xù)命令延遲的踩坑經(jīng)歷對(duì)比Pipeline、事務(wù)與Lua腳本的適用場(chǎng)景。開(kāi)篇故事上個(gè)月我們做用戶積分批量更新1000個(gè)用戶的積分要同時(shí)刷新。同事寫(xiě)了個(gè)for循環(huán)每個(gè)用戶一條命令發(fā)到Redis。上線后Redis的QPS翻了十倍網(wǎng)絡(luò)延遲也上去了。我改成Pipeline后1000條命令一次發(fā)完執(zhí)行時(shí)間從800ms降到50ms。Redis單線程處理命令但網(wǎng)絡(luò)往返的延遲可以優(yōu)化。Pipeline、事務(wù)、發(fā)布訂閱是go-redis的三個(gè)進(jìn)階能力用好了性能能提升一個(gè)量級(jí)。這篇我把三者的原理和使用場(chǎng)景講清楚。一、Pipeline管道批量操作Pipeline把多條命令打包一次網(wǎng)絡(luò)往返發(fā)到Redis結(jié)果一次性拿回來(lái)。100條命令從100次往返變成1次。packagemainimport(contextfmttimegithub.com/redis/go-redis/v9)// Pipeline基本用法funcpipelineBasic(ctx context.Context,rdb*redis.Client){pipe:rdb.Pipeline()// 注冊(cè)命令不立即執(zhí)行返回Cmd用于取結(jié)果cmd1:pipe.Set(ctx,key1,value1,0)cmd2:pipe.Get(ctx,key1)// Exec一次性發(fā)送所有命令pipe.Exec(ctx)fmt.Println(key1設(shè)置結(jié)果:,cmd1.Val())fmt.Println(key1的值:,cmd2.Val())}// Pipeline vs 普通操作性能對(duì)比f(wàn)uncpipelineBenchmark(ctx context.Context,rdb*redis.Client){constcount1000// 普通方式:逐條執(zhí)行start:time.Now()fori:0;icount;i{rdb.Set(ctx,fmt.Sprintf(normal:%d,i),x,0)}fmt.Printf(普通方式 %d條: %v\n,count,time.Since(start))// Pipeline方式:批量執(zhí)行starttime.Now()pipe:rdb.Pipeline()fori:0;icount;i{pipe.Set(ctx,fmt.Sprintf(pipe:%d,i),x,0)}pipe.Exec(ctx)fmt.Printf(Pipeline %d條: %v\n,count,time.Since(start))// Pipeline通常快10-20倍}TxPipeline是事務(wù)型Pipeline命令包在MULTI/EXEC里執(zhí)行保證原子性。// TxPipeline:事務(wù)型Pipeline命令要么全成功要么全失敗functxPipelineExample(ctx context.Context,rdb*redis.Client){pipe:rdb.TxPipeline()pipe.Incr(ctx,counter)pipe.Set(ctx,flag,done,0)pipe.Expire(ctx,counter,10*time.Minute)pipe.Exec(ctx)}二、事務(wù)(MULTI/EXEC/WATCH)Redis事務(wù)是一組命令的順序執(zhí)行中間不會(huì)被其他客戶端打斷。但Redis事務(wù)不支持回滾某條命令出錯(cuò)后面的照樣執(zhí)行。WATCH實(shí)現(xiàn)樂(lè)觀鎖的場(chǎng)景很典型。扣庫(kù)存時(shí)先WATCH庫(kù)存key讀取當(dāng)前值如果大于0就扣減。如果在WATCH和EXEC之間庫(kù)存被別人改了事務(wù)自動(dòng)失敗需要重試。// WATCH實(shí)現(xiàn)樂(lè)觀鎖扣庫(kù)存funcdeductStock(ctx context.Context,rdb*redis.Client,productIDstring)error{stockKey:fmt.Sprintf(stock:%s,productID)maxRetry:3fori:0;imaxRetry;i{// Watch監(jiān)視stockKey被修改則事務(wù)失敗err:rdb.Watch(ctx,func(tx*redis.Tx)error{stock,err:tx.Get(ctx,stockKey).Int()iferrredis.Nil{returnfmt.Errorf(商品不存在)}iferr!nil{returnerr}ifstock0{returnfmt.Errorf(庫(kù)存不足)}// 事務(wù)內(nèi)扣減庫(kù)存pipe:tx.TxPipeline()pipe.Decr(ctx,stockKey)pipe.HIncrBy(ctx,sales,productID,1)_,errpipe.Exec(ctx)returnerr},stockKey)iferrnil{returnnil// 成功}// 事務(wù)沖突則重試iferrredis.TxFailedErr{fmt.Printf(第%d次重試\n,i1)continue}returnerr}returnfmt.Errorf(重試次數(shù)用完)}三、發(fā)布訂閱(Pub/Sub)Pub/Sub是Redis內(nèi)置的消息廣播機(jī)制。發(fā)布者往channel發(fā)消息所有訂閱了該channel的客戶端都能收到。適合實(shí)時(shí)通知和聊天室。packagemainimport(contextfmttimegithub.com/redis/go-redis/v9)// 訂閱者:監(jiān)聽(tīng)channelfuncsubscriber(ctx context.Context,rdb*redis.Client,namestring){pubsub:rdb.Subscribe(ctx,chat_room)deferpubsub.Close()// 循環(huán)接收消息formsg:rangepubsub.Channel(){fmt.Printf([%s] %s: %s\n,name,msg.Channel,msg.Payload)}}funcmain(){rdb:redis.NewClient(redis.Options{Addr:localhost:6379})deferrdb.Close()ctx:context.Background()// 啟動(dòng)兩個(gè)訂閱者gosubscriber(ctx,rdb,客戶端A)gosubscriber(ctx,rdb,客戶端B)time.Sleep(time.Second)// 等訂閱者就緒// 發(fā)布者發(fā)消息fori:0;i5;i{rdb.Publish(ctx,chat_room,fmt.Sprintf(消息%d,i1))time.Sleep(time.Second)}time.Sleep(2*time.Second)}Pub/Sub有個(gè)特點(diǎn)消息發(fā)出去沒(méi)人訂閱就丟了Redis不保存歷史消息。需要消息可靠性用Stream或外部消息隊(duì)列。四、獨(dú)家踩坑:Pipeline中混入耗時(shí)操作這個(gè)坑比較隱蔽。我們?cè)谝粋€(gè)Pipeline里塞了50條命令其中有一條是KEYS *。結(jié)果整個(gè)Pipeline的執(zhí)行時(shí)間從20ms飆到3秒后面所有命令都被阻塞了。// 問(wèn)題代碼:Pipeline中混入O(N)耗時(shí)命令funcbadPipeline(ctx context.Context,rdb*redis.Client){pipe:rdb.Pipeline()pipe.Set(ctx,k1,v1,0)pipe.Set(ctx,k2,v2,0)// KEYS *是O(N)操作阻塞Redis主線程// Redis單線程執(zhí)行期間所有后續(xù)命令都等著pipe.Keys(ctx,*)pipe.Get(ctx,k1)// 要等KEYS *執(zhí)行完才能返回pipe.Get(ctx,k2)pipe.Exec(ctx)}排查過(guò)程比較曲折。看網(wǎng)絡(luò)延遲0.3ms正常然后看Redis的SLOWLOG發(fā)現(xiàn)KEYS *執(zhí)行時(shí)間有2-3秒。Redis是單線程Pipeline里命令逐條執(zhí)行一個(gè)慢命令阻塞整個(gè)Pipeline。修復(fù)方案很簡(jiǎn)單。把KEYS *從Pipeline拿出來(lái)單獨(dú)執(zhí)行更好的做法是用SCAN替代游標(biāo)式遍歷不阻塞主線程。// 修復(fù)后:慢命令獨(dú)立執(zhí)行用SCAN替代KEYSfuncfixedPipeline(ctx context.Context,rdb*redis.Client){// Pipeline只放輕量級(jí)命令pipe:rdb.Pipeline()pipe.Set(ctx,k1,v1,0)pipe.Set(ctx,k2,v2,0)pipe.Get(ctx,k1)pipe.Get(ctx,k2)pipe.Exec(ctx)// SCAN游標(biāo)式遍歷不阻塞varcursoruint64for{result,newCursor,_:rdb.Scan(ctx,cursor,*,100).Result()fmt.Println(掃描到:,result)cursornewCursorifcursor0{break// 遍歷完成}}}經(jīng)驗(yàn)就是Pipeline里只放O(1)或O(log N)的輕量命令。KEYS、FLUSHALL、大范圍SORT這些重操作要獨(dú)立執(zhí)行或用SCAN替代。五、對(duì)比分析特性Pipeline事務(wù)(MULTI/EXEC)Lua腳本原子性無(wú)命令間可插入其他客戶端命令有順序執(zhí)行不被打斷有整個(gè)腳本原子執(zhí)行網(wǎng)絡(luò)往返1次1次1次條件邏輯不支持不支持WATCH只做沖突檢測(cè)支持腳本內(nèi)可寫(xiě)if/for錯(cuò)誤回滾不涉及不回滾出錯(cuò)繼續(xù)執(zhí)行腳本報(bào)錯(cuò)不回滾已執(zhí)行部分適用場(chǎng)景批量讀寫(xiě)無(wú)依賴需要原子性的簡(jiǎn)單操作復(fù)雜條件判斷的原子操作調(diào)試難度低中高選擇思路很直接。批量讀寫(xiě)無(wú)依賴用Pipeline。需要原子性但邏輯簡(jiǎn)單用事務(wù)。邏輯復(fù)雜且必須原子執(zhí)行用Lua腳本。日常開(kāi)發(fā)Pipeline用得最多事務(wù)次之Lua腳本用在扣庫(kù)存這類需要條件判斷的場(chǎng)景。總結(jié)與預(yù)告Pipeline是Redis性能優(yōu)化第一手段把N次網(wǎng)絡(luò)往返壓成1次。事務(wù)配合WATCH能實(shí)現(xiàn)樂(lè)觀鎖但Redis事務(wù)不支持回滾。Pub/Sub適合實(shí)時(shí)廣播不保證消息到達(dá)需要可靠消息用Stream。下一篇講Redis緩存策略深入緩存穿透、擊穿、雪崩三種經(jīng)典問(wèn)題的解決方案。