作者: 技术架构组
关键词: DDD、CQRS、事件驱动、FastAPI、Clean Architecture、Python架构
适用读者: 高级后端工程师、架构师
一、为什么后端架构对高级工程师至关重要?
当我们从"能工作的代码"迈向"能演进的系统",架构就不再是锦上添花,而是生存基石。初级工程师关注功能实现,高级工程师关注系统熵增——每一次业务迭代,要么降低复杂度,要么加速腐烂。
在后端领域,Python 常被误解为"只适合写脚本和原型"。但当你面对日均千万级请求、数百个业务实体、跨团队协作的领域模型时,架构模式决定了下半场的胜负。本文将结合实际项目经验,深入探讨 DDD(Domain-Driven Design)、CQRS 与事件驱动架构在 Python/FastAPI 生态中的落地实践。
二、分层架构与整洁架构
2.1 传统分层架构的困境
经典的三层架构(Controller → Service → Repository)在小项目中高效直接,但在复杂业务中暴露出严重问题:
# ❌ 反模式:Service 层承载了太多职责
class OrderService:
def create_order(self, user_id, items):
# 校验用户
user = self.user_repo.get(user_id)
if not user.is_active:
raise ValueError("用户已禁用")
# 计算价格
total = sum(item.price for item in items)
# 创建订单
order = Order(user_id=user_id, total=total, status="pending")
# 发送邮件
self.email_service.send(order)
# 记录审计
self.audit_logger.log(f"订单创建: {order.id}")
return order
这段代码的问题在于:业务规则散落在各个基础设施调用之间,无法单独测试、无法复用、无法在变更时评估影响面。
2.2 整洁架构(Clean Architecture)的核心理念
Robert C. Martin 提出的整洁架构强调依赖方向由外向内——领域层不依赖任何外部框架或数据库。用 Python 表达:
📁 src/
├── domain/ ← 核心业务逻辑(零外部依赖)
├── application/ ← 用例编排(依赖 domain)
├── infrastructure/ ← 数据库、消息队列等(依赖 domain)
└── interfaces/ ← API 路由、CLI 等(依赖 application)
核心原则:Infrastructure 实现 Domain 定义的接口,Interfaces 调用 Application 的用例。
三、领域驱动设计(DDD)核心概念
3.1 实体(Entity)
实体有唯一标识,其属性可以变化但身份不变。在 Python 中,我们通过 __eq__ 和 __hash__ 确保实体的身份语义:
# domain/entities/order.py
from __future__ import annotations
from uuid import UUID, uuid4
from datetime import datetime
from dataclasses import dataclass, field
@dataclass
class Order:
id: UUID = field(default_factory=uuid4)
customer_id: UUID
status: str = "pending"
total_amount: float = 0.0
created_at: datetime = field(default_factory=datetime.utcnow)
def __eq__(self, other):
if not isinstance(other, Order):
return NotImplemented
return self.id == other.id
def __hash__(self):
return hash(self.id)
# 领域行为:实体携带业务方法
def confirm(self):
if self.status != "pending":
raise DomainError(f"无法确认状态为 {self.status} 的订单")
self.status = "confirmed"
def cancel(self):
if self.status in ("shipped", "delivered"):
raise DomainError("已发货订单无法取消")
self.status = "cancelled"
3.2 值对象(Value Object)
值对象没有身份,由属性值定义相等性。不可变是其核心特征:
# domain/value_objects/money.py
@dataclass(frozen=True)
class Money:
amount: float
currency: str = "CNY"
def __add__(self, other: Money) -> Money:
if self.currency != other.currency:
raise ValueError("货币类型不一致")
return Money(self.amount + other.amount, self.currency)
def __mul__(self, multiplier: float) -> Money:
return Money(self.amount * multiplier, self.currency)
3.3 聚合(Aggregate)与聚合根
聚合是数据一致性的边界。聚合根是外部访问聚合的唯一入口:
# domain/aggregates/cart.py
@dataclass
class Cart(AggregateRoot): # 聚合根
id: UUID = field(default_factory=uuid4)
customer_id: UUID
items: list[CartItem] = field(default_factory=list)
def add_item(self, product_id: UUID, quantity: int, price: Money):
# 业务校验:数量不能为负
if quantity <= 0:
raise DomainError("数量必须大于0")
# 查找是否已存在同商品
existing = next((i for i in self.items if i.product_id == product_id), None)
if existing:
existing.quantity += quantity
else:
self.items.append(CartItem(
product_id=product_id,
quantity=quantity,
unit_price=price,
))
# 注册领域事件——稍后详述
self.register_event(CartItemAdded(
cart_id=self.id,
product_id=product_id,
quantity=quantity,
))
def total(self) -> Money:
return sum(
(item.unit_price * item.quantity for item in self.items),
Money(0),
)
@dataclass
class CartItem(Entity): # 聚合内部实体,不对外暴露
id: UUID = field(default_factory=uuid4)
product_id: UUID
quantity: int
unit_price: Money
3.4 仓库(Repository)
仓库为领域层提供集合风格的聚合存取接口,不暴露数据库细节:
# domain/repositories/order_repository.py
from abc import ABC, abstractmethod
class OrderRepository(ABC):
@abstractmethod
def save(self, order: Order) -> None:
...
@abstractmethod
def get_by_id(self, order_id: UUID) -> Order | None:
...
@abstractmethod
def get_by_customer(self, customer_id: UUID) -> list[Order]:
...
基础设施层实现此接口:
# infrastructure/repositories/sqlalchemy_order_repo.py
class SqlAlchemyOrderRepository(OrderRepository):
def __init__(self, session: AsyncSession):
self._session = session
async def save(self, order: Order) -> None:
orm_model = OrderORM.from_domain(order)
self._session.add(orm_model)
await self._session.flush()
async def get_by_id(self, order_id: UUID) -> Order | None:
orm_model = await self._session.get(OrderORM, order_id)
return orm_model.to_domain() if orm_model else None
3.5 领域服务(Domain Service)
当某个业务逻辑不属于任何实体或值对象时,用领域服务来表达:
# domain/services/pricing_service.py
class PricingService:
"""计算订单价格的领域服务——横跨多个聚合"""
def calculate_discount(self, order: Order, customer: Customer) -> Money:
if customer.is_vip and order.total_amount > Money(100):
return order.total_amount * Decimal("0.1")
if order.total_amount > Money(500):
return order.total_amount * Decimal("0.05")
return Money(0)
四、CQRS:命令查询职责分离
4.1 为什么需要 CQRS?
传统 CRUD 模式下,读取和写入共用同一模型。随着业务复杂化,问题浮现:
- 写模型需要封装业务规则、维护不变条件
- 读模型需要灵活的投影、高效的查询、DTO 定制
CQRS 将模型一分为二:Command 改变状态,Query 读取状态。
4.2 Command 侧实现
# application/commands/create_order.py
from dataclasses import dataclass
from domain.entities.order import Order
from domain.repositories.order_repository import OrderRepository
from domain.repositories.cart_repository import CartRepository
@dataclass
class CreateOrderCommand:
cart_id: UUID
customer_id: UUID
class CreateOrderHandler:
def __init__(self,
order_repo: OrderRepository,
cart_repo: CartRepository,
pricing: PricingService):
self._order_repo = order_repo
self._cart_repo = cart_repo
self._pricing = pricing
async def handle(self, cmd: CreateOrderCommand) -> Order:
cart = await self._cart_repo.get_by_id(cmd.cart_id)
if not cart:
raise DomainError("购物车不存在")
discount = self._pricing.calculate_discount(cart, ...)
order = Order(
customer_id=cmd.customer_id,
total_amount=cart.total() - discount,
)
await self._order_repo.save(order)
return order
4.3 Query 侧实现
查询侧直接面向展现需求,使用轻量级 DTO:
# application/queries/order_queries.py
from dataclasses import dataclass, asdict
@dataclass
class OrderSummaryDTO:
order_id: str
customer_name: str
total: float
status: str
item_count: int
created_at: str
class OrderQueryService:
"""读模型——直接使用优化过的查询,不走领域模型"""
def __init__(self, db_session: AsyncSession):
self._session = db_session
async def get_customer_orders(
self, customer_id: UUID, page: int, size: int
) -> list[OrderSummaryDTO]:
query = text("""
SELECT o.id AS order_id,
c.name AS customer_name,
o.total_amount AS total,
o.status,
COUNT(oi.id) AS item_count,
o.created_at
FROM orders o
JOIN customers c ON c.id = o.customer_id
LEFT JOIN order_items oi ON oi.order_id = o.id
WHERE o.customer_id = :cid
GROUP BY o.id, c.name, o.total_amount, o.status, o.created_at
ORDER BY o.created_at DESC
LIMIT :lim OFFSET :off
""")
result = await self._session.execute(query, {
"cid": customer_id, "lim": size, "off": (page - 1) * size,
})
return [OrderSummaryDTO(**row) for row in result.mappings()]
4.4 CQRS 在 FastAPI 中的集成
# interfaces/api/orders.py
from fastapi import APIRouter, Depends
from application.commands.create_order import CreateOrderCommand, CreateOrderHandler
from application.queries.order_queries import OrderQueryService
router = APIRouter(prefix="/orders", tags=["orders"])
@router.post("", status_code=201)
async def create_order(
cmd: CreateOrderCommand,
handler: CreateOrderHandler = Depends(get_create_order_handler),
):
order = await handler.handle(cmd)
return {"order_id": str(order.id), "status": order.status}
@router.get("")
async def list_orders(
customer_id: UUID,
page: int = 1,
size: int = 20,
query_service: OrderQueryService = Depends(get_order_query_service),
):
items = await query_service.get_customer_orders(customer_id, page, size)
return {"items": items, "page": page, "size": size}
五、事件驱动架构
5.1 领域事件(Domain Event)
领域事件是"已发生事实"的不可变记录,用过去时命名:
# domain/events/order_events.py
from dataclasses import dataclass, field
from datetime import datetime
from uuid import UUID, uuid4
@dataclass(frozen=True)
class OrderConfirmed:
event_id: UUID = field(default_factory=uuid4)
order_id: UUID
confirmed_at: datetime = field(default_factory=datetime.utcnow)
@dataclass(frozen=True)
class OrderShipped:
event_id: UUID = field(default_factory=uuid4)
order_id: UUID
tracking_number: str
shipped_at: datetime = field(default_factory=datetime.utcnow)
5.2 事件总线(Event Bus)
事件总线负责将领域事件发布给订阅者,解耦事件产生者和处理者:
# infrastructure/event_bus.py
from collections import defaultdict
from typing import Callable, Type
Handler = Callable[[object], None]
class InMemoryEventBus:
def __init__(self):
self._handlers: dict[type, list[Handler]] = defaultdict(list)
def register(self, event_type: Type, handler: Handler):
self._handlers[event_type].append(handler)
async def publish(self, event: object):
for handler in self._handlers.get(type(event), []):
await handler(event)
# 实际生产可用 Redis Streams / RabbitMQ / Kafka 实现
class RedisStreamEventBus:
def __init__(self, redis_client: Redis):
self._redis = redis_client
self._handlers: dict[str, list[Handler]] = defaultdict(list)
async def publish(self, stream: str, event: object):
await self._redis.xadd(stream, {
"type": type(event).__name__,
"payload": json.dumps(asdict(event), default=str),
})
async def consume(self, stream: str, group: str, consumer: str):
messages = await self._redis.xreadgroup(group, consumer, {stream: ">"})
for msg_id, data in messages[stream].items():
yield data
5.3 聚合根中的事件注册
聚合根基类管理领域事件的注册与清空:
# domain/aggregates/base.py
from typing import Protocol
class AggregateRoot(Protocol):
"""聚合根协议——所有聚合根必须实现"""
_events: list[object]
def register_event(self, event: object):
self._events.append(event)
def pull_events(self) -> list[object]:
events = self._events[:]
self._events.clear()
return events
# 在 Repository 中自动发布事件
class EventSourcingOrderRepository(OrderRepository):
def __init__(self, session: AsyncSession, event_bus: EventBus):
self._session = session
self._event_bus = event_bus
async def save(self, order: Order) -> None:
orm_model = OrderORM.from_domain(order)
self._session.add(orm_model)
# 发布聚合根累积的领域事件
for event in order.pull_events():
await self._event_bus.publish(event)
5.4 事件订阅者处理跨领域关注点
# application/subscribers/order_subscribers.py
class OrderEventSubscribers:
"""事件订阅者——处理跨边界的副作用"""
def __init__(self, bus: EventBus, email_svc, audit_svc, inventory_svc):
bus.register(OrderConfirmed, self.on_order_confirmed)
bus.register(OrderShipped, self.on_order_shipped)
async def on_order_confirmed(self, event: OrderConfirmed):
# 发送确认邮件(非核心业务流程)
await self.email_svc.send_confirmation(event.order_id)
# 记录审计日志
self.audit_svc.log(f"订单确认: {event.order_id}", event)
async def on_order_shipped(self, event: OrderShipped):
# 扣减库存
await self.inventory_svc.release_reserved(event.order_id)
# 通知物流
await self.email_svc.send_shipping_notification(
event.order_id, event.tracking_number
)
六、FastAPI 完整项目结构
以下是一个生产级 DDD 项目目录结构,融合了分层架构与模块化设计:
📁 src/
├── __init__.py
│
├── domain/ # 🔵 领域层(零外部依赖)
│ ├── __init__.py
│ ├── entities/
│ │ ├── __init__.py
│ │ ├── order.py
│ │ └── customer.py
│ ├── value_objects/
│ │ ├── __init__.py
│ │ ├── money.py
│ │ └── address.py
│ ├── aggregates/
│ │ ├── __init__.py
│ │ ├── base.py # AggregateRoot 基类
│ │ └── cart.py
│ ├── events/
│ │ ├── __init__.py
│ │ └── order_events.py
│ ├── services/
│ │ ├── __init__.py
│ │ └── pricing_service.py
│ ├── repositories/
│ │ ├── __init__.py
│ │ ├── order_repository.py # 抽象接口
│ │ └── cart_repository.py
│ └── exceptions.py # DomainError 等
│
├── application/ # 🟢 应用层(用例编排)
│ ├── __init__.py
│ ├── commands/
│ │ ├── __init__.py
│ │ ├── create_order.py
│ │ └── cancel_order.py
│ ├── queries/
│ │ ├── __init__.py
│ │ └── order_queries.py
│ └── subscribers/
│ ├── __init__.py
│ └── order_subscribers.py
│
├── infrastructure/ # 🟡 基础设施层
│ ├── __init__.py
│ ├── database/
│ │ ├── __init__.py
│ │ ├── models/ # SQLAlchemy ORM 模型
│ │ │ ├── __init__.py
│ │ │ └── order_orm.py
│ │ └── migrations/ # Alembic 迁移
│ ├── repositories/
│ │ ├── __init__.py
│ │ └── sqlalchemy_order_repo.py
│ ├── message_queue/
│ │ ├── __init__.py
│ │ └── redis_event_bus.py
│ └── external/
│ ├── __init__.py
│ └── payment_gateway.py
│
├── interfaces/ # 🟠 接口层
│ ├── __init__.py
│ ├── api/
│ │ ├── __init__.py
│ │ ├── app.py # FastAPI 应用
│ │ ├── orders.py
│ │ └── customers.py
│ ├── di/ # 依赖注入配置
│ │ ├── __init__.py
│ │ └── container.py
│ └── cli/
│ ├── __init__.py
│ └── seed_data.py
│
└── main.py # 启动入口
七、依赖注入与 IoC 容器
7.1 为什么需要 DI?
DDD 的核心依赖规则——领域层定义接口,基础设施层实现——如果没有 DI 容器,你将不得不在各处手动 new 实现,导致层间耦合。
7.2 基于 FastAPI Depends 的 DI 方案
# interfaces/di/container.py
from functools import lru_cache
from sqlalchemy.ext.asyncio import AsyncSession, create_async_engine
from domain.repositories.order_repository import OrderRepository
from infrastructure.repositories.sqlalchemy_order_repo import SqlAlchemyOrderRepository
from infrastructure.event_bus import InMemoryEventBus
from application.commands.create_order import CreateOrderHandler
class Container:
def __init__(self):
self._engine = create_async_engine(
"postgresql+asyncpg://user:pass@localhost/db",
)
self._event_bus = InMemoryEventBus()
self._init_subscribers()
def _init_subscribers(self):
from application.subscribers.order_subscribers import OrderEventSubscribers
OrderEventSubscribers(self._event_bus, ..., ..., ...)
async def get_session(self) -> AsyncSession:
async with AsyncSession(self._engine) as session:
yield session
def get_order_repo(self, session: AsyncSession = Depends(get_session)):
return SqlAlchemyOrderRepository(session)
def get_create_order_handler(
self,
order_repo: OrderRepository = Depends(get_order_repo),
cart_repo: CartRepository = Depends(get_cart_repo),
pricing: PricingService = Depends(...),
):
return CreateOrderHandler(order_repo, cart_repo, pricing)
@lru_cache()
def get_container() -> Container:
return Container()
在 FastAPI 路由中:
# interfaces/api/orders.py
from interfaces.di.container import get_container
@router.post("")
async def create_order(
cmd: CreateOrderCommand,
handler = Depends(lambda: get_container().get_create_order_handler()),
):
return await handler.handle(cmd)
这种方案的优点:零框架侵入、完全可测试(替换 Container 中的实现即可 Mock)。
八、从单体到微服务的架构演进策略
架构演进不是一蹴而就的重写,而是一个持续的、有节奏的拆解过程。
8.1 演进路线图
阶段一 阶段二 阶段三
┌──────────┐ ┌──────────────┐ ┌────────────────┐
│ │ │ │ │ │
│ 单体 │ ──▶ │ 模块化单体 │ ──▶ │ 按限界上下文 │
│ DDD │ │ (模块间 │ │ 拆分微服务 │
│ 项目 │ │ 通过事件 │ │ │
│ 结构 │ │ 通信) │ │ Order Service │
│ │ │ │ │ Payment Svc │
│ │ │ │ │ Notification │
└──────────┘ └──────────────┘ └────────────────┘
8.2 事件驱动在演进中的关键作用
- 阶段一:事件总线在进程内传递事件(InMemoryEventBus)
- 阶段二:切换到 Redis Streams 或 RabbitMQ,模块间通过消息解耦
- 阶段三:将单个模块独立为服务,事件总线变为跨服务消息中间件
关键洞察:由于你在阶段一就遵循了 DDD 的限界上下文划分和事件驱动设计,每个模块的边界天然清晰——拆分的成本远低于从零重构。
九、架构决策记录(ADR)最佳实践
ADR(Architecture Decision Record)是团队架构共识的载体。每个重要决策都应记录。
ADR 模板示例
# ADR-001:选择 CQRS + 事件溯源作为订单核心架构
## 状态
✅ 已采纳 · 2026-03-15
## 背景
订单模块需要完整的操作审计和状态回溯能力,
同时团队需要支持未来拆分为独立微服务。
## 决策
在订单限界上下文内采用 CQRS + 事件溯源模式:
- Command 侧:使用 DDD 聚合 + SQLAlchemy
- Query 侧:使用专用读模型(物化视图 + Redis Cache)
- 事件存储:PostgreSQL + Debezium → Kafka
## 权衡
- ✅ 可审计:所有状态变更可溯源
- ✅ 可演进:事件是自然的微服务拆分边界
- ❌ 复杂度增加:最终一致性引入补偿事务需求
- ❌ 学习成本:团队需要理解事件溯源模式
## 替代方案
1. **传统 CRUD**:简单但无审计、难拆分 → 被否决
2. **仅 CQRS**:满足查询优化,但无法回溯状态 → 不满足审计需求
## 影响
- 订单写入吞吐量下降约 15%(事件序列化开销)
- 读取延迟降低约 40%(专用读模型优化)
- 需要引入 Saga 模式处理跨服务事务
ADR 管理要点
- 存储位置:在项目根目录建立
docs/adr/,与代码同仓库 - 更新机制:当决策被推翻时(ADR 状态设为"已废弃"),保留原始记录
- 评审节点:所有涉及模块边界、数据流、技术选型的变更,必须附带 ADR
- 数量控制:不记录 trivial 决策(如"用 f-string 还是 format"),聚焦架构级决策
十、总结
| 模式 | 解决的问题 | Python 落地要点 |
|---|---|---|
| DDD | 业务复杂度膨胀、领域逻辑散落 | 边界在于 domain/ 目录零外部依赖 |
| CQRS | 读模型与写模型冲突 | Command 用聚合,Query 用 DTO+SQL |
| 事件驱动 | 跨边界耦合、异步协作 | 聚合根注册事件,Repository 自动发布 |
| DI 容器 | 层间依赖倒置 | FastAPI Depends + 自定义 Container 函数 |
| ADR | 架构决策失忆、团队认知断层 | 代码仓库内/docs/adr/ 轻量模板 |
后端架构没有银弹,但 DDD + CQRS + 事件驱动的组合,为 Python 后端团队提供了一条可验证、可演进、可协作的架构路径。当你下一次面对复杂业务时,不妨从限界上下文开始划分,让领域模型驱动你的代码结构——而不是让框架或数据库替你决定。
延伸阅读
- 《领域驱动设计:软件核心复杂性应对之道》—— Eric Evans
- 《实现领域驱动设计》—— Vaughn Vernon
- Clean Architecture: A Craftsman's Guide —— Robert C. Martin
- FastAPI 官方文档:https://fastapi.tiangolo.com/
评论