Redis 流 Stream 实现消息队列
FreeGuideOnline
最新
2026-07-11
bash
向名为 mystream 的流中添加一条消息,字段为 sensor_id 和 temperature
XADD mystream * sensor_id "1234" temperature "19.8"
返回值就是自动生成的消息 ID,如 "1672531214822-0"
- `*` 表示让 Redis 自动生成 ID,也可以手动指定,但需保证 ID 大于最后一条。
- 一条消息可包含多个键值对。
**限制流长度**:通过 `MAXLEN` 可以近似控制流的大小,防止内存无限增长。
```bash
XADD mystream MAXLEN ~ 1000 * data "log entry"
使用 ~ 表示近似裁剪,Redis 会在合适的时机进行高效删除。
3.2 读取消息 (XREAD)
# 从最早的消息开始读取,不带阻塞
XREAD STREAMS mystream 0
# 只读取新消息(ID 用 $),并且阻塞 1000 毫秒
XREAD BLOCK 1000 STREAMS mystream $
0代表从头开始,$代表只接收最新消息。BLOCK可使消费者在没有消息时等待指定毫秒数,实现类似阻塞队列的效果。
3.3 创建消费者组 (XGROUP CREATE)
# 创建消费者组 mygroup,从 mystream 的头部开始消费(0)
XGROUP CREATE mystream mygroup 0 MKSTREAM
# 或者从最新消息开始消费(只消费新消息)
XGROUP CREATE mystream mygroup $ MKSTREAM
MKSTREAM会在流不存在时自动创建空流。- 后续可通过
XGROUP SETID调整消费位置。
查看消费者组信息:
XINFO GROUPS mystream
3.4 消费者组读取消息 (XREADGROUP)
# 消费者 consumer1 从 mygroup 中读取最多 2 条未处理的消息,阻塞 2000 毫秒
XREADGROUP GROUP mygroup consumer1 BLOCK 2000 COUNT 2 STREAMS mystream >
>表示只读取从未被分配给该组任何消费者的新消息。- 也可使用特定 ID 或
0来读取已分配但未确认的“待处理消息”(用于重试)。
典型流程:先通过 XREADGROUP 获取消息,处理完成后确认。
4. 消息确认与重试机制
4.1 确认消息处理完成 (XACK)
XACK mystream mygroup 1672531214822-0
处理完消息后必须执行 XACK,否则该消息会一直保留在消费者的 Pending 列表中,且即使消费者重启,消息仍然会被再次分配。
4.2 查看待处理消息 (XPENDING)
# 查看整个负载的概要信息
XPENDING mystream mygroup
# 详细信息:从最旧开始读取 10 条
XPENDING mystream mygroup - + 10
返回每个待处理消息的 ID、被分配的消费者、空闲时间(毫秒)等。
4.3 重新认领超时消息 (XCLAIM / XAUTOCLAIM)
如果某个消费者宕机且消息长时间未确认,可由其他消费者接管:
# 手动认领指定 ID 的消息,最小空闲时间 3600000 毫秒(1 小时)
XCLAIM mystream mygroup consumer2 3600000 1672531214822-0
Redis 6.2 起更推荐使用 XAUTOCLAIM,它会自动扫描并认领符合条件的消息:
XAUTOCLAIM mystream mygroup consumer2 3600000 0-0 COUNT 25
XAUTOCLAIM 返回被认领的消息以及下一次扫描的起始 ID,能更安全地实现自动重试。
5. 流的管理与优化
5.1 限制流长度 (XADD MAXLEN / XTRIM)
除了在 XADD 时使用 MAXLEN,也可以事后修剪:
# 保留最新的 1000 条消息
XTRIM mystream MAXLEN ~ 1000
# 精确裁剪(可能略有性能开销)
XTRIM mystream MINID ~ 1672531200000-0
5.2 查看流信息 (XINFO)
XINFO STREAM mystream:长度、基数、最后生成 ID 等。XINFO GROUPS mystream:所有消费者组信息。XINFO CONSUMERS mystream mygroup:组内每个消费者的 Pending 数量。
5.3 删除消息 (XDEL)
XDEL mystream 1672531214822-0
- 通常只需依赖 MAXLEN 自动删除,手动删除用于纠正错误。
6. 实战:构建一个简单的任务队列
使用 Python 和 redis-py 库演示典型的生产者-消费者组模式。
安装依赖:pip install redis
生产者代码 (producer.py):
import redis
import time
import json
r = redis.Redis(host='localhost', port=6379, decode_responses=True)
stream_name = 'tasks'
def add_task(data):
msg_id = r.xadd(stream_name, {'data': json.dumps(data), 'status': 'pending'})
print(f"Added task {msg_id}")
return msg_id
# 模拟不断产生任务
for i in range(10):
add_task({'task_id': i, 'payload': f'work_{i}'})
time.sleep(0.5)
消费者组初始化 (init_group.py):
import redis
r = redis.Redis(decode_responses=True)
try:
r.xgroup_create('tasks', 'workers', id='0', mkstream=True)
print("Consumer group 'workers' created.")
except redis.exceptions.ResponseError as e:
if 'BUSYGROUP' in str(e):
print("Group already exists.")
else:
raise
工人类 (worker.py):
import redis
import time
import json
r = redis.Redis(decode_responses=True)
stream = 'tasks'
group = 'workers'
consumer_name = f'worker-{os.getpid()}'
def process_message(msg_id, data):
print(f"[{consumer_name}] Processing {msg_id}: {data}")
# 模拟业务处理
time.sleep(1)
# 处理成功,进行确认
r.xack(stream, group, msg_id)
print(f"[{consumer_name}] Acknowledged {msg_id}")
# 启动消费循环
while True:
try:
# 每次读取最多 2 条未分配的新消息,阻塞 5000 毫秒
messages = r.xreadgroup(group, consumer_name, {stream: '>'}, count=2, block=5000)
if not messages:
# 超时,可以在这里执行清理或健康检查
continue
# messages 的结构: [[stream_name, [(msg_id, {field: value}), ...]]]
for stream_name, msg_list in messages:
for msg_id, msg_data in msg_list:
# 处理消息
process_message(msg_id, msg_data.get('data'))
except Exception as e:
print(f"Error: {e}")
time.sleep(1)