概念Redis Stream 是 Redis 5.0 引入的持久化消息队列数据结构。你可以把它想象成一个只能追加的日志文件每条消息都有一个自动生成的、按时间排序的唯一 ID格式如 1633000000000-0前半是毫秒时间戳后半是序列号它的核心特性有序消息按加入顺序严格排列持久化消息写入后会落到 Redis 的数据文件中重启不丢失消费者组Consumer Group多个消费者可以组成一个组分工消费消息每条消息只被一个消费者处理ACK 机制消费者处理完消息后需要发送 ACK 确认未确认的消息可以被重新消费阻塞读取消费者可以阻塞等待新消息到来使用场景拍卖行竞拍一个拍卖品一个Stream。性能拍卖品有数百个创建上百个 Stream 完全没问题Redis 轻松应对。但如果是几万、几十万个商品就需要优化了。原因Redis 本身非常轻量每个 Stream 的底层是一个基数树Radix Tree结构即使有上万个 StreamRedis 的内存开销和管理成本也很低。具体来说内存开销每个空的 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 个 Streamimporthashlibdefget_stream_key(auction_id,shard_count1024): 将 auction_id 映射到 1024 个 Stream 中的一个 同一个 auction_id 永远映射到同一个 Stream保证顺序 shardauction_id%shard_countreturnfbid_stream_shard:{shard}# 使用示例stream_keyget_stream_key(12345)# 总是返回 bid_stream_shard:xxx优点Stream 数量固定1024 个无论商品多少都不会膨胀同一个商品的出价永远进入同一个 Stream顺序得到保证消费者数量可控1024 个或更少可以多个 Stream 共用一个消费者方案 2动态消费者池classConsumerPool:管理固定数量的消费者每个消费者负责多个 Streamdef__init__(self,shard_count1024,consumer_per_shard1):self.shard_countshard_count self.consumers[]defstart(self):# 启动 1024 个消费者协程每个负责一个分片forshard_idinrange(self.shard_count):stream_keyfbid_stream_shard:{shard_id}consumerShardConsumer(shard_id,stream_key)self.consumers.append(consumer)asyncio.create_task(consumer.run())方案 3批量读取进一步优化如果某些分片流量很低可以让一个消费者负责多个分片使用 XREADGROUP 同时读取多个 Stream# 一个消费者同时监听多个 StreamRedis 支持多 Stream 读取streams{bid_stream_shard:0:,bid_stream_shard:1:,bid_stream_shard:2:,# ... 最多 1024 个}resultawaitredis.xreadgroup(group_namebid_group,consumer_nameconsumer_0,streamsstreams,count10,block1000)这样我们可以用少量消费者如 32 个处理全部 1024 个分片大幅降低资源消耗。
郑州网站建设
网页设计
企业官网