深入浅出Redis:使用Stream实现消息队列
mhr18 2024-11-12 11:26 27 浏览 0 评论
1 介绍
我们前面文章介绍了如何使用List实现消息队列么,但是我们也看到很多局限性,如下:
- 不支持消息确认机制,没有很好的ACK应答
- 不支持消息回溯,无法排查问题和做消息分析
- List遵循FIFO机制,所以存在消息堆积的风险。
- 查询效率低,作为线性结构,List中定位一个数据需要进行遍历,O(N)的时间复杂度。
- 不存在消费组(Consumer Group)的概念,无法进行分组消费和批量消费
Redis中有三种消息队列模式:
名称 | 简要说明 |
List | 不支持消息确认机制(Ack),不支持消息回朔 |
pubSub | 不支持消息确认机制(Ack),不支持消息回朔,不支持消息持久化 |
stream | 支持消息确认机制(Ack),支持消息回朔,支持消息持久化,支持消息阻塞 |
可以看出,作为Redis 5.0 引入的专门为消息队列设计的数据类型,Stream 功能更加健全,更适合做消息队列分发。
Stream 可以包含 0个 到 n个元素的有序队列,并根据ID的大小进行排序。
Stream类型消息队列的具备以下命令特点:
- 可以序列化生成消息ID,方便索引、排序
- 消息可回朔
- 支持Consumer Groups 消费组:多消费者消息争抢,加快消费速度
- 可以阻塞读取消息和非阻塞读取消息
- 没有消息漏读风险
- 有ACK消息确认机制,保证消息至少被消费一次
- 支持多播模式:可以让队列从逻辑上分组进行隔离消费
这些特性,基本达到了一个消息中间件的基本能力,比如:
- 类似 Kafka 的 Consumer Groups 的概念,它也具备了消费组的能力。
- 类似 Rocket MQ的持久化能力,以及高可用的文件存储机制,它也具备了消息的持久化和主从复制机制,可以记录访问位置,方便后续其他时间段继续访问,避免数据丢失。
详细的stream操作见官网文档:https://redis.io/docs/data-types/streams-tutorial/
2 XADD 消息写入
即讲消息添加到队列中,语法如下:
# 队列名称后面的队列id如果用 * 号表示 ,这代表让 Redis 为插入的消息自动生成唯一 序列化ID,当然也可以自己指定。
# 后面可以包含多个键值对,代表多个消息元素
XADD 队列名称 队列id key1 value1 [key2 value2 ....]
注释比较清楚,以下举例说明:
> xadd stream_user * user_id 1 user_name brand age 18
"1680926230000-0"
不指配*,这可以直接指定顺序Id
> XADD stream_user 0-1 user_name lili
0-1
> XADD stream_user 0-2 user_name brand
0-2
> XADD stream_user 0-* user_name candy
0-3
队列的消息ID 由两部分组成:
- 毫秒级别的当前时间的时间戳;
- 顺序编号。从 0 为起始值,用于区分同一时间内产生的多个ID,如果同一个时间戳内生成多ID,按序号顺序增长,这种方式可解决顺序识别和时间回拨问题。
通过这种时间戳 + 顺序编号的模式,变成数据Append的模式,这种流式记性数据顺序推送的方式符合MQ的基本消费逻辑,也为后面的有序性消费提供基本条件。
2 XREAD 消息阅读
即讲消息从队列中读取出来(消费),语法如下:
# COUNT:指的是对于每个Steam流中最多读取几个元素;
# BLOCK:当配置时阻塞读取,队列中没有消息即阻塞等待, 单位是ms,0 表示无限等待,类似MQ中的订阅,等待新消息出现。
# key:表示stream的名称
# ID:消息 id,读取消息的时候可以指定Id,并且指定某个Id的第一条甚至第n条开始读取,图中0-0 则表示从队列Id为0的队列的第1个元素开始读取。
XREAD [COUNT count] [BLOCK milliseconds] STREAMS key [key ...] ID [ID ...]
注释比较清楚,以下举例说明:
XREAD COUNT 1 BLOCK 0 STREAMS stream_user 0-0
1) 1) "stream_user"
2) 1) 1) "1680926230000-0"
2) 1) "user_name"
2) "brand"
3) "age"
4) "18"
如何顺序性消费:我们每次读取之后都会返回消息Id和序号,比如上面的 1680926230000-0,所以在下一次调用的时候,可以用上一次返回的ID序号作为参数,就可以从指定位置上进行消费。
问题:XREAD之后数据并没有删除,所以没记住读取的位置,下次可能重复阅读,造成重复消费。所以需要消费确认机制(即ACK)。
3 消费者组模式(Consumer Group)
消息队列很重要的一个能力就是分组消费(Consumer Group),无论是Kafka 还是 RabbitMQ。他允许队列从逻辑上进行分组来保证隔离消费。
这是典型的多播模,如下图所示:
它有如下特点:
- Redis Stream 实际结构是一个链式的队列,一个消息由消息Id和消息内容组成,消息Id具有唯一性;
- 消费组的状态是独立的,像图中的GroupA、GroupB、GroupC,Stream 消息可以被这几个组消费;
- 同时一个消费者组可以有多个消费者,但是他们的竞选关系,任意消费者消费之后就会导致 last_deliverd_id 偏移,这样避免了重复消费。
- 每个消费者都携带pending_ids 变量,记录读取但还未消费(未被ack)的消息,来保证消息有且仅有一次被消费。
消费组实现的消息队列主要有3类指令,如下:
- XGROUP:用于创建消费群组,包括注销和其他管理职能。
- XREADGROUP:消费者群组,通过这些组从流中有序读取数据。
- XACK:通过该命令,消费者将处理的完的消息标记为已正确完成。
3.1 写入队列数据
咱们先做一下数据准备,创建队列,并往里面写入一些数据,如下:
> xadd stream_user * user_id 1 user_name brand age 18
"1681126033000-0"
> xadd stream_user * user_id 2 user_name jay age 19
"1681126222000-0"
> xadd stream_user * user_id 3 user_name candy age 20
"1681126235000-0"
> xadd stream_user * user_id 4 user_name lili age 21
"1681126251000-0"
> xadd stream_user * user_id 5 user_name hanry age 22
"1681126263000-0"
3.2 创建消费者群组
这个的做法就是在队列中创建消费者组,然后指定消费的位置。
语法如下:
# stream_name:队列名称
# consumer_group:消费者组
# msgIdStartIndex:消息Id开始位置
# msgIdStartIndex:消息Id结束位置
xgroup create stream_name consumer_group msgIdStartIndex-msgIdStartIndex
下面是具体实现示例,为队列 stream_user 创建了消费组1(consumer_group1)和 消费组2(consumer_group2):
> xgroup create stream_user consumer_group1 0-0
OK
> xgroup create stream_user consumer_group2 0-0
OK
3.3 读取消费组信息
消费队列消息的语法如下:
# groupName: 消费者群组名
# consumerName: 消费者名称
# COUNT number: count 消费个数
# BLOCK ms: 阻塞时间,如果为 0 则代表无线阻塞
# streamName: 队列名称
# id: 消息消费ID
# []:代表可选参数
# `>`:放在命令参数的最后面,表示从尚未被消费的消息开始读取;
XREADGROUP GROUP groupName consumerName [COUNT number] [BLOCK ms] STREAMS streamName [stream ...] id [id ...]
实现示例:消费组 consumer_group1 的消费者 consumer1 从 stream_user 中以阻塞的方式读取一条消息:
XREADGROUP GROUP consumer_group1 consumer1 COUNT 1 BLOCK 0 STREAMS stream_user >
1) 1) "stream_user"
2) 1) 1) "1681126033000-0"
2) 1) "user_id"
2) "1"
3) "user_name"
4) "brand"
5) "age"
6) "18"
这边需要主意的是,同一个消费组内,消息只能单次消费,如果被消费组内消费过了,就不会被同组的其他消费组读取到。
如下:
XREADGROUP GROUP consumer_group1 consumer2 COUNT 1 BLOCK 0 STREAMS stream_user >
1) 1) "stream_user"
2) 1) 1) "1681126222000-0"
2) 1) "user_id"
2) "2"
3) "user_name"
4) "jay"
5) "age"
6) "19"
上面 user_name 为 brand 的数据已经被consumer1消费了,所以consumer2 就读不到了,只能读取到下一条 user_name 为 jay 的数据。
多个消费者可以达到流量分摊的目的,为大业务流量的场景做负载和分流。如下图,多个消费者相对平均的进行消息消费。
3.4 XPENDING 检查已读取但未ACK的数据
有时候会出现这种情况,就是消费者组或者消费者发生了故障,甚至整个消费者都故障重启了,那么如何避免消息丢失呢,那就是将读取到的但是还没消费的数据进行暂存。
Redis在Stream内部实现了一个待决队列(pending List),消费者读取之后且没有进行ACK的数据都保存在这里。
这种情况就是:
- 消费者使用 XREADGROUP 读取消息
- 读取完成之后,发生故障或者异常,没有给 Stream 发送 XACK 命令,消息依然保留在Stream 的 pending List中。
比如查看 stream_user 中的 消费组 consumer_group1 中各个消费者已读取未确认的消息信息:
XPENDING stream_user consumer_group1
1) (integer) 2 # 未确认消息条数
2) "1681126235000-0" # consumer_group1 消费组中所有消费者读取的最小ID
3) "1681126251000-0" # consumer_group1 消费组中所有消费者读取的最大ID
4) 1) 1) "consumer1"
2) "1"
2) 1) "consumer2"
2) "1"
3.5 消息消费完成之后确认(ACK)
正如3.4中所说的相关内容消费完之后,需要 ACK 通知 Streams,然后Stream除消息。否则就会造成消息变成待决队列中,可能造成重复消费的情况。
执行命令语法如下:
# XACK stream_name group_name ID [ID ...]
# stream_name:队列名称
# group_name:消费组名称
# ID:消费ID,可多选
XACK stream_user consumer_group1 1681126235000-0 1681126251000-0
(integer) 2
ack的本意就是对消费完成的消息进行确认,业务处理没有问题之后进行一个check的过程,代表这个消息已经被消费完了。流程如下:
4 使用Redission实现Stream队列能力
4.1 添加maven依赖 和 配置基本连接
# maven信息
<dependency>
<groupId>org.redisson</groupId>
<artifactId>redisson-spring-boot-starter</artifactId>
<version>3.16.8</version>
</dependency>
# 基本配置
spring:
application:
name: redission_test
redis:
host: x.x.x.x
port: 6379
ssl: false
password: xxxx.xxxx
4.2 Java程序实现
@Slf4j
@Service
public class StreamQueueService {
@Autowired
private RedissonClient redissonClient;
/**
* 生产消息内容
*
* @param msg
* @return
*/
@Override
public void produceMsg(String msg) {
RStream<String, String> stream = redissonClient.getStream("stream_user");
stream.add("user_id", "1");
stream.add("user_name", "brand");
stream.add("age", "18");
}
/**
* 消费消息内容
*/
@Override
public void consumeMessage() {
// 根据队列名称获取消息队列
RStream<String, String> stream = redissonClient.getStream("stream_user");
// 创建消费者小组
stream.createGroup("consumer_group1", StreamMessageId.ALL);
// 消费者读取消息
Map<StreamMessageId, Map<String, String>> msgs
= stream.readGroup("consumer_group1", "consumer1");
for (Map.Entry<StreamMessageId, Map<String, String>> entry : msgs.entrySet()) {
Map<String, String> msg = entry.getValue();
log.info(msg);
// todo:处理消息的业务逻辑代码
stream.ack("consumer_group1", entry.getKey());
}
}
}
5 总结
相对List,Stream的能力有比较大的提升:
- 支持消息确认机制(ACK应答确认)
- 支持消息回溯,方便排查问题和做消息分析
- 存在消费组(Consumer Group)的概念,可以进行分组消费和批量消费,可以负载多个消费实例
为帮助开发者们提升面试技能、有机会入职BATJ等大厂公司,特别制作了这个专辑——这一次整体放出。
大致内容包括了: Java 集合、JVM、多线程、并发编程、设计模式、Spring全家桶、Java、MyBatis、ZooKeeper、Dubbo、Elasticsearch、Memcached、MongoDB、Redis、MySQL、RabbitMQ、Kafka、Linux、Netty、Tomcat等大厂面试题等、等技术栈!
欢迎大家关注公众号【Java烂猪皮】,回复【666】,获取以上最新Java后端架构VIP学习资料以及视频学习教程,然后一起学习,一文在手,面试我有。
每一个专栏都是大家非常关心,和非常有价值的话题,如果我的文章对你有所帮助,还请帮忙点赞、好评、转发一下,你的支持会激励我输出更高质量的文章,非常感谢!
相关推荐
- Docker集群管理之Docker Compose
-
前言:在上一篇《Docker集群管理之DockerMachine》中,我们通过源码分析了解了DockerMachine的工作原理,使用者可以通过DockerMachine的一条命令在任意支持的平...
- 使用Dockerfile build镜像
-
Docker映像可以看作是Docker容器的压缩包,包含了应用程序以及运行应用程序所需的依赖,容器是映像的运行时实例。一般构建镜像都是使用dockerfile进行构建而不是dockercommit,...
- 自建私有云相册:Docker一键部署Immich,照片视频备份利器
-
自建私有云相册:Docker一键部署Immich,照片视频备份利器前言随着人们手机、PC、平板等电子产品多样,我们拍摄和保存的照片和视频数量也在不断增加。如何高效地管理和备份这些珍贵的记忆成为了一个重...
- docker容器的使用以及部署mysql
-
首先什么是docker官方:翻译:Docker是一个用于开发、发布和运行应用程序的开放平台。Docker使您能够将应用程序与基础架构分离,以便您可以快速交付软件。使用Docker,您可以像管理应...
-
- 自建Docker镜像加速服务,免费且简单,服务器VPS、NAS皆可用
-
写在前面:流程十分简单,有手就行,还请耐心看完。本文的实例仅做演示用,不久后将会删除,有需要的各位请自行搭建。免费实例如果15分钟内未收到入站流量,Render会关闭实例的网络服务。Render会在下次收到处理请求时重新启动该服务。Ren...
-
2025-05-24 15:40 mhr18
- 用了8年的方式-用 Docker 瞬间搭建本地开发环境
-
有些时候我们需要在本地搭开发环境,比如平时学习新技术的时候。或者有时候公司的项目需要在本地建一套类似的,方便调试修改。开发环境可能包括MySQL、Redis、Nginx、MQ、Elasticsea...
- 使用dockerfile构建docker镜像
-
准备工作购买vps使用ssh工具连接上1、更新系统aptupdate-y2、一键安装Dockercurl-fsSLhttps://get.docker.com-oget-docker.sh...
- 快速搭建 SpringCloud 微服务开发环境的脚手架
-
本文适合有SpringBoot和SpringCloud基础知识的人群,跟着本文可使用和快速搭建SpringCloud项目。本文作者:HelloGitHub-秦人HelloGitHub推出...
- Docker Hub最全详解(图文全面总结)
-
DockerHubDockerHub是一个由Docker公司负责维护的公共注册中心,它包含了超过15000多个可用来下载和构建容器的Docker镜像。DockerHub作用Docker好比一个代...
- Docker 命令详解
-
dockerimages—查看本地镜像命令dockerimages说明列出本地已下载的所有镜像及其标签、ID、大小等信息。适用场景查看本地镜像资源、准备删除或管理镜像时。常见用法docker...
- Kylin安装Dify
-
cd/mntgitclonehttps://github.com/langgenius/dify.gitcp/mnt/dify/docker/.env.example/mnt/dif...
- kali下对Docker的详细安装
-
Docker是渗透测试中必学不可的一个容器工具,在其中,我们能够快速创建、运行、测试以及部署应用程序。如,我们对一些漏洞进行本地复现时,可以使用Docker快速搭建漏洞环境,完成复现学习。注:本教程仅...
- 银河麒麟V10使用Docker方式部署应用
-
现在越来越多的企业级应用需要运行在国产化环境中,而银河麒麟V10是目前我碰到的最常用的服务器,在银河麒麟上部署应用有两种方式:使用二进制文件编译部署和使用Docker。关于使用二进制文件的方式...
- Docker入门到精通超详细教程,Docker全家桶实战攻略
-
大家好,我是各位双生的武魂、随身老爷爷。从看到这篇内容开始,你就是被选定的天命骚年,将承担起学完docker教程的使命,本使命为单向契约,你可选择YES或者选择YES。正式学习之前,我先给大家做一下d...
- 【Docker 新手入门指南】第一章:前言
-
一、基本介绍Docker介绍Docker是基于Go语言开发的开源容器化平台,旨在实现“一次镜像,处处运行”。它通过将应用程序及其依赖环境(代码、运行时、系统工具、系统库等)打包成一个轻量级、可移...
你 发表评论:
欢迎- 一周热门
- 最近发表
- 标签列表
-
- oracle位图索引 (63)
- oracle批量插入数据 (62)
- oracle事务隔离级别 (53)
- oracle 空为0 (50)
- oracle主从同步 (55)
- oracle 乐观锁 (51)
- redis 命令 (78)
- php redis (88)
- redis 存储 (66)
- redis 锁 (69)
- 启动 redis (66)
- redis 时间 (56)
- redis 删除 (67)
- redis内存 (57)
- redis并发 (52)
- redis 主从 (69)
- redis 订阅 (51)
- redis 登录 (54)
- redis 面试 (58)
- 阿里 redis (59)
- redis 搭建 (53)
- redis的缓存 (55)
- lua redis (58)
- redis 连接池 (61)
- redis 限流 (51)