找回密码
立即注册
搜索
热搜: Java Python Linux Go
发回帖 发新帖

4722

积分

0

好友

612

主题
发表于 1 小时前 | 查看: 3| 回复: 0

聚合不是把“订单相关的一切”装进一个 Java 对象;它定义了一次本地提交必须同时成立的业务事实。边界过大,锁、加载和缓存一起膨胀;边界过小,不变量就会在并发下被绕过。

很多 DDD 项目在低流量时看不出问题:@Transactional 包住订单、库存、优惠券和支付,接口也能跑通。压力上来后,慢事务占满连接、乐观锁冲突重试、消息漏发或重复消费,最后才发现“强一致”没有换来正确性。

本文用下单支付场景回答四件事:

  • 怎么用不变量而不是对象关系划定聚合;
  • 怎样让订单提交和异步事件不丢失;
  • 并发、超时、重复投递时,各系统分别知道什么、如何恢复;
  • 何时需要 CQRS、Saga、缓存或分片,何时还不需要。

文中的表名、阈值、Topic 和代码均为示例;请按自己的数据库、消息中间件和合规要求调整。

1. 先找不变量,不要先画对象图

订单有一个典型的聚合内不变量:

totalAmount = Σ(orderItem.subtotal)

订单项增删、价格快照变化与总价更新必须在同一提交中成立。它们适合放在 Order 聚合内:

Order
├── OrderItem
├── Money
└── AddressSnapshot

PaymentInventoryLogisticsInvoice 虽然与订单有关,却有独立生命周期和故障边界。它们应是不同聚合,通过 ID 和事件协作:

Order --OrderPaid--> Inventory
   \--OrderPaid--> Fulfillment

判断边界时依次问:

  1. 修改 A 时,B 的哪个事实必须在同一原子操作中成立?
  2. 如果 B 暂后处理,用户看到短暂不一致是否可接受?谁负责收敛?
  3. A、B 是否会被不同请求高频独立修改,或拥有无限增长的集合?

前两问都回答“必须同时成立”时,才有理由放进同一聚合。业务关联、同一页面展示、同一数据库表都不是理由。

一个反例:巨型订单为何会放大故障

下面的对象图看似完整,实际上把不同的写热点绑到同一个版本和事务上:

Order
├── items
├── paymentAttempts
├── coupons
├── shipmentTracks (持续增长)
└── operationLogs  (持续增长)

地址修改和订单项修改本可并行,却会争抢同一个 Order.version;加载订单详情又可能触发集合的 JOIN、Lazy Loading 或 N+1。把整个图序列化为 order:{id} 缓存后,一条物流轨迹也会使整份订单缓存失效。

这不是 ORM 的“小问题”,而是聚合边界把查询、并发与缓存粒度错误地绑在了一起。

2. 写模型小而完整,读模型按页面组装

建议的写侧边界如下:

Order Aggregate        Payment Aggregate       Inventory Aggregate
├── OrderItem          ├── PaymentAttempt      └── StockReservation
├── totalAmount        └── providerReference       └── reservedQuantity
└── status

跨聚合只存标识,不持有实体引用:

// L3:写模型的边界示例
public final class Order {
    private final OrderId id;
    private final CustomerId customerId;
    private final List<OrderItem> items;
    private Money totalAmount;
    private OrderStatus status;
    private long version;

    public void addItem(ProductSnapshot product, int quantity) {
        if (status != OrderStatus.DRAFT) {
            throw new DomainException("ORDER_NOT_EDITABLE");
        }
        if (quantity <= 0) throw new DomainException("INVALID_QUANTITY");

        OrderItem item = OrderItem.of(product, quantity);
        items.add(item);
        totalAmount = totalAmount.add(item.subtotal()); // 不变量在唯一写入口维护
    }

    public OrderPaidEvent markPaid(PaymentId paymentId, Instant occurredAt) {
        if (status != OrderStatus.PENDING_PAYMENT) {
            throw new DomainException("ORDER_NOT_PAYABLE");
        }
        status = OrderStatus.PAID;
        return new OrderPaidEvent(EventId.newId(), id, paymentId,
            items.stream().map(OrderItem::toPurchaseLine).toList(), occurredAt);
    }
}

事件契约显式表达消费者需要的领域事实,并带版本;不要直接发送 ORM Entity:

public record OrderPaidEvent(
    EventId eventId,
    OrderId orderId,
    PaymentId paymentId,
    List<PurchaseLine> lines,
    Instant occurredAt,
    int schemaVersion
) {}

public record PurchaseLine(SkuId skuId, int quantity) {}

应用服务只负责加载、调用、保存和事务编排;它不应通过 setter 拼接领域状态:

@Transactional
public void recordPaymentSucceeded(RecordPaymentSucceeded command) {
    Order order = orderRepository.mustFind(command.orderId());
    OrderPaidEvent event = order.markPaid(command.paymentId(), clock.instant());
    orderRepository.save(order);
    outboxRepository.append(OutboxEvent.from(event));
}

这里的本地事务保护的是:orders.status = PAID 与“有一条待发布的 OrderPaid 事实”同时存在。它保护库存、物流或通知已经完成。

订单详情页则不必从 Order 聚合递归加载。可以从专用投影 order_detail_view 查询订单、支付状态和物流摘要;投影允许延迟,写模型不允许绕过不变量。

3. 状态不是枚举装饰:先写清合法迁移

订单支付过程中最危险的状态是 UNKNOWN:调用方超时不等于支付失败。必须把“已知事实”和“尚未知晓”存下来,不能据此直接换渠道重扣。

DRAFT -> PENDING_PAYMENT -> PAID -> FULFILLING -> COMPLETED
                    |           
                    +-> PAYMENT_UNKNOWN -> PAID | PAYMENT_FAILED
PENDING_PAYMENT -> CANCELLED

约束如下:

状态 允许动作 明确禁止
PENDING_PAYMENT 发起一次支付尝试 同订单并发发起第二次尝试
PAYMENT_UNKNOWN 查询原支付方、等待 webhook、恢复扫描 换渠道重扣、直接取消为失败
PAID 发出履约事件 再次扣款
PAYMENT_FAILED 在明确失败后由用户重试 使用旧 payment attempt 重放

支付尝试必须保存可查询的事实:payment_id、本系统幂等键、支付方交易号(若已获得)、请求时间、状态、最后查询时间和错误码。它是 webhook、定时恢复任务与人工排障的共同依据。

远程支付的失败路径

以“支付方调用超时”为例,正确性不来自一句“稍后查询”,而来自完整闭环:

  1. 触发点:HTTP 客户端超时;本系统不知道支付方是否已扣款。
  2. 本地事实PaymentAttempt 已以唯一幂等键创建,订单为 PAYMENT_UNKNOWN;支付方结果未知。
  3. 保护:订单版本 CAS 和 payment_attempt.order_id + active_attempt 唯一约束,拒绝第二笔扣款;webhook 与恢复任务只能做条件迁移。
  4. 恢复:任务按 next_check_at 查询相同的支付方交易号/幂等键;查询仍失败则退避重试并告警,不能把未知伪造成失败。
  5. 验证:注入“支付方已成功但响应丢失”的故障;确认只产生一笔扣款、webhook 与任务竞争后仅一次迁移、订单最终收敛。

外部副作用的幂等键应由业务操作生成并在重试中复用,例如 pay:{orderId}:{attemptNo};不要用每次 HTTP 调用的新 UUID。

4. 跨聚合一致性:本地事务 + Outbox,而不是大事务

把订单、库存、优惠券、钱包全部包进一个事务,在单库低并发时或许能工作;一旦跨库、跨服务,普通数据库事务并不能覆盖所有资源,长锁和同步调用反而扩大故障面。

对“订单支付后通知库存”的简单最终一致场景,最小可靠链路是:

Order transaction                 asynchronous delivery
UPDATE orders(status=PAID)        claim outbox -> broker ack -> mark published
INSERT outbox_event(PENDING)  -->  Kafka(orderId as key) --> inventory consumer

订单更新与 Outbox 插入在同一个数据库事务提交。Broker 暂时不可用时,事件留在数据库等待发布;这解决了“数据库已提交、消息完全没留下”的双写窗口。

Outbox 数据模型与抢占

以下 DDL 为 MySQL 风格示例。字段的重点不是表有多全,而是明确事实来源、写入者和恢复条件:

CREATE TABLE outbox_event (
  id             BIGINT PRIMARY KEY AUTO_INCREMENT,
  event_id       CHAR(26) NOT NULL,
  aggregate_id   VARCHAR(64) NOT NULL,
  event_type     VARCHAR(128) NOT NULL,
  payload        JSON NOT NULL,
  status         VARCHAR(16) NOT NULL, -- PENDING, PUBLISHING, PUBLISHED
  retry_count    INT NOT NULL DEFAULT 0,
  next_retry_at  DATETIME NOT NULL,
  lease_until    DATETIME NULL,
  last_error_code VARCHAR(64) NULL,
  created_at     DATETIME NOT NULL,
  published_at   DATETIME NULL,
  UNIQUE KEY uk_outbox_event_id (event_id),
  KEY idx_outbox_dispatch (status, next_retry_at, lease_until)
);
  • 业务事务写入 PENDING,发布器只读取已提交的行;payload 是领域事实,不是 JPA Entity 序列化快照。
  • 多个 Pod 通过条件更新领取发送权;lease_until 使崩溃实例遗留的 PUBLISHING 能被恢复扫描重新领取。
  • PUBLISHED 仅在收到 Broker ack 后写入。ack 丢失时可能重发,因此消费者仍必须幂等。
-- 发布器领取一条可重试事件;受影响行数为 1 才拥有租约
UPDATE outbox_event
SET status = 'PUBLISHING', lease_until = :leaseUntil
WHERE id = :id
AND ((status = 'PENDING' AND next_retry_at <= :now)
OR (status = 'PUBLISHING' AND lease_until < :now));

发布流程(L2 伪代码,关键失败分支已展开):

for each candidate selected in a small batch:
  if conditional-claim(candidate) == 0: continue
  send(eventType, aggregateId as partition key, payload)
  if broker ack:
      UPDATE ... SET status=PUBLISHED, published_at=now WHERE id=? AND status=PUBLISHING
  else:
      UPDATE ... SET status=PENDING, retry_count=retry_count+1,
          next_retry_at=backoff(retry_count), last_error_code=? WHERE id=? AND status=PUBLISHING

批量大小、租约时长和退避上限取决于消息大小、Broker 延迟与数据库负载,需压测后设定,不能把示例数字照搬线上。

消费端:至少一次投递,业务幂等

不要把“只发一次”当作整个链路的承诺。事件可能在 ack 丢失、消费者重平衡或进程崩溃后重复到达。消费者在同一事务中先登记去重事实,再修改自己的聚合:

CREATE TABLE consumed_event (
  consumer_name VARCHAR(64) NOT NULL,
  event_id      CHAR(26) NOT NULL,
  processed_at  DATETIME NOT NULL,
  PRIMARY KEY (consumer_name, event_id)
);
@Transactional
public void onOrderPaid(OrderPaidEvent event) {
    boolean firstDelivery = consumedEventRepository.insertIfAbsent(
        "inventory-service", event.eventId(), clock.instant());
    if (!firstDelivery) return;

    for (PurchaseLine line : event.lines()) {
        int changed = jdbc.update(
            """
            UPDATE stock
            SET available = available - :quantity, version = version + 1
            WHERE sku_id = :skuId AND available >= :quantity
            """, line.skuId(), line.quantity());
        if (changed != 1) throw new RetryableBusinessException("STOCK_NOT_AVAILABLE");
        reservationRepository.create(event.orderId(), line.skuId(), line.quantity());
    }
}

这里的条件更新同时是库存不变量的保护位置:available >= quantity。若库存不足,抛出可处理异常让消费事务回滚,去重行也会回滚;之后可重试或由流程状态机转为补偿,而不是吞掉事件。

同一订单需要有序的事件,应以 orderId 作为 Kafka key,使其落在同一分区。这个保证只在单分区范围内有效,不是全局顺序。

5. 并发控制不由聚合自动完成

聚合能定义谁有权修改事实,但不能替数据库处理竞争。库存扣减不要使用“先查库存、Java 判断、再更新”的读改写模式;并发窗口会使两个请求都看到同一库存。

优先将判断放进单条条件更新:

UPDATE inventory
SET available_stock = available_stock - :quantity,
    version = version + 1
WHERE sku_id = :skuId
AND available_stock >= :quantity
AND version = :expectedVersion;

affected_rows = 1 才表示成功。失败时应用必须返回明确的竞争或库存不足结果,有限次数重读重试;不能无限重试,也不能在失败后写成功状态。

对不同热点,选择不同机制:

约束 优先选择 代价与边界
普通库存 数据库条件更新 / CAS 简单、强约束;热点会有冲突
少量强竞争资源 短时悲观锁 等待与死锁风险,禁止持锁远程调用
极高热点、可预留 库存预扣或分段 引入过期释放与对账
缓存侧计数 原子操作/Lua 仍需落库事实与故障恢复

并发测试应覆盖不变量,而不是只测接口返回:

@Test
void onlyOneRequestCanConsumeTheLastUnit() throws Exception {
    // Given: sku-1 available=1;两个线程从同一屏障同时开始
    // When : 并发执行条件扣减
    // Then : 恰好一个 affected_rows=1;最终 available=0,且从不为负
}

还应演练两类竞争:两个 Publisher 领取同一 Outbox 行,及 webhook 与支付恢复任务同时抵达同一个 PaymentAttempt。两者都应以唯一约束或条件状态迁移收敛到一个事实。

6. 查询、缓存与 ORM:服务访问模式,不服务对象图

写侧 Repository 只加载维护不变量需要的 Order。复杂查询交给 Query Repository 或物化读模型:

OrderDetail API -> order_detail_view
                     ├── order status
                     ├── payment summary
                     └── shipment summary

这避免为了展示详情而加载订单、用户、支付、物流、优惠券的完整实体图。JPA 中默认使用 LAZY 并不等于可以忽略 N+1;应为具体查询使用 projection、显式 join 或 EntityGraph,测量 SQL 数与扫描行数。

缓存也按读取契约拆分,例如 order:summary:{id}order:status:{id},而非默认缓存一个不断变大的聚合 JSON。缓存失效或陈旧只能影响读取体验,不能成为支付、库存等写侧事实来源。

当读查询已经影响写库,或需要搜索、分页、跨聚合视图时,再由事件构建读模型。投影的延迟需要向产品明确;用户刚支付时,详情页短暂显示旧物流状态是可接受的,支付状态本身仍应从权威写侧或带版本的读模型确认。

7. 何时需要 Saga,何时不需要

OrderPaid -> reserve inventory 这类单一后续动作,Outbox、重试和幂等通常足够。

若流程有多个可失败步骤,且每一步都有明确的反向业务动作,例如“锁库存 → 扣款 → 建物流”,才需要状态机/Saga:

RESERVING_STOCK -> CHARGING -> CREATING_SHIPMENT -> COMPLETED
       |               |                |
       +---------------+----------------+-> COMPENSATING -> CANCELLED

补偿不是技术性回滚:已通知用户、已发货或已向银行扣款时,反向动作分别是取消预留、退款或人工处置。把它们建成可观察状态,优于用 2PC 试图把跨服务网络调用塞进一个长事务。2PC 在受控场景可用,但不应成为高并发业务的默认架构。

8. 上线、验证与回滚:让值班人员能做决定

聚合重构上线不是换一层代码。先用历史订单和故障样本验证状态机、事件 schema 与投影,再以确定性方式灰度(例如按订单 ID 哈希固定分桶)。影子计算可以记录新旧决策差异,但绝不能执行第二次扣款、扣库存或发货。

建议发布路径:

DRAFT
  -> schema / 状态迁移 / 唯一约束校验
  -> 历史样本模拟(无候选、重复、未知支付)
  -> 小流量确定性灰度
  -> 观察 Outbox、消费与业务收敛
  -> 全量;异常则回滚到 Last Known Good 配置/版本

数据库迁移须向后兼容:先加表、列和索引,再发布可同时读旧写新的代码;确认无旧版本实例后才清理旧字段。回滚应用版本不等于回滚业务事实;已发布事件继续按 event schema 兼容消费,必要时启动恢复任务而不是删除记录。

最小观测面与处置动作

信号 先看什么 典型处置
P99 与连接池同时升高 慢 SQL、事务时长、下游等待 缩短事务;移出同步远程调用;限流
乐观锁冲突升高 热点 aggregate ID、写操作类型 合并重复写;拆独立生命周期数据;有限重试
Outbox 积压增长 最老 created_at、Broker 错误、租约卡死 扩/修 Publisher;恢复过期租约;不丢弃事件
消费延迟或 DLQ 增长 event type、错误码、重复率 修复可重试故障;对不可处理事件走业务补偿/人工队列
PAYMENT_UNKNOWN 超阈值 provider、最后查询时间、webhook 失败率 触发查询恢复;必要时降级支付入口;禁止重扣

日志和指标应关联 order_idpayment_idevent_id、事件版本、trace ID 与错误码;订单金额、支付凭证、地址等按数据分级脱敏。报警阈值必须来自基线与业务时限,本文不虚构一个通用数值。

9. 从单体演进,不要一次引入所有组件

一个可控的演进顺序是:

  1. 先明确 Order、Payment、Inventory 的事实归属和不变量;消除跨聚合对象引用。
  2. 将聚合内写操作收敛为短本地事务,移除事务内远程调用。
  3. 当需要可靠异步通知时引入 Outbox,并实现 Broker ack、租约恢复与消费幂等。
  4. 当查询拖累写侧时增加读模型;当出现多步骤且需补偿的流程时再引入 Saga。
  5. 只有数据量或热点已经证明单库、单分区不足时,再评估分区、分片、预扣等更复杂方案。

Kubernetes 扩容、Redis、Kafka、ES、Saga 都不能修复错误的聚合边界。它们各自引入新的容量、陈旧性和恢复责任,应在简单方案确实失效时再引入。

10. 交付前检查

聚合与状态

  • 每个聚合都有明确的本地不变量和唯一修改入口。
  • 独立生命周期、无限集合或高频独立写入的数据没有被塞进巨型聚合。
  • 状态迁移、非法迁移和 UNKNOWN 的事实来源均已定义。

并发与消息

  • 唯一约束、条件更新/CAS 的位置明确;竞争失败有确定行为和并发测试。
  • 业务变更与 Outbox 事件在同一事务提交;发送 ack 前不标记 PUBLISHED
  • Publisher 崩溃租约可回收;消费者按 consumer + event_id 幂等。
  • 需要顺序的事件使用聚合 ID 分区键,并接受单分区边界。

运行与恢复

  • 有历史样本、重复消息、ack 丢失、支付未知和库存竞争的演练。
  • 能从 Outbox、消费去重表、支付尝试记录还原事实并安全恢复。
  • 监控能区分锁竞争、慢查询、积压、Broker 故障和业务拒绝;日志不泄露敏感数据。
  • 版本发布可灰度、可观察、可回滚;回滚不会抹掉已经发生的业务事实。

结语

成熟的聚合设计最终不是“对象划分得漂亮”,而是能清楚回答:谁拥有这个事实?哪个不变量必须一起提交?并发时由谁拒绝错误写入?远程结果未知时系统如何避免二次副作用?失败后又从哪里恢复?

这些答案落在状态、约束、事件和运行步骤上,DDD 才从建模术语变成一个能承受真实故障的工程边界。对这类围绕 后端与架构 展开的一致性、分布式系统与高并发问题,最终都要回到可恢复的工程边界上去验证。




上一篇:set_false_path 别乱用:FPGA 时序约束与 CDC 跨时钟域避坑指南
下一篇:Trellis 14.4k Star 给 AI 编程工具装「项目记忆」:这 3 类人先别碰
您需要登录后才可以回帖 登录 | 立即注册

手机版|小黑屋|网站地图|云栈社区 ( 苏ICP备2022046150号-2 )

GMT+8, 2026-9-13 20:01 , Processed in 0.335290 second(s), 40 queries , Gzip On.

Powered by Discuz! X3.5

© 2025-2026 云栈社区.

快速回复 返回顶部 返回列表