架构模式——什么时候该用 CQRS / Architecture Patterns and When to Use CQRS
📅 创建时间:2026-05-08 🏷️ 标签:#CQRS #ES #Saga #熔断降级 #舱壁模式 #微服务拆分 📚 前置知识:[[00-backend-overview]] [[04-distributed-system]](分布式事务) 📚 相关知识:[[09-mysql-optimization]](读写分离) [[05-middleware]](熔断降级)
场景:你的查询把事务系统拖垮了
┌─────────────────────────────────────────────────────────────┐
│ │
│ 产品经理:新需求——在订单列表页展示: │
│ 用户昵称、商品名称、商品图片、优惠信息、 │
│ 物流状态、仓储位置、推荐商品... │
│ │
│ 现有表:orders, users, products, logistics, warehouses │
│ │
│ SQL: │
│ SELECT o.*, u.name, p.name, p.image, l.status, │
│ w.location, r.product_id │
│ FROM orders o │
│ LEFT JOIN users u ON o.user_id = u.id │
│ LEFT JOIN order_items oi ON o.id = oi.order_id │
│ LEFT JOIN products p ON oi.product_id = p.id │
│ LEFT JOIN logistics l ON o.id = l.order_id │
│ LEFT JOIN warehouses w ON p.warehouse_id = w.id │
│ LEFT JOIN recommendations r ON p.id = r.product_id │
│ │
│ 优化后耗时:3.5 秒。 │
│ │
│ 但这只是单个订单。如果是订单列表(20 条)呢? │
│ │
└─────────────────────────────────────────────────────────────┘这一章,我们理解高级架构模式:CQRS、事件溯源、Saga、熔断降级。
第1节:CQRS——读写分离到极致
核心思想
┌─────────────────────────────────────────────────────────────┐
│ CQRS 核心思想 │
├─────────────────────────────────────────────────────────────┤
│ │
│ Command Query Responsibility Segregation │
│ 命令查询职责分离 │
│ │
│ 核心:一个方法不应该既是 Command(写)又是 Query(读) │
│ │
│ 传统架构: │
│ OrderService: │
│ createOrder() → 写 │
│ getOrder() → 读 │
│ listOrders() → 读 │
│ → 同一个数据模型同时服务读写 → 互相影响 │
│ │
│ CQRS 架构: │
│ Command Side: 写入 │
│ OrderCommandService: │
│ createOrder() → 写入 MySQL │
│ cancelOrder() → 写入 MySQL │
│ │
│ Query Side:读取 │
│ OrderQueryService: │
│ getOrder() → 读取预聚合好的 Read Model │
│ listOrders() → 读取 Read Model(JOIN 结果) │
│ │
│ 关键:写入和读取使用不同的数据模型 │
│ │
└─────────────────────────────────────────────────────────────┘CQRS 实现
┌─────────────────────────────────────────────────────────────┐
│ CQRS 数据流 │
├─────────────────────────────────────────────────────────────┤
│ │
│ Write Path(命令路径): │
│ 用户 ──▶ API ──▶ Command Handler ──▶ MySQL (Write Model) │
│ │ │
│ ▼ │
│ Event Bus(事件总线) │
│ │ │
│ ▼ │
│ Read Path(查询路径): │
│ 用户 ◀── API ◀── Query Handler ◀── Read Model (预聚合表) │
│ │
│ 同步/异步: │
│ • 同步:Event Bus = 同一数据库事务 │
│ • 异步:Event Bus = Kafka / RabbitMQ │
│ → 立即返回(更快),但有短暂不一致 │
│ │
└─────────────────────────────────────────────────────────────┘CQRS 实战:订单查询优化
python
# Write Model(命令端):订单表
class OrderRepository:
def create_order(self, user_id, items):
order = Order(
id=uuid(),
user_id=user_id,
status="pending",
created_at=now()
)
for item in items:
order.add_item(OrderItem(
product_id=item["product_id"],
quantity=item["quantity"]
))
self.db.add(order)
self.db.commit()
# 发布事件
self.event_bus.publish(OrderCreatedEvent(order))
# Read Model(查询端):预聚合表
class OrderQueryService:
def get_order_detail(self, order_id):
# 直接查预聚合表,不用 JOIN
return self.db.query("""
SELECT * FROM v_order_detail WHERE order_id = %s
""", order_id)
def list_user_orders(self, user_id, page):
return self.db.query("""
SELECT order_id, status, total, created_at
FROM orders
WHERE user_id = %s
ORDER BY created_at DESC
LIMIT 20 OFFSET %s
""", user_id, (page - 1) * 20)
# View(预聚合视图),由事件同步更新
class OrderViewSync:
def on_order_created(self, event):
# 更新订单详情视图(join 结果)
self.db.execute("""
INSERT INTO v_order_detail (order_id, user_name, product_names, ...)
VALUES (%s, %s, %s, ...)
""", ...)
# 更新用户订单列表
self.db.execute("""
INSERT INTO orders_summary (...)
VALUES (...)
""", ...)第2节:Saga 模式——长流程的编排
场景:订单创建涉及多个服务
┌─────────────────────────────────────────────────────────────┐
│ │
│ 订单创建流程: │
│ │
│ 1. 扣减库存(Inventory Service) │
│ 2. 创建订单(Order Service) │
│ 3. 扣减余额(Account Service) │
│ 4. 发送通知(Notification Service) │
│ │
│ 如果步骤 3 失败了? │
│ → 库存已经扣了,余额没扣 → 数据不一致 │
│ → 需要回滚:库存加回去 │
│ │
└─────────────────────────────────────────────────────────────┘Saga 编排 vs Saga 编舞
┌─────────────────────────────────────────────────────────────┐
│ Saga 编排 vs 编舞 │
├─────────────────────────────────────────────────────────────┤
│ │
│ Saga 编排器(Choreography): │
│ │
│ OrderService ──▶ InventoryService ──▶ AccountService ──▶ Notification │
│ ◀─────────────────────────────── ◀────────────── │
│ 补偿事件 │
│ │
│ 服务 A 执行 → 发布成功事件 → 服务 B 执行 → 发布成功事件 │
│ 服务 B 失败 → 发布失败事件 → 服务 A 补偿(回滚) │
│ │
│ 优点:服务间松耦合 │
│ 缺点:流程不清晰,调试困难 │
│ │
│ Saga 编排器(Orchestration): │
│ │
│ ┌──────────────────────────┐ │
│ │ Saga Orchestrator │ │
│ │ (订单编排器) │ │
│ └───────────┬──────────────┘ │
│ │ │
│ ┌──────────┼──────────┐ │
│ ▼ ▼ ▼ │
│ Inventory Account Notification │
│ Service Service Service │
│ │
│ 编排器统一管理流程: │
│ 编排器调用库存服务 → 成功 → 调用账户服务 → 失败 → │
│ 调用库存服务补偿(加回库存) │
│ │
│ 优点:流程清晰,便于调试 │
│ 缺点:编排器是中心,存在单点问题 │
│ │
└─────────────────────────────────────────────────────────────┘Saga 实现
python
class CreateOrderSaga:
def execute(self, user_id, items):
steps = [
Step("reserve_inventory", # 步骤名
lambda ctx: self.inventory_service.reserve(ctx), # 正向操作
lambda ctx: self.inventory_service.release(ctx)), # 补偿操作
Step("create_order",
lambda ctx: self.order_service.create(ctx),
lambda ctx: self.order_service.cancel(ctx)),
Step("charge_account",
lambda ctx: self.account_service.charge(ctx),
lambda ctx: self.account_service.refund(ctx)),
]
context = {"user_id": user_id, "items": items, "order_id": None}
completed = []
for step in steps:
try:
result = step.forward(context)
context.update(result)
completed.append(step.name)
except Exception as e:
# 补偿已完成的步骤
for name in reversed(completed):
step_to_rollback = next(s for s in steps if s.name == name)
step_to_rollback.compensate(context)
raise OrderCreationFailed(f"Step {step.name} failed: {e}")
return context["order_id"]第3节:熔断降级——当服务不可用时
场景:库存服务挂了,但订单还是要能下
┌─────────────────────────────────────────────────────────────┐
│ │
│ 库存服务不可用。 │
│ │
│ 方案 A(无降级): │
│ → 订单服务调用库存服务 → 超时 → 订单创建失败 │
│ → 用户买不到东西 │
│ │
│ 方案 B(有降级): │
│ → 订单服务调用库存服务 → 超时 → 返回"库存待确认" │
│ → 订单创建成功,但标记为"库存待确认" │
│ → 后台人工/自动补偿 │
│ │
└─────────────────────────────────────────────────────────────┘三种降级策略
python
# 方案 1:返回默认值(最常见)
def get_stock(product_id):
try:
return stock_service.get_stock(product_id)
except ServiceUnavailable:
# 降级:返回默认值
return 0 # 或者返回"库存充足"的乐观值
# 方案 2:返回缓存数据
def get_stock(product_id):
try:
stock = stock_service.get_stock(product_id)
cache.set(f"stock:{product_id}", stock, ttl=3600) # 同时缓存
return stock
except ServiceUnavailable:
# 降级:返回缓存数据
cached = cache.get(f"stock:{product_id}")
if cached:
return cached
return -1 # 标记为异常
# 方案 3:返回兜底数据(静态/历史数据)
def get_recommendations(user_id):
try:
return recommendation_service.get(user_id)
except ServiceUnavailable:
# 降级:返回热门商品
return hot_products.get_top(10)第4节:舱壁模式——隔离故障
问题:一个服务拖垮了整个系统
┌─────────────────────────────────────────────────────────────┐
│ │
│ 线程池 A(100 线程):处理所有请求 │
│ │
│ 请求类型: │
│ • 订单查询(需要调用库存服务) │
│ • 用户查询(MySQL 直接查) │
│ • 商品查询(Elasticsearch) │
│ │
│ 库存服务挂了: │
│ → 100 个线程全部等待库存服务 │
│ → 线程池耗尽 │
│ → 用户查询、商品查询也全部超时 │
│ → 整个系统崩溃 │
│ │
└─────────────────────────────────────────────────────────────┘舱壁模式:用不同的线程池处理不同请求
┌─────────────────────────────────────────────────────────────┐
│ 舱壁模式 │
├─────────────────────────────────────────────────────────────┤
│ │
│ 线程池 A(20 线程):处理订单查询 │
│ 线程池 B(10 线程):处理用户查询 │
│ 线程池 C(30 线程):处理商品查询 │
│ │
│ 库存服务挂了: │
│ → 线程池 A 的 20 个线程等待库存服务 │
│ → 线程池 B、线程池 C 完全不受影响 │
│ → 用户查询、商品查询正常响应 │
│ │
│ 舱壁的作用:故障被隔离在一个隔舱内,不蔓延 │
│ │
└─────────────────────────────────────────────────────────────┘python
from concurrent.futures import ThreadPoolExecutor
# 舱壁模式:不同的服务用不同的线程池
order_pool = ThreadPoolExecutor(max_workers=20, thread_name_prefix="order-")
user_pool = ThreadPoolExecutor(max_workers=10, thread_name_prefix="user-")
goods_pool = ThreadPoolExecutor(max_workers=30, thread_name_prefix="goods-")
class OrderService:
def __init__(self, inventory_client, user_client):
self.inventory = inventory_client
self.user = user_client
def get_order_detail(self, order_id):
# 订单查询用订单线程池
return order_pool.submit(self._fetch_order, order_id).result()
def _fetch_order(self, order_id):
# 调用库存服务和用户服务
order = self.db.query("SELECT * FROM orders WHERE id = %s", order_id)
inventory = order_pool.submit(self.inventory.get_stock, order["product_id"]).result()
user = user_pool.submit(self.user.get, order["user_id"]).result()
return {**order, "inventory": inventory, "user": user}第5节:微服务拆分——什么时候该拆
问题:什么时候该拆微服务
┌─────────────────────────────────────────────────────────────┐
│ │
│ 很多公司的问题不是"微服务太多了" │
│ 而是"还没到那个规模就开始拆了" │
│ │
│ 单体 → 微服务 容易。 │
│ 微服务 → 单体 很难(拆分容易,合并难)。 │
│ │
└─────────────────────────────────────────────────────────────┘拆分的信号
┌─────────────────────────────────────────────────────────────┐
│ 应该拆微服务的信号 │
├─────────────────────────────────────────────────────────────┤
│ │
│ ✅ 技术栈差异大 │
│ → 用户服务用 Node.js,订单服务用 Go │
│ → 不同团队,不同部署周期,不同性能需求 │
│ │
│ ✅ 独立扩展需求 │
│ → 商品服务是 CPU 密集(需要多 CPU) │
│ → 用户服务是 IO 密集(需要高并发) │
│ │
│ ✅ 发布频率差异大 │
│ → 用户服务每天发布 10 次 │
│ → 订单服务每周发布 1 次(牵一发而动全身) │
│ │
│ ✅ 故障隔离 │
│ → 订单服务挂了 → 用户登录正常 │
│ │
│ ❌ 不要拆的信号: │
│ → 团队规模 < 10 人(拆了管不过来) │
│ → 数据强关联(JOIN 频繁) │
│ → 还没遇到单体性能瓶颈 │
│ │
└─────────────────────────────────────────────────────────────┘升华:架构模式的选择
┌─────────────────────────────────────────────────────────────┐
│ 架构模式选择指南 │
├─────────────────────────────────────────────────────────────┤
│ │
│ CQRS: │
│ → 读压力远大于写压力(90% vs 10%) │
│ → 查询逻辑复杂,JOIN 多表 │
│ → 报表和事务需要分离 │
│ │
│ Saga: │
│ → 长流程(5+ 个步骤) │
│ → 每个步骤可能失败,需要补偿 │
│ → 不适合需要强一致性的场景(金融转账) │
│ │
│ 熔断降级: │
│ → 所有服务都应该做(关键依赖) │
│ → 核心接口 → 降级返回默认值 │
│ → 非核心接口 → 直接返回错误 │
│ │
│ 舱壁模式: │
│ → 调用多个外部服务的接口 │
│ → 不同服务的超时时间差异大 │
│ │
│ 一句话总结: │
│ 架构模式是工具,不是目的。 │
│ 先让系统跑起来,遇到问题了再引入对应的模式。 │
│ │
└─────────────────────────────────────────────────────────────┘"AI 可查 vs 必须理解"清单
AI 可查:
✅ CQRS 的具体框架实现(Axon / Eventuate / NestJS CQRS)
✅ Saga 编排器的具体实现代码
✅ 各种降级策略的具体返回值设计
必须理解:
🔴 CQRS 和传统读写分离的本质区别
🔴 Saga 编排 vs 编舞的 trade-off
🔴 为什么 Saga 不适合强一致性场景
🔴 舱壁模式的隔离原理
🔴 什么时候该拆微服务,什么时候不该拆学习状态:🟡 开始学习