Redis 5.0 引入了 Stream 数据类型,它像一个只追加的日志结构,支持按消息 ID 进行范围读取。XREAD 命令是读取 Stream 消息的基础手段,可以从某个游标位置向后获取数据,也可以阻塞等待新消息到达。与传统的 List 配合 BLPOP 相比,XREAD 的游标控制更灵活,消息 ID 也提供了明确的位置标记。

一、XREAD 命令基础与消息 ID 游标
XREAD 命令的标准语法如下:
XREAD [COUNT count] [BLOCK milliseconds] STREAMS key [key ...] ID [ID ...]
其中 COUNT 用于限制单次读取的消息数量,BLOCK 用于设置阻塞毫秒数,STREAMS 后面依次给出键名和对应的读取起点 ID。ID 参数可以是具体消息 ID、0 或 $。0 表示从 Stream 头部开始读取,$ 表示从当前最大 ID 之后开始读取,也就是只读取执行命令之后新到达的消息。举例来说:
# 写入两条消息 XADD mystream * sensor_id 1 temp 23.5 XADD mystream * sensor_id 2 temp 24.0 # 从头部读取所有消息 XREAD COUNT 10 STREAMS mystream 0
返回结果是一个数组结构,每个元素对应一个 Stream,包含 Stream 名称和消息列表。每条消息由消息 ID 和字段-值对组成。消息 ID 的格式为毫秒时间戳-序号,例如 1717000000000-0。这个 ID 是 Redis 自动生成的,既保证了时间上的大致顺序,也解决了同一毫秒内的冲突。掌握 ID 游标是使用 XREAD 的关键:下次读取时可以把上次处理完的最后一条消息 ID 作为起点,实现增量消费。
不过需要注意,XREAD 只是读取消息,并不会把消息从 Stream 中删除。即使消息被读取后,Stream 的长度不会因此缩小,其他客户端用同样的 ID 起点依然能读到相同的内容。这种特性让 XREAD 非常适合广播式消费场景,但在任务队列里就可能带来重复处理的问题,后面会专门讨论。
另外,当需要同时读取多个 Stream 时,STREAMS 后面的 ID 数量必须与键数量一致,并且按顺序对应。如果某个 Stream 不存在,使用 0 读取时会被当作空 Stream 返回,不会报错。
二、XREAD 阻塞读取与 COUNT 参数
如果希望在消息到达时立即得到通知,而不是轮询,可以使用 BLOCK 参数。BLOCK 后面跟毫秒数,0 表示无限期阻塞。例如下面的命令会阻塞 5 秒,等待 mystream 中出现新消息:
# 阻塞 5 秒等待新消息 XREAD COUNT 5 BLOCK 5000 STREAMS mystream $
当阻塞期间有新消息 XADD 写入 mystream,XREAD 会立即返回,并带上新消息。如果 5 秒内没有新消息,则返回空值(nil)。实际使用中,建议设置一个有限的阻塞时间,比如 5000 或 10000 毫秒,这样客户端可以在超时后做一些健康检查或重新建立连接。
COUNT 参数用于限制本次读取的最大消息数量。在非阻塞模式下,COUNT 可以避免一次性拉取太多消息导致内存压力。在阻塞模式下,COUNT 的作用依然有效,一旦有消息到达,Redis 会返回最多 COUNT 条新消息,不会等到凑满 COUNT 条才返回。例如设置 COUNT 10,但阻塞期间只写入了 3 条消息,那么命令会立即返回这 3 条消息。
需要注意,BLOCK 与 $ 组合时,只会读取阻塞开始之后的新消息,不会读取阻塞之前已经存在但尚未被消费的历史消息。如果业务希望先把历史遗留消息处理完,再进入阻塞等待新消息,就需要先以 0 或某个具体 ID 为起点做一次非阻塞读取,处理完后再用 $ 进入阻塞循环。
三、XREAD 的局限:没有消费确认与重复读取问题
XREAD 最大的局限在于它没有消费确认机制,也没有消费者组的概念。用 XREAD 读取消息后,Stream 中的消息仍然保留,其他客户端用相同的起点 ID 读取时,可以拿到完全一样的数据。例如两个客户端都执行 XREAD COUNT 10 STREAMS mystream 0,它们会得到相同的 10 条消息。这对于日志采集、事件广播等场景很合适,因为多个消费者都需要各自独立地看到完整数据。
但在任务队列场景中,这就会造成重复消费。如果多个消费者都从 0 开始读取,每个任务可能被处理多次。即使只使用一个消费者,如果在处理完消息后没有及时记录游标而崩溃,重启后必须从之前的 ID 重新读取,也会导致部分消息重复。因此,使用 XREAD 构建的任务队列需要业务层自行保证幂等性,例如通过唯一业务 ID 去重、把已处理的消息 ID 写入另一个 Set 等方式。
如果业务不能容忍重复消费,或者需要多个消费者分组处理同一批消息,就应该考虑消费组模式,使用 XREADGROUP、XACK 等命令。消费组会为每个消费者维护独立的待处理列表,只有收到 XACK 确认后消息才会从待处理列表移除,从而实现更可靠的投递语义。XREAD 更适合无需确认、允许重复读取的日志管道或通知流。
四、实践:用 XREAD 构建轻量消息读取循环
下面是一段 Python 伪代码,展示如何用 Redis 客户端循环读取消息,并维护游标避免重复处理已经处理过的记录:
import redis
r = redis.Redis(host='127.0.0.1', port=6379, decode_responses=True)
last_id = '0'
while True:
try:
resp = r.xread({'mystream': last_id}, count=10, block=5000)
if not resp:
continue
for stream, messages in resp:
for msg_id, fields in messages:
print(f'处理消息 {msg_id}: {fields}')
# 实际业务处理逻辑
last_id = msg_id
except Exception as e:
print(f'处理异常: {e}')
# 可以在此处记录日志,稍后重试
这段代码的核心是使用变量 last_id 保存上次处理完成的最后一条消息 ID。每次调用 xread 时,都从该 ID 之后开始读取。这样正常情况下不会重复处理消息,但如果进程在处理完消息后、更新 last_id 之前崩溃,那么重启后仍然会从旧的 ID 开始读取,从而重复处理一条消息。因此,业务处理逻辑最好设计成幂等操作。
此外,循环中使用 block=5000 阻塞等待 5 秒,超时后返回空列表,然后继续下一轮循环。这样既减少了空轮询对 CPU 的消耗,又能在超时后处理一些维护任务。还需要注意网络异常、Redis 服务重启等情况,客户端应当具备自动重连机制,一些成熟的 Redis 客户端库会自动处理。
如果希望进一步提高可靠性,可以把游标持久化到 Redis 或本地文件中,以便进程重启后恢复。但即便持久化了游标,也只能保证“至少一次”投递,仍需要业务幂等来兜底。相比之下,XREADGROUP 配合 XACK 可以在 Redis 层面管理待确认列表,实现更完善的消息确认。根据实际场景选择合适的读取方式,是使用 Redis Stream 时的第一步。
总的来说,XREAD 是读取 Stream 消息的基础命令,掌握消息 ID 游标、阻塞参数和 COUNT 限制,就能构建出简单高效的消息读取循环。但它并不负责消息确认,理解这一局限有助于在可靠性和实现成本之间做出权衡。
Redis XREADStream流消息读取修改时间:2026-09-19 17:48:22