redis数据结构学习——stream
概念
Redis Stream 是 Redis 5.0 引入的持久化消息队列数据结构。你可以把它想象成一个只能追加的日志文件,每条消息都有一个自动生成的、按时间排序的唯一 ID(格式如 1633000000000-0,前半是毫秒时间戳,后半是序列号)
它的核心特性:
有序:消息按加入顺序严格排列
持久化:消息写入后会落到 Redis 的数据文件中,重启不丢失
消费者组(Consumer Group):多个消费者可以组成一个组,分工消费消息,每条消息只被一个消费者处理
ACK 机制:消费者处理完消息后需要发送 ACK 确认,未确认的消息可以被重新消费
阻塞读取:消费者可以阻塞等待新消息到来
使用场景
拍卖行竞拍,一个拍卖品一个Stream。
性能
拍卖品有数百个,创建上百个 Stream 完全没问题,Redis 轻松应对。但如果是几万、几十万个商品,就需要优化了。
原因:
Redis 本身非常轻量,每个 Stream 的底层是一个基数树(Radix Tree)结构,即使有上万个 Stream,Redis 的内存开销和管理成本也很低。具体来说:
内存开销:每个空的 Stream 约占几百字节到 1KB 的内存,100 个就是 100KB 左右,微不足道。
连接开销:消费者通过 XREADGROUP 阻塞读取,100 个 Stream 意味着 100 个阻塞连接。Redis 单机支持数万连接,100 个完全没压力。
CPU 开销:每个阻塞读取在有消息到达时才会唤醒,空闲时几乎零 CPU。
什么时候会有问题?
如果拍卖行有 10 万个活跃商品,每个商品一个 Stream,就会出现:
连接数爆炸:10 万个阻塞连接,虽然 Redis 能扛,但网络和系统资源浪费严重。
消费者线程爆炸:如果每个 Stream 单开一个线程消费,10 万个线程直接打垮服务器。
Stream 元数据内存:10 万个 Stream 的元数据可能占用数百 MB 内存。
运维复杂:监控、排查问题困难。
优化方案:分片 + 合并消费
对于大规模场景,我们不会"一个商品一个 Stream",而是采用分片策略:
方案 1:哈希分片(推荐)
将商品 ID 哈希到固定数量的 Stream 中,例如 1024 个 Stream:
importhashlibdefget_stream_key(auction_id,shard_count=1024):""" 将 auction_id 映射到 1024 个 Stream 中的一个 同一个 auction_id 永远映射到同一个 Stream,保证顺序 """shard=auction_id%shard_countreturnf"bid_stream_shard:{shard}"# 使用示例stream_key=get_stream_key(12345)# 总是返回 bid_stream_shard:xxx优点:
Stream 数量固定(1024 个),无论商品多少都不会膨胀
同一个商品的出价永远进入同一个 Stream,顺序得到保证
消费者数量可控(1024 个或更少,可以多个 Stream 共用一个消费者)
方案 2:动态消费者池
classConsumerPool:"""管理固定数量的消费者,每个消费者负责多个 Stream"""def__init__(self,shard_count=1024,consumer_per_shard=1):self.shard_count=shard_count self.consumers=[]defstart(self):# 启动 1024 个消费者协程,每个负责一个分片forshard_idinrange(self.shard_count):stream_key=f"bid_stream_shard:{shard_id}"consumer=ShardConsumer(shard_id,stream_key)self.consumers.append(consumer)asyncio.create_task(consumer.run())方案 3:批量读取(进一步优化)
如果某些分片流量很低,可以让一个消费者负责多个分片,使用 XREADGROUP 同时读取多个 Stream:
# 一个消费者同时监听多个 Stream(Redis 支持多 Stream 读取)streams={"bid_stream_shard:0":">","bid_stream_shard:1":">","bid_stream_shard:2":">",# ... 最多 1024 个}result=awaitredis.xreadgroup(group_name="bid_group",consumer_name="consumer_0",streams=streams,count=10,block=1000)这样我们可以用少量消费者(如 32 个)处理全部 1024 个分片,大幅降低资源消耗。
