免费获取学习方案
ARTICLE DETAIL

资讯详情

深耕编程基础知识与建站技术分享的一线实战洞察。

Redis队列与阻塞队列:从原理到落地的完整指南

Redis队列与阻塞队列:从原理到落地的完整指南 最近手上的订单异步处理项目刚好改造成基于Redis的队列方案前后踩了不少坑从最初随手用List塞消息到后来换成Stream做可靠消费再到排查重复消费问题一路下来积累了不少值得记录的东西。顺手看了看热搜词发现“Redis队列和阻塞队列”“Redis Stream如何拉取消息”“线程池的阻塞队列选择”这些词热度一直不低说明很多人都在同一条路上摸爬滚打过。这篇文章就结合我的实际改造过程把Redis队列和阻塞队列从原理、选型到落地一次讲清楚。1. 先把概念理清楚Redis队列和阻塞队列分别是什么1.1 一分钟认清Redis里的“队列”形态Redis本身不是为消息队列设计的但它提供了List、Stream、Sorted Set这些数据结构恰好能被用来组成不同特性的队列。最基础的是List队列思路就是左进右出或者右进左出。比如我用LPUSH把任务塞到列表左边消费者用RPOP从右边取先入先出这就是一个标准的队列模型。之前很多老项目都是这么玩的简单直接几行命令就能跑起来。但随着使用加深你会感受到它的局限性消息没有确认机制消费者拿走即删如果处理过程中宕机了消息就丢了。后来Redis 5.0引入了Stream类型专门为消息队列这种场景设计支持消费者组、消息确认、持久化、Pending Entries ListPEL在可靠性上比List方案高一个档次。这也是我最终选择的方案后面会详细讲。从Redis 6.2开始官方也持续增强了Stream能力新增了XAUTOCLAIM等命令来处理长时间未确认的消息。整体趋势很明显Redis官方希望Stream能承担更重的消息队列职责。顺带一提这个队列和Windows那个“打印队列”报错完全是两码事别混淆。1.2 “阻塞队列”的两种理解很多人搞混了我在团队里问过一圈发现大家对“阻塞队列”这四个字的理解完全不一样。一种是Redis自带的阻塞读取命令比如BRPOP、BLPOP。它们的特点是在列表为空时客户端连接会一直挂在那里等待直到有新消息进来或者超过设定的超时时间。这种阻塞不是Redis端的阻塞而是“客户端的行为被挂起”本质上是一种高效的长轮询机制避免消费者空转疯狂请求Redis。另一种是Java并发包里的BlockingQueue像ArrayBlockingQueue、LinkedBlockingQueue、SynchronousQueue它们解决的是JVM内部线程之间的协作问题。线程池的任务队列就是典型场景生产者线程放任务如果队列满了就阻塞消费者线程取任务如果队列空了也阻塞。这两种概念经常被放到一起聊是因为它们解决的问题非常相似都是为了削峰填谷、解耦生产者和消费者。但作用域完全不同一个是跨进程、跨服务的分布式队列一个是单机进程内的线程协调器。明白这一点后面看各种技术方案就不会懵。1.3 什么时候该用Redis队列什么时候该用线程池阻塞队列我的判断标准很简单先看生产者和消费者是不是在同一个进程里。如果是同一个JVM内的线程间协作比如一个线程池接收任务另一个线程池处理任务直接用Java的BlockingQueue就够了压根不需要引入Redis。你想象一下两个线程就在一个进程里消息还要先写到Redis再读回来一次序列化加一次网络往返白白增加好几毫秒延迟完全没有必要。如果生产者和消费者分布在不同的服务节点上那就必须用Redis队列或者真正的消息队列中间件了。比如用户下单后订单服务要把消息推送给积分服务、短信服务这几个服务部署在不同的机器上进程内的BlockingQueue连消息都传不过去这时候就需要一个集中式的队列来做中转。Redis队列适合消息量中等、对可靠性要求不是极端苛刻、又不想额外引入Kafka/RabbitMQ这类重组件的场景。2. 基于Redis List的队列最简单但坑也不少2.1 核心命令拆解LPUSH和BRPOP的正确用法List队列的核心就是一对命令生产者LPUSH消费者BRPOP。为什么一左一右而不是两个都用同方向很简单如果都用LPUSH那消费端也拿到的是最新消息就变成栈结构了后进先出这不是队列。我项目的早期版本就是这样的// 生产者往队列左边塞消息 stringRedisTemplate.opsForList().leftPush(order:queue, orderJson); // 消费者从队列右边阻塞取出 String order stringRedisTemplate.opsForList() .rightPop(order:queue, 5, TimeUnit.SECONDS); if (order ! null) { process(order); }rightPop的第二个参数是超时时间这里设置5秒意思是如果队列里没有消息最多阻塞5秒超过就返回null。如果传0在Redis命令层面表示永久阻塞但Spring Data Redis里传0我建议慎用因为可能在部分版本上导致非预期行为我一般习惯给个明确超时再配合循环。2.2 阻塞读取的实现细节和参数选择BRPOP命令还支持一次监听多个key比如BRPOP order:queue notify:queue sms:queue 5Redis会从左到右依次检查这些列表哪个先有数据就弹出哪个。这个特性可以用来做简单的优先级效果把高优先级的队列放在前面低优先级的放在后面。在Java里用Spring Data Redis实现多个队列阻塞读取可以这样写ListString queueKeys Arrays.asList(high:queue, normal:queue, low:queue); ListString result stringRedisTemplate.opsForList() .rightPop(queueKeys, 5, TimeUnit.SECONDS); if (result ! null) { // result.get(0) 是队名result.get(1) 是弹出来的消息 String queueName result.get(0); String message result.get(1); }这里有个经验点阻塞时间不要设置太长也不要太短。太长了Redis连接长时间占用如果客户端有连接池限制可能把连接池耗尽太短了比如100ms消费者会频繁空轮询白白消耗CPU和网络资源。我个人的经验值是3到5秒再配合外部循环既能快速响应新消息又不会太消耗资源。2.3 可靠性这个方案最大的坑用List做队列最大的问题就是没有确认机制。BRPOP一执行消息立刻从列表里消失了如果消费者在process(order)这一步崩溃这条消息就永久丢失了。我早期做的一个内部通知系统就栽在这里。消费者把用户通知消息从队列里取出来后在组装模版的环节抛了异常消息已经没了重试都无从谈起。后来查了半天发现是消息里有一批脏数据导致模版解析失败。要缓解这个问题可以用BRPOPLPUSH命令它在弹出消息的同时把消息内容存到另一个“备份队列”里处理成功后手动删除备份。Redis 6.2以后官方推荐用BLMOVE替代BRPOPLPUSH功能一样但是更灵活。// 从order:queue弹出消息同时备份到order:queue-backup String message stringRedisTemplate.opsForList() .rightPopAndLeftPush(order:queue, order:queue-backup, 5, TimeUnit.SECONDS); try { process(message); // 处理成功后从备份队列删除 stringRedisTemplate.opsForList().remove(order:queue-backup, 1, message); } catch (Exception e) { // 处理失败消息还在备份队列里可以后续补偿 }这样相当于实现了一个简易的“未确认消息”机制。但代码复杂度上来了而且备份队列里的消息没有过期时间如果一直处理失败会越积越多还得额外写一套扫描补偿脚本。这也是我后来彻底切换到Stream方案的根本原因。3. Spring Boot集成Redis Stream生产级队列方案3.1 为什么Stream比List更适合做生产环境队列Stream用起来比List复杂一点但换来的可靠性完全值得尤其适合订单、支付这类不允许丢消息的业务。Stream有几个核心特性我一个个说。第一是消息持久化。Stream里的消息会存在Redis内存中根据配置可以通过AOF和RDB持久化到磁盘。Redis重启后消息还能恢复这是List完全不具备的。第二是消费者组。同一个Stream可以被多个消费者组订阅组内消息被竞争消费组间消息互相隔离。这个模型和Kafka的消费者组非常像如果你用过Kafka上手会很快。第三是消息确认机制。消费者取到消息后消息不会立刻被删除而是进入Pending Entries ListPEL。处理成功后消费者发XACK命令确认消息才会被真正标记为已处理。如果消费者处理失败或者崩溃了消息留在PEL里后续可以用XAUTOCLAIM重新认领。下面是我理解的两个方案对比能力维度List方案Stream方案消息持久化随Redis持久化但无独立结构独立数据结构支持AOF/RDB消费确认无取出即删XACK确认失败可重领消费者组不支持原生支持消息追溯不支持支持按ID范围读取实现复杂度低中适合场景允许少量丢消息的非核心业务订单、支付等核心链路3.2 生产者端用StreamRecords发送消息在Spring Boot项目中生产端代码很简单。我用StringRedisTemplate操作Stream消息体直接放JSON字符串。Autowired private StringRedisTemplate stringRedisTemplate; public void sendOrderMessage(String orderJson) { stringRedisTemplate.opsForStream().add( StreamRecords.newRecord() .ofObject(orderJson) .withStreamKey(order:stream) ); }ofObject方法可以接收任意对象内部会通过RedisSerializer序列化。我建议直接存JSON字符串这样在Redis Desktop Manager或者Redis Insight里排查问题时能直接看到原始内容不用反序列化。如果直接存Java对象默认JDK序列化会存成二进制排障的时候你会崩溃。Stream支持自动生成消息ID默认是毫秒时间戳加序号比如1735689600000-0。也可以自己指定ID但一般不建议自动生成的ID天然有时序性后面按ID范围拉消息会很方便。3.3 消费者端消费者组和手动ACK消费者端的核心是组的概念。首次创建一个消费者组的命令是XGROUP CREATE order:stream group-a 0在Spring Boot里可以通过StreamOperations来创建StreamOperationsString, Object, Object streamOps stringRedisTemplate.opsForStream(); try { streamOps.createGroup(order:stream, group-a); } catch (RedisSystemException e) { // 分组已存在时会报错这里不做处理即可 }消费的代码我用的是StreamReadOptions手动读的方式ListMapRecordString, Object, Object records streamOps.read( Consumer.from(group-a, consumer-1), StreamReadOptions.empty() .count(10) .block(Duration.ofSeconds(5)), StreamOffset.create(order:stream, ReadOffset.lastConsumed()) ); for (MapRecordString, Object, Object record : records) { try { // 拿到消息内容 MapObject, Object value record.getValue(); // 你的业务处理逻辑 processOrder(value.get(orderJson).toString()); // 处理成功后确认消息 streamOps.acknowledge(order:stream, group-a, record.getId()); } catch (Exception e) { log.error(处理消息失败消息ID: {}, record.getId(), e); // 不ack消息会停留在PEL中后续可重试 } }这里的ReadOffset.lastConsumed()表示读取本组中上次未确认的消息及之后的新消息通常配合使用。在Spring Data Redis中ReadOffset.lastConsumed()对应的就是语义ReadOffset.from(0)则是从头读。有个关键点成功处理的标志是acknowledge调用成功而不是processOrder执行完。所以ack一定要放在try块里而且是处理成功之后。我之前见过有人把ack写在processOrder前面结果消息处理失败还是被确认了等于白加了一套机制。3.4 重复消费问题为什么Stream也会重复很多人以为用了Stream的ack机制就不会重复消费了其实恰恰相反Stream的ack保证的是“消息不丢”不是“消息不重”。什么时候会重复最常见的情况是消费者取到消息执行完业务逻辑但是还没来及ackJVM就宕机了。这时候消息在PEL里还是未确认状态组内其他消费者通过XAUTOCLAIM或者重启后再次lastConsumed()读到这条消息于是又处理了一遍。所以消息队列的“至少一次投递”语义在Stream里体现得淋漓尽致。要解决重复处理唯一可靠的办法是让消费逻辑具备幂等性。我常用的方案是在业务表里加一个message_id字段并建唯一索引。处理消息时先把messageId插入业务表如果插入冲突说明处理过了直接跳过。另外也可以用Redis的SETNX做一个轻量幂等标记String messageId record.getId().getValue(); Boolean firstProcess stringRedisTemplate.opsForValue() .setIfAbsent(order:processed: messageId, 1, Duration.ofHours(24)); if (Boolean.FALSE.equals(firstProcess)) { log.info(消息重复投递忽略处理messageId: {}, messageId); return; }这条命令的意思是如果key不存在就设置成功返回true如果key已经存在就返回false。24小时后key自动过期既防止重复又不会让Redis被历史消息ID占满。4. 队列的两个常见变种延迟队列与线程池阻塞队列选型4.1 用Redis ZSet实现一个延迟队列延迟队列在业务里很常见比如订单超时自动关闭、定时任务调度。用Redis的话最经典的方案是基于ZSetscore存的是任务预期执行的时间戳。入队的时候把任务内容放在ZSet的member里score设置为当前时间加上延迟时间// 订单下单后设置15分钟后自动关闭 long executeTime System.currentTimeMillis() 15 * 60 * 1000; redisTemplate.opsForZSet().add(order:delay, orderId, executeTime);消费端用一个线程循环扫描每次取出score小于等于当前时间的任务while (true) { SetString expiredOrders redisTemplate.opsForZSet() .rangeByScore(order:delay, 0, System.currentTimeMillis(), 0, 100); for (String orderId : expiredOrders) { // 先移除再执行防止多个消费者重复拿到同一个任务 Long removed redisTemplate.opsForZSet() .remove(order:delay, orderId); if (removed ! null removed 0) { // 只有真正移除成功的人才有权利执行 closeOrder(orderId); } } Thread.sleep(500); }这里的设计思想是“先删除后执行”。因为ZSet的rangeByScore只是读取不会把元素移出去如果两个消费者同时扫描到同一个任务就会重复执行。用remove的返回值做判断返回1说明这个任务是你删掉的你才有资格执行等于用Redis的原子操作实现了一个简单的分布式锁。这是一个比较实用的技巧但也要注意如果执行closeOrder失败任务已经从ZSet中删除了需要额外记录失败日志或者放入重试队列。4.2 线程池的阻塞队列怎么选一份基于实践的选择表再来看Java进程内的阻塞队列这块的热度在热搜词里排得很靠前说明很多人在做线程池调优。Java里常见的阻塞队列有这么几个阻塞队列特性典型使用场景ArrayBlockingQueue有界数组结构必须指定容量需要控制任务堆积量的场景LinkedBlockingQueue链表结构默认无界可指定容量线程池默认队列但要注意OOM风险SynchronousQueue不存储元素线程间直接交接CachedThreadPool默认队列PriorityBlockingQueue无界按优先级出队定时任务的调度场景DelayQueue无界延迟到期才可取定时任务、订单超时检查具体到Executors框架newFixedThreadPool用的是无界的LinkedBlockingQueue任务堆积不会让线程池拒绝任务但极可能导致内存溢出newCachedThreadPool用的是SynchronousQueue每一个任务都尝试创建新线程空闲线程60秒回收适合大量短任务。我的个人建议是生产环境尽量不要用Executors的快捷方法而是自己new ThreadPoolExecutor并显式指定有界队列比如new ArrayBlockingQueue(1000)然后配一个合适的拒绝策略。原因很简单有界队列可以作为一个天然的背压机制队列满了以后拒绝策略会触发告警让你能意识到系统处理能力已经跟不上了。无界队列反而会把这个信号藏起来直到内存爆掉才后悔莫及。4.3 Redis分布式锁和队列的联动坑说到分布式锁热搜词里也有我多说两句。有些人为了防止队列消息被多个消费者同时处理会先加一把Redis分布式锁再执行业务逻辑。这个想法本身没错但要注意锁的粒度。选择用整个Stream或者队列的key做锁粒度太粗会导致同一时间只有一个消费者在处理消息整个队列就串行化了吞吐量大打折扣。正确做法是按消息里面的业务维度加锁比如订单号、用户ID这样不同用户的订单能被不同消费者并发处理同时同一个订单不会被两个人同时处理。锁的实现可以用Redisson的RLock它是开箱即用的支持看门狗自动续期不用担心锁超时导致的并发问题。RLock lock redissonClient.getLock(order:lock: orderId); if (lock.tryLock(3, TimeUnit.SECONDS)) { try { processOrder(orderId); } finally { lock.unlock(); } }这里有个坑用tryLock时第三个参数leaseTime如果省略Redisson会启动一个后台看门狗线程默认每10秒检查一次如果锁还在使用中就自动续期到30秒防止业务没执行完锁先过期了。如果自己传了leaseTime看门狗就失效了一定要根据业务耗时设置一个足够长的过期时间。5. 常见问题与排查实录5.1 消息推不出去、消费不动的排查路径我遇到过几次“消费者一直收不到消息”的情况一开始都在怀疑Redis配置最后发现原因五花八门。一个是消费者组的ReadOffset用错了。如果用了ReadOffset.from(0)消费者只会读取Stream创建以来的所有历史消息而且每次都是从最早的开始读读完了不会自动更新位点看起来就像消息没进来其实是卡在重复读老消息上了。正确做法是用lastConsumed()表示从当前消费位点继续往后读。另一个是Stream的key不存在。如果生产者和消费者启动顺序不一致消费者先启动了Stream还没创建XREADGROUP默认不会自动创建Stream会直接报错。解决办法是在消费者启动时先进行XGROUP CREATE或者用Redis 7.0新增的XREADGROUP ... MKSTREAM选项自动创建。在Spring Data Redis里我是在消费逻辑里判断Stream类型是否存在if (!stringRedisTemplate.hasKey(order:stream)) { stringRedisTemplate.opsForStream().createGroup(order:stream, group-a); }5.2 分页查询慢怎么用Redis优化热搜词里有“分页查询慢怎么用redis优化”这和队列其实是同一个父话题缓存治理。我顺便分享一个实际案例。之前有个查询接口数据量到了几百万行MySQL分页越翻越慢尤其是查到后面几页LIMIT 200000, 20要扫全索引耗时超过3秒。后来我把数据id列表全量放到了Redis的Sorted Set里score用业务排序字段比如创建时间分页用ZRANGEBYSCORE配合LIMIT来做// 建立索引缓存 redisTemplate.opsForZSet().add(order:page:list, orderId, createTime); // 分页查询从第1000条开始取20条 SetString ids redisTemplate.opsForZSet() .rangeByScore(order:page:list, 0, System.currentTimeMillis(), 1000, 20);这样MySQL只需要按id批量查20条完整记录不再需要深度分页扫描接口耗时从3秒降到了100毫秒以内。这不是队列的内容但思路是一样的把热点数据前置到Redis用Redis的表达力解决数据库的瓶颈。5.3 队列积压了怎么处理队列积压是消息队列运营中最常见的故障。我在一次活动大促时遇到过上游订单量激增Stream里堆积了上百万条消息消费者处理不过来。排查思路是这样先看积压幅度用XLEN order:stream查看消息总量再用XINFO GROUPS order:stream查看每个消费者组的落后数量。如果落后数量在持续增长说明消费速度小于生产速度单纯加消费者不一定有用得看瓶颈在哪。瓶颈通常是下游依赖比如消费逻辑里要调用外部API外部接口响应慢导致整个消费链路卡住。解决办法是拉大消费者的批量拉取数量比如count从10调到100减少网络往返次数同时检查下游接口是否能承受更大并发必要时做超时降级。还有一次是Redis的慢查询导致的。消费逻辑在Redis里执行了一个KEYS order:*命令直接拉垮了Redis性能队列消费全部阻塞。后来我把KEYS换成SCAN并用Stream消息里的业务维度做了缓存索引问题才解决。这里也提醒一下生产环境KEYS通配符命令能不用就别用尤其像KEYS ekyc_pic_*这种全量扫描数据量大时Redis会卡在原地所有请求都跟着排队。5.4 常见问题速查表现象可能原因解决措施消费者取不到新消息ReadOffset用了from(0)改用lastConsumed()消息处理失败后丢失业务异常未被catch或去掉ack捕获异常不ack走重试同一消息被多次处理宕机导致ack未发消费逻辑幂等或记录messageId唯一键队列积压持续增长下游依赖慢拉取数量太少增大count优化下游逻辑List队列消息丢失无确认机制消费者崩溃用Stream替代或BRPOPLPUSH备份消费线程池满线程池无界/拒绝策略不合理有界队列明确的拒绝策略Redis响应突然变慢操作了大量的KEYS命令换SCAN搞定具体的key查询5.5 性能优化补充消费端批量处理最后分享一批量处理的小技巧。默认情况下消费者每次读一条消息处理一条在消息量大时性能上不去因为单条处理的方式浪费了大量Redis网络往返。我现在的做法是每次拉取50到100条消息攒批处理处理完统一ack。ListMapRecordString, Object, Object records streamOps.read( Consumer.from(group-a, consumer-1), StreamReadOptions.empty().count(100).block(Duration.ofSeconds(2)), StreamOffset.create(order:stream, ReadOffset.lastConsumed()) );批量拉取之后要注意一个细节批里如果有某几条消息处理失败不要整批不ack否则成功的那几条也会被反复消费。我的做法是逐条try-catch处理成功的单独记到一个list里统一ack失败的记录日志后放到一个本地的重试队列稍后再compensate。6. 写到最后一些与工具有关的经验聊了很多架构和代码最后说点实际工具层面的事。排查Redis队列问题的时候命令行虽然能用但在看Stream积压情况、消息字段内容的时候确实费劲。我用过Redis Desktop Manager也用过Redis Insight个人更习惯Redis Insight它在Stream类型的可视化上做得更直观点开一个Stream就能看到消息列表、消费者组、PEL待确认消息数排查问题效率提升明显。另外Windows上做本地开发的话Redis官方其实不提供Windows版本可以下载微软维护的Win版本或者用WSL跑Linux版这个在热搜词里也出现了。下载时注意别装到那些套壳的推广站点尽量找官方GitHub Release或者验证过的国内镜像。我见过有人为了装个Redis给电脑装了一堆全家桶纯粹是给自己挖坑。最后再分享一个体会队列选型不要一步到位追求最复杂的方案也不要一直停留在最原始的List方案。选型的核心是看业务对可靠性的容忍度。非核心的日志上报、链路追踪用List甚至直接发Redis Pub/Sub都行但一旦涉及订单、支付、库存这种每一条消息都不能丢的业务Stream的消费者组加ack机制是底线。我之前在List方案上吃了亏之后现在统一把核心业务的消息全部切到了Stream配套加上幂等校验线上基本上没再出现过消息丢失导致的资损问题。这套组合拳打下来可靠性和复杂度之间能达到一个比较舒服的平衡。
返回列表