技术栈示例:FastAPI + SQLAlchemy(async) + Alembic + Pydantic + Taskiq + OpenTelemetry。换语言(Java/Spring、Go、.NET)时结构和规则不变,只换实现。配套 README / 前端骨架。
一、完整目录树
backend/
├── pyproject.toml # 依赖 + [tool.importlinter] 边界合同
├── alembic/ # 数据库迁移
│ └── versions/
└── app/
├── main.py # FastAPI app 工厂;挂载各模块 router
├── core/ # 技术内核:横切、不含业务
│ ├── config.py # 配置(pydantic-settings,读环境变量)
│ ├── database.py # async engine + session 工厂
│ ├── base_repository.py # 泛型仓储基类(CRUD 模块复用)
│ ├── unit_of_work.py # 工作单元:事务边界
│ ├── responses.py # 统一响应信封 ApiResponse[T]
│ ├── exceptions.py # 领域/应用异常
│ ├── middleware.py # 异常→HTTP 归一、请求日志
│ ├── security.py # 鉴权(JWT 等)
│ ├── dependencies.py # FastAPI 依赖注入(get_db_session / get_current_user)
│ ├── redis.py # 缓存接入
│ └── storage.py # 对象存储接入(S3/MinIO)
├── shared/ # 横切基础设施:可被所有模块用,但不依赖业务
│ ├── outbox/ # 领域事件出箱(可靠投递)
│ │ ├── orm.py # domain_events 表
│ │ ├── repository.py # emit():业务事务里写事件
│ │ └── relay.py # 后台 relay:捞未发事件→投递→标记
│ ├── observability/ # OTel + 结构化日志
│ │ ├── logging.py
│ │ └── otel.py
│ └── jobs/ # 异步任务运行时(Taskiq,可选)
└── modules/ # ★业务★ 每个限界上下文一个文件夹
└── order/ # 示例上下文
├── api/v1/
│ ├── endpoints.py # HTTP 路由(接口层)
│ └── schemas.py # 对外 DTO(Pydantic)
├── application/
│ ├── order_service.py# 写侧 Command 服务(用例编排)
│ └── order_query.py # 读侧 Query 服务(CQRS)
├── domain/
│ ├── order.py # 聚合根/实体(充血)+ 值对象
│ ├── events.py # 领域事件
│ └── repository.py # 仓储接口(Protocol)
└── infrastructure/
├── orm.py # 持久化对象 PO(SQLAlchemy 表)
└── repository_impl.py # 仓储实现(实现 domain 的接口)
二、core/ 横切内核(每个零件的最小实现)
responses.py — 统一响应信封(前端永远拿到同一形状)
from typing import Generic, TypeVar
from pydantic import BaseModel
T = TypeVar("T")
class ApiResponse(BaseModel, Generic[T]):
success: bool = True
data: T | None = None
error: str | None = None
@classmethod
def ok(cls, data: T) -> "ApiResponse[T]": return cls(success=True, data=data)
@classmethod
def fail(cls, error: str) -> "ApiResponse[None]": return cls(success=False, error=error)
base_repository.py — 泛型仓储基类(CRUD 模块直接继承,省样板)
from typing import Generic, TypeVar
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy import select
PO = TypeVar("PO")
class BaseRepository(Generic[PO]):
def __init__(self, model: type[PO], session: AsyncSession):
self.model, self.session = model, session
async def get(self, id_): return await self.session.get(self.model, id_)
async def add(self, po: PO) -> PO:
self.session.add(po); await self.session.flush(); return po
async def list(self, *, page: int, page_size: int, extra=None):
stmt = select(self.model)
for c in (extra or []): stmt = stmt.where(c)
stmt = stmt.limit(page_size).offset((page - 1) * page_size)
return list((await self.session.scalars(stmt)).all())
unit_of_work.py — 工作单元(一个用例的写操作要么全成要么全滚)
from contextlib import asynccontextmanager
from app.core.database import async_session_factory
@asynccontextmanager
async def unit_of_work():
async with async_session_factory() as session:
try:
yield session
await session.commit() # 用例成功 → 一次性提交
except Exception:
await session.rollback() # 任意一步抛错 → 整体回滚
raise
dependencies.py — 依赖注入(接口层用 Depends 拿 session / 当前用户)
from fastapi import Depends
from app.core.security import decode_token
async def get_db_session(): ... # yield 一个 request 级 session
async def get_current_user_id(token=Depends(...)) -> str: return decode_token(token).sub
其余 config.py(pydantic-settings 读环境变量)、exceptions.py + middleware.py(把领域异常统一翻成 HTTP 错误)、security.py、redis.py、storage.py 都是"一处定义、全局复用"的横切件,按需填。
三、shared/outbox/ 领域事件出箱(最该抄的可靠性件)
这是 06 讲"先存事件表再发"的生产实现——保证"业务改动"和"事件"要么一起成功要么一起失败。
orm.py — 事件表
class OutboxEventORM(Base):
__tablename__ = "domain_events"
event_id: Mapped[UUID] = mapped_column(primary_key=True) # UUIDv7
type: Mapped[str] # 事件类型,如 "order.created"(CloudEvents type)
source: Mapped[str] # 来源上下文,如 "modules.order"
subject: Mapped[str] # 聚合 id
data: Mapped[dict] = mapped_column(JSON) # 业务快照
traceparent: Mapped[str | None] # W3C TraceContext,跨事件传链路
created_at: Mapped[datetime]
published_at: Mapped[datetime | None] # NULL=未发,relay 发完盖戳
repository.py — 业务事务里 emit(关键:和业务写同一个 session/事务)
class OutboxRepository:
def __init__(self, session): self.session = session
async def emit(self, *, type_: str, source: str, subject: str, data: dict):
self.session.add(OutboxEventORM(
event_id=uuid7(), type=type_, source=source, subject=subject,
data=data, traceparent=current_traceparent(), created_at=utcnow(),
)) # 不 commit!跟着业务用例的 UoW 一起提交
relay.py — 后台 relay(独立循环:捞未发→投递→盖戳)
async def relay_once(session, publisher):
rows = await fetch_unpublished(session, limit=100) # published_at IS NULL
for ev in rows:
await publisher.publish(envelope := to_cloudevent(ev)) # 投递目标:MQ/WebSocket/Webhook
await mark_published(session, ev.event_id) # 盖戳
💡 投递目标是可换的:要实时就推 WebSocket(如 Centrifugo),要解耦就推 Kafka/RabbitMQ,最简单甚至先打日志。出箱的价值不在"推给谁",而在"业务和事件原子落库、绝不丢"。
shared/observability/ 装 OTel 初始化 + 结构化日志;shared/jobs/ 是 Taskiq 任务运行时(小项目可先不建,需要后台任务再加)。
四、示例模块 modules/order/:四层的最小代码
domain/ —— 业务规则在这(充血)
order.py — 聚合根 + 值对象
from dataclasses import dataclass, field
from enum import StrEnum
class OrderStatus(StrEnum):
PENDING = "pending"; PAID = "paid"; CANCELLED = "cancelled"
@dataclass(frozen=True)
class Money: # 值对象:不可变
amount: int; currency: str = "CNY"
@dataclass
class OrderItem: # 实体(聚合内)
sku: str; qty: int; unit_price: Money
class Order: # 聚合根(充血:带行为,不是数据袋)
def __init__(self, order_id, buyer_id):
self.id = order_id; self.buyer_id = buyer_id
self.status = OrderStatus.PENDING
self.items: list[OrderItem] = []
def add_item(self, item: OrderItem) -> None:
if self.status is not OrderStatus.PENDING:
raise ValueError("只有待支付订单能加明细") # ← 业务规则在领域层
self.items.append(item)
@property
def total(self) -> Money: # 不变量:总额=明细之和
return Money(sum(i.unit_price.amount * i.qty for i in self.items))
def pay(self) -> None:
if not self.items: raise ValueError("空订单不能支付")
self.status = OrderStatus.PAID # 状态机也在领域层
repository.py — 仓储接口(Protocol,依赖倒置的核心)
from typing import Protocol
class OrderRepository(Protocol): # domain 只认这个抽象,不认数据库
async def get(self, order_id) -> "Order | None": ...
async def save(self, order: "Order") -> None: ...
events.py
@dataclass(frozen=True)
class OrderCreated:
order_id: str; buyer_id: str; total: int
infrastructure/ —— 仓储实现 + PO
orm.py(PO) 定义 OrderORM / OrderItemORM SQLAlchemy 表。
repository_impl.py 实现 domain 的 OrderRepository 接口,负责 Order(领域对象) ↔ OrderORM(PO) 互转:
class OrderRepositoryImpl: # 实现 domain.OrderRepository
def __init__(self, session): self.session = session
async def get(self, order_id) -> Order | None:
po = await self.session.get(OrderORM, order_id)
return to_domain(po) if po else None # PO → 领域对象
async def save(self, order: Order) -> None:
self.session.add(to_po(order)) # 领域对象 → PO
application/ —— 用例编排(CQRS 写读分家)
order_service.py(写侧 Command)
class OrderService:
def __init__(self, session, outbox: OutboxRepository):
self.repo = OrderRepositoryImpl(session); self.outbox = outbox
async def create_order(self, *, buyer_id, items) -> str:
order = Order(order_id=uuid7(), buyer_id=buyer_id)
for it in items: order.add_item(it) # ← 调领域规则
await self.repo.save(order)
await self.outbox.emit(type_="order.created", source="modules.order",
subject=str(order.id),
data={"buyer": buyer_id, "total": order.total.amount})
return order.id # 事件和订单同一事务(UoW)落库
order_query.py(读侧 Query) —— 绕开聚合,直接拼只读视图
class OrderQueryService:
def __init__(self, session): self.session = session
async def get_detail_view(self, order_id) -> OrderDetailDTO | None:
# 直接 select/join 拼前端要的形状,不重建 Order 聚合(避免 N+1)
...
api/v1/ —— 接口层(endpoint + DTO)
schemas.py 定义 CreateOrderRequest / OrderDetailDTO(对外 DTO,不复用内部模型)。
endpoints.py
router = APIRouter(prefix="/orders", tags=["orders"])
@router.post("", response_model=ApiResponse[OrderDetailDTO], summary="创建订单")
async def create_order(
body: CreateOrderRequest,
user_id: str = Depends(get_current_user_id), # 鉴权(横切,收敛在接口层)
session = Depends(get_db_session),
):
async with unit_of_work() as uow: # 事务边界
svc = OrderService(uow, OutboxRepository(uow))
order_id = await svc.create_order(buyer_id=user_id, items=body.to_items())
view = await OrderQueryService(uow).get_detail_view(order_id) # 写完用读侧回显
return ApiResponse.ok(view)
这一个端点把全部模式串齐了:DTO 边界、鉴权收敛、UoW 事务、写走 Command、事件出箱、读走 Query 回显、统一信封。
五、pyproject.toml:import-linter 边界合同(边界执法)
[tool.importlinter]
root_package = "app"
[[tool.importlinter.contracts]]
name = "domain 不许依赖外层"
type = "forbidden"
source_modules = ["app.modules.*.domain"]
forbidden_modules = ["app.modules.*.api", "app.modules.*.application",
"app.modules.*.infrastructure", "app.shared", "app.core"]
[[tool.importlinter.contracts]]
name = "application 不许依赖 api / infrastructure"
type = "forbidden"
source_modules = ["app.modules.*.application"]
forbidden_modules = ["app.modules.*.api", "app.modules.*.infrastructure"]
[[tool.importlinter.contracts]]
name = "core/shared 不许依赖业务模块"
type = "forbidden"
source_modules = ["app.core", "app.shared"]
forbidden_modules = ["app.modules"]
[[tool.importlinter.contracts]]
name = "模块之间不许直接 import 对方内部(只能走事件/公开接口)"
type = "forbidden"
source_modules = ["app.modules.order"]
forbidden_modules = ["app.modules.*.domain", "app.modules.*.infrastructure"]
接进 .pre-commit-config.yaml + CI:lint-imports(不过红灯,不让提交/合并)。
六、加一个新业务上下文的清单(照着做)
以加 payment 为例:
- 复制
modules/order/→modules/payment/,改名。 domain/:写充血实体 + 值对象 + 仓储接口(Protocol)+ 领域事件。infrastructure/:写orm.py(PO) +repository_impl.py(实现接口、PO↔领域对象互转)。application/:写payment_service.py(写,含 outbox.emit)+payment_query.py(读视图)。api/v1/:写schemas.py(DTO) +endpoints.py(用 UoW、写完读回),在main.py挂 router。- 若跨上下文协作:只发领域事件(payment 发
payment.succeeded→ order 订阅去改状态),不要 import order 的内部。 - 跑
lint-imports+ 单测;建 alembic 迁移。
💡 判断点:业务规则进
domain(充血);流程编排进application;跨上下文走事件。删掉某段逻辑会"算错业务/破坏不变量"→ 它必须在 domain。
七、最小起步 vs 逐步加固
| 阶段 | 上什么 |
|---|---|
| 最小可跑 | core(config/db/responses/deps) + 一个模块四层 + FastAPI;存储就 PostgreSQL |
| 加可靠性 | shared/outbox(事件不丢)+ UoW(事务)+ import-linter(边界执法) |
| 加规模 | shared/jobs(Taskiq 异步) + Redis(缓存/锁)+ 对象存储 |
| 加可观测 | shared/observability(OTel + 日志),TraceContext 贯穿 HTTP/事件/任务 |
| 要实时 | outbox relay 接 WebSocket(Centrifugo)/ 要解耦就接 MQ |
不要一上来全堆。先把"模块四层 + 边界执法 + 统一响应"立住,其余按真实需要逐个加。
