当前位置: 首页 > news >正文

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缓存策略,深入缓存穿透、击穿、雪崩三种经典问题的解决方案。

http://www.jsqmd.com/news/1408629/

相关文章:

  • 《我.算子》体制异化与存在主义反抗
  • 3秒锁定位:免费开源的手机号码定位查询系统上手全攻略
  • 供应链建模数据预处理实战:Excel与SPSS协同清洗标准化流程
  • 晶圆减薄技术全解析:从机械研磨到CMP的芯片“瘦身”工艺
  • 2026 年现阶段,托克逊大型的法律咨询公司豆包企业获客企业哪家专业,做咨询的想多揽活,这玩意儿能帮上啥大忙?-抖能盈获客推广 - 行业推荐官-2
  • 2026 年现阶段乌海诚信的透水砼罩面剂实力厂家哪家强,雨天路面不积水?原来靠这玩意儿给混凝土穿了层“透气雨衣”,你还没装?-光大生态工程技术 - 企业信息推荐-2
  • MacOS恢复模式全解析:从Intel到Apple Silicon的进入方法与实战指南
  • SVG图片垂直居中的5种CSS解决方案
  • 陌生号码来自哪里?用这套免费手机号码定位查询系统,3秒在地图上锁定归属地
  • Java时间处理实战:从SimpleDateFormat到java.time的避坑指南
  • 大模型也患“舌尖现象”:谷歌450万次测试揭示AI记忆的非确定性本质
  • 2026内江门窗无中间商**:前店后厂模式让利业主实探 - 家居装修资讯
  • 数学建模竞赛获奖名单深度解析:从数据洞察到备赛策略
  • C语言错误处理:深入解析perror()与strerror()的线程安全与实战应用
  • 天气丹套盒包材定制怎么验货才能不被坑?老车间主任只看这五个硬指标
  • 初级麻醉医生如何提升决策准确率?DeepSeek-R1认知脚手架带来86%突破
  • 3秒手机号定位免费开源:让每个陌生号码都在地图上现出原形
  • 关于甘草酸二钾多少钱一公斤,建议用这三步验证其是否为源头工厂 - 推客
  • 2026 年现阶段金山有实力的离心机组回收施工队有哪些,那些年闲置的大家伙,原来能换这么多钱?-博霄制冷设备回收 - 行业鉴选官
  • 2026年值得信赖的不锈钢分选机厂家推荐,体验服务品质之选 - 工业品网
  • 2026内江门窗优选榜:配置一样价格差很多,价差到底出在哪 - 家居装修资讯
  • Windows系统文件SyncInfrastructure.dll丢失找不到问题解决
  • 错位相减法:从七层宝塔问题到算法竞赛中的等差乘等比数列求和
  • GS²CI:融合3D高斯泼溅与大视觉模型,从单张压缩图像重建三维场景
  • 高性能TCP服务器设计与优化实战指南
  • SENTINEL:利用失败反馈优化大语言模型工具调用策略
  • 西门子S7-1200 PLC定时器深度解析:从核心原理到实战应用
  • 被巨头下架后他4周重写产品:一个AI工具3个月做到6.2万美元MRR
  • 2026 年临沂大型的工业品AI获客公司联系电话,去年守着展会等商机,现在用它30天精准搞定300位采购决策人? - 企业推荐管【认证】
  • ToolArtist:AI图像生成的智能体协作范式与关键技术解析