Redis Streams 消息队列入门与实战
Raj Kumar | 2026-08-27T20:54:52 | Database
详解 Redis Streams 数据结构,通过消费者组实现可靠的消息队列,对比 Kafka 和 RabbitMQ 的适用场景。
# Redis Streams 消息队列入门与实战 ## 什么是 Redis Streams? Redis Streams 是 Redis 5.0 引入的日志型数据结构,支持消费者组、消息确认、持久化,可作为轻量级消息队列。 ## 基础操作 ```bash # 发送消息 XADD mystream * name "order-created" data '{"orderId":1001}' # 返回消息 ID: 1693456789012-0 # 读取消息(从头开始) XRANGE mystream - + COUNT 10 # 读取最新消息(阻塞等待) XREAD BLOCK 5000 STREAMS mystream $ ``` ## 消费者组 ```bash # 创建消费者组(从最新消息开始消费) XGROUP CREATE mystream order-processors $ MKSTREAM # 消费者读取消息 XREADGROUP GROUP order-processors consumer-1 COUNT 5 BLOCK 2000 STREAMS mystream > # 确认消息已处理 XACK mystream order-processors 1693456789012-0 ``` ## Python 客户端示例 ```python import redis import json import time r = redis.Redis(host='localhost', port=6379, decode_responses=True) # 生产者 def publish_order(order_data): msg_id = r.xadd('orders', { 'event': 'order_created', 'data': json.dumps(order_data), }) print('Published: {}'.format(msg_id)) return msg_id # 消费者 def consume_orders(group, consumer_name): # 确保消费者组存在 try: r.xgroup_create('orders', group, id='0', mkstream=True) except redis.exceptions.ResponseError: pass # 组已存在 while True: messages = r.xreadgroup( group, consumer_name, {'orders': '>'}, count=10, block=5000, ) for stream, msgs in messages: for msg_id, data in msgs: try: order = json.loads(data['data']) print('Processing order: {}'.format(order)) # ... 业务处理 ... r.xack('orders', group, msg_id) except Exception as e: print('Error: {}'.format(e)) # 消息不 ACK,会在 XPENDING 中保留 # 查看未确认的消息 def check_pending(group): pending = r.xpending('orders', group) print('Pending messages: {}'.format(pending['pending'])) return pending ``` ## 消息重试(处理死信) ```python def claim_stale_messages(group, consumer, min_idle_ms=60000): # 获取超时未确认的消息 pending = r.xpending_range('orders', group, '-', '+', 10) stale_ids = [ p['message_id'] for p in pending if p['time_since_delivered'] > min_idle_ms ] if stale_ids: # 认领超时消息 claimed = r.xclaim('orders', group, consumer, min_idle_ms, stale_ids) print('Claimed {} stale messages'.format(len(claimed))) ``` ## 与 Kafka/RabbitMQ 对比 | 特性 | Redis Streams | Kafka | RabbitMQ | |------|--------------|-------|----------| | 持久化 | RDB/AOF | 磁盘 | 磁盘 | | 吞吐量 | 中 (万/秒) | 极高 (百万/秒) | 中 | | 消费者组 | 支持 | 支持 | 不原生 | | 运维复杂度 | 低 | 高 | 中 | | 适用场景 | 轻量级队列 | 大数据流 | 企业集成 | ## 总结 Redis Streams 适合中小规模消息队列场景,已经用 Redis 的项目无需引入额外组件。数据量大或高可靠性要求选 Kafka。