I'm Aron

Redis实战:Redisson 分布式锁与 Stream 消息队列

3526 字
18 分钟
Redis实战:Redisson 分布式锁与 Stream 消息队列

本文整理两个常见的 Redis 实战场景:使用 Redisson 实现分布式锁,以及使用 Redis Stream 构建支持消息持久化、阻塞读取、消费者组和消息确认的消息队列。

一、Redisson 分布式锁#

Redisson 是 Redis 的 Java 客户端之一。它在 Redis 命令之上封装了分布式锁、集合、信号量等常用组件,可以减少手动编写锁续期和解锁脚本的工作。

1. 引入依赖#

<dependency>
<groupId>org.redisson</groupId>
<artifactId>redisson</artifactId>
<version>3.13.6</version>
</dependency>

2. 配置 Redisson 客户端#

@Configuration
public class RedissonConfig {
@Bean
public RedissonClient redissonClient() {
// 配置
Config config = new Config();
config.useSingleServer().setAddress("redis://10.211.55.3:6379").setPassword("123123");
// 创建RedissonClient对象
return Redisson.create(config);
}
}

3. 在业务中使用锁#

这段代码以用户 ID 作为锁粒度,保证同一用户的下单请求串行执行。不同用户使用不同的锁,因此仍然可以并发下单。

@Resource
private RedissonClient redissonClient;
Long userId = UserHolder.getUser().getId();
// 创建锁对象
// SimpleRedisLock lock = new SimpleRedisLock("order" + userId, stringRedisTemplate);
RLock lock = redissonClient.getLock("lock:order:" + userId);
// 获取锁
boolean isLock = lock.tryLock();
// 判断是否获取锁成功
if(!isLock){
// 获取锁失败,返回错误或重试
return Result.fail("不允许重复下单");
}
try{
// 获取代理对象(事物)
IVoucherOrderService proxy = (IVoucherOrderService) AopContext.currentProxy();
return proxy.createVoucherOrder(voucherId);
} finally {
// 释放锁
lock.unlock();
}

4. 使用分布式锁时需要注意什么#

  • tryLock() 会立即尝试获取锁,失败后直接返回;需要等待时可以使用带等待时间的重载方法。
  • 未显式指定租约时间时,Redisson 的看门狗机制会在持锁线程仍然存活时自动续期。
  • 解锁必须放在 finally 中,避免业务异常导致锁无法释放。
  • 只有锁的持有线程才能解锁。复杂流程中可在解锁前使用 lock.isHeldByCurrentThread() 检查。
  • 分布式锁不能替代数据库约束。防止重复下单时,数据库唯一索引仍是最后一道保障。
  • 通过 AopContext.currentProxy() 获取代理对象,需要启用代理暴露;否则可能抛出异常。也可以通过拆分 Service 避免同类内部调用导致事务失效。

二、Redis Stream 基础#

Redis Stream 是 Redis 5.0 引入的数据类型,适合保存按时间追加的事件。一个 Stream 可以理解为一条只向尾部追加的日志,每条消息称为一个 Entry。

Stream 的基本结构如下:

Stream key
├── 1720000000000-0 name=jack age=21
├── 1720000000100-0 name=rose age=20
└── 1720000000100-1 name=tom age=22

每条 Entry 包含两部分:

  • 唯一消息 ID,例如 1720000000000-0
  • 一个或多个 field value 键值对。

消息 ID 的格式是 毫秒时间戳-序列号。同一毫秒内产生多条消息时,序列号会递增。通常使用 * 让 Redis 自动生成 ID。

三、使用 XADD 发送消息#

1. 基本语法#

XADD key [NOMKSTREAM] [MAXLEN | MINID [= | ~] threshold [LIMIT count]]
* | ID field value [field value ...]

常用参数:

  • key:Stream 的名称。
  • NOMKSTREAM:Stream 不存在时不自动创建;默认会自动创建。
  • MAXLEN:限制 Stream 的最大长度,避免消息无限堆积。
  • MINID:删除小于指定 ID 的旧消息。
  • ~:近似裁剪,通常比精确裁剪 = 的性能更好。
  • *:让 Redis 自动生成消息 ID。
  • field value:消息内容,可以包含多个键值对。

2. 添加一条消息#

XADD users * name jack age 21

返回值是新消息的 ID:

"1644805700523-0"

3. 限制队列长度#

XADD users MAXLEN ~ 10000 * name jack age 21

该命令在添加消息的同时,把 Stream 大致控制在 10000 条以内。~ 是近似裁剪,实际长度可能略大于阈值。

常用的查看命令:

XLEN users
XRANGE users - + COUNT 10
XREVRANGE users + - COUNT 10
XINFO STREAM users

四、使用 XREAD 读取消息#

XREAD 直接读取 Stream,不使用消费者组。它适合单消费者读取,或者多个消费者都需要接收同一份消息的广播式场景。

1. 基本语法#

XREAD [COUNT count] [BLOCK milliseconds] STREAMS key [key ...] ID [ID ...]

阻塞 1 秒,等待 users 中的新消息:

XREAD COUNT 1 BLOCK 1000 STREAMS users $

如果 1 秒内没有新消息,命令返回 nilBLOCK 0 表示一直阻塞,直到有消息到达。

2. 起始 ID 的含义#

  • 0-0:从最早的消息开始读取。
  • 某个具体 ID:读取 ID 大于它的消息,适合接着上一次的位置继续读取。
  • $:只等待执行命令之后到达的新消息。

3. 为什么循环使用 $ 可能漏读#

如果每次循环都执行下面的命令:

XREAD COUNT 1 BLOCK 2000 STREAMS users $

那么 $ 每次都会重新代表“当前 Stream 的最后一条消息”。消费者处理上一条消息期间,如果又有多条消息到达,下次读取会从新的队尾开始等待,这些已经到达的消息可能被跳过。

更稳妥的做法是保存最后一次成功读取的消息 ID,并在下一次读取时传入该 ID:

lastId = 0-0
循环:
messages = XREAD COUNT 10 BLOCK 2000 STREAMS users lastId
依次处理 messages
lastId = 本批最后一条消息的 ID

如果进程重启后仍要从原位置继续,还需要把 lastId 持久化。Redis 本身不会为普通 XREAD 保存消费者进度。

4. XREAD 的特点#

  • 支持按 ID 回溯历史消息。
  • 支持阻塞读取。
  • 多个客户端可以读取同一条消息,彼此之间不会自动分流。
  • Redis 不记录每个读取者的确认状态。
  • 使用方式不当会跳过消息,因此不适合直接承担需要可靠消费的复杂业务。

五、消费者组 Consumer Group#

消费者组将多个消费者组织成一个组,共同消费同一个 Stream。它主要解决三个问题:消息分流、消费进度记录和消息确认。

1. 消息分流#

同一消费者组中的多个消费者会竞争组内的新消息。一条新消息通常只交给组内的一个消费者处理,从而提高整体吞吐量。

不同消费者组之间互不影响。同一条消息可以分别被多个消费者组各消费一次。

2. 消费进度#

消费者组维护最后投递位置,用于记录哪些新消息已经投递给该组。消费者宕机重启后,消费者组可以继续从组的进度读取新消息。

3. Pending Entries List#

消费者读取消息后,如果没有使用 NOACK,该消息会进入 Pending Entries List,简称 PEL。PEL 记录消息 ID、所属消费者、空闲时间和投递次数。

业务处理成功后,需要执行 XACK 确认消息。确认成功后,消息才会从 PEL 中移除。注意,XACK 只会修改消费者组的确认状态,不会删除 Stream 中的原始消息。

六、管理消费者组#

1. 创建消费者组#

XGROUP CREATE key groupName ID [MKSTREAM]

示例:

XGROUP CREATE users user-group 0-0 MKSTREAM

参数说明:

  • key:Stream 名称。
  • groupName:消费者组名称。
  • ID:消费者组的起始位置。
  • 0-0:从 Stream 中现有的第一条消息开始消费。
  • $:忽略当前已有消息,只消费创建消费者组之后的新消息。
  • MKSTREAM:Stream 不存在时自动创建一个空 Stream。

2. 其他管理命令#

# 删除消费者组
XGROUP DESTROY users user-group
# 手动创建消费者
XGROUP CREATECONSUMER users user-group consumer-1
# 删除消费者
XGROUP DELCONSUMER users user-group consumer-1
# 查看消费者组和消费者信息
XINFO GROUPS users
XINFO CONSUMERS users user-group

消费者通常不必手动创建。第一次执行 XREADGROUP 时,如果消费者名称不存在,Redis 会自动创建该消费者。

删除消费者前必须先处理其 Pending 消息。直接执行 XGROUP DELCONSUMER 会删除该消费者对应的 Pending 记录,未确认消息将无法再通过原 PEL 恢复。

七、使用 XREADGROUP 消费消息#

1. 基本语法#

XREADGROUP GROUP group consumer
[COUNT count]
[BLOCK milliseconds]
[NOACK]
STREAMS key [key ...] ID [ID ...]

读取新消息:

XREADGROUP GROUP user-group consumer-1 COUNT 10 BLOCK 2000 STREAMS users >

参数说明:

  • group:消费者组名称。
  • consumer:当前消费者名称;不存在时会自动创建。
  • COUNT:本次最多返回的消息数量。
  • BLOCK:没有消息时的最长等待时间,单位为毫秒。
  • STREAMS:后面依次提供 Stream 名称和对应的起始 ID。
  • >:读取从未投递给当前消费者组的新消息。

2. > 和具体 ID 的区别#

  • 使用 > 时,读取尚未投递给消费者组的新消息。
  • 使用 00-0 等具体 ID 时,读取已经分配给当前消费者、但尚未确认的 Pending 消息。
  • 具体 ID 不会直接读取其他消费者名下的 Pending 消息。接管其他消费者的消息需要使用 XCLAIMXAUTOCLAIM

3. 谨慎使用 NOACK#

NOACK 表示消息被读取后不进入 PEL。它更接近“无需确认”,而不是业务成功后的“自动确认”。如果消费者读取后立刻宕机,消息不会留在 PEL 中供后续重试。

对订单、支付、库存等可靠性要求高的业务,不建议使用 NOACK

八、消息确认与异常恢复#

1. 确认消息#

业务处理成功后执行:

XACK users user-group 1644805700523-0

推荐顺序是:

读取消息 -> 执行业务 -> 业务成功提交 -> XACK

如果业务处理失败,不要确认消息,让它继续保留在 PEL 中,等待重试或被其他消费者接管。

2. 查看 Pending 消息#

XPENDING users user-group
XPENDING users user-group - + 10
XPENDING users user-group - + 10 consumer-1

第一条命令查看 PEL 摘要,后两条命令查看 Pending 消息明细。

3. 接管长时间未处理的消息#

Redis 6.2 及以上可以使用 XAUTOCLAIM

XAUTOCLAIM users user-group consumer-2 60000 0-0 COUNT 10

这表示由 consumer-2 接管空闲时间至少达到 60000 毫秒的 Pending 消息。较早版本可以先用 XPENDING 找到消息 ID,再使用 XCLAIM 接管。

4. 推荐的可靠消费流程#

1. 使用 XREADGROUP ... STREAMS users > 阻塞读取新消息
2. 执行业务逻辑
3. 业务成功后执行 XACK
4. 业务失败时保留 Pending 状态,并记录失败原因
5. 定时扫描空闲时间过长的 Pending 消息
6. 使用 XAUTOCLAIM 或 XCLAIM 接管并重试
7. 超过最大重试次数后写入自定义死信 Stream,并人工或定时处理

消费者组配合 PEL 和 ACK 可以实现至少一次投递,但不能天然实现严格的恰好一次。消费者可能在“业务已成功、ACK 尚未发送”时宕机,恢复后同一条消息会再次处理,因此业务必须具备幂等性。

常见幂等方案:

  • 使用消息 ID 作为幂等键,并记录已处理消息。
  • 利用数据库唯一索引阻止重复写入。
  • 使用业务单号检查当前状态,只允许合法的状态迁移。
  • 把业务写入和消费记录放在同一个数据库事务中。

九、XREAD 与 XREADGROUP 对比#

对比项XREADXREADGROUP
消费模式多个读取者可以重复读取同一消息组内消费者竞争新消息
消费进度客户端自行保存最后 IDRedis 保存消费者组投递进度
消息确认不支持支持 XACK
Pending 记录不支持支持 PEL
故障恢复需要应用自行实现可通过 XPENDINGXCLAIMXAUTOCLAIM 恢复
适用场景简单读取、广播、历史查询任务分发、并行消费、可靠消息处理

不能简单地认为 XREADGROUP “绝对不会漏消息”。可靠性仍取决于是否正确使用 ACK、是否恢复 Pending 消息、是否误用 NOACK、是否过早裁剪 Stream,以及业务是否具备幂等性。

十、List、Pub/Sub 与 Stream 对比#

能力ListPub/SubStream
消息持久化支持不支持支持
阻塞读取支持支持订阅等待支持
多消费者分流支持,但需自行设计不适用,订阅者各收一份消费者组原生支持
消息确认无原生 ACK不支持支持
消息回溯不方便不支持按 ID 支持
消费进度应用自行维护不维护消费者组维护
典型场景简单任务队列实时通知、允许丢失的广播可靠任务队列、事件流

选择建议:

  • 只需要非常简单的先进先出任务队列,可以考虑 List。
  • 只关心在线订阅者,并且允许离线期间丢消息,可以使用 Pub/Sub。
  • 需要持久化、回溯、消费者组、确认和故障恢复时,优先考虑 Stream。

十一、生产实践建议#

  • 使用 MAXLEN ~ 或定期 XTRIM 控制 Stream 长度,防止内存无限增长。
  • 每个消费者实例使用唯一且稳定的消费者名称,便于观察和故障接管。
  • 使用 BLOCK 阻塞读取,避免无消息时高频轮询消耗 CPU。
  • 业务成功后再执行 XACK,不要把“读取成功”当成“处理成功”。
  • 定时监控 PEL 数量、消息空闲时间、投递次数和消费者状态。
  • 对重复消费做好幂等处理,这是至少一次投递模型中的必要设计。
  • 设计最大重试次数和死信 Stream,避免异常消息永久占用重试资源。
  • 裁剪或删除消息前确认消费组的处理进度。过早清理可能导致 PEL 中只剩消息 ID,而消息正文已经不存在。

总结#

Redisson 分布式锁解决的是多个应用实例之间的并发互斥问题,Redis Stream 解决的是消息的持久化、分发与消费进度管理问题。二者可以组合使用,但职责不同:锁用于保护临界区,Stream 用于解耦生产者和消费者。

学习 Redis Stream 时,需要重点掌握以下关系:

XADD 负责生产消息
XREAD 负责不带消费者组的读取
XGROUP 负责创建和管理消费者组
XREADGROUP 负责组内分流和消费
XPENDING 负责查看未确认消息
XACK 负责确认处理完成
XCLAIM / XAUTOCLAIM 负责接管超时未完成的消息

评论区

文章目录