Redis Stream消息队列与消费者组实战应用

Redis Stream数据结构与核心概念

Redis Stream是Redis 5.0引入的持久化消息队列数据结构,底层基于Radix Tree实现,支持时间序写入、消费者组、消息确认和Pending列表等核心功能。相比List实现的消息队列,Stream提供了更完整的消息语义:每条消息拥有全局唯一的毫秒时间戳+序列号ID、支持多消费者组独立消费、未确认消息可被重新分配。

Stream的核心数据结构由Entry组成,每个Entry是一个field-value对的有序集合:

# 写入消息
XADD orders * user_id 1001 amount 99.9 status pending
# 返回: 1723445566123-0 (毫秒时间戳-序列号)

# 读取消息
XRANGE orders - +
# 1) 1) "1723445566123-0"
#    2) 1) "user_id"
#       2) "1001"
#       3) "amount"
#       4) "99.9"
#       5) "status"
#       6) "pending"

消费者组与消息分发机制

消费者组(Consumer Group)让多个消费者协作处理同一个Stream的消息,每条消息只会被组内一个消费者处理。XREADGROUP命令从指定消费者组读取消息,last-delivered-id记录组内已投递的最后一条消息ID。

# 创建消费者组
# MKSTREAM: 如果Stream不存在则自动创建
# $符号: 从最新消息开始消费(0表示从头开始)
XGROUP CREATE orders order-processors $ MKSTREAM

# 消费者读取消息
# GROUP: 指定消费者组和消费者名称
# COUNT: 每次最多读取条数
# BLOCK: 阻塞等待毫秒数(0表示不阻塞)
XREADGROUP GROUP order-processors consumer-1 COUNT 10 BLOCK 5000 STREAMS orders >

XREADGROUP的stream参数使用>表示只接收新消息(未被任何消费者读取过的),如果使用0或其他ID则返回Pending列表中的消息。这种设计支持”拉取新消息”和”重试未确认消息”两种模式。

消息确认与Pending列表管理

消费者处理完消息后必须调用XACK确认,否则消息会留在Pending列表中。Pending列表记录了每条已投递但未确认的消息,包含消费者名称、投递时间和已投递次数。

# 确认消息已处理
XACK orders order-processors 1723445566123-0

# 查看Pending列表
XPENDING orders order-processors
# 返回: 待确认消息数、最小ID、最大ID、各消费者的待确认数

# 查看Pending消息详情
XPENDING orders order-processors - + 10

当消费者宕机或长时间未确认消息时,需要将Pending消息转移给其他消费者处理。XPENDING命令的idle参数可筛选空闲时间超过阈值的消息:

# 声明转移:将空闲超过60秒的Pending消息转给当前消费者
XAUTOCLAIM orders order-processors consumer-2 60000 0 COUNT 10

XAUTOCLAIM是Redis 6.2引入的原子操作,替代了早期的XCLAIM+XPENDING组合方案。它同时完成消息筛选和所有权转移,返回被转移的消息列表。

Python消费者实现与容错处理

生产环境的消费者需要处理网络断连、消息处理失败、消费者宕机等异常场景。以下是一个完整的消费者实现:

import redis
import time
import logging

logger = logging.getLogger(__name__)

class StreamConsumer:
    def __init__(self, host, stream, group, consumer_name):
        self.client = redis.Redis(host=host, decode_responses=True)
        self.stream = stream
        self.group = group
        self.consumer = consumer_name
        self._ensure_group()

    def _ensure_group(self):
        """确保消费者组存在"""
        try:
            self.client.xgroup_create(
                self.stream, self.group, id='0', mkstream=True
            )
        except redis.ResponseError as e:
            if 'BUSYGROUP' not in str(e):
                raise

    def consume(self, batch_size=10, block_ms=5000):
        """消费消息主循环"""
        while True:
            try:
                messages = self.client.xreadgroup(
                    self.group, self.consumer,
                    {self.stream: '>'},
                    count=batch_size, block=block_ms
                )
                if not messages:
                    continue
                for stream, entries in messages:
                    for msg_id, fields in entries:
                        self._process_message(msg_id, fields)
            except redis.ConnectionError:
                logger.warning('Redis连接断开,5秒后重连')
                time.sleep(5)
            except Exception as e:
                logger.error(f'消费异常: {e}')
                time.sleep(1)

    def _process_message(self, msg_id, fields):
        """处理单条消息"""
        try:
            self.handle(fields)
            self.client.xack(self.stream, self.group, msg_id)
            logger.info(f'消息处理成功: {msg_id}')
        except Exception as e:
            logger.error(f'消息处理失败: {msg_id}, {e}')

    def handle(self, fields):
        """业务处理逻辑,子类重写"""
        raise NotImplementedError

    def reclaim_stuck(self, idle_ms=120000, count=10):
        """回收卡住的消息"""
        result = self.client.xautoclaim(
            self.stream, self.group, self.consumer,
            idle_ms, '0', count=count
        )
        if result[1]:
            logger.info(f'回收了 {len(result[1])} 条卡住的消息')

Stream容量规划与内存优化

Stream的消息会持续累积占用内存,需要设置上限策略。MAXLEN参数限制Stream的最大长度,超过后自动删除最旧的消息。近似修剪(~符号)比精确修剪性能更好:

# 写入时限制Stream长度(近似修剪,性能优先)
XADD orders MAXLEN ~ 100000 * user_id 1001 amount 99.9

# 手动修剪
XTRIM orders MAXLEN ~ 100000

容量规划经验值:每条Entry占用约100-200字节(含field-value对和元数据),10万条消息约占10-20MB内存。消费者组的Pending列表同样占用内存,处理延迟高的消费者组会显著增加内存消耗。

Stream适用场景的边界判断:消息量在每秒1万以下、不需要严格的事务保证、消费者数量在百级以内——Redis Stream完全胜任。当日吞吐量超过亿级或需要exactly-once语义时,应考虑Kafka或RocketMQ。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/redisstream-xiao-xi-dui-lie-yu-xiao-fei-zhe-zu-shi-zhan/

(0)
小编小编
上一篇 7小时前
下一篇 7小时前

相关推荐