最近手上的订单异步处理项目刚好改造成基于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 List(PEL),在可靠性上比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实现多个队列阻塞读取,可以这样写:
List<String> queueKeys = Arrays.asList("high:queue", "normal:queue", "low:queue"); List<String> 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 List(PEL)。处理成功后,消费者发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来创建:
StreamOperations<String, Object, Object> streamOps = stringRedisTemplate.opsForStream(); try { streamOps.createGroup("order:stream", "group-a"); } catch (RedisSystemException e) { // 分组已存在时会报错,这里不做处理即可 }消费的代码,我用的是StreamReadOptions手动读的方式:
List<MapRecord<String, 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 (MapRecord<String, Object, Object> record : records) { try { // 拿到消息内容 Map<Object, 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保证的是“消息不丢”,不是“消息不重”。
什么时候会重复?最常见的情况是:消费者取到消息,执行完业务逻辑,但是还没来及ack,JVM就宕机了。这时候消息在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的话,最经典的方案是基于ZSet,score存的是任务预期执行的时间戳。
入队的时候把任务内容放在ZSet的member里,score设置为当前时间加上延迟时间:
// 订单下单后,设置15分钟后自动关闭 long executeTime = System.currentTimeMillis() + 15 * 60 * 1000; redisTemplate.opsForZSet().add("order:delay", orderId, executeTime);消费端用一个线程循环扫描,每次取出score小于等于当前时间的任务:
while (true) { Set<String> 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条 Set<String> 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。
List<MapRecord<String, 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,配套加上幂等校验,线上基本上没再出现过消息丢失导致的资损问题。这套组合拳打下来,可靠性和复杂度之间能达到一个比较舒服的平衡。