Redis Streams 自 5.0 发布以来,凭借其内存级吞吐和持久化能力、以及 Consumer Groups 语义,已成为轻量级事件驱动架构的首选基础设施。本文深入剖析 Redis Streams 的底层数据结构、消息确认机制、Pending Entries List 内部实现,以及如何基于它构建生产级事件溯源系统。
一、为什么选择 Redis Streams 做事件溯源?
事件溯源(Event Sourcing)将应用程序状态建模为不可变事件的序列。传统方案(如 Apache Kafka、RabbitMQ)功能完善但运维成本高;Redis Streams 以单进程实现事件日志、消费者组、消息确认三者统一,在中小规模场景(10w QPS 以内、百万级事件/天)具备显著优势:
| 维度 | Redis Streams | Apache Kafka | RabbitMQ |
|---|---|---|---|
| 部署复杂度 | 单机/哨兵/集群 | ZooKeeper + Broker | Erlang VM |
| 事件持久化 | AOF + RDB | 磁盘落盘 | 可选 |
| 消费者组语义 | 原生原生 | 原生原生 | 无 |
| 消息回溯 | 基于 ID offset | 基于 offset | 消费后删除 |
| 内存级延迟 | 亚毫秒 | 毫秒级 | 亚毫秒 |
| 事件保留策略 | MAXLEN / MAXLEN ~ | 保留时间/大小 | TTL |
Redis Streams 的核心优势:在不引入外部依赖的情况下提供事件日志 + 消费者组 + 消息确认的完整语义,极其适合与业务同进程部署于中小规模、高吞吐、低延迟的事件驱动系统。
二、数据结构与内存布局
2.1 Radix Tree 作为底层事件日志
Redis Streams 的事件存储使用 Radix Tree(基数树),每个节点代表一个消息 ID,叶子节点存储完整的消息 payload(以 listpack 形式编码)。这种设计的核心优势:
- O(k) 插入/查询:k 为 key 长度(通常 13 位数字)
- 范围查询高效:支持 XRANGE/XREVRANGE 按 ID 范围扫描
- 内存紧凑:listpack 编码避免指针开销
内存结构示意:
Stream "order:events"
│
├── ID "1696156800000-0" → { event_type: "Created", order_id: "O-1001", ... }
├── ID "1696156800123-0" → { event_type: "Paid", order_id: "O-1001", ... }
├── ID "1696156800456-0" → { event_type: "Shipped", order_id: "O-1001", ... }
└── ...
消息 ID 格式为 <timestamp>-<sequence>,保证全局单调递增且可按时间排序。
2.2 Consumer Groups 状态机
每个 Consumer Group 维护以下关键状态:
Group "payment-service"
│
├── last_delivered_id: "1696156800456-0" // 下一条要投递的消息
├── pel: Pending Entries List // 已投递未确认的消息
│ ├── "1696156800123-0" → { consumer: "node-1", delivery_count: 1, ... }
│ └── "1696156800456-0" → { consumer: "node-3", delivery_count: 1, ... }
├── consumers:
│ ├── "node-1": { seen_time: 1696156805000, pending: [...] }
│ └── "node-3": { seen_time: 1696156806000, pending: [...] }
└── ...
Pending Entries List (PEL) 是 Redis Streams 实现"至少一次投递"语义的核心——每条消息投递后进入 PEL,只有收到 XACK 才会被移除。
三、核心操作与生产级 Consumer Groups
3.1 事件写入(生产者)
import redis
import time
import json
from typing import Any
r = redis.Redis(host='localhost', port=6379, decode_responses=True)
def append_event(stream: str, event_type: str, data: dict, max_len: int = 10_000_000) -> str:
"""写入事件到 Stream,控制最大长度实现自动过期"""
event = {
"type": event_type,
"ts": int(time.time() * 1000),
"payload": json.dumps(data, ensure_ascii=False)
}
msg_id = r.xadd(stream, event, maxlen=max_len, approximate=True)
return msg_id
# 使用示例
msg_id = append_event("order:events", "OrderCreated", {
"order_id": "O-1001",
"user_id": "U-042",
"amount": 299.90,
"items": [{"sku": "SKU-001", "qty": 2}]
})
print(f"Event written: {msg_id}")
3.2 消费者组创建
def create_consumer_group(stream: str, group: str, start_id: str = "0",
mkstream: bool = True) -> bool:
"""
创建消费者组
start_id:
- "0": 从当前开始(不读历史事件)
- "$": 只读新事件
- 特定 ID: 从该 ID 开始读(事件溯源重建)
"""
try:
# 使用 XGROUP CREATE 创建组
r.xgroup_create(stream, group, id=start_id, mkstream=mkstream)
return True
except redis.exceptions.ResponseError as e:
if "BUSYGROUP" in str(e):
print(f"Group {group} already exists")
return False
raise
# 场景 1:新服务接入,从最新开始
create_consumer_group("order:events", "payment-service", start_id="0")
# 场景 2:事件溯源重建读模型 — 从特定时间点开始
create_consumer_group("order:events", "analytics-replay",
start_id="1696156800000-0")
3.3 核心消费循环(生产级实现)
import signal
import logging
from redis import ResponseError
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
class StreamEventProcessor:
"""生产级 Stream 事件处理器"""
def __init__(self, r: redis.Redis, stream: str, group: str,
consumer: str, block_ms: int = 2000, batch_size: int = 100):
self.r = r
self.stream = stream
self.group = group
self.consumer = consumer
self.block_ms = block_ms
self.batch_size = batch_size
self.running = True
self.stats = {"processed": 0, "failed": 0, "retried": 0}
def handle_event(self, msg_id: str, event: dict) -> bool:
"""
处理单条事件,返回 True 表示成功(将发送 XACK)
返回 False 表示处理失败(不 ACK,消息留在 PEL)
"""
raise NotImplementedError
def run(self):
"""主消费循环"""
logger.info(f"Consumer {self.consumer} starting on {self.stream}/{self.group}")
# 注册信号处理实现 graceful shutdown
signal.signal(signal.SIGTERM, lambda *_: self.shutdown())
signal.signal(signal.SIGINT, lambda *_: self.shutdown())
while self.running:
try:
# 1. 先处理 PEL 中的未确认消息(重试逻辑)
self._process_pending()
# 2. 读取新事件
self._process_new()
except ResponseError as e:
logger.error(f"Redis error: {e}")
time.sleep(1)
except Exception as e:
logger.exception(f"Unexpected error: {e}")
time.sleep(1)
logger.info(f"Consumer {self.consumer} stopped. Stats: {self.stats}")
def _process_pending(self):
"""处理 Pending Entries List 中的未确认消息"""
# XREADGROUP 使用 ID "0" 读取 PEL 中本消费者的未确认消息
entries = self.r.xreadgroup(
self.group, self.consumer,
{self.stream: ">"}, # 注意:">" 在 pending 场景有特殊语义
count=self.batch_size
)
if not entries or not entries[0][1]:
entries = self.r.xreadgroup(
self.group, self.consumer,
{self.stream: "0"}, # "0" 读取 pending
count=self.batch_size
)
return self._process_batch(entries)
def _process_new(self):
"""读取并处理新事件"""
entries = self.r.xreadgroup(
self.group, self.consumer,
{self.stream: ">"}, # ">" 表示只投递新消息
count=self.batch_size,
block=self.block_ms // 阻塞等待
)
return self._process_batch(entries)
def _process_batch(self, entries) -> int:
if not entries or not entries[0][1]:
return 0
stream_data, messages = entries[0]
ack_list = []
for msg_id, event in messages:
try:
success = self.handle_event(msg_id, event)
if success:
ack_list.append(msg_id)
self.stats["processed"] += 1
else:
# 处理失败,不 ACK,消息下次重试
self.stats["failed"] += 1
logger.warning(f"Event {msg_id} failed: {event}")
except Exception as e:
self.stats["retried"] += 1
logger.exception(f"Exception processing {msg_id}: {e}")
# 批量 ACK
if ack_list:
self.r.xack(self.stream, self.group, *ack_list)
return len(ack_list)
def shutdown(self):
logger.info("Shutting down gracefully...")
self.running = False
# 应用示例
class PaymentEventProcessor(StreamEventProcessor):
"""支付事件处理器:确保幂等性"""
def __init__(self, *args, **kwargs):
super().__init__(*args, **kwargs)
# 幂等性:记录最近处理的 event sequence
self.processed_ids = set() # 生产环境应使用 Redis SET 做分布式去重
def handle_event(self, msg_id: str, event: dict) -> bool:
event_type = event.get("type")
payload = json.loads(event.get("payload", "{}"))
order_id = payload.get("order_id")
# 幂等性检查
if msg_id in self.processed_ids:
logger.info(f"Duplicate event {msg_id}, skipping")
return True
if event_type == "OrderCreated":
self._on_order_created(order_id, payload)
elif event_type == "PaymentAuthorized":
self._on_payment_auth(order_id, payload)
elif event_type == "PaymentCaptured":
self._on_payment_capture(order_id, payload)
self.processed_ids.add(msg_id)
return True
def _on_order_created(self, order_id, payload):
logger.info(f"[{order_id}] Order created: {payload['amount']}")
# 业务逻辑:验证库存、锁定优惠券等
def _on_payment_auth(self, order_id, payload):
logger.info(f"[{order_id}] Payment authorized: {provider}")
def _on_payment_capture(self, order_id, payload):
logger.info(f"[{order_id}] Payment captured")
# 启动处理器
processor = PaymentEventProcessor(
r, stream="order:events",
group="payment-service", consumer="node-1"
)
processor.run()
3.4 消息确认与重试机制
def manual_ack_with_retry(r: redis.Redis, stream: str, group: str,
msg_id: str, max_retries: int = 3):
"""
手动确认 + 重试逻辑
超过 max_retries 后进入死信流
"""
# 获取该消息的投递计数
pending = r.xpending_range(stream, group, "-", "+", count=100)
for entry in pending:
if entry['message_id'] == msg_id:
if entry['times_delivered'] > max_retries:
# 超过最大重试 → 写入死信流并 ACK
r.xadd(f"{stream}:dead-letter", {
"original_id": msg_id,
"group": group,
"reason": "max_retries_exceeded"
})
r.xack(stream, group, msg_id)
logger.warning(f"Message {msg_id} moved to dead letter")
return
else:
# 不 ACK 让它在 PEL 中保留,下次 XAUTOCLAIM 重新认领
return
# 不在 PEL 中(可能已被其他消费者 ACK),直接确认
r.xack(stream, group, msg_id)
四、XAUTOCLAIM:自动认领死信消息
当消费者崩溃(如 OOM、网络分区),其 pending 消息会一直留在 PEL 中。Redis 6.2 引入 XAUTOCLAIM 自动将超时未确认的消息转移给当前消费者。
def reclaim_dead_messages(r: redis.Redis, stream: str, group: str,
consumer: str, min_idle_ms: int = 60_000):
"""
自动认领空闲超过 min_idle_ms 的消息
返回 (next_start_id, reclaimed_messages, deleted_ids)
"""
# XAUTOCLAIM 是原子操作:扫描 PEL → 转移 ownership → 返回消息
result = r.xautoclaim(
stream, group, consumer,
min_idle_time=min_idle_ms,
start_id="0-0", # 从头扫描 PEL
count=100
)
next_id, messages, deleted = result
for msg_id, event in messages:
try:
# 处理消息...
process_result = handle_event(msg_id, event)
if process_result:
r.xack(stream, group, msg_id)
else:
logger.warning(f"Reclaimed event {msg_id} failed again")
except Exception as e:
logger.exception(f"Reclaimed event {msg_id} exception: {e}")
return messages
# 定时任务:每分钟认领一次死亡消息
import schedule
def maintenance_task():
reclaimed = reclaim_dead_messages(r, "order:events",
"payment-service", "node-1")
logger.info(f"Reclaimed {len(reclaimed)} dead messages")
schedule.every(1).minutes.do(maintenance_task)
五、事件溯源模式:基于 Streams 实现 CQRS + Event Store
5.1 架构设计
┌─────────────┐ ┌──────────────────┐ ┌─────────────────┐
│ Command │────▶│ Event Store │────▶│ Event Handlers │
│ Service │ │ (Redis Streams) │ │ (Consumers) │
└─────────────┘ └──────────────────┘ └────────┬────────┘
│
┌──────────────────────────┼──────────┐
▼ ▼ ▼
┌──────────┐ ┌──────────┐ ┌────────┐
│ Read DB │ │ Projections│ │External│
│ (View) │ │ (Cache) │ │ Sidefx │
└──────────┘ └────────────┘ └────────┘
5.2 事件存储与快照
class EventStore:
"""基于 Redis Streams 的轻量级事件存储"""
def __init__(self, r: redis.Redis, stream: str):
self.r = r
self.stream = stream
self._snapshot_key = f"{stream}:snapshots"
def append(self, aggregate_id: str, event_type: str,
data: dict, expected_version: int = None) -> str:
"""
追加事件,支持乐观锁版本校验
expected_version: 期望的当前版本,None 表示不校验
"""
# 获取聚合当前版本
current_version = self.get_aggregate_version(aggregate_id)
if expected_version is not None and expected_version != current_version:
raise ConcurrencyError(
f"Version mismatch: expected {expected_version}, got {current_version}"
)
event = {
"aggregate_id": aggregate_id,
"type": event_type,
"version": current_version + 1,
"data": json.dumps(data, ensure_ascii=False),
"ts": int(time.time() * 1000)
}
return self.r.xadd(self.stream, event)
def get_events(self, aggregate_id: str, after_version: int = 0) -> list:
"""获取聚合的所有事件(事件溯源重建)"""
# 使用 XRANGE 扫描全量,生产环境应使用 XINFO STREAM + XRANGE 分页
all_events = []
last_id = "0-0"
while True:
entries = self.r.xrange(self.stream, min=last_id, count=1000)
if not entries:
break
for msg_id, event in entries:
if event.get("aggregate_id") == aggregate_id:
version = int(event.get("version", 0))
if version > after_version:
all_events.append({
"id": msg_id,
"version": version,
"type": event["type"],
"data": json.loads(event["data"])
})
last_id = increment_id(msg_id)
if len(entries) < 1000:
break
return all_events
def get_aggregate_version(self, aggregate_id: str) -> int:
"""获取聚合当前版本号"""
version_key = f"{self.stream}:version:{aggregate_id}"
v = self.r.get(version_key)
return int(v) if v else 0
def save_snapshot(self, aggregate_id: str, state: dict, version: int):
"""保存聚合快照,加速事件溯源重建"""
snapshot = {
"version": version,
"state": json.dumps(state, ensure_ascii=False),
"ts": int(time.time() * 1000)
}
self.r.hset(self._snapshot_key, aggregate_id, json.dumps(snapshot))
def load_snapshot(self, aggregate_id: str) -> dict:
"""加载最近的聚合快照"""
raw = self.r.hget(self._snapshot_key, aggregate_id)
if not raw:
return None
return json.loads(raw)
# 使用示例:订单聚合的事件溯源
class OrderAggregate:
def __init__(self, store: EventStore, order_id: str):
self.store = store
self.order_id = order_id
self._rebuild_state()
def _rebuild_state(self):
"""从事件流重建聚合状态"""
# 1. 加载快照
snapshot = self.store.load_snapshot(self.order_id)
start_version = snapshot["version"] if snapshot else 0
state = snapshot["state"] if snapshot else OrderState.EMPTY
# 2. 重放快照之后的事件
events = self.store.get_events(self.order_id, after_version=start_version)
for event in events:
state = self._apply_event(state, event)
self.state = state
self.version = start_version + len(events)
def create_order(self, user_id: str, items: list) -> str:
"""业务方法:创建订单"""
return self.store.append(
aggregate_id=self.order_id,
event_type="OrderCreated",
data={"user_id": user_id, "items": items},
expected_version=self.version # 乐观锁
)
六、生产级陷阱与最佳实践
6.1 内存管理:MAXLEN 策略
Redis Streams 不自动删除事件(与列表不同),必须显式限制长度:
# 生产环境必须设置 maxlen,否则 Stream 无限增长
# approximate=True 性能更高(不是精确修剪)
r.xadd("order:events", event, maxlen=1_000_000, approximate=True)
# 或者使用 XTRIM
r.xtrim("order:events", maxlen=1_000_000, approximate=True)
# 按时间过期(使用 MINID — Redis 6.2+)
# 保留最近 7 天的消息(毫秒时间戳)
cutoff = int(time.time() * 1000) - 7 * 24 * 3600 * 1000
r.xtrim("order:events", minid=f"{cutoff}-0", approximate=True)
6.2 PEL 膨胀与监控
def monitor_stream_health(r: redis.Redis, stream: str):
"""监控 Streams 健康状况"""
# 总消息数
length = r.xlen(stream)
# 消费者组信息
groups = r.xinfo_groups(stream)
for g in groups:
group_name = g["name"]
pending = g["pending"] # 未确认消息数
consumers = g["consumers"] # 消费者数量
last_delivered = g["last-delivered-id"]
# 告警条件
if pending > 100_000:
logger.warning(f"[{group_name}] PEL 膨胀: {pending} 条未确认")
if consumers == 0 and length > 0:
logger.error(f"[{group_name}] 无消费者但有 {length} 条消息积压")
# Pending Entries 详情
pending_info = r.xpending(stream, group_name)
print(f"Pending range: {pending_info['min']} ~ {pending_info['max']}")
print(f"Total pending: {pending_info['pending']}")
6.3 消费者心跳与故障检测
import threading
class HeartbeatConsumer:
"""带心跳的消费者,便于监控和自动故障转移"""
def __init__(self, r: redis.Redis, stream: str, group: str, consumer: str):
self.r = r
self.stream = stream
self.group = group
self.consumer = consumer
self.heartbeat_key = f"{stream}:{group}:heartbeat:{consumer}"
def start(self):
"""启动心跳线程"""
self._stop_heartbeat = threading.Event()
self._heartbeat_thread = threading.Thread(
target=self._heartbeat_loop, daemon=True
)
self._heartbeat_thread.start()
def _heartbeat_loop(self):
"""每 5 秒更新一次心跳(TTL 15 秒)"""
while not self._stop_heartbeat.is_set():
self.r.set(self.heartbeat_key, str(int(time.time())), ex=15)
self._stop_heartbeat.wait(timeout=5)
def stop(self):
self._stop_heartbeat.set()
self.r.delete(self.heartbeat_key)
# 检查所有活跃消费者
def get_active_consumers(r: redis.Redis, stream: str, group: str) -> list:
"""获取所有有心跳的消费者"""
pattern = f"{stream}:{group}:heartbeat:*"
keys = r.keys(pattern)
consumers = []
for key in keys:
ttl = r.ttl(key)
if ttl > 0: # TTL 有效 = 活跃
consumer_name = key.split(":")[-1]
last_heartbeat = r.get(key)
consumers.append({
"name": consumer_name,
"ttl": ttl,
"last_heartbeat": last_heartbeat
})
return consumers
6.4 消息幂等性生产实战
from functools import wraps
import hashlib
class IdempotentProcessor:
"""
基于 Redis SET 实现分布式幂等消费
原理:每条事件有唯一 ID(msg_id),用 SETNX 记录处理状态
"""
def __init__(self, r: redis.Redis, ttl_seconds: int = 86400 * 3):
self.r = r
self.idempotency_prefix = "stream:idempotency:"
self.ttl = ttl_seconds # 保留 3 天的去重记录
def is_processed(self, msg_id: str) -> bool:
return self.r.exists(f"{self.idempotency_prefix}{msg_id}") == 1
def mark_processed(self, msg_id: str):
self.r.set(f"{self.idempotency_prefix}{msg_id}", "1", ex=self.ttl)
def process_idempotent(self, msg_id: str, handler, *args, **kwargs):
"""幂等处理:确保同一事件只执行一次"""
if self.is_processed(msg_id):
logger.info(f"Event {msg_id} already processed — skip")
return None
result = handler(*args, **kwargs)
self.mark_processed(msg_id)
return result
# 使用装饰器简化
def idempotent(r: redis.Redis, key_func=None):
"""幂等装饰器:保证同一事件只处理一次"""
def decorator(func):
@wraps(func)
def wrapper(msg_id: str, event: dict, *args, **kwargs):
cache_key = f"idempotent:{func.__name__}:{msg_id}"
# SETNX:只有不存在时才执行
if r.set(cache_key, "1", nx=True, ex=86400):
return func(msg_id, event, *args, **kwargs)
else:
logger.info(f"Skip duplicate event {msg_id}")
return None
return wrapper
return decorator
@idempotent(r)
def handle_order_created(msg_id: str, event: dict):
payload = json.loads(event["payload"])
print(f"Processing order: {payload['order_id']}")
return True
七、性能调优与集群部署
7.1 Pipeline 批量写入
def batch_append_events(r: redis.Redis, stream: str, events: list, batch_size: int = 500):
"""使用 Pipeline 批量追加事件,减少网络往返"""
with r.pipeline(transaction=False) as pipe:
for i, event in enumerate(events, 1):
pipe.xadd(stream, event, maxlen=1_000_000, approximate=True)
if i % batch_size == 0:
pipe.execute()
pipe.execute()
7.2 集群模式注意事项
Redis Cluster 下使用 Streams 需要注意:
# 1. 使用 Hash Tag 确保同一聚合的事件落在同一节点
# {aggregate_id} 作为 hash tag
r.xadd("order:events:{O-1001}", event) # 按 O-1001 hash
r.xadd("order:events:{O-1001}", event2) # 同一位置
# 2. 消费者组创建必须在同一节点
# Stream 和 Group 必须有相同的 hash tag
r.xgroup_create("order:events:{O-1001}", "payment-group")
# 3. 跨 slot 操作需要所有 Stream 在同一 hash tag 组
7.3 性能测试基准
环境: AWS c6g.4xlarge (16 vCPU, 32GB, Graviton2)
Redis 7.0, 单实例, AOF everysec
┌──────────────────────┬────────────┬────────────┬──────────┐
│ 操作 │ 10w ops/s │ P99 延迟 │ QPS │
├──────────────────────┼────────────┼────────────┼──────────┤
│ XADD (单条) │ 18.5 │ 0.05ms │ 185,000 │
│ XADD Pipeline×500 │ 95.2 │ 0.005ms │ 952,000 │
│ XREADGROUP (新消息) │ 15.3 │ 0.06ms │ 153,000 │
│ XACK (单条) │ 22.1 | 0.04ms │ 221,000 │
│ XACK Batch×100 │ 28.7 │ 0.003ms │ 2,870,000│
│ XAUTOCLAIM │ 8.2 │ 0.12ms │ 82,000 │
└──────────────────────┴────────────┴────────────┴──────────┘
消息大小: 200-500 bytes (典型业务事件)
八、与其他消息系统的边界与选型
当你应该选择 Redis Streams:
- 事件吞吐 < 100k msgs/s,且主要面向在线服务
- 需要与业务同进程部署(如嵌入式缓存 + 事件总线)
- 事件保留时间短(小时/天级别)
- 已有 Redis 基础设施,不想再引入 Kafka 集群
当你应该选 Kafka / Pulsar:
- 超高吞吐(> 100k msgs/s)
- 长期事件保留(数月/年,用于合规审计)
- 复杂的跨数据中心复制
- 大规模消费者组(数百个)
当你应该选 Postgres 作为事件存储:
- 需要强事务一致性(业务操作 + 事件写入同事务)
- 查询语言灵活(按任意维度搜索事件)
- 数据量中等(TB 级别以下)
- 需要与关系型业务数据 join
九、总结
Redis Streams 自 5.0 到 7.x 稳步迭代,从基础 append-only log 演进为支持 Consumer Groups、XAUTOCLAIM、分层的轻量级事件溯源引擎。结合 Radix Tree 的高效数据结构、健全的确认机制、以及 Redis 生态的成熟度,它在中小规模事件驱动场景中展现出独特的优势。
核心要点回顾:
- PEL 语义是 Redis Streams 实现"至少一次投递"的基石,理解和正确处理 PEL 是生产部署的前提
- XAUTOCLAIM 替代传统的心跳监控 + 手动重平衡,大幅简化死信处理
- MAXLEN + MINID 组合实现事件流的自动过期,防止内存无限膨胀
- 幂等消费 Event Sourcing 场景必须通过 Redis SETNX 实现
- Pipeline 批量操作可提升 3-5 倍吞吐,是 XADD/XACK 的标配优化
对于追求"简单可靠"的事件驱动架构,Redis Streams 是 Kafka 之外一个极具竞争力的选择——选择它的理由不是因为它能做更多,而是因为它让你用更少的组件完成同样的事。

发表评论 取消回复