2.9 流 (Streams) (Redis 5.0+) Redis Streams (Redis 5.0+): 现代数据流的强大引擎 Redis 一直以来以其高性能的键值存储和丰富的数据类型而闻名。自 5.0 版本开始,Redis 引入了一种全新的数据类型—— Streams (流),为开发者提供了处理实时数据流的强大工具。Streams 的出现填补了 Redis 在处理高吞吐量、持久化、可消费的消息流方面的空白,使其在消息队列、事件溯源、实时分析等领域拥有了更广阔的应用前景。 Streams 的诞生背景:弥补 Redis 的短板 在 Redis 5.
Redis 一直以来以其高性能的键值存储和丰富的数据类型而闻名。自 5.0 版本开始,Redis 引入了一种全新的数据类型—— Streams (流),为开发者提供了处理实时数据流的强大工具。Streams 的出现填补了 Redis 在处理高吞吐量、持久化、可消费的消息流方面的空白,使其在消息队列、事件溯源、实时分析等领域拥有了更广阔的应用前景。
在 Redis 5.0 之前,开发者在处理消息队列和数据流时,通常会选择以下几种方式:
List (列表): List 可以作为简单的消息队列使用,支持生产者使用 LPUSH 或 RPUSH 添加消息,消费者使用 LPOP 或 RPOP 获取消息。然而,List 存在一些明显的不足:
无持久化消费组: List 无法追踪哪些消息已经被哪些消费者处理,难以实现可靠的消费组模式。
阻塞式读取的局限: BLPOP 和 BRPOP 可以实现阻塞式读取,但只能由单个消费者阻塞在一个 List 上,难以扩展到多个消费者。
消息丢失风险: 如果消费者在处理消息过程中崩溃,且消息未被持久化,则可能导致消息丢失。
Pub/Sub (发布/订阅): Pub/Sub 提供了发布和订阅消息的机制,适用于实时广播场景。但 Pub/Sub 也存在一些限制:
无持久化: 消息发布后,如果没有订阅者在线,则消息会直接丢失。
无消息回溯: 订阅者只能接收到订阅之后发布的新消息,无法回溯历史消息。
无消费者状态管理: Pub/Sub 不维护消费者状态,难以实现复杂的消费逻辑和消息确认机制。
面对这些局限性,Redis 需要一种更强大的数据类型来应对现代应用中日益增长的实时数据处理需求。Streams 正是在这样的背景下应运而生,它借鉴了消息队列和日志系统的设计思想,提供了持久化、有序、可消费的消息流,并引入了消费者组的概念,极大地提升了 Redis 在数据流处理方面的能力。
要深入理解 Streams,首先需要掌握其核心概念:
Stream Key (流键): 与 Redis 其他数据类型一样,Stream 也是通过键来标识的。Stream Key 就是用来访问和操作特定流的键名。
Entry (条目): Stream 中存储的最小单元是 Entry,每个 Entry 代表一条消息或一个事件。一个 Entry 由以下部分组成:
ID (ID): 每个 Entry 都有一个唯一的 ID,用于标识和排序 Entry。ID 由毫秒级时间戳和序列号组成,例如 1678886400000-1。ID 可以由 Redis 自动生成,也可以由用户自定义 (但通常不推荐自定义 ID)。
Field-Value Pairs (字段-值对): Entry 的内容以键值对的形式存储,类似于 Hash 数据类型。一个 Entry 可以包含多个字段-值对。
Consumer Group (消费者组): 这是 Streams 最重要的概念之一。消费者组允许多个消费者协作消费同一个 Stream 中的消息,实现负载均衡和消息的可靠消费。
Group Name (组名): 消费者组的名称,用于唯一标识一个消费者组。
Consumers (消费者): 属于同一个消费者组的客户端。
Last Delivered ID (最后投递 ID): 消费者组维护一个 Last Delivered ID,用于记录组内消费者最后消费的消息 ID。新加入消费者组的消费者将从 Last Delivered ID 之后的消息开始消费。
Pending Entries List (PEL) (待处理条目列表): 每个消费者组都维护一个 PEL,用于跟踪已经投递给消费者但尚未被确认 (acknowledged) 的消息。当消费者崩溃或超时未确认消息时,PEL 中的消息可以被其他消费者重新消费,保证消息的可靠性。
Consumer (消费者): 客户端程序,通过消费者组来消费 Stream 中的消息。每个消费者在消费者组内需要有一个唯一的名称。
Pending Entries List (PEL) (待处理条目列表): 如上所述,PEL 是消费者组的核心机制,用于保证消息的可靠消费。PEL 记录了已经投递给消费者但尚未被确认的消息 ID、投递给哪个消费者以及投递时间。
总结 Streams 的核心特点:
持久化: Streams 中的数据会被持久化存储到 Redis 数据库中,即使 Redis 服务器重启,数据也不会丢失。
有序性: Streams 中的 Entry 按照 ID 的顺序严格排序,保证消息的有序性。
可消费组: 通过消费者组,可以实现多个消费者并行消费同一个 Stream,提高消费效率,并保证消息的可靠消费。
消息确认机制: 消费者需要显式地确认 (acknowledge) 已经处理的消息,确保消息不会被重复处理,并支持消息的重试和故障恢复。
阻塞式读取: 消费者可以使用阻塞式命令等待新消息的到来,避免轮询,提高效率。
接下来,我们将通过实际的 Redis 命令和 Python 代码示例,详细介绍 Streams 的常用操作。我们将使用 redis-py Python 客户端库进行演示。
XADD 命令XADD 命令用于向 Stream 中添加新的 Entry。
命令语法:
XADD key [NOMKSTREAM] [MAXLEN|MAXLEN~ [threshold]] *|ID field value [field value ...]
参数解释:
key: Stream 的键名。
NOMKSTREAM: 可选参数,如果 Stream 不存在,则不创建 Stream,直接返回错误。
MAXLEN [threshold]: 可选参数,用于限制 Stream 的最大长度。
MAXLEN: 指定 Stream 的最大长度。当 Stream 的 Entry 数量超过最大长度时,旧的 Entry 会被删除。
~ [threshold]: 近似修剪。Redis 会删除至少 threshold 个旧的 Entry,但实际删除的数量可能更多,以提高性能。
*: 使用 * 表示让 Redis 自动生成 Entry ID。这是最常用的方式。
ID: 用户自定义的 Entry ID。不推荐使用,除非有特殊需求。
field value [field value ...]: Entry 的字段-值对。
Redis-CLI 示例:
> XADD mystream * sensor-id 123 temperature 25.5 humidity 60 1678886400000-0 > XADD mystream * sensor-id 124 temperature 26.0 humidity 62 1678886400001-0
以上命令向名为 mystream 的 Stream 中添加了两个 Entry。每个 Entry 包含 sensor-id, temperature, humidity 三个字段。Redis 自动生成了 Entry ID。
Python 代码示例:
import redis r = redis.Redis(host='localhost', port=6379, db=0) # 添加一个 Entry entry_id = r.xadd('mystream', {'sensor-id': 123, 'temperature': 25.5, 'humidity': 60}) print(f"Added entry with ID: {entry_id}") # 添加另一个 Entry entry_id = r.xadd('mystream', {'sensor-id': 124, 'temperature': 26.0, 'humidity': 62}) print(f"Added entry with ID: {entry_id}") # 使用 MAXLEN 限制 Stream 长度 entry_id = r.xadd('capped_stream', {'data': 'value'}, maxlen=100, approximate=True) # approximate=True 等价于 MAXLEN~ print(f"Added entry to capped stream with ID: {entry_id}")
XRANGE 和 XREVRANGE 命令XRANGE 命令用于读取指定 ID 范围内的 Entry,按照 ID 升序排列。XREVRANGE 命令与 XRANGE 类似,但按照 ID 降序排列。
命令语法:
XRANGE key start end [COUNT count] XREVRANGE key end start [COUNT count]
参数解释:
key: Stream 的键名。
start: 起始 ID。可以使用 - 表示最小值,+ 表示最大值。
end: 结束 ID。可以使用 - 表示最小值,+ 表示最大值。
COUNT count: 可选参数,限制返回的 Entry 数量。
Redis-CLI 示例:
> XRANGE mystream - + 1) 1) "1678886400000-0" 2) 1) "sensor-id" 2) "123" 3) "temperature" 4) "25.5" 5) "humidity" 6) "60" 2) 1) "1678886400001-0" 2) 1) "sensor-id" 2) "124" 3) "temperature" 4) "26.0" 5) "humidity" 6) "62" > XRANGE mystream 1678886400000-0 + COUNT 1 1) 1) "1678886400000-0" 2) 1) "sensor-id" 2) "123" 3) "temperature" 4) "25.5" 5) "humidity" 6) "60" > XREVRANGE mystream + - COUNT 1 1) 1) "1678886400001-0" 2) 1) "sensor-id" 2) "124" 3) "temperature" 4) "26.0" 5) "humidity" 6) "62"
Python 代码示例:
# 读取所有 Entry entries = r.xrange('mystream', min='-', max='+') print("All entries:", entries) # 读取指定范围的 Entry entries = r.xrange('mystream', min='1678886400000-0', max='+') print("Entries from ID 1678886400000-0:", entries) # 读取最新的 Entry (使用 XREVRANGE) entries = r.xrevrange('mystream', max='+', min='-', count=1) print("Latest entry:", entries)
XREAD 命令XREAD 命令用于从一个或多个 Stream 中读取新的 Entry。它可以实现阻塞式读取,等待新消息的到来。
命令语法:
XREAD [COUNT count] [BLOCK milliseconds] [STREAMS key [key ... ] ID [ID ...]]
参数解释:
COUNT count: 可选参数,限制每个 Stream 返回的 Entry 数量。
BLOCK milliseconds: 可选参数,指定阻塞等待的时间 (毫秒)。如果设置为 0,则表示永不超时。
STREAMS key [key ... ] ID [ID ... ]: 指定要读取的 Stream 键名和起始 ID。
key [key ... ]: 要读取的 Stream 键名,可以指定多个 Stream。
ID [ID ... ]: 每个 Stream 对应的起始 ID。可以使用 $ 表示读取最新的 Entry 之后的新消息,0 表示从 Stream 的开头开始读取。
Redis-CLI 示例:
# 读取 mystream 中 ID 大于 0 的所有新消息 (实际上就是从头开始读,因为ID不可能小于0) > XREAD STREAMS mystream 0 1) 1) "mystream" 2) 1) 1) "1678886400000-0" 2) 1) "sensor-id" 2) "123" 3) "temperature" 4) "25.5" 5) "humidity" 6) "60" 2) 1) "1678886400001-0" 2) 1) "sensor-id" 2) "124" 3) "temperature" 4) "26.0" 5) "humidity" 6) "62" # 读取 mystream 中最新的消息之后的新消息 (使用 $) 并阻塞 5 秒 > XREAD BLOCK 5000 STREAMS mystream $ (等待新消息...) # 在另一个 Redis-CLI 窗口中添加新消息 > XADD mystream * new-data value 1678886405000-0 # 阻塞的 XREAD 命令立即返回新消息 1) 1) "mystream" 2) 1) 1) "1678886405000-0" 2) 1) "new-data" 2) "value"
Python 代码示例:
# 读取 mystream 中 ID 大于 0 的所有新消息 entries = r.xread(streams={'mystream': '0'}) print("Read entries from beginning:", entries) # 读取 mystream 中最新的消息之后的新消息并阻塞 5 秒 entries = r.xread(streams={'mystream': '$'}, block=5000) print("Read new entries with blocking:", entries) # 同时读取多个 Stream entries = r.xread(streams={'mystream': '0', 'another_stream': '0'}) print("Read entries from multiple streams:", entries)
XREADGROUP 命令XREADGROUP 命令是 Streams 最核心的命令之一,用于从消费者组中读取消息。
命令语法:
XREADGROUP GROUP groupname consumername [COUNT count] [BLOCK milliseconds] [NOACK] STREAMS key [key ...] ID [ID ...]
参数解释:
GROUP groupname consumername: 指定消费者组名称和消费者名称。
COUNT count: 可选参数,限制每个 Stream 返回的 Entry 数量。
BLOCK milliseconds: 可选参数,指定阻塞等待的时间 (毫秒)。
NOACK: 可选参数,表示自动确认消息。不建议使用,因为它会降低消息的可靠性。
STREAMS key [key ... ] ID [ID ... ]: 指定要读取的 Stream 键名和起始 ID。
key [key ... ]: 要读取的 Stream 键名,可以指定多个 Stream。
ID [ID ... ]: 每个 Stream 对应的起始 ID。
>: 表示读取消费者组中尚未被消费的新消息。这是最常用的方式。
$: 表示从 Stream 的末尾开始读取,通常用于新创建的消费者组。
0: 表示从 Stream 的开头开始读取 (谨慎使用,可能导致重复消费)。
具体 ID: 可以指定具体的 ID 开始读取,用于回溯历史消息。
首次使用消费者组需要先创建消费者组:
命令语法:
XGROUP CREATE key groupname ID [$ | 0 | entry-id] [MKSTREAM]
参数解释:
key: Stream 的键名。
groupname: 消费者组名称。
ID: 消费者组的起始 ID。
$: 表示从 Stream 的末尾开始消费新消息。
0: 表示从 Stream 的开头开始消费所有消息 (谨慎使用)。
entry-id: 指定具体的 Entry ID 作为起始 ID。
MKSTREAM: 可选参数,如果 Stream 不存在,则创建 Stream。
Redis-CLI 示例:
# 创建消费者组 mygroup,从 Stream 末尾开始消费新消息 > XGROUP CREATE mystream mygroup $ MKSTREAM OK # 消费者 consumer1 从 mygroup 消费者组读取 mystream 的新消息 (>) 并阻塞 5 秒 > XREADGROUP GROUP mygroup consumer1 BLOCK 5000 COUNT 1 STREAMS mystream > (等待新消息...) # 在另一个 Redis-CLI 窗口中添加新消息 > XADD mystream * data-field new-message 1678886410000-0 # 阻塞的 XREADGROUP 命令立即返回新消息 1) 1) "mystream" 2) 1) 1) "1678886410000-0" 2) 1) "data-field" 2) "new-message" # 消费者 consumer2 也从 mygroup 消费者组读取 mystream 的新消息 (>) 并阻塞 5 秒 > XREADGROUP GROUP mygroup consumer2 BLOCK 5000 COUNT 1 STREAMS mystream > (等待新消息...) # 再次添加新消息 > XADD mystream * another-field another-message 1678886415000-0 # 这次新消息被 consumer2 消费 (负载均衡) 1) 1) "mystream" 2) 1) 1) "1678886415000-0" 2) 1) "another-field" 2) "another-message"
Python 代码示例:
# 创建消费者组 mygroup,从 Stream 末尾开始消费新消息 try: r.xgroup_create(name='mystream', groupname='mygroup', id='$', mkstream=True) except redis.exceptions.ResponseError as e: if 'BUSYGROUP Consumer Group name already exists' not in str(e): raise # 如果不是组已存在错误,则抛出异常 print("Consumer group already exists, skipping creation.") # 消费者 consumer1 从 mygroup 消费者组读取 mystream 的新消息并阻塞 5 秒 entries = r.xreadgroup(groupname='mygroup', consumername='consumer1', streams={'mystream': '>'}, block=5000, count=1) print("Consumer1 read entries:", entries) # 消费者 consumer2 从 mygroup 消费者组读取 mystream 的新消息并阻塞 5 秒 entries = r.xreadgroup(groupname='mygroup', consumername='consumer2', streams={'mystream': '>'}, block=5000, count=1) print("Consumer2 read entries:", entries)
XACK 命令XACK 命令用于消费者确认已经成功处理了从消费者组中读取的消息。
命令语法:
XACK key groupname ID [ID ...]
参数解释:
key: Stream 的键名。
groupname: 消费者组名称。
ID [ID ... ]: 要确认的消息 ID,可以确认多个消息。
Redis-CLI 示例:
# 假设 consumer1 消费了消息 1678886410000-0,需要进行确认 > XACK mystream mygroup 1678886410000-0 (integer) 1
Python 代码示例:
# 假设 consumer1 读取到 entries,获取消息 ID 并确认 if entries: stream_name, messages = entries[0] for message_id, message_data in messages: # ... 处理消息 ... ack_count = r.xack('mystream', 'mygroup', message_id) print(f"Acknowledged message ID: {message_id}, ACK count: {ack_count}")
重要提示: 务必在消费者成功处理消息后调用 XACK 命令进行确认。 如果消费者在处理消息过程中崩溃或未及时确认,消息将保留在 PEL 中,并可能被其他消费者重新消费,保证消息的可靠性。
XPENDING 命令XPENDING 命令用于查看消费者组的 PEL,了解哪些消息已经被投递但尚未被确认。
命令语法:
XPENDING key groupname [start end count] [consumername]
参数解释:
key: Stream 的键名。
groupname: 消费者组名称。
start end count: 可选参数,用于分页查看 PEL。
start: 起始消息 ID,可以使用 - 表示最小值。
end: 结束消息 ID,可以使用 + 表示最大值。
count: 返回的消息数量。
consumername: 可选参数,只查看指定消费者的 PEL。
Redis-CLI 示例:
> XPENDING mystream mygroup 1) (integer) 1 # PEL 中待处理消息总数 2) "1678886410000-0" # PEL 中最早的消息 ID 3) "1678886410000-0" # PEL 中最新的消息 ID 4) 1) 1) "consumer1" # 消费者名称 2) "1" # 该消费者在 PEL 中的消息数量 > XPENDING mystream mygroup - + 10 1) 1) "1678886410000-0" 2) "consumer1" 3) (integer) 123456789 # 消息投递时间戳 (毫秒) 4) (integer) 1 # 消息被投递的次数 (重试次数)
Python 代码示例:
# 获取 PEL 概览信息 pending_info = r.xpending('mystream', 'mygroup') print("Pending info:", pending_info) # 获取 PEL 详细信息 pending_entries = r.xpending_range('mystream', 'mygroup', '-', '+', 10) print("Pending entries:", pending_entries) # 获取指定消费者的 PEL 详细信息 pending_entries_consumer = r.xpending_range('mystream', 'mygroup', '-', '+', 10, consumername='consumer1') print("Pending entries for consumer1:", pending_entries_consumer)
XCLAIM 命令XCLAIM 命令用于从 PEL 中认领消息。当消费者崩溃或超时未确认消息时,其他消费者可以使用 XCLAIM 命令将 PEL 中的消息转移到自己名下,并重新处理。
命令语法:
XCLAIM key groupname consumername min-idle-time ID [ID ...] [RETRYCOUNT count] [FORCE] [JUSTID]
参数解释:
key: Stream 的键名。
groupname: 消费者组名称。
consumername: 认领消息的消费者名称。
min-idle-time: 最小空闲时间 (毫秒)。只有 PEL 中消息的空闲时间超过这个值,才会被认领。空闲时间是指消息最后一次被投递到现在的时间间隔。
ID [ID ... ]: 要认领的消息 ID,可以认领多个消息。
RETRYCOUNT count: 可选参数,设置消息的重试计数器。
FORCE: 可选参数,强制认领消息,即使消息的空闲时间没有达到 min-idle-time。
JUSTID: 可选参数,只返回认领的消息 ID,不返回消息内容。
Redis-CLI 示例:
# 消费者 consumer2 认领 PEL 中空闲时间超过 60 秒的消息 > XCLAIM mystream mygroup consumer2 60000 1678886410000-0 1) 1) "1678886410000-0" 2) 1) "data-field" 2) "new-message"
Python 代码示例:
# 消费者 consumer2 认领 PEL 中空闲时间超过 60 秒的消息 claimed_entries = r.xclaim('mystream', 'mygroup', 'consumer2', 60000, ['1678886410000-0']) print("Claimed entries:", claimed_entries) # 消费者 consumer3 认领 PEL 中所有空闲时间超过 60 秒的消息 pending_entries_all = r.xpending_range('mystream', 'mygroup', '-', '+', 1000) # 获取 PEL 中所有消息 (假设不超过 1000 条) claimable_ids = [entry[0] for entry in pending_entries_all if entry[3] > 60000] # 筛选空闲时间超过 60 秒的消息 ID if claimable_ids: claimed_entries = r.xclaim('mystream', 'mygroup', 'consumer3', 60000, claimable_ids) print("Consumer3 claimed entries:", claimed_entries) else: print("No claimable entries found.")
XDEL 命令XDEL 命令用于从 Stream 中删除指定的 Entry。通常情况下,不建议手动删除 Stream 中的 Entry,因为 Streams 的设计目标是持久化和有序的消息流。 删除操作可能会破坏消息的完整性和顺序性。
命令语法:
XDEL key ID [ID ...]
参数解释:
key: Stream 的键名。
ID [ID ... ]: 要删除的 Entry ID,可以删除多个 Entry。
Redis-CLI 示例:
> XDEL mystream 1678886400000-0 (integer) 1
Python 代码示例:
# 删除指定的 Entry delete_count = r.xdel('mystream', '1678886400000-0') print(f"Deleted {delete_count} entry.")
XINFO 命令XINFO 命令用于获取 Stream 的各种信息,例如 Stream 的长度、消费者组信息、消费者信息等。
命令语法:
XINFO [STREAM key | GROUPS key | CONSUMERS key groupname]
参数解释:
STREAM key: 获取 Stream 的基本信息。
GROUPS key: 获取 Stream 的消费者组信息。
CONSUMERS key groupname: 获取指定消费者组的消费者信息.
Redis-CLI 示例:
> XINFO STREAM mystream 1) "length" 2) (integer) 2 3) "groups-len" 4) (integer) 1 5) "first-entry" 6) 1) "1678886400001-0" 2) 1) "sensor-id" 2) "124" 3) "temperature" 4) "26" 5) "humidity" 6) "62" 7) "last-entry" 8) 1) "1678886415000-0" 2) 1) "another-field" 2) "another-message" 9) "max-deleted-entry-id" 10) "0-0" 11) "entries-added" 12) (integer) 3 > XINFO GROUPS mystream 1) 1) "name" 2) "mygroup" 3) "consumers" 4) (integer) 2 5) "pending" 6) (integer) 1 7) "last-delivered-id" 8) "1678886415000-0" > XINFO CONSUMERS mystream mygroup 1) 1) "name" 2) "consumer1" 3) "pending" 4) (integer) 1 5) "idle" 6) (integer) 123456789
Python 代码示例:
# 获取 Stream 信息 stream_info = r.xinfo_stream('mystream') print("Stream info:", stream_info) # 获取消费者组信息 groups_info = r.xinfo_groups('mystream') print("Groups info:", groups_info) # 获取指定消费者组的消费者信息 consumers_info = r.xinfo_consumers('mystream', 'mygroup') print("Consumers info:", consumers_info)