加载中...

技术栈示例: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.pyredis.pystorage.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 为例:

  1. 复制 modules/order/modules/payment/,改名。
  2. domain/:写充血实体 + 值对象 + 仓储接口(Protocol)+ 领域事件。
  3. infrastructure/:写 orm.py(PO) + repository_impl.py(实现接口、PO↔领域对象互转)。
  4. application/:写 payment_service.py(写,含 outbox.emit)+ payment_query.py(读视图)。
  5. api/v1/:写 schemas.py(DTO) + endpoints.py(用 UoW、写完读回),在 main.py 挂 router。
  6. 若跨上下文协作:只发领域事件(payment 发 payment.succeeded → order 订阅去改状态),不要 import order 的内部
  7. 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

不要一上来全堆。先把"模块四层 + 边界执法 + 统一响应"立住,其余按真实需要逐个加。

公告栏
这是我的个人知识库。
记录技术,也记录生活 —— 读过的、试过的、想明白的,都堆在这儿。
最新文章
网站资讯
文章数目 :
5
已运行时间 :
本站总字数 :
15.7k
本站访客数 :
本站总访问量 :
最后更新时间 :
全局知识图谱
当前页面 已访问 文章 标签
ESC 关闭 · 滚轮缩放 · 拖拽移动 · Ctrl+G 开关