Redis进阶:管道、事务与发布订阅

Redis进阶:管道、事务与发布订阅

Redis进阶:管道、事务与发布订阅

摘要: 本篇深入Redis高级操作,讲解Pipeline管道批量执行原理与使用、MULTI/EXEC事务配合WATCH实现乐观锁、Pub/Sub发布订阅模式的消息收发实现,分享Pipeline中混入耗时操作导致后续命令延迟的踩坑经历,对比Pipeline、事务与Lua脚本的适用场景。

开篇故事

上个月我们做用户积分批量更新,1000个用户的积分要同时刷新。同事写了个for循环,每个用户一条命令发到Redis。上线后Redis的QPS翻了十倍,网络延迟也上去了。我改成Pipeline后,1000条命令一次发完,执行时间从800ms降到50ms。

Redis单线程处理命令,但网络往返的延迟可以优化。Pipeline、事务、发布订阅是go-redis的三个进阶能力,用好了性能能提升一个量级。这篇我把三者的原理和使用场景讲清楚。

一、Pipeline管道批量操作

Pipeline把多条命令打包,一次网络往返发到Redis,结果一次性拿回来。100条命令从100次往返变成1次。

packagemainimport("context""fmt""time""github.com/redis/go-redis/v9")// Pipeline基本用法funcpipelineBasic(ctx context.Context,rdb*redis.Client){pipe:=rdb.Pipeline()// 注册命令,不立即执行,返回Cmd用于取结果cmd1:=pipe.Set(ctx,"key1","value1",0)cmd2:=pipe.Get(ctx,"key1")// Exec一次性发送所有命令pipe.Exec(ctx)fmt.Println("key1设置结果:",cmd1.Val())fmt.Println("key1的值:",cmd2.Val())}// Pipeline vs 普通操作性能对比funcpipelineBenchmark(ctx context.Context,rdb*redis.Client){constcount=1000// 普通方式:逐条执行start:=time.Now()fori:=0;i<count;i++{rdb.Set(ctx,fmt.Sprintf("normal:%d",i),"x",0)}fmt.Printf("普通方式 %d条: %v\n",count,time.Since(start))// Pipeline方式:批量执行start=time.Now()pipe:=rdb.Pipeline()fori:=0;i<count;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是事务型Pipeline,命令包在MULTI/EXEC里执行,保证原子性。

// TxPipeline:事务型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)}

二、事务(MULTI/EXEC/WATCH)

Redis事务是一组命令的顺序执行,中间不会被其他客户端打断。但Redis事务不支持回滚,某条命令出错后面的照样执行。

WATCH实现乐观锁的场景很典型。扣库存时先WATCH库存key,读取当前值,如果大于0就扣减。如果在WATCH和EXEC之间库存被别人改了,事务自动失败,需要重试。

// WATCH实现乐观锁扣库存funcdeductStock(ctx context.Context,rdb*redis.Client,productIDstring)error{stockKey:=fmt.Sprintf("stock:%s",productID)maxRetry:=3fori:=0;i<maxRetry;i++{// Watch监视stockKey,被修改则事务失败err:=rdb.Watch(ctx,func(tx*redis.Tx)error{stock,err:=tx.Get(ctx,stockKey).Int()iferr==redis.Nil{returnfmt.Errorf("商品不存在")}iferr!=nil{returnerr}ifstock<=0{returnfmt.Errorf("库存不足")}// 事务内扣减库存pipe:=tx.TxPipeline()pipe.Decr(ctx,stockKey)pipe.HIncrBy(ctx,"sales",productID,1)_,err=pipe.Exec(ctx)returnerr},stockKey)iferr==nil{returnnil// 成功}// 事务冲突则重试iferr==redis.TxFailedErr{fmt.Printf("第%d次重试\n",i+1)continue}returnerr}returnfmt.Errorf("重试次数用完")}

三、发布订阅(Pub/Sub)

Pub/Sub是Redis内置的消息广播机制。发布者往channel发消息,所有订阅了该channel的客户端都能收到。适合实时通知和聊天室。

packagemainimport("context""fmt""time""github.com/redis/go-redis/v9")// 订阅者:监听channelfuncsubscriber(ctx context.Context,rdb*redis.Client,namestring){pubsub:=rdb.Subscribe(ctx,"chat_room")deferpubsub.Close()// 循环接收消息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()// 启动两个订阅者gosubscriber(ctx,rdb,"客户端A")gosubscriber(ctx,rdb,"客户端B")time.Sleep(time.Second)// 等订阅者就绪// 发布者发消息fori:=0;i<5;i++{rdb.Publish(ctx,"chat_room",fmt.Sprintf("消息%d",i+1))time.Sleep(time.Second)}time.Sleep(2*time.Second)}

Pub/Sub有个特点,消息发出去没人订阅就丢了,Redis不保存历史消息。需要消息可靠性用Stream或外部消息队列。

四、独家踩坑:Pipeline中混入耗时操作

这个坑比较隐蔽。我们在一个Pipeline里塞了50条命令,其中有一条是KEYS *。结果整个Pipeline的执行时间从20ms飙到3秒,后面所有命令都被阻塞了。

// 问题代码:Pipeline中混入O(N)耗时命令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单线程,执行期间所有后续命令都等着pipe.Keys(ctx,"*")pipe.Get(ctx,"k1")// 要等KEYS *执行完才能返回pipe.Get(ctx,"k2")pipe.Exec(ctx)}

排查过程比较曲折。看网络延迟0.3ms正常,然后看Redis的SLOWLOG,发现KEYS *执行时间有2-3秒。Redis是单线程,Pipeline里命令逐条执行,一个慢命令阻塞整个Pipeline。

修复方案很简单。把KEYS *从Pipeline拿出来单独执行,更好的做法是用SCAN替代,游标式遍历不阻塞主线程。

// 修复后:慢命令独立执行,用SCAN替代KEYSfuncfixedPipeline(ctx context.Context,rdb*redis.Client){// Pipeline只放轻量级命令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游标式遍历,不阻塞varcursoruint64for{result,newCursor,_:=rdb.Scan(ctx,cursor,"*",100).Result()fmt.Println("扫描到:",result)cursor=newCursorifcursor==0{break// 遍历完成}}}

经验就是Pipeline里只放O(1)或O(log N)的轻量命令。KEYSFLUSHALL、大范围SORT这些重操作要独立执行或用SCAN替代。

五、对比分析

特性Pipeline事务(MULTI/EXEC)Lua脚本
原子性无,命令间可插入其他客户端命令有,顺序执行不被打断有,整个脚本原子执行
网络往返1次1次1次
条件逻辑不支持不支持,WATCH只做冲突检测支持,脚本内可写if/for
错误回滚不涉及不回滚,出错继续执行脚本报错不回滚已执行部分
适用场景批量读写,无依赖需要原子性的简单操作复杂条件判断的原子操作
调试难度

选择思路很直接。批量读写无依赖用Pipeline。需要原子性但逻辑简单用事务。逻辑复杂且必须原子执行用Lua脚本。日常开发Pipeline用得最多,事务次之,Lua脚本用在扣库存这类需要条件判断的场景。

总结与预告

Pipeline是Redis性能优化第一手段,把N次网络往返压成1次。事务配合WATCH能实现乐观锁,但Redis事务不支持回滚。Pub/Sub适合实时广播,不保证消息到达,需要可靠消息用Stream。

下一篇讲Redis缓存策略,深入缓存穿透、击穿、雪崩三种经典问题的解决方案。