Redis Streams
提示
来自deepseek解释
原文链接:https://redis.io/docs/latest/develop/data-types/streams/
代码示例说明
下方的代码示例展示如何在不同的编程语言和客户端库中执行相同的操作:
- Redis CLI:Redis 命令行界面
- C# (同步):StackExchange.Redis 同步客户端
- C# (异步):StackExchange.Redis 异步客户端
- Go:go-redis 客户端
- Java (同步 - Jedis):Jedis 同步客户端
- Java (异步 - Lettuce):Lettuce 异步客户端
- Java (响应式 - Lettuce):Lettuce 响应式/流式客户端
- JavaScript (Node.js):node-redis 客户端
- PHP:Predis 客户端
- Python:redis-py 客户端
- Rust (同步):redis-rs 同步客户端
- Rust (异步):redis-rs 异步客户端
每个代码示例均以不同语言展示相同的基本操作。具体语法和模式会因语言和客户端库而异,但底层的 Redis 命令和行为保持一致。
流命令摘要
本组共 30 条命令:
| 命令 | 摘要 | 复杂度 | 引入版本 |
|---|---|---|---|
| XACK | 返回消费者组成员成功确认的消息数量。 | O(1) 每条消息 ID。 | 5.0.0 |
| XACKDEL | 为流消费者组确认并删除一条或多条消息。 | O(1) 每条消息 ID。 | 8.2.0 |
| XADD | 向流中追加新消息。若 key 不存在则创建。 | O(1) 添加新条目时,O(N) 修剪时... | 5.0.0 |
| XAUTOCLAIM | 更改或获取消费者组中消息的所有权,如同消息已投递给消费者组成员。 | O(1) 如果 COUNT 较小。 | 6.2.0 |
| XCFGSET | 设置流的 IDMP 配置参数。 | O(1) | 8.6.0 |
| XCLAIM | 更改或获取消费者组中消息的所有权,如同消息已投递给消费者组成员。 | O(log N),N 为流中消息数... | 5.0.0 |
| XDEL | 返回从流中删除消息后的数量。 | O(1) 每条待删除的流条目... | 5.0.0 |
| XDELEX | 从流中删除一条或多条条目。 | O(1) 每条待删除的流条目... | 8.2.0 |
| XGROUP | 消费者组命令的容器。 | 取决于子命令。 | 5.0.0 |
| XGROUP CREATE | 创建消费者组。 | O(1) | 5.0.0 |
流命令摘要(第2部分)
| 命令 | 摘要 | 复杂度 | 引入版本 |
|---|---|---|---|
| XGROUP CREATECONSUMER | 在消费者组中创建一个消费者。 | O(1) | 6.2.0 |
| XGROUP DELCONSUMER | 从消费者组中删除消费者。 | O(1) | 5.0.0 |
| XGROUP DESTROY | 销毁消费者组。 | O(N),N 为组中待处理条目数... | 5.0.0 |
| XGROUP HELP | 返回关于不同子命令的帮助文本。 | O(1) | 5.0.0 |
| XGROUP SETID | 设置消费者组的最后投递 ID。 | O(1) | 5.0.0 |
| XIDMPRECORD | 用于在现有流消息上设置 IDMP 元数据的内部命令。 | O(1) | 8.6.2 |
| XINFO | 流内省命令的容器。 | 取决于子命令。 | 5.0.0 |
| XINFO CONSUMERS | 返回消费者组中的消费者列表。 | O(1) | 5.0.0 |
| XINFO GROUPS | 返回流的消费者组列表。 | O(1) | 5.0.0 |
| XINFO HELP | 返回关于不同子命令的帮助文本。 | O(1) | 5.0.0 |
流命令摘要(第3部分)
| 命令 | 摘要 | 复杂度 | 引入版本 |
|---|---|---|---|
| XINFO STREAM | 返回流的信息。 | O(1) | 5.0.0 |
| XLEN | 返回流中的消息数量。 | O(1) | 5.0.0 |
| XNACK | 将已认领的消息释放回组的 PEL,但不确认,使其可重新投递。 | O(1) 每条消息 ID。 | 8.8.0 |
| XPENDING | 返回流消费者组待处理条目列表的信息与条目。 | O(N),N 为返回的元素数... | 5.0.0 |
| XRANGE | 按 ID 范围返回流中的消息。 | O(N),N 为返回的元素数... | 5.0.0 |
| XREAD | 返回多个流中 ID 大于指定值的消息。若无消息则阻塞。 | 不适用 | 5.0.0 |
| XREADGROUP | 为组中的消费者返回流的新消息或历史消息。若无消息则阻塞。 | 每个流 O(M),M 为返回的元素数... | 5.0.0 |
| XREVRANGE | 按 ID 范围逆序返回流中的消息。 | O(N),N 为返回的元素数... | 5.0.0 |
| XSETID | 用于复制流值的内部命令。 | O(1) | 5.0.0 |
| XTRIM | 从流开头删除消息。 | O(N),N 为被驱逐的条目数... | 5.0.0 |
Redis 流是一种类似于追加日志的数据结构,但同时实现了若干操作以克服典型追加日志的一些限制,例如 O(1) 时间复杂度的随机访问以及复杂的消费策略(如消费者组)。 您可以使用流来实时记录并同时同步事件。 Redis 流的使用场景示例包括:
- 事件溯源(例如跟踪用户操作、点击等)
- 传感器监控(例如现场设备的读数)
- 通知(例如将每个用户的通知记录存储在独立的流中)
Redis 为每个流条目生成唯一 ID。 您可以使用这些 ID 稍后检索对应的条目,或者读取并处理流中所有后续条目。请注意,由于这些 ID 与时间相关,此处显示的 ID 可能与您自己 Redis 实例中看到的有所不同。
Redis 流支持多种修剪策略(防止流无限增长)以及多种消费策略(参见 XREAD、XREADGROUP 和 XRANGE)。从 Redis 8.2 开始,XACKDEL、XDELEX、XADD 和 XTRIM 命令提供了对流操作如何与多个消费者组交互的细粒度控制,简化了不同应用程序之间消息处理的协调。
从 Redis 8.6 开始,Redis 流支持幂等消息处理(最多一次生产),以在使用至少一次投递模式时防止重复条目。此功能支持带自动去重的可靠消息提交。详见幂等消息处理。
示例
- 当我们的赛车手通过检查点时,我们为每位赛车手添加一个流条目,包含车手姓名、速度、位置和位置 ID:
命令: XADD
复杂度: XADD: O(1)
可用客户端: Redis CLI、C#、Go、Java(同步 - Jedis)、JavaScript(Node.js)、Python、Ruby、Rust(异步)、Rust(同步)
Java(同步 - Jedis)
StreamEntryID res1 = jedis.xadd("race:france",new HashMap<String,String>(){{put("rider","Castilla");put("speed","30.2");put("position","1");put("location_id","1");}} , XAddParams.xAddParams());
System.out.println(res1); // >>> 1701760582225-0
StreamEntryID res2 = jedis.xadd("race:france",new HashMap<String,String>(){{put("rider","Norem");put("speed","28.8");put("position","3");put("location_id","1");}} , XAddParams.xAddParams());
System.out.println(res2); // >>> 1701760582225-1
StreamEntryID res3 = jedis.xadd("race:france",new HashMap<String,String>(){{put("rider","Prickett");put("speed","29.7");put("position","2");put("location_id","1");}} , XAddParams.xAddParams());
System.out.println(res3); // >>> 1701760582226-0- 从 ID
1692632086370-0开始读取两条流条目:
命令: XRANGE
复杂度: XRANGE: O(N)
可用客户端: Redis CLI、C#、Go、Java(同步 - Jedis)、JavaScript(Node.js)、Python、Ruby、Rust(异步)、Rust(同步)
Java(同步 - Jedis)
List<StreamEntry> res4 = jedis.xrange("race:france","1701760582225-0","+",2);
System.out.println(res4); // >>> [1701760841292-0 {rider=Castilla, speed=30.2, location_id=1, position=1}, 1701760841292-1 {rider=Norem, speed=28.8, location_id=1, position=3}]- 从流末尾开始读取最多 100 条新流条目,若无新条目则阻塞最多 300 毫秒:
难度: 中级
命令: XREAD
可用客户端: Redis CLI、C#、Go、Java(同步 - Jedis)、JavaScript(Node.js)、Python、Ruby、Rust(异步)、Rust(同步)
Java(同步 - Jedis)
List<Map.Entry<String, List<StreamEntry>>> res5= jedis.xread(XReadParams.xReadParams().block(300).count(100),new HashMap<String,StreamEntryID>(){{put("race:france",new StreamEntryID());}});
System.out.println(
res5
); // >>> [race:france=[1701761996660-0 {rider=Castilla, speed=30.2, location_id=1, position=1}, 1701761996661-0 {rider=Norem, speed=28.8, location_id=1, position=3}, 1701761996661-1 {rider=Prickett, speed=29.7, location_id=1, position=2}]]性能
向流中添加条目是 O(1) 复杂度。 访问任意单个条目是 O(n),其中 n 是 ID 的长度。由于流 ID 通常较短且长度固定,这实际上等同于常量时间查找。 有关原因的详细信息,请注意流是作为基数树实现的。
简言之,Redis 流提供了高效的插入和读取操作。 各命令的时间复杂度详见对应文档。
流基础
流是一种仅追加的数据结构。基本的写入命令 XADD 将新条目追加到指定流中。
每个流条目由一个或多个字段-值对组成,类似于字典或 Redis 哈希:
命令: XADD
复杂度: XADD: O(1)
可用客户端: Redis CLI、C#、Go、Java(同步 - Jedis)、JavaScript(Node.js)、Python、Ruby、Rust(异步)、Rust(同步)
Java(同步 - Jedis)
StreamEntryID res6 = jedis.xadd("race:france",new HashMap<String,String>(){{put("rider","Castilla");put("speed","29.9");put("position","2");put("location_id","1");}} , XAddParams.xAddParams());
System.out.println(res6); // >>> 1701762285679-0上述 XADD 调用向键为 race:france 的流中添加了一个条目 rider: Castilla, speed: 29.9, position: 1, location_id: 2,使用自动生成的条目 ID,即命令返回的 1692632147973-0。第一个参数是键名 race:france,第二个参数是条目 ID,用于标识流中的每个条目。但此处我们传入 *,因为希望服务器自动生成新 ID。每个新 ID 都是单调递增的,即每个新条目的 ID 都会高于所有过去的条目。服务器自动生成 ID 几乎总是您想要的,显式指定 ID 的情况极为罕见。我们稍后会详细讨论。每个流条目都有 ID,这与日志文件相似,日志文件中的行号或文件内字节偏移量可用于标识特定条目。回到我们的 XADD 示例,键名和 ID 之后是组成流条目的字段-值对。
可以使用 XLEN 命令获取流中的条目数量:
命令: XLEN
复杂度: XLEN: O(1)
可用客户端: Redis CLI、C#、Go、Java(同步 - Jedis)、JavaScript(Node.js)、Python、Ruby、Rust(异步)、Rust(同步)
Java(同步 - Jedis)
long res7 = jedis.xlen("race:france");
System.out.println(res7); // >>> 4条目 ID
XADD 命令返回的条目 ID 唯一标识流中的每个条目,由两部分组成:
<毫秒时间>-<序列号>毫秒时间部分实际上是生成流 ID 的本地 Redis 节点的本地时间,但如果当前毫秒时间小于前一个条目的时间,则使用前一个条目的时间,因此即使时钟回拨,单调递增 ID 属性仍然成立。序列号用于在同一毫秒内创建的条目。由于序列号为 64 位宽,实际上在同一毫秒内可生成的条目数没有限制。
这种 ID 格式初看可能有些奇怪,细心的读者可能会问为什么 ID 中要包含时间。原因是 Redis 流支持按 ID 进行范围查询。由于 ID 与条目生成的时间相关,这使我们能够轻松地进行时间范围查询。我们将在介绍 XRANGE 命令时看到这一点。
如果用户出于某种原因需要与时间无关的递增 ID,而是要与另一个外部系统 ID 关联,如前所述,XADD 命令可以接受显式 ID,而不是触发自动生成的 * 通配符 ID,如下例所示:
难度: 高级
命令: XADD
复杂度: XADD: O(1)
可用客户端: Redis CLI、C#、Go、Java(同步 - Jedis)、JavaScript(Node.js)、Python、Ruby、Rust(异步)、Rust(同步)
Java(同步 - Jedis)
StreamEntryID res8 = jedis.xadd("race:usa", new HashMap<String,String>(){{put("racer","Castilla");}},XAddParams.xAddParams().id("0-1"));
System.out.println(res8); // >>> 0-1
StreamEntryID res9 = jedis.xadd("race:usa", new HashMap<String,String>(){{put("racer","Norem");}},XAddParams.xAddParams().id("0-2"));
System.out.println(res9); // >>> 0-2请注意,最小 ID 为 0-1,并且命令不会接受等于或小于之前 ID 的 ID:
难度: 高级
命令: XADD
复杂度: XADD: O(1)
可用客户端: Redis CLI、C#、Go、Java(同步 - Jedis)、JavaScript(Node.js)、Python、Ruby、Rust(异步)、Rust(同步)
Java(同步 - Jedis)
try {
StreamEntryID res10 = jedis.xadd("race:usa", new HashMap<String,String>(){{put("racer","Prickett");}},XAddParams.xAddParams().id("0-1"));
System.out.println(res10); // >>> 0-1
}
catch (JedisDataException e){
System.out.println(e); // >>> ERR The ID specified in XADD is equal or smaller than the target stream top item
}如果您运行的是 Redis 7 或更高版本,还可以提供仅包含毫秒部分的显式 ID。此时,ID 的序列部分将自动生成。使用以下语法:
难度: 中级
命令: XADD
复杂度: XADD: O(1)
可用客户端: Redis CLI、C#、Go、Java(同步 - Jedis)、JavaScript(Node.js)、Python、Ruby、Rust(异步)、Rust(同步)
Java(同步 - Jedis)
StreamEntryID res11 = jedis.xadd("race:usa", new HashMap<String,String>(){{put("racer","Norem");}},XAddParams.xAddParams().id("0-*"));
System.out.println(res11);从流中获取数据
现在我们能够通过 XADD 向流中追加条目了。然而,虽然追加数据到流中非常直观,但查询流以提取数据的方式就不那么直观了。继续用日志文件作类比,一种显而易见的方式是模仿 Unix 命令 tail -f,即开始监听以获取追加到流中的新消息。请注意,与 Redis 阻塞列表操作(其中某个元素会到达单个正在以 弹出样式 阻塞的客户端,如 BLPOP)不同,对于流,我们希望多个消费者都能看到追加到流中的新消息(就像多个 tail -f 进程可以看到追加到日志中的内容一样)。用传统术语来说,我们希望流能够将消息 扇出 给多个客户端。
然而,这只是其中一种访问模式。我们还可以从完全不同的角度来看待流:不仅作为消息系统,还可以作为 时间序列存储。在这种情况下,获取新追加的消息可能也有用,但另一种自然的查询模式是按时间范围获取消息,或者使用游标迭代消息以逐步检查所有历史记录。这无疑是另一种有用的访问模式。
最后,如果从消费者的角度来看流,我们可能希望以另一种方式访问流,即将流作为可分区给多个消费者的消息流,这些消费者正在处理此类消息,从而消费者组中的每个消费者只能看到单个流中到达的消息的子集。这样,就可以在不同消费者之间扩展消息处理,而无需单个消费者处理所有消息:每个消费者将获得不同的消息进行处理。这基本上就是 Kafka (TM) 使用消费者组所做的。通过消费者组读取消息是另一种从 Redis 流读取的有趣模式。
Redis 流通过不同的命令支持上述所有三种查询模式。接下来的部分将逐一介绍,从最简单、最直接的范围查询开始。
范围查询:XRANGE 和 XREVRANGE
要通过范围查询流,我们只需指定两个 ID:start 和 end。返回的范围将包含具有 start 或 end 作为 ID 的元素,因此范围是包含性的。两个特殊 ID - 和 + 分别表示最小和最大的可能 ID。
命令: XRANGE
复杂度: XRANGE: O(N)
可用客户端: Redis CLI、C#、Go、Java(同步 - Jedis)、JavaScript(Node.js)、Python、Ruby、Rust(异步)、Rust(同步)
Java(同步 - Jedis)
List<StreamEntry> res12 = jedis.xrange("race:france","-","+");
System.out.println(
res12
); // >>> [1701764734160-0 {rider=Castilla, speed=30.2, location_id=1, position=1}, 1701764734160-1 {rider=Norem, speed=28.8, location_id=1, position=3}, 1701764734161-0 {rider=Prickett, speed=29.7, location_id=1, position=2}, 1701764734162-0 {rider=Castilla, speed=29.9, location_id=1, position=2}]每个返回的条目是一个包含两个元素的数组:ID 和字段-值对列表。我们已经提到条目 ID 与时间有关,因为 - 左侧部分是创建该流条目的本地节点在创建时的 Unix 毫秒时间(但请注意,流通过完整的 XADD 命令进行复制,因此副本将与主节点具有相同的 ID)。这意味着我可以使用 XRANGE 查询时间范围。为此,我可以省略 ID 的序列部分:如果省略,在范围的起始部分将假定为 0,在结束部分将假定为最大可用序列号。这样,仅使用两个毫秒 Unix 时间进行查询,就能获得在该时间范围内生成的所有条目(包含性)。例如,如果我想查询一个 2 毫秒的时间段,可以这样:
难度: 中级
命令: XRANGE
复杂度: XRANGE: O(N)
可用客户端: Redis CLI、C#、Go、Java(同步 - Jedis)、JavaScript(Node.js)、Python、Ruby、Rust(异步)、Rust(同步)
Java(同步 - Jedis)
List<StreamEntry> res13 = jedis.xrange("race:france",String.valueOf(System.currentTimeMillis()-1000),String.valueOf(System.currentTimeMillis()+1000));
System.out.println(
res13
); // >>> [1701764734160-0 {rider=Castilla, speed=30.2, location_id=1, position=1}, 1701764734160-1 {rider=Norem, speed=28.8, location_id=1, position=3}, 1701764734161-0 {rider=Prickett, speed=29.7, location_id=1, position=2}, 1701764734162-0 {rider=Castilla, speed=29.9, location_id=1, position=2}]我在此范围内只有一个条目。但在真实数据集中,我可能查询数小时的范围,或者在两毫秒内就有很多条目,返回的结果可能非常大。因此,XRANGE 支持一个可选的 COUNT 选项。通过指定计数,我可以只获取前 N 个条目。如果需要更多,可以获取最后返回的 ID,将序列部分加 1,然后再次查询。让我们看下面的示例。假设流 race:france 中有 4 个条目。要开始迭代,每次获取 2 个条目,我从完整范围开始,但 COUNT 设为 2。
难度: 中级
命令: XRANGE
复杂度: XRANGE: O(N)
可用客户端: Redis CLI、C#、Go、Java(同步 - Jedis)、JavaScript(Node.js)、Python、Ruby、Rust(异步)、Rust(同步)
Java(同步 - Jedis)
List<StreamEntry> res14 = jedis.xrange("race:france","-","+",2);
System.out.println(res14); // >>> [1701764887638-0 {rider=Castilla, speed=30.2, location_id=1, position=1}, 1701764887638-1 {rider=Norem, speed=28.8, location_id=1, position=3}]要继续迭代接下来的两个条目,我需要获取最后返回的 ID,即 1692632094485-0,并在其前面加上前缀 (。得到的排他范围区间,即本例中的 (1692632094485-0,现在可以作为下一次 XRANGE 调用的新 start 参数:
难度: 中级
命令: XRANGE
复杂度: XRANGE: O(N)
可用客户端: Redis CLI、C#、Go、Java(同步 - Jedis)、JavaScript(Node.js)、Python、Ruby、Rust(异步)、Rust(同步)
Java(同步 - Jedis)
List<StreamEntry> res15 = jedis.xrange("race:france",String.valueOf(System.currentTimeMillis()-1000)+"-0","+",2);
System.out.println(res15); // >>> [1701764887638-0 {rider=Castilla, speed=30.2, location_id=1, position=1}, 1701764887638-1 {rider=Norem, speed=28.8, location_id=1, position=3}]现在我们已从流中获取了全部 4 个条目(流原本只有 4 个条目),如果尝试获取更多,将得到一个空数组:
难度: 中级
命令: XRANGE
复杂度: XRANGE: O(N)
可用客户端: Redis CLI、C#、Go、Java(同步 - Jedis)、JavaScript(Node.js)、Python、Ruby、Rust(异步)、Rust(同步)
Java(同步 - Jedis)
List<StreamEntry> res16 = jedis.xrange("race:france",String.valueOf(System.currentTimeMillis()+1000)+"-0","+",2);
System.out.println(res16); // >>> []由于 XRANGE 的复杂度为 O(log(N)) 用于定位,然后 O(M) 用于返回 M 个元素,在较小的 COUNT 下,命令具有对数时间复杂度,这意味着迭代的每一步都很快。因此 XRANGE 实际上也是 流迭代器,不需要 XSCAN 命令。
XREVRANGE 是 XRANGE 的逆序版本,返回逆序元素,因此 XREVRANGE 的一个实际用途是检查流中的最后一项:
命令: XREVRANGE
复杂度: XREVRANGE: O(N)
可用客户端: Redis CLI、C#、Go、Java(同步 - Jedis)、JavaScript(Node.js)、Python、Ruby、Rust(异步)、Rust(同步)
Java(同步 - Jedis)
List<StreamEntry> res17 = jedis.xrevrange("race:france","+","-",1);
System.out.println(res17); // >>> [1701765218592-0 {rider=Castilla, speed=29.9, location_id=1, position=2}]注意 XREVRANGE 命令的 start 和 stop 参数顺序是相反的。
使用 XREAD 监听新条目
当我们不想通过范围访问流中的条目时,通常想要的是 订阅 到达流的新条目。这个概念可能类似于 Redis Pub/Sub(订阅频道)或 Redis 阻塞列表(等待键获取新元素),但在消费流的方式上存在根本差异:
- 一个流可以有多个客户端(消费者)等待数据。默认情况下,每个新条目将投递给 每个 正在等待该流数据的消费者。这种行为不同于阻塞列表(每个消费者获得不同元素)。然而,能够 扇出 给多个消费者类似于 Pub/Sub。
- 在 Pub/Sub 中,消息是 即发即弃 且从不存储,而在使用阻塞列表时,消息被客户端接收时会从列表中 弹出(实际上被移除),流的工作方式则根本不同。所有消息都会无限期地追加到流中(除非用户明确要求删除条目):不同的消费者通过记住最后接收消息的 ID 来了解从自身角度看哪些是新消息。
- 流的消费者组提供了 Pub/Sub 或阻塞列表无法实现的控制级别,例如同一流的不同组、对已处理条目的显式确认、检查待处理条目的能力、认领未处理消息的能力,以及每个客户端仅能看到自己私有历史消息的一致性历史可见性。
用于监听流中新消息的命令称为 XREAD。它比 XRANGE 稍微复杂一些,因此我们将从简单形式开始,然后再介绍完整命令布局。
命令: XREAD
可用客户端: Redis CLI、C#、Go、Java(同步 - Jedis)、JavaScript(Node.js)、Python、Ruby、Rust(异步)、Rust(同步)
Java(同步 - Jedis)
List<Map.Entry<String, List<StreamEntry>>> res18= jedis.xread(XReadParams.xReadParams().count(2),new HashMap<String,StreamEntryID>(){{put("race:france",new StreamEntryID());}});
System.out.println(
res18
); // >>> [race:france=[1701765384638-0 {rider=Castilla, speed=30.2, location_id=1, position=1}, 1701765384638-1 {rider=Norem, speed=28.8, location_id=1, position=3}]]以上是 XREAD 的非阻塞形式。注意 COUNT 选项不是必须的,事实上命令唯一必须的选项是 STREAMS,它指定一个键列表以及调用消费者对每个流已看到的最大 ID,以便命令只向客户端提供 ID 大于指定值的消息。
在上述命令中,我们写了 STREAMS race:france 0,因此我们想要流 race:france 中所有 ID 大于 0-0 的消息。如示例所示,命令返回键名,因为实际上可以同时调用此命令读取多个不同流。例如,我可以写:STREAMS race:france race:italy 0 0。注意在 STREAMS 选项之后,我们需要提供键名,然后是 ID。因此 STREAMS 选项必须始终是最后一个选项。 任何其他选项必须放在 STREAMS 选项之前。
除了 XREAD 可以同时访问多个流,并且我们可以指定自己拥有的最后 ID 以仅获取更新的消息之外,在这种简单形式下,该命令与 XRANGE 并没有太大区别。然而有趣的是,我们可以通过指定 BLOCK 参数将 XREAD 变为 阻塞命令:
> XREAD BLOCK 0 STREAMS race:france $注意在上述示例中,除了去掉 COUNT,我还指定了新的 BLOCK 选项,超时时间为 0 毫秒(表示永不超时)。此外,我没有为流 race:france 传递普通 ID,而是传递了特殊 ID $。这个特殊 ID 表示 XREAD 应将流 race:france 中已存储的最大 ID 作为最后 ID,因此我们将只收到 新 消息,从开始监听之时起。这类似于 Unix 命令 tail -f 的某些方面。
注意当使用 BLOCK 选项时,不必使用特殊 ID $。我们可以使用任何有效 ID。如果命令能够立即满足请求而不阻塞,则会立即返回,否则会阻塞。通常如果我们想从新条目开始消费流,我们从 ID $ 开始,之后继续使用最后接收消息的 ID 进行下一次调用,依此类推。
XREAD 的阻塞形式也能够监听多个流,只需指定多个键名。如果请求可以同步满足(因为至少有一个流中存在大于对应指定 ID 的元素),则返回结果。否则,命令将阻塞,并在第一个有新数据(根据指定 ID)的流上返回其条目。
与阻塞列表操作类似,从等待数据的客户端角度来看,阻塞流读取是 公平 的,因为语义是 FIFO 风格。第一个为给定流阻塞的客户端将是新条目可用时第一个被唤醒的客户端。
XREAD 除了 COUNT 和 BLOCK 之外没有其他选项,因此它是一个非常基本的命令,特定用途是将消费者挂接到一个或多个流。更强大的消费流功能可通过消费者组 API 实现,但通过消费者组读取由另一个命令 XREADGROUP 实现,将在下一节介绍。
消费者组
当任务是从不同客户端消费同一个流时,XREAD 已经提供了一种 扇出 给 N 个客户端的方式,也可能使用副本提供更多读取可扩展性。但在某些问题中,我们想要的不是向许多客户端提供相同的消息流,而是向许多客户端提供来自同一流的不同 子集 消息。一个明显的用例是处理速度慢的消息:拥有 N 个不同的工作进程,每个工作进程接收流的不同部分,这样我们可以通过将不同消息路由给准备好执行更多工作的不同工作进程来扩展消息处理。
实际上,假设我们有三个消费者 C1、C2、C3 和一个包含消息 1、2、3、4、5、6、7 的流,我们希望按以下图示分配消息:
1 -> C1
2 -> C2
3 -> C3
4 -> C1
5 -> C2
6 -> C3
7 -> C1为了实现这一点,Redis 使用了一个称为 消费者组 的概念。需要特别指出的是,Redis 消费者组在实现上与 Kafka (TM) 消费者组毫无关系。但在功能上相似,因此我决定保留 Kafka (TM) 的术语,因为它最初推广了这一概念。
消费者组就像一个从流中获取数据的 伪消费者,实际上为多个消费者提供服务,并提供以下保证:
- 每条消息投递给不同的消费者,因此同一消息不可能投递给多个消费者。
- 消费者在消费者组内通过名称标识,该名称是客户端(实现消费者的程序)必须选择的区分大小写的字符串。这意味着即使在断开连接后,流消费者组也会保留所有状态,因为客户端会再次声明自己是同一个消费者。但这也意味着客户端有责任提供唯一标识符。
- 每个消费者组都有 从未消费的第一条 ID 的概念,因此当消费者请求新消息时,它只会提供以前未投递的消息。
- 然而,消费消息需要使用特定命令进行显式确认。Redis 将此确认解释为:消息已正确处理,因此可以从消费者组中移除。
- 消费者组跟踪所有当前待处理的消息,即已投递给消费者组中某个消费者但尚未确认已处理的消息。借助此功能,在访问流的历史记录时,每个消费者 只会看到投递给自己的消息。
从某种意义上说,消费者组可以被想象为关于流的一些 状态:
+----------------------------------------+
| consumer_group_name: mygroup |
| consumer_group_stream: somekey |
| last_delivered_id: 1292309234234-92 |
| |
| consumers: |
| "consumer-1" with pending messages |
| 1292309234234-4 |
| 1292309234232-8 |
| "consumer-42" with pending messages |
| ... (and so forth) |
+----------------------------------------+从这个角度来看,很容易理解消费者组能做什么,如何能够仅为消费者提供其待处理消息的历史记录,以及请求新消息的消费者如何只能获得 ID 大于 last_delivered_id 的消息。同时,如果您将消费者组视为 Redis 流的辅助数据结构,那么显然单个流可以有多个消费者组,每个组包含不同的消费者集合。实际上,同一个流甚至可能同时有通过 XREAD 读取的客户端和通过 XREADGROUP 在不同消费者组中读取的客户端。
现在是时候深入了解基本的消费者组命令了。它们包括:
XGROUP用于创建、销毁和管理消费者组。XREADGROUP用于通过消费者组从流中读取。XACK允许消费者将待处理消息标记为已正确处理。XNACK允许消费者将待处理消息释放回组而不确认,使其立即可供其他消费者重新投递。XACKDEL将确认和删除合并为单个原子操作,并增强了对消费者组引用的控制。
创建消费者组
假设我已经有一个类型为流的键 race:france,要创建一个消费者组,只需执行:
命令: XGROUP
复杂度: XGROUP: 取决于子命令
可用客户端: Redis CLI、C#、Go、Java(同步 - Jedis)、JavaScript(Node.js)、Python、Ruby、Rust(异步)、Rust(同步)
Java(同步 - Jedis)
String res19 = jedis.xgroupCreate("race:france","france_riders",StreamEntryID.LAST_ENTRY,false);
System.out.println(res19); // >>> OK如上所示,创建消费者组时需要指定一个 ID,在示例中就是 $。这是因为消费者组除了其他状态外,还必须知道第一个消费者连接时应该提供什么消息,即组刚创建时的 最后消息 ID。如果我们像这样提供 $,那么组中的消费者只会收到从现在开始到达流中的新消息。如果指定 0,则消费者组将消费流历史中的 所有 消息作为开始。当然,您可以指定任何其他有效 ID。您知道的是,消费者组将开始投递 ID 大于您指定 ID 的消息。因为 $ 表示流中当前最大的 ID,指定 $ 的效果是只消费新消息。
XGROUP CREATE 还支持通过可选的 MKSTREAM 子命令(作为最后一个参数)自动创建流(如果不存在):
难度: 中级
命令: XGROUP
复杂度: XGROUP: 取决于子命令
可用客户端: Redis CLI、C#、Go、Java(同步 - Jedis)、JavaScript(Node.js)、Python、Ruby、Rust(异步)、Rust(同步)
Java(同步 - Jedis)
String res20 = jedis.xgroupCreate("race:italy","italy_riders",StreamEntryID.LAST_ENTRY,true);
System.out.println(res20); // >>> OK现在消费者组已创建,我们可以立即使用 XREADGROUP 命令通过消费者组读取消息。我们将使用名为 Alice 和 Bob 的消费者来演示系统如何向 Alice 或 Bob 返回不同的消息。
XREADGROUP 与 XREAD 非常相似,并提供相同的 BLOCK 选项,否则是一个同步命令。但是有一个 必须 始终指定的选项 GROUP,它有两个参数:消费者组名称和尝试读取的消费者名称。COUNT 选项也受支持,与 XREAD 中的相同。
我们将向 race:italy 流添加车手并尝试使用消费者组读取:
可用客户端: Redis CLI、C#、Go、Java(同步 - Jedis)、JavaScript(Node.js)、Python、Ruby、Rust(异步)、Rust(同步)
Java(同步 - Jedis)
StreamEntryID id1 = jedis.xadd("race:italy", new HashMap<String,String>(){{put("rider","Castilaa");}},XAddParams.xAddParams());
StreamEntryID id2 = jedis.xadd("race:italy", new HashMap<String,String>(){{put("rider","Royce");}},XAddParams.xAddParams());
StreamEntryID id3 = jedis.xadd("race:italy", new HashMap<String,String>(){{put("rider","Sam-Bodden");}},XAddParams.xAddParams());
StreamEntryID id4 = jedis.xadd("race:italy", new HashMap<String,String>(){{put("rider","Prickett");}},XAddParams.xAddParams());
StreamEntryID id5 = jedis.xadd("race:italy", new HashMap<String,String>(){{put("rider","Norem");}},XAddParams.xAddParams());
List<Map.Entry<String, List<StreamEntry>>> res21 = jedis.xreadGroup("italy_riders","Alice", XReadGroupParams.xReadGroupParams().count(1),new HashMap<String,StreamEntryID>(){{put("race:italy",StreamEntryID.UNRECEIVED_ENTRY);}});
System.out.println(res21); // >>> [race:italy=[1701766299006-0 {rider=Castilaa}]]XREADGROUP 的回复与 XREAD 类似。注意上面的 GROUP <group-name> <consumer-name>。它表示我想使用消费者组 mygroup 从流中读取,我是消费者 Alice。每次消费者对消费者组执行操作时,都必须指定其名称,唯一标识该组内的这个消费者。
命令行中还有一个非常重要的细节,在必须的 STREAMS 选项之后,为键 race:italy 请求的 ID 是特殊 ID >。这个特殊 ID 仅在消费者组上下文中有效,表示:尚未投递给其他消费者的消息。
这几乎总是您想要的,但也可以指定真实 ID,例如 0 或任何其他有效 ID,在这种情况下,我们向 XREADGROUP 请求的是 待处理消息的历史记录,并且在这种情况下,永远不会看到组中的新消息。因此,基本上 XREADGROUP 根据我们指定的 ID 具有以下行为:
- 如果 ID 是特殊 ID
>,则命令只返回尚未投递给其他消费者的新消息,并且作为副作用,会更新消费者组的 最后 ID。 - 如果 ID 是任何其他有效数字 ID,则命令将允许我们访问 待处理消息的历史记录。即,已投递给此指定消费者(由提供的名称标识)且尚未通过
XACK确认的消息集。
我们可以立即测试此行为,指定 ID 为 0,不指定 COUNT 选项:我们将只看到唯一的待处理消息,即关于 Castilla 的消息:
难度: 中级
命令: XREADGROUP
可用客户端: Redis CLI、C#、Go、Java(同步 - Jedis)、JavaScript(Node.js)、Python、Ruby、Rust(异步)、Rust(同步)
Java(同步 - Jedis)
List<Map.Entry<String, List<StreamEntry>>> res22 = jedis.xreadGroup("italy_riders","Alice", XReadGroupParams.xReadGroupParams().count(1),new HashMap<String,StreamEntryID>(){{put("race:italy",new StreamEntryID());}});
System.out.println(res22); // >>> [race:italy=[1701766299006-0 {rider=Castilaa}]]然而,如果我们将消息确认为已处理,它将不再属于待处理消息历史,因此系统将不再报告任何内容:
命令: XACK, XREADGROUP
复杂度: XACK: O(1)
可用客户端: Redis CLI、C#、Go、Java(同步 - Jedis)、JavaScript(Node.js)、Python、Ruby、Rust(异步)、Rust(同步)
Java(同步 - Jedis)
long res23 = jedis.xack("race:italy","italy_riders",id1);
System.out.println(res23); // >>> 1
List<Map.Entry<String, List<StreamEntry>>> res24 = jedis.xreadGroup("italy_riders","Alice", XReadGroupParams.xReadGroupParams().count(1),new HashMap<String,StreamEntryID>(){{put("race:italy",new StreamEntryID());}});
System.out.println(res24); // >>> [race:italy=[]]如果您还不了解 XACK 的工作原理,没关系,关键是已处理的消息不再是我们能访问的历史的一部分。
现在轮到 Bob 读取一些内容:
难度: 中级
命令: XREADGROUP
可用客户端: Redis CLI、C#、Go、Java(同步 - Jedis)、JavaScript(Node.js)、Python、Ruby、Rust(异步)、Rust(同步)
Java(同步 - Jedis)
List<Map.Entry<String, List<StreamEntry>>> res25 = jedis.xreadGroup("italy_riders","Bob", XReadGroupParams.xReadGroupParams().count(2),new HashMap<String,StreamEntryID>(){{put("race:italy",StreamEntryID.UNRECEIVED_ENTRY);}});
System.out.println(res25); // >>> [race:italy=[1701767632261-1 {rider=Royce}, 1701767632262-0 {rider=Sam-Bodden}]]Bob 请求最多两条消息,并通过同一个组 mygroup 读取。结果 Redis 只报告 新 消息。如您所见,“Castilla”消息并未投递给 Bob,因为它已经投递给 Alice,因此 Bob 获得了 Royce 和 Sam-Bodden,依此类推。
这样,Alice、Bob 以及组中的任何其他消费者都可以从同一流中读取不同的消息,读取各自尚未处理的消息历史,或将消息标记为已处理。这允许创建不同的拓扑结构和语义来消费流中的消息。
需要记住几点:
- 消费者在首次被提及时会自动创建,无需显式创建。
- 即使使用
XREADGROUP,您也可以同时从多个键读取,但为此,您需要在每个流中创建同名的消费者组。这不是常见需求,但值得一提的是,该功能在技术上是可用的。 XREADGROUP是一个 写命令,因为尽管它从流中读取,但消费者组作为读取的副作用被修改,因此它只能在主节点上调用。
以下是一个使用 Ruby 语言编写的消费者实现示例,使用消费者组。Ruby 代码旨在让任何有经验的程序员都能读懂,即使不懂 Ruby:
require 'redis'
if ARGV.length == 0
puts "Please specify a consumer name"
exit 1
end
ConsumerName = ARGV[0]
GroupName = "mygroup"
r = Redis.new
def process_message(id,msg)
puts "[#{ConsumerName}] #{id} = #{msg.inspect}"
end
$lastid = '0-0'
puts "Consumer #{ConsumerName} starting..."
check_backlog = true
while true
# 根据迭代选择 ID:第一次我们要读取待处理消息,以防崩溃后恢复。
# 一旦消费完历史记录,我们就可以开始获取新消息。
if check_backlog
myid = $lastid
else
myid = '>'
end
items = r.xreadgroup('GROUP',GroupName,ConsumerName,'BLOCK','2000','COUNT','10','STREAMS',:my_stream_key,myid)
if items == nil
puts "Timeout!"
next
end
# 如果收到空回复,表示我们正在消费历史记录且历史记录已空。开始消费新消息。
check_backlog = false if items[0][1].length == 0
items[0][1].each{|i|
id,fields = i
# 处理消息
process_message(id,fields)
# 确认消息已处理
r.xack(:my_stream_key,GroupName,id)
$lastid = id
}
end如您所见,这里的思路是从消费历史记录(即我们的待处理消息列表)开始。这很有用,因为消费者可能在之前崩溃,所以在重启时我们希望重新读取已投递给我们的但尚未确认的消息。请注意,我们可能会多次处理同一条消息(至少在消费者故障的情况下,但 Redis 持久化和复制也存在限制,请参阅关于此主题的专门章节)。
一旦历史记录消费完毕,并且我们得到空消息列表,我们就可以切换到使用 > 特殊 ID 来消费新消息。