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)