系列:Han Menu 外卖系统实践 · 从第一篇开始

上一篇:第6篇 · 下一篇:第8篇

系列:Han Menu 外卖系统实践 · P6

面向读者:已经理解订单、支付、事务和可靠事件,希望继续学习实时通知、长连接认证、统计口径与数据投影的开发者。

本篇依据 P6 提交 82c2b8c 编写,承接 P5 提交 ecb4228。关键代码直接摘自该提交;注明省略的片段不是独立完整文件。设计图表达职责与执行顺序,不代表已经实现了新的前端工程。

P6 完成顾客催单、员工通知、WebSocket、经营统计、XLSX 导出和投影重建。Flutter、新管理端页面及多实例实时广播仍不属于本篇的已实现范围。

1. 支付和履约已经能跑,为什么还需要通知与报表

上一篇已经能处理真实支付结果、取消、退款、接单和配送。但业务运行起来后,会遇到新的问题:

  • 顾客已经付款,员工怎样及时知道有新订单?
  • 员工浏览器没打开,订单通知是不是就丢了?
  • 顾客催单能否无限点击,能否催别人的订单?
  • 员工退出登录后,已经建立的 WebSocket 还能继续收消息吗?
  • 今天完成的订单,是否一定是今天创建的订单?
  • 今天退款了一笔昨天的付款,营业额应该直接减掉吗?
  • 统计投影损坏后,能否重建,又会不会在重建期间显示半份报表?

这些问题分属不同边界,却有共同点:业务事实已经存在,其他功能需要可靠地理解和使用它们。

本篇将沿着两条主线展开:

主线 需要建立的能力
让员工知道发生过什么 持久化消息、实时提示、断线补查、个人阅读进度
让经营者理解业务情况 明确指标口径、独立统计投影、同快照查询、原子重建

WebSocket 只负责其中一段网络通信,图表也只负责其中一层展示。先把数据含义和恢复流程设计清楚,再选择传输与展示工具。

2. 为什么 notification 和 reporting 应当独立

通知关注“有什么事实值得提醒员工”。统计关注“怎样按明确口径组织历史和当前事实”。

它们不应该直接取得修改订单状态或重新判断付款成功的权限。

flowchart LR
  O[ordering 订单事实] --> OE[公开订单事件]
  C[customer 注册事实] --> CE[公开注册事件]
  P[payment 收退款事实] --> PE[公开支付事件]
  OE --> N[notification 持久化通知]
  OE --> R[reporting 统计投影]
  CE --> R
  PE --> R
  APIs[ordering/customer/payment 快照API] --> R
  I[identity 公开会话与权限能力] --> N
  I --> R

这张图表示事实和能力的使用关系,不表示业务模块之间互相发送 HTTP。当前仍是同一个 Spring Boot 应用、同一个数据库事务管理体系。

P6 的模块边界可以概括为:

模块 自己保存什么 不能做什么
ordering 催单次数、催单时间、权威订单状态 把浏览器是否在线作为订单提交条件
notification 消息、投递尝试、个人阅读进度、连接票据 擅自修改订单或把发送成功当作员工阅读
reporting 最小业务事实投影、控制版本、统计查询结果 直接访问其他模块业务表,或用报表修造支付事实
identity 员工会话与账号有效性 因为已经握手就永久认可旧权限

notification 和 reporting 的表由各自模块维护。它们通过公开事件或 API 获取业务事实,不导入订单 JPA 实体,也不直接查询 ordering_order

3. 在增加消费者前,先把事件含义定义清楚

P5 已有 PaymentResult 和 RefundResult。P6 继续增加几种不同职责的事实:

事件 表达什么 主要用途
OrderChanged 已保存版本的完整订单统计快照 更新统计投影
OrderReady 真实付款使订单进入待接单状态 员工来单提醒
OrderReminderRaised 一次合法催单已经登记 员工催单提醒
CustomerRegistered 顾客注册标识与时间 顾客增长统计
PaymentResult / RefundResult 已确认收款或退款结果 资金事实投影

不要用一个泛化的 OrderUpdated 事件,让每个消费者自己猜测“这是来单、催单,还是普通版本变化”。

3.1 已付款不一定应该产生新的待接单提示

P6 在付款结果应用之后,只对真实的 UNPAID → PAID 转换登记来单事件。实际关键片段如下:

save(order);
if (previous == Order.Status.UNPAID && order.status() == Order.Status.PAID) {
  events.ready(order.id(), order.lifecycle().paidAt());
}

如果取消后的迟到付款正在触发退款补偿,订单不会重新进入 PAID,也不应该给员工制造“请接新单”的提示。

事件描述业务意义,不能只看是否收到了某种技术回调。

3.2 统计事件使用已经持久化的版本

OrderEvents.changed() 的实际实现:

void changed(Order previous) {
  var current = orders.find(previous.id()).orElseThrow();
  if (current.version() > previous.version()) {
    events.publishEvent(new OrderChanged(OrderFactsService.snapshot(current)));
  }
}

外层业务用例先保存订单,再按已 flush 的版本读取快照。如果确实产生新版本,才发布 OrderChanged。

这种顺序让消费者可以比较源版本,而不是依靠事件抵达时间来猜测哪一份订单更新。

事件登记仍与业务状态同事务;后面发生回滚时,不留下脱离业务结果的消息。

4. 顾客催单是一项业务行为,不是一个自由文本消息接口

新增接口:

POST /api/v1/orders/<订单UUID>/reminders
Authorization: Bearer <顾客令牌>
Content-Type: application/json

{"version": 3}

这里的版本只是示例,应使用本人订单实际版本。请求不能提交正文、目标员工或任意顾客 ID。

4.1 状态、版本和冷却窗口属于订单规则

public void remind(long expectedVersion, Instant now) {
  requireVersion(expectedVersion);
  if (status != Status.PAID && status != Status.ACCEPTED && status != Status.DELIVERING) {
    conflict("当前订单不能催单");
  }
  if (lastRemindedAt != null && now.isBefore(lastRemindedAt.plusSeconds(60))) {
    conflict("两次催单至少间隔六十秒");
  }
  reminderCount = Math.incrementExact(reminderCount);
  lastRemindedAt = Objects.requireNonNull(now);
}

当前允许 PAID、ACCEPTED、DELIVERING 催单。未支付、取消处理中、退款中、已取消和已完成都不允许。

冷却规则为两次成功催单至少间隔六十秒:

t_{next}\ge t_{last}+60\text{秒}

它与前端按钮防抖不同。即使两台设备同时发请求,后端仍然要保护该规则。

4.2 应用服务在订单行锁里完成判断和保存

public OrderViews.Detail remind(CustomerIdentity identity, UUID id, long version) {
  customers.lockActive(identity);
  var order = ownLocked(identity, id);
  order.remind(version, clock.instant());
  save(order);
  events.reminder(order);
  return OrderViews.detail(orders.find(id).orElseThrow());
}

顾客身份重验、订单归属、订单版本和冷却窗口共同构成这次用例。状态保存与催单事件一起提交。

成功催单会推进订单版本,因此响应必须带回最新订单详情,不能让客户端继续用旧版本处理下一步操作。

4.3 稳定事件 ID 避免重复建立通知

催单事件的实际登记方法:

void reminder(Order order) {
  events.publishEvent(
      new OrderReminderRaised(
          key("reminder:" + order.id() + ":" + order.reminderCount()),
          order.id(),
          order.reminderCount(),
          order.lastRemindedAt()));
}

其中 key() 使用订单标识和催单次数生成稳定 UUID。重复投递同一次催单,得到相同通知标识;下一次合法催单则是另一个标识。

这是业务去重键,不是认证凭证。是否允许催单,仍由前面的身份和聚合规则决定。

5. 通知持久化、实时投递、员工阅读,是三种不同状态

假设顾客付款成功时,所有员工都离线:

  • 来单事实应该被保存。
  • 实时投递暂时没有接收者。
  • 不能因此认为员工已经看到了消息。

P6 将这三个层次分开:

flowchart TD
  E[已提交业务事件] --> N[持久化通知 Notice]
  N --> D[实时投递与尝试轨迹]
  N --> F[HTTP按游标补查]
  D --> W[当前合法WebSocket连接]
  F --> U[员工处理完整消息页]
  U --> R[该员工的Receipt阅读进度]

Notice.status=DELIVERED 不会自动修改 Receipt。

判断本次实时投递是否成功的代码也很明确:

record Outcome(int sent, int failed) {
  /** 全部当前合法连接发送成功才完成本次实时提示任务. */
  public boolean successful() {
    return sent > 0 && failed == 0;
  }

  /** 返回不包含网络细节或身份的固定失败分类. */
  public String failure() {
    return failed > 0 ? "SEND_FAILED" : "NO_SUBSCRIBERS";
  }
}

至少成功发送到一个当前合法连接,而且没有发送失败,才算本次实时提示投递成功。

它不是“所有员工都在线并已阅读”,也不是浏览器已经处理消息的端到端确认。即使实时提示耗尽,持久化消息仍然可补查。

6. 游标最容易踩的坑:递增编号不等于提交顺序

6.1 为什么直接使用数据库序列可能漏消息

假设消息 ID 使用一个提前分配的递增序列:

时刻 事务 A 事务 B 客户端
T1 分配序号 10,尚未提交
T2 分配序号 11 并提交
T3 仍未提交 读到 11,保存游标 11
T4 提交序号 10 后续只查大于 11 的消息

序号 10 就被跳过了。

PostgreSQL 序列的变更不是普通业务行的事务回滚语义,因此不能仅凭 sequence/serial 就推断消息按提交顺序可见。PostgreSQL 事务隔离说明

6.2 用同事务控制行串行分配消息游标

P6 的真实追加实现:

public void append(UUID id, UUID orderId, Notice.Kind kind, Instant occurredAt, Instant now) {
  var feed = feeds.findLockedById(1).orElseThrow();
  if (notices.existsById(id)) {
    return;
  }
  feed.sequence = Math.incrementExact(feed.sequence);
  notices.saveAndFlush(
      NoticeEntity.from(Notice.create(id, feed.sequence, orderId, kind, occurredAt, now)));
}

所有追加都先锁 notification_feed 中同一条控制行:

  1. 检查稳定事件 ID 是否已经生成通知。
  2. 在托管控制行上递增 sequence。
  3. 保存该 sequence 对应的消息。
  4. 两者共同提交,或共同回滚。

第二个追加事务必须等第一个提交或回滚后再分配。因此客户端不会先越过一个尚未提交的较小消息游标。

这里的 sequence 是通知持久化提交顺序。它不是跨模块事件的全局发生时间排序;较早发生的业务事件仍然可能较晚被消费和追加。

这份保证依赖所有消息追加遵守同一控制行协议。它也意味着追加存在一个串行协调点,这是当前单店场景的明确取舍。

7. 补查分页不能跳到另外查询的“全局最大值”

NotificationService.feed() 的实际实现:

public Feed feed(StaffIdentity actor, long after, int limit) {
  authorization.requireStaff(actor);
  if (after < 0 || limit < 1 || limit > 100 || after > repository.head()) {
    throw new NotificationException(NotificationException.Reason.INVALID_INPUT, "通知游标或分页范围不合法");
  }
  var all = repository.after(after, limit + 1);
  var items = all.stream().limit(limit).map(NotificationService::view).toList();
  return new Feed(
      items, items.isEmpty() ? after : items.getLast().sequence(), all.size() > limit);
}

它多读取一项,用于判断 hasMore,但返回至多 limit 项。

nextCursor 只取本页真正返回的最后一个序号。如果没有返回消息,就保持原 after。

为什么不在查询结束后再查一下全局 head,并把它作为 nextCursor?因为分页和查询 head 之间可能有新消息提交。如果没有把这些消息返回,却把游标推到它们后面,客户端就会漏读。

例如查询 limit=2,返回 10、11,期间 12、13 也提交了。当前页游标仍然应该是 11,下一页再查询大于 11 的内容。

服务端的消息排序、分页边界和客户端推进规则必须互相配合;仅仅提供一个 after 参数还不够。

8. 实时投递的重试,需要自己的状态机和尝试标识

Notice 使用四种状态:

stateDiagram-v2
  [*] --> PENDING
  PENDING --> IN_FLIGHT: 领取并生成尝试UUID
  IN_FLIGHT --> DELIVERED: 当前尝试发送成功
  IN_FLIGHT --> PENDING: 失败且仍可重试
  IN_FLIGHT --> EXHAUSTED: 连续五次失败
  IN_FLIGHT --> IN_FLIGHT: 处理窗口到期后重新领取
  EXHAUSTED --> PENDING: 管理员按版本重投

8.1 领取时生成独立的 activeAttemptId

public UUID claim(Instant now) {
  if (nextAttemptAt == null
      || nextAttemptAt.isAfter(now)
      || status == Status.DELIVERED
      || status == Status.EXHAUSTED) {
    return null;
  }
  status = Status.IN_FLIGHT;
  attempts = Math.incrementExact(attempts);
  activeAttemptId = UUID.randomUUID();
  nextAttemptAt = now.plusSeconds(60);
  return activeAttemptId;
}

这里与 P5 仅使用时间窗口的任务领取相比,多了一个明确的投递尝试 UUID。

如果 A 领取后卡住,六十秒后 B 重新领取,B 会得到新的 activeAttemptId。

8.2 迟到结果必须证明自己仍属于当前尝试

public boolean finish(UUID attemptId, boolean success, String failure, Instant now) {
  if (status != Status.IN_FLIGHT || !Objects.equals(activeAttemptId, attemptId)) {
    return false;
  }
  activeAttemptId = null;
  if (success) {
    status = Status.DELIVERED;
    nextAttemptAt = null;
    lastFailure = null;
    failures = 0;
  } else {
    failures++;
    lastFailure = failure;
    status = failures >= 5 ? Status.EXHAUSTED : Status.PENDING;
    nextAttemptAt =
        status == Status.EXHAUSTED ? null : now.plusSeconds(Math.min(60, 5L << (failures - 1)));
  }
  return true;
}

A 后来返回成功时,如果它的 attemptId 已经不是当前值,就不能把 B 正在处理的通知改成 DELIVERED。

这个令牌保护的是本地结果写入。它不能撤回 A 已经发送出去的网络帧,因此客户端仍然必须接受重复消息。

8.3 失败次数与累计尝试次数不是同一个指标

当前普通失败的等待时间依次为 5、10、20、40 秒。第五次连续失败后进入 EXHAUSTED,不再安排第六次自动投递。

第1次失败 → 5秒后重试
第2次失败 → 10秒后重试
第3次失败 → 20秒后重试
第4次失败 → 40秒后重试
第5次失败 → EXHAUSTED

管理员可以用当前通知版本重投已耗尽消息。本轮 failures 会清零,但累计 attempts 和历史轨迹仍然保留。

没有在线订阅者记录 NO_SUBSCRIBERS;发送故障记录 SEND_FAILED。耗尽不会删除消息,也不会阻止 HTTP 补查。

9. 网络发送仍然放在数据库事务之外

P5 的三段式处理在这里继续适用:领取、发送、登记结果。

NotificationDispatcher 使用类级 Propagation.NEVER,发送方法如下:

public void deliver(UUID id) {
  var claimed = notifications.claim(id);
  if (claimed.isEmpty()) {
    return;
  }
  var notice = claimed.orElseThrow();
  NotificationPush.Outcome outcome;
  try {
    outcome = push.send(notice);
  } catch (RuntimeException exception) {
    outcome = new NotificationPush.Outcome(0, 1);
  }
  notifications.finish(id, notice.activeAttemptId(), outcome);
}

notifications.claim() 在短事务中登记尝试;push.send() 在事务外执行;notifications.finish() 再进入另一段短事务。

这让浏览器断线不会回滚订单,也不会让数据库锁一直等待慢连接。

仓储还会把被新领取替代的旧尝试标为 EXPIRED,并保存结束时间与固定失败原因。返回的旧结果不覆盖已经结束的轨迹和新通知状态。

当前消息扫描最多读取一百条候选,发送和空闲连接复验分别按五秒 fixedDelay 调度。P6 配置四线程维护调度器,避免一个慢任务阻塞全部本地维护。

fixedDelay 不是严格的五秒端到端时限;线程繁忙、数据库或网络延迟仍会影响实际处理时间。

10. 个人阅读进度为什么也需要版本

员工甲确认看到消息 20,不代表员工乙已经看到。P6 按员工 UUID 保存独立 Receipt,而不是在通知表上放一个全店通用的 read=true

领域方法如下:

public Receipt acknowledge(long target, long expectedVersion, long head, Instant now) {
  if (version != expectedVersion) {
    throw new NotificationException(NotificationException.Reason.VERSION_CONFLICT, "阅读进度已更新");
  }
  if (target < sequence || target > head) {
    throw new NotificationException(
        NotificationException.Reason.INVALID_INPUT, "阅读游标不能倒退或超出消息范围");
  }
  return target == sequence ? this : new Receipt(employeeId, target, version, now);
}

它同时要求:

sequence_{old}\le sequence_{target}\le sequence_{committedHead}

以及客户端携带当前阅读进度版本。

两个设备同时推进同一员工的进度时,旧版本不能覆盖新版本。读取不到持久化记录时,返回 sequence=0/version=0;第一次真正前移后持久化版本推进为 1,思路与 P3 首次加购相似。

一个标量游标还隐含了产品约定:确认序号 20,表示此前应处理的消息已经形成连续前缀。它不适合直接表达“只读了 20,但 18、19 明确未读”的任意集合。

如果未来需要逐消息未读状态,就要重新设计阅读模型,不能继续用一个最高序号掩盖中间空洞。

11. 浏览器 WebSocket 为什么先换一次性票据

普通浏览器 WebSocket 构造方式不能像 fetch 一样随意设置业务 Authorization 请求头。把长期员工 Bearer 放进 URL 查询参数又容易进入访问日志和复制链接。

P6 采用两步:

  1. 先通过带员工 Bearer 的普通 HTTP 接口申请票据。
  2. 用固定业务子协议和一次性票据建立连接。

示意调用如下,它是接口用法示例,不是 P7 管理端工程:

const socket = new WebSocket(streamUrl, [
  ticket.protocol,
  `ticket.${ticket.ticket}`
]);

业务协议固定为 han-menu.notifications.v1,票据形如 hmw_...,只在签发响应中出现。

11.1 票据怎样签发

public Ticket issue(String bearer) {
  var proof = sessions.authenticate(bearer).orElseThrow(StreamTicketService::invalid);
  byte[] bytes = new byte[32];
  random.nextBytes(bytes);
  String secret = "hmw_" + BaseEncoding.base64Url().omitPadding().encode(bytes);
  Instant expires = clock.instant().plusSeconds(30);
  repository.issue(
      new StreamTicket(
          proof.sessionHash(),
          digest(secret),
          proof.employeeId(),
          proof.securityVersion(),
          expires,
          null),
      clock.instant());
  return new Ticket(secret, expires, "han-menu.notifications.v1");
}

32 个随机字节提供 256 位随机量,再使用 Base64URL 表达。数据库保存 SHA-256 摘要,同时绑定原员工会话摘要、员工 ID 和账号安全版本。

每个原会话只保留最新票据。重新申请会让之前尚未使用的票据失效;握手失败后也应申请新票据,不把已经消费的票据当成可重复登录密码。

11.2 短有效期还不够,还必须只能消费一次

public StaffSessions.Proof consume(String value) {
  if (value == null || !value.matches("hmw_[A-Za-z0-9_-]{43}")) {
    throw invalid();
  }
  var ticket = repository.lockTicket(digest(value)).orElseThrow(StreamTicketService::invalid);
  var proof =
      new StaffSessions.Proof(
          ticket.employeeId(), ticket.securityVersion(), ticket.sessionHash());
  if (!sessions.active(proof)) {
    throw invalid();
  }
  ticket.consume(clock.instant());
  repository.consume(ticket);
  return proof;
}

票据行锁让两个并发握手不能同时消费同一份票据。领域 StreamTicket.consume() 还会拒绝已消费或已到期状态。

三十秒是票据可用期,不是员工长连接的会话有效期。票据消费后,连接仍绑定原员工会话,必须继续复验。

票据依然是凭证。放进子协议并不会让它变成可以公开记录的普通字符串,代理与应用日志都不应完整记录携带票据的请求头。

12. 握手放行、Origin 检查和身份认证不能混为一谈

通知 WebSocket 有单独安全链,让精确 GET 路径进入握手处理。这里的 permitAll 不代表匿名员工可以建立连接。

真正升级之前,还要通过 Origin 检查和票据消费。实际握手方法如下:

public boolean beforeHandshake(
    ServerHttpRequest request,
    ServerHttpResponse response,
    WebSocketHandler handler,
    Map<String, Object> attributes) {
  if (request.getURI().getRawQuery() != null) {
    response.setStatusCode(HttpStatus.UNAUTHORIZED);
    return false;
  }
  var protocols =
      request.getHeaders().getOrEmpty("Sec-WebSocket-Protocol").stream()
          .flatMap(value -> Arrays.stream(value.split(",")))
          .map(String::strip)
          .toList();
  var secrets = protocols.stream().filter(value -> value.startsWith("ticket.")).toList();
  if (!protocols.contains(NotificationSocket.PROTOCOL) || secrets.size() != 1) {
    response.setStatusCode(HttpStatus.UNAUTHORIZED);
    return false;
  }
  try {
    attributes.put(NotificationSocket.PROOF, tickets.consume(secrets.getFirst().substring(7)));
    return true;
  } catch (com.hanserwei.hanmenu.notification.domain.NotificationException exception) {
    response.setStatusCode(HttpStatus.UNAUTHORIZED);
    return false;
  } catch (RuntimeException exception) {
    response.setStatusCode(HttpStatus.SERVICE_UNAVAILABLE);
    return false;
  }
}

几个设计点值得分开理解:

  • 有 URL 查询参数直接拒绝,避免以查询参数承载凭证。
  • 必须声明固定业务协议,且只能提供一个 ticket. 子协议值。
  • 握手前重验原会话,票据不能使已退出账号恢复有效。
  • 服务端只回显固定业务子协议,不回显票据字符串。
  • Origin 先检查,之后才消费票据。

默认同源;如有跨来源部署需求,配置明确的 HTTP/HTTPS Origin 列表,不使用 *

Origin 是浏览器来源约束,不是身份凭证。没有 Origin 的非浏览器客户端仍然需要同样的有效票据。Spring 的处理器注册与握手拦截机制可参考 WebSocket 服务端文档

13. 已经连接成功,不意味着权限永久有效

握手成功只能说明那一刻认证通过。员工可能随后退出、会话到期、被停用或修改密码。

P6 通过 identity 的公开能力检查原会话:

public boolean active(Proof proof) {
  return proof != null
      && sessions
          .findActive(proof.sessionHash(), clock.instant())
          .filter(
              account ->
                  account.id().equals(proof.employeeId())
                      && account.securityVersion() == proof.securityVersion())
          .isPresent();
}

WebSocket 适配器在每次发送和响应 ping 前调用这项能力,另外周期检查空闲连接。

private boolean authorized(Connection connection) {
  try {
    if (connection.session().isOpen() && sessions.active(connection.proof())) {
      return true;
    }
  } catch (RuntimeException exception) {
    // 认证依赖故障时也关闭连接,不能沿用旧权限。
    close(connection.session());
    return false;
  }
  close(connection.session());
  return false;
}

数据库认证依赖故障时也关闭连接,不能继续沿用曾经验证过的权限。

这不是把数据库事务维持到消息发送结束。每次身份检查使用自己的短只读事务,返回后再写网络。

连接侧还有明确的资源边界:

限制 当前值
每员工连接数量 最多 3 个
当前进程连接总数 最多 64 个
接收文本消息 最多 1024 字节,仅接受 ping
并发发送保护参数 5 秒发送时限、64 KiB 缓冲
空闲认证检查配置 5 秒 fixedDelay

使用 ConcurrentWebSocketSessionDecorator 管理发送,避免多个发送者直接并发写同一会话。非法消息或失效身份按策略关闭连接,不把 WebSocket 变成另一个任意业务命令入口。

14. 客户端恢复顺序,决定会不会漏消息

实时消息可以重复、乱序,也可能在连接中断时丢失。因此客户端不能收到一帧 sequence=20,就直接把个人进度推进到 20。

假设序号 19 的实时发送失败,20 的发送成功。直接确认 20,会让后续从 20 开始的补查跳过 19。

推荐流程如下:

sequenceDiagram
  participant B as 员工客户端
  participant W as WebSocket
  participant H as 通知HTTP接口
  B->>H: 用员工Bearer申请一次性票据
  B->>W: 持票据建立连接
  W-->>B: READY
  B->>H: 读取本人receipt
  loop 按已确认游标连续补查
    B->>H: GET notifications?after=sequence
    H-->>B: items、nextCursor、hasMore
    B->>B: 按id/sequence去重并完成整页处理确认
    B->>H: 按receipt版本前移至本页nextCursor
  end
  W-->>B: 新消息实时提示
  B->>H: 再次补查,而非直接跳到实时帧序号

先建立连接,再读取个人进度并补查,有助于避免“补查完成与连接建立之间”的通知空窗。连接期间到来的实时帧应该合并触发补查,不并发启动一堆互相覆盖的游标写入。

出现 409 时,重读本人阅读进度再恢复;不要猜测应该提交哪个更大的版本或游标。

这里的“整页处理确认”应由客户端产品明确约定。单纯接到网络帧不等于用户阅读,服务端不会自动替员工确认。

通知描述的是发生过的事实。员工真正接单或配送前,仍须读取当前订单及版本,不能因为看到一条旧来单消息就绕过 P5 的状态机。

15. 经营统计从“能查询”开始,还必须回答“查询的是什么”

reporting 为自己的查询保存最小事实投影,而不是让每张图表临时跨模块拼表。

主要事实如下:

投影 保存内容 更新方式
订单 当前源版本、状态、成交金额、关键时间、支付/退款引用、成交行 按源订单版本推进
商品摘要 商品 UUID、种类、最近观察到的订单快照名称 供销量分组显示
注册 顾客 UUID 与注册时间 按稳定 ID 去重
收款 支付 UUID、订单引用、确认金额与付款时间 按稳定 ID 去重
退款 退款 UUID、支付/订单引用、确认金额与确认时间 按稳定 ID 去重

这些投影没有地址、手机号、昵称或令牌。统计需要顾客标识来表达事实关系,不需要复制顾客个人资料。

源码中的公开订单快照契约:

/** 统计专用的无个人资料订单快照,重建通过业务 API 读取而非访问订单表. */
public interface OrderFacts {
  /** 按数据库 UUID 顺序进行键集分页,每批最多二百条;跨批一致性由调用方只读事务保证. */
  List<Snapshot> after(UUID cursor, int limit);

  /** 订单的版本化统计事实,不含收货信息、令牌或客户端价格. */
  record Snapshot(
      UUID id,
      UUID customerId,
      long version,
      String status,
      BigDecimal total,
      Instant createdAt,
      Instant paidAt,
      Instant completedAt,
      Instant cancelledAt,
      UUID paymentId,
      String refundStatus,
      UUID refundId,
      List<Line> lines) {
    /** 固定成交明细快照. */
    public Snapshot {
      lines = List.copyOf(lines);
    }
  }

  /** 销量统计只需要商品标识、历史名称、数量与成交单价. */
  record Line(
      UUID id, UUID productId, String name, String kind, int quantity, BigDecimal unitPrice) {}
}

应用层再将它转换成 reporting 自己的 ReportingFacts.Order,避免领域层依赖其他模块的契约类型。

这是面向查询组织数据的实践,不等于必须部署一个独立“读微服务”。当前投影仍在同一个模块化单体里。

16. 幂等统计不能每收到一次事件就加一次营业额

16.1 订单更新使用完整快照和源版本

如果重复收到一次完成订单事件,简单执行 turnover += total 会重复统计。

P6 保存订单投影,再依据投影做聚合。重复与旧版本快照不更新当前行:

var existing = orders.findById(value.id());
if (existing.isPresent()) {
  var entity = existing.orElseThrow();
  if (!value.newerThan(entity.sourceVersion)) {
    return false;
  }
  if (entity.total.compareTo(value.total()) != 0
      || !entity.customerId.equals(value.customerId())
      || !entity.createdAt.equals(value.createdAt())) {
    conflict();
  }
  entity.apply(value);
}

这是仓储方法的已有行分支,省略新订单投影的建立过程。

newerThan() 使用严格大于。版本 5 已经落库后,版本 4 的迟到事件不会把状态改回去;版本 5 重放也不会增加统计修订号。

这个策略依赖源模块提供完整、可信的版本化快照。它不是用事件时间戳简单选“最近一条”。

16.2 已确认资金事实按稳定 ID 去重

收款事实的实际保存代码:

public boolean receipt(ReportingFacts.Receipt value) {
  var existing = receipts.findById(value.id());
  if (existing.isPresent()) {
    var entity = existing.orElseThrow();
    if (!new ReportingFacts.Receipt(entity.id, entity.orderId, entity.amount, entity.paidAt)
        .equals(value)) {
      conflict();
    }
    return false;
  }
  var entity = new ReceiptFactEntity();
  entity.id = value.id();
  entity.orderId = value.orderId();
  entity.amount = value.amount();
  entity.paidAt = value.paidAt();
  entity.businessDate = BusinessTime.date(value.paidAt());
  receipts.save(entity);
  return true;
}

同一支付 ID 再到达时,需要内容一致,才可以作为重复事实忽略;相同标识却携带不同金额或时间,会触发冲突。

退款允许先于收款事件到达,因为异步消费者的执行顺序不能被想当然地保证。投影暂时保留退款事实,后续由对账查询指出缺失关系,而不是把事实直接丢掉。

16.3 只有真正应用新事实,才推进修订号

private void apply(BooleanSupplier action) {
  var state = repository.lock();
  if (action.getAsBoolean()) {
    state.applied(clock.instant());
    repository.saveState(state);
  }
}

正常事件更新先取得投影控制行锁。这个协调点稍后也会用于原子重建。

revision 的变化意味着接受了新的投影事实,不是“收到了一次 HTTP 查询”或“又重投了一次重复事件”。

17. 时间精度和经营时区会影响幂等与统计结果

17.1 UTC 时刻与经营日期是两个概念

数据库和事件使用 UTC Instant,经营日期统一按 Asia/Shanghai 转换。

例如:

2026-09-18T15:59:59Z → 上海 2026-09-18 23:59:59
2026-09-18T16:00:00Z → 上海 2026-09-19 00:00:00

如果直接按 UTC 日期截取,第二笔事实会被放错经营日。

项目不跟随运行机器默认时区改变统计口径:

/** 单店经营日期统一按上海时区,源事件与数据库快照统一至 PostgreSQL 微秒精度. */
public final class BusinessTime {
  /** 当前单店经营时区,不跟随服务器默认时区改变口径. */
  public static final ZoneId ZONE = ZoneId.of("Asia/Shanghai");

  private BusinessTime() {}

  /** 对齐 PostgreSQL 时间舍入,避免事件原始纳秒与重建快照产生假冲突. */
  public static Instant canonical(Instant value) {
    return value == null
        ? null
        : Instant.ofEpochSecond(value.getEpochSecond(), ((value.getNano() + 500L) / 1000L) * 1000L);
  }

  /** 将 UTC 事实时间映射为固定经营日期. */
  public static LocalDate date(Instant value) {
    return value == null ? null : canonical(value).atZone(ZONE).toLocalDate();
  }
}

17.2 为什么还有微秒舍入

Java 事件可能保留纳秒,而数据库源快照已经按 PostgreSQL 时间精度保存。

例如原事件带 ...123456789 纳秒部分,数据库重读得到对应的微秒舍入结果。如果拿未经规范化的时间直接比较,不可变事实可能被误判为冲突。

canonical() 在事件和重建路径统一精度,还允许舍入进位到下一秒。这里使用的是舍入规则,不是简单把末尾三位截掉。

规范化的目的不是人为调整业务发生时间,而是让同一存储事实经过不同传递路径后具有一致表示。

17.3 查询日期范围也要有上界

from/to 都包含首尾日期,当前最多 366 天,年份限制在 1970—9999。这样可以补齐空日期并生成有界导出,而不接受任意大区间把整个历史搬进内存。

18. 营业额、完成率与净收款,不能共用一套日期口径

先定义集合,再写公式,会比先写一个 sum(total) 更稳妥。

对经营日 d:

  • S_d:在 d 创建的订单。
  • C_d:当前已完成,并且 completedAt 属于 d 的订单。
  • K_d:S_d 中当前状态已完成的订单。
  • R_d:付款时间属于 d 的已确认收款。
  • F_d:退款确认时间属于 d 的已确认退款。

当前口径为:

Turnover(d)=\sum_{o\in C_d}total(o)
CompletionRate(d)= \begin{cases} 0.00,&|S_d|=0\\ 100\times\dfrac{|K_d|}{|S_d|},&|S_d|>0 \end{cases}
NetReceived(d)=\sum_{r\in R_d}amount(r)-\sum_{f\in F_d}amount(f)

18.1 用跨日例子理解差异

以下都是教学数据,时刻按上海经营时间表示,并假定查询时 C 的退款已确认:

订单 创建 付款 完成/退款确认 金额
A 9月18日 23:55 9月18日 23:56 9月19日 00:10 完成 20.00
B 9月19日 09:00 9月19日 09:01 9月19日 09:30 完成 30.00
C 9月19日 10:00 9月19日 10:01 9月20日 09:10 全额退款确认 10.00

按完成日与资金发生日分别统计:

经营日 新建订单数 当日完成数 营业额 已确认收款 已确认退款 净收款
9月18日 1 0 0.00 20.00 0.00 20.00
9月19日 2 2 50.00 40.00 0.00 40.00
9月20日 0 0 0.00 0.00 10.00 -10.00

9月19日创建的群组是 B、C,其中只有 B 当前已完成,所以完成率为 50.00%。不能使用“当日完成 2 / 当日新建 2”得出 100%。

同样,9月20日净收款为负,并不说明代码计算错误;这一天确认退款,却没有新的收款。

18.2 群组状态不是历史截面状态

completedCohort、cancelledCohort 表示“这些创建日期的订单,在当前投影快照中是什么状态”。

未来订单继续完成或取消,同一历史创建群组的指标会变化。它不是“截至那个历史日期晚上,系统当时看到的状态”回放。

平均订单金额则使用完成日口径的营业额除以完成订单数,零分母返回 0.00。完成率与均价均保留两位小数。

18.3 顾客增长也要说明含义

newCustomers 按注册日统计;累计顾客包括查询开始日期之前的注册事实,也包含后来停用的账号。

它不是日活、付款人数、留存率,也不应被图表标题包装成这些尚未实现的指标。

19. 用数据库做分组聚合,再在 Java 中补齐空日期

聚合查询只操作 reporting 自己的表。完成日营业额的实际查询如下:

private List<ReportData.CompletionDay> completionDays(ReportPeriod period) {
  var cb = entityManager.getCriteriaBuilder();
  var query = cb.createTupleQuery();
  var root = query.from(OrderFactEntity.class);
  var date = root.<LocalDate>get("completedDate");
  query
      .multiselect(date, cb.count(root), cb.sum(root.<BigDecimal>get("total")))
      .where(ReportingSpecifications.completed(period).toPredicate(root, query, cb))
      .groupBy(date)
      .orderBy(cb.asc(date));
  return entityManager.createQuery(query).getResultList().stream()
      .map(
          row ->
              new ReportData.CompletionDay(
                  row.get(0, LocalDate.class), number(row, 1), row.get(2, BigDecimal.class)))
      .toList();
}

其中 ReportingSpecifications.completed(period) 组合当前 COMPLETED 状态与 completedDate 日期条件。分组、计数和金额求和由数据库完成。

查询返回的是有界日汇总,不是把所有历史订单加载到 Java 后再循环相加。

随后领域模型才补齐没有数据的日期:

public static List<Day> days(ReportPeriod period, Aggregates data) {
  var orders = index(data.orders(), OrderDay::date);
  var completions = index(data.completions(), CompletionDay::date);
  var receipts = index(data.receipts(), CashDay::date);
  var refunds = index(data.refunds(), CashDay::date);
  var customers = index(data.customers(), CustomerDay::date);
  var result = new java.util.ArrayList<Day>();
  long cumulative = data.customersBefore();
  for (int offset = 0; offset < period.days(); offset++) {
    var date = period.from().plusDays(offset);
    var order = orders.getOrDefault(date, new OrderDay(date, 0, 0, 0));
    var completed = completions.getOrDefault(date, new CompletionDay(date, 0, zero()));
    var incoming = receipts.getOrDefault(date, new CashDay(date, 0, zero()));
    var outgoing = refunds.getOrDefault(date, new CashDay(date, 0, zero()));
    long added = customers.getOrDefault(date, new CustomerDay(date, 0)).added();
    cumulative = Math.addExact(cumulative, added);
    result.add(
        new Day(
            date,
            order.submitted(),
            order.completedCohort(),
            order.cancelledCohort(),
            completed.completed(),
            completed.turnover(),
            incoming.amount(),
            outgoing.amount(),
            added,
            cumulative));
  }
  return List.copyOf(result);
}

这里的循环次数与日期数有关,最多 366;它不会为每一个日期再发送一次数据库查询。

累计顾客从 customersBefore 开始,而不是每次查询都把起日当成系统注册的第一天。

一日和一年查询次数一致,说明消除了逐日往返,不代表一年数据的数据库扫描成本与一天完全相同。数据量、索引与执行计划仍会影响实际耗时。

20. 商品改名以后,销量应该分成两种商品吗

P6 按商品 UUID 合并销量,不按名称分组识别业务商品。

同一个商品先叫“招牌面”,后来叫“招牌牛肉面”,不能因为显示名称变化就分裂成两个销量对象。

实际销量查询:

public List<ReportData.Sale> sales(ReportPeriod period, int limit) {
  var cb = entityManager.getCriteriaBuilder();
  var query = cb.createTupleQuery();
  var root = query.from(LineFactEntity.class);
  var product = root.join("product");
  var quantity = cb.sumAsLong(root.<Integer>get("quantity"));
  var amount = cb.sum(root.<BigDecimal>get("subtotal"));
  query
      .multiselect(root.get("productId"), root.get("kind"), product.get("name"), quantity, amount)
      .where(ReportingSpecifications.completedSales(period).toPredicate(root, query, cb))
      .groupBy(root.get("productId"), root.get("kind"), product.get("name"))
      .orderBy(cb.desc(quantity), cb.desc(amount), cb.asc(root.get("productId")));
  return entityManager.createQuery(query).setMaxResults(limit).getResultList().stream()
      .map(
          row ->
              new ReportData.Sale(
                  row.get(0, UUID.class),
                  row.get(1, String.class),
                  row.get(2, String.class),
                  number(row, 3),
                  row.get(4, BigDecimal.class)))
      .toList();
}

这里的 product 关系指向 reporting 自己的商品摘要表,没有跨到 catalog 业务表。

几个具体口径:

  • 只统计当前 COMPLETED 订单,并按订单完成日期筛选。
  • 数量来自成交明细,金额使用历史 subtotal,不读取当前目录价。
  • 名称使用投影中最近订单快照观察到的商品名称,不保证只取当前查询区间内已完成订单的名称。
  • 排序按数量降序、金额降序、UUID 升序。
  • 套餐作为套餐商品计件,不再重复累加其组成菜品。

商品从目录删除后,历史成交明细和统计投影仍然存在。这个能力来自 P4 的历史快照,也说明前几阶段模型边界会影响后续统计是否可靠。

21. 资金对账为什么要双向检查

只从收款表出发检查“有没有订单”,发现不了“订单已经显示付款,但收款事实漏投影了”。

P6 保留两个方向:

检查 范围 提示的关系问题
missingOrders 日期区间内收款 收款没有对应订单投影
paymentMismatches 日期区间内收款 金额或支付引用与订单不匹配
missingPayments 日期区间内退款 退款缺少原收款事实
refundMismatches 日期区间内退款 与原收款的订单/金额不匹配
ordersMissingReceipts 全店当前 已付款订单没有匹配收款事实
ordersMissingRefunds 全店当前 已退款订单没有匹配退款事实
pendingRefunds 全店当前 订单仍处在待退款状态

后面三个没有按日期区间截断,所以不能把所有结果都标成“本月差异”。

反向检查“已付款订单缺少收款”的实际查询:

private long ordersMissingReceipts() {
  var cb = entityManager.getCriteriaBuilder();
  var query = cb.createQuery(Long.class);
  var order = query.from(OrderFactEntity.class);
  var source = query.subquery(UUID.class);
  var receipt = source.from(ReceiptFactEntity.class);
  source
      .select(receipt.get("id"))
      .where(
          cb.equal(receipt.get("id"), order.get("paymentId")),
          cb.equal(receipt.get("orderId"), order.get("id")));
  query
      .select(cb.count(order))
      .where(cb.isNotNull(order.get("paidAt")), cb.not(cb.exists(source)));
  return entityManager.createQuery(query).getSingleResult();
}

这仍是本模块内部投影之间的关系检查,不是直接读取支付宝账单或在生产资金账户上对账。

异步滞后也可能暂时形成差异。例如退款事件先到、收款事件尚未应用,此时 missingPayments 会暂时大于零。

差异指标也可能重叠,不能把它们简单求和当作互不重复的异常订单总数。报表只指出问题,不修改订单、补造付款或把退款受理当作已支出。

22. 同一张报表里的数据,也需要来自同一个快照

如果先查询总额,随后恰好发生投影重建,再查询日账,用户可能看到两个不同代际的数据。

P6 把报表组成放进同一个只读一致性事务:

@Service
@Transactional(readOnly = true, isolation = Isolation.REPEATABLE_READ)
public class ReportQueries {
  // 依赖、构造器及其他方法省略。
}

实际加载方法如下:

public ReportWorkbook load(StaffIdentity actor, ReportPeriod period, int salesLimit) {
  authorization.requireAdministrator(actor);
  if (salesLimit < 1 || salesLimit > 100) {
    throw new ReportingException(ReportingException.Reason.INVALID_INPUT, "销量条数须为 1 至 100");
  }
  var state = repository.state();
  state.requireReady();
  var days = ReportData.days(period, repository.aggregate(period));
  return new ReportWorkbook(
      period,
      metadata(state),
      ReportData.summary(days),
      days,
      repository.sales(period, salesLimit),
      repository.reconcile(period));
}

先检查管理员权限和投影是否就绪,再加载日账、汇总、销量和对账,附上控制版本、generation、revision、更新时间与重建时间。

PostgreSQL 的 Repeatable Read 让事务内连续查询使用同一个一致性快照;它与默认 Read Committed 的逐语句快照不同。PostgreSQL Repeatable Read

这保证单次报表不会混合重建前后的投影数据,不保证异步投影已经追到所有最新源业务事件。updatedAt 也不是一个证明“全部事件已经处理完”的全局水位。

JSON 和 XLSX 复用同一加载方式与统计规则。两次独立请求之间若发生新事件,结果仍可能变化,应结合元数据解释,不能承诺它们在不同时刻永远逐字相同。

工作台与财务报表的权限不同

普通 STAFF 和 ADMIN 可以查看工作台待办;财务指标、导出和投影维护只允许当前有效管理员。

工作台的 PAID、ACCEPTED、DELIVERING、CANCELLING、REFUNDING 待办统计覆盖所有日期,避免昨天的未完成订单今天消失。另行提供当前经营日的创建和完成数量。

门店状态和商品上下架数量来自 shop/catalog 的公开实时摘要 API。它们是当前业务摘要,不应伪装成财务历史,也不把展示缓存当作业务事实。

23. 投影可以重建,前提是源事实有公开读取契约

P6 之前已经存在订单、顾客和支付数据,但过去并没有向 P6 发布全部事件。

因此,首次初始化不能只等新事件。它需要通过各模块公开 API 补齐已有事实。

23.1 重建是对投影的维护,不是重新执行历史业务

ProjectionMaintenance 的事务定义:

@Service
@Transactional(isolation = Isolation.REPEATABLE_READ, timeout = 180)
public class ProjectionMaintenance {
  // 依赖、构造器与公开入口省略。
}

公开入口先锁投影控制行,管理员重建还要检查当前控制版本。内部重建方法如下:

private void rebuild(ProjectionState state) {
  repository.clear();
  copy(
      cursor -> orders.after(cursor, 200),
      OrderFacts.Snapshot::id,
      value -> repository.order(FactMapper.order(value)));
  copy(
      cursor -> customers.after(cursor, 200),
      CustomerFacts.Registration::id,
      value -> repository.customer(new ReportingFacts.Customer(value.id(), value.createdAt())));
  copy(
      cursor -> payments.receiptsAfter(cursor, 200),
      PaymentFacts.Receipt::id,
      value ->
          repository.receipt(
              new ReportingFacts.Receipt(
                  value.id(), value.orderId(), value.amount(), value.paidAt())));
  copy(
      cursor -> payments.refundsAfter(cursor, 200),
      PaymentFacts.Refund::id,
      value ->
          repository.refund(
              new ReportingFacts.Refund(
                  value.id(),
                  value.paymentId(),
                  value.orderId(),
                  value.amount(),
                  value.confirmedAt())));
  state.rebuilt(clock.instant());
  repository.saveState(state);
}

它读取订单、注册、已确认收款和已确认退款快照,不回放原来的“下单命令”“付款命令”或“通知发送命令”。

所以重建不会再次扣款、再次创建订单,也不会为所有历史已付款订单批量制造新的来单提醒。

23.2 使用键集分页,控制持久化上下文大小

private <T> void copy(Function<UUID, List<T>> reader, Function<T, UUID> id, Consumer<T> writer) {
  UUID cursor = null;
  while (true) {
    var batch = reader.apply(cursor);
    for (T value : batch) {
      writer.accept(value);
    }
    repository.finishBatch();
    if (batch.size() < 200) {
      return;
    }
    cursor = id.apply(batch.getLast());
  }
}

每批最多二百条,游标取本批最后一个 UUID。比较和排序在源数据库中保持一致,不把随机 UUID 当作业务时间。

每批结束执行:

entityManager.flush();
entityManager.clear();

flush 将待写数据发送给数据库,clear 释放持久化上下文里已管理对象的引用。二者都不是提交,也不会释放当前事务持有的数据库锁。

这里批次限制的是单批数据与 ORM 管理对象规模,不表示整个重建只复制二百条,或整个重建事务只持续一个批次。

23.3 为什么当前跨模块读取可以保持同一数据库快照

公开 API 是同一个应用内的 Java 调用,参与相同事务管理器和当前 Repeatable Read 事务。因此不同源模块的分页读取可以使用同一数据库快照。

如果未来将这些接口换成远程微服务 HTTP,这个前提就不再成立。需要重新设计源版本、水位、快照切点或增量追赶协议,不能只保留 @Transactional 就宣称有跨服务一致性。

24. 重建期间,读者与新事件分别会发生什么

24.1 删除旧投影和建立新投影一起提交

当前方案在同一个事务内清理旧投影、分批建立新投影,最后才执行 state.rebuilt() 并保存控制状态。

它没有提前提交一次“清空”,也没有把每个批次独立提交。

普通只读查询通过 MVCC 继续看到之前已经提交的完整数据。重建中途失败,删除、新投影和新控制状态全部回滚,旧代际保留。

generation 表示成功重建的代际;revision 表示接受新事实或完成重建的修订进度;version 是控制行的 ORM 并发版本。它们不是三个可以随意互换的游标。

24.2 新源业务可以继续提交,投影更新等待控制锁

sequenceDiagram
  participant R as 重建事务
  participant P as reporting投影
  participant O as 订单源业务
  participant E as 可靠事件消费者
  participant Q as 报表读者
  R->>P: 锁控制行,建立一致性快照
  R->>P: 删除并分批写入投影,尚未提交
  Q->>P: 读取已提交旧代际
  P-->>Q: 旧的完整结果
  O->>O: 提交新的订单变化及事件登记
  E->>P: 应用事件,等待控制行锁
  R->>P: 保存新代际并提交
  E->>P: 获得锁,按源版本补上新变化

新业务事实没有因为重建快照看不到它,就永久丢失。源业务已经同事务登记事件,重建后消费者继续应用。

如果事件中的版本已经被重建快照覆盖,就作为重复或旧事实忽略;如果它更新,就推进投影。

24.3 首次初始化失败不能显示一张假的“全零报表”

public void requireReady() {
  if (!initialized) {
    throw new ReportingException(ReportingException.Reason.NOT_READY, "统计投影尚未初始化,请重建后查询");
  }
}

初始化成功前返回 503 REPORTING_NOT_READY。管理员仍可查询投影状态、按版本重建;启动初始化失败也会保留未就绪状态并在后续周期重试。

只有 initialized=true 的空投影,才表示已经建立基线但没有相应业务数据。

控制行竞争、事务超时或数据库隔离冲突也可能让重建失败。失败时应重读状态并重新执行整个维护用例,而不是从已经回滚事务的某个分页中间位置继续。Repeatable Read 不等同于任意竞争下绝不会失败。PostgreSQL 事务隔离与重试

当前使用一个控制协调点和最长 180 秒的重建事务,适合现阶段单店场景。它不是已经实现双表在线切换、无限数据量导入或零成本重建的大型数仓方案。

25. XLSX 导出:金额保真和公式隔离比“能下载”更重要

P6 使用 Apache POI 5.5.1,输出四张表:口径与汇总、经营日账、商品销量、资金对账。

25.1 先结束查询事务,再生成压缩文件

应用入口:

public byte[] export(StaffIdentity actor, LocalDate from, LocalDate to) {
  return exporter.export(queries.load(actor, new ReportPeriod(from, to), 100));
}

queries.load() 通过另一个 Bean 的代理完成快照查询并结束事务,随后 exporter 生成文件。导出适配器还声明 Propagation.NEVER,防止把压缩过程意外放进数据库事务。

日期最多 366 行,销量最多 100 行。它不是把全历史明细无限写入一个巨大的工作簿。

25.2 金额写成十进制文本,是一种明确取舍

单元格写入方法:

private void row(
    org.apache.poi.ss.usermodel.Sheet sheet, int index, CellStyle style, Object... values) {
  var row = sheet.createRow(index);
  for (int column = 0; column < values.length; column++) {
    var cell = row.createCell(column);
    Object value = values[column];
    if (value instanceof BigDecimal amount) {
      cell.setCellValue(amount.toPlainString());
    } else if (value instanceof Number number) {
      cell.setCellValue(number.doubleValue());
    } else {
      cell.setCellValue(String.valueOf(value));
    }
    if (style != null) {
      cell.setCellStyle(style);
    }
  }
}

BigDecimal 使用 toPlainString() 作为字符串写入,百分比也按精确十进制文本保存;普通数量等 Number 仍按数值写入。

好处是金额不会在写入 Excel 数值时悄悄经过二进制浮点转换。代价是这些金额单元格不是可以直接当作原生数值任意计算的格式,汇总已由服务端计算,用户再做公式分析时需要理解文本类型。

不能同时宣传“这里金额作为文本保证表示精确”,又承诺它与所有 Excel 数值公式完全无差别。

25.3 像公式的商品名称仍然是名称

如果商品名称是:

=HYPERLINK("https://invalid.example")

代码仍然调用字符串版本的 setCellValue,没有调用 setCellFormula。导出测试使用真实 POI 重新读取并检查它是普通字符串。

字符串单元格与公式单元格的 API 分别承担不同职责,可参考 Apache POI XSSF 指南。这个做法针对当前 XLSX 导出,不能直接推导所有 CSV 或其他表格格式都自动安全。

文件包含冻结表头、筛选和口径/投影元数据。财务导出要求管理员,响应使用标准 XLSX MIME、附件文件名与 Cache-Control: no-store

26. 从接口实操到测试,再留下演进问题

26.1 一组可以对应到代码的操作顺序

先使用有效员工会话,读取个人进度和消息:

GET /api/v1/notifications/receipt
Authorization: Bearer <员工令牌>
GET /api/v1/notifications?after=0&limit=50
Authorization: Bearer <员工令牌>

处理完整页面后,以实际游标和阅读进度版本确认:

PUT /api/v1/notifications/receipt
Authorization: Bearer <员工令牌>
Content-Type: application/json

{"sequence": 12, "version": 1}

这些数字需要替换。它们分别是消息提交游标和个人进度版本,不是订单版本或通知投递版本。

申请 WebSocket 票据:

POST /api/v1/notifications/stream-tickets
Authorization: Bearer <员工令牌>

票据响应带 Cache-Control: no-store。在有效期内连接通知 stream,观察 READY,再按第 14 节补查。断线重连时重新申请票据,不复用已经消费的那一份。

管理员可以查询报表:

GET /api/v1/reports/operations?from=2026-09-18&to=2026-09-20
Authorization: Bearer <管理员令牌>

需要重建时,先读取 /api/v1/reports/projection 的控制版本:

POST /api/v1/reports/projection/rebuild
Authorization: Bearer <管理员令牌>
Content-Type: application/json

{"version": 7}

同样不要直接照抄示例版本。重建是受控维护操作,不应在普通页面每次刷新时自动执行。

26.2 其余 HTTP 能力与权限

资源 主要权限与用途
POST /orders/{id}/reminders 顾客本人,按订单版本催单
GET /notifications/{id}/attempts 管理员查看最近最多 100 条投递轨迹
POST /notifications/{id}/redelivery 管理员按通知版本重投已耗尽消息
GET /workspace 有效 ADMIN/STAFF,工作台待办与当前门店/目录摘要
GET /reports/sales 管理员,完成订单销量,limit 1—100
GET /reports/reconciliation 管理员,区间资金事实与全店反向检查
GET /reports/export 管理员,有界 XLSX 导出

表内路径均以前缀 /api/v1 开始。JSON 错误沿用 RFC 9457、code 和 traceId;握手失败发生在 HTTP 升级之前,建立后的连接撤销使用 WebSocket 关闭语义。

26.3 有效测试应当复现那些“不按理想顺序发生”的场景

提交顺序游标测试会让第一个消息事务暂不提交,第二个追加同时开始,再检查客户端没有读到能越过第一条消息的更大游标。

重建失败测试则验证旧数据和旧代际都保留:

void failedRebuildRollsBackDeletionAndKeepsOldGenerationVisible() {
  metrics();
  var before = reports.operations(admin, FROM, TO);
  long generation = queries.projection(admin).generation();
  doThrow(new IllegalStateException("模拟源快照失败"))
      .when(AopTestUtils.<CustomerFacts>getUltimateTargetObject(sourceCustomers))
      .after(any(), anyInt());
  assertThatThrownBy(() -> maintenance.rebuild(admin, queries.projection(admin).version()))
      .isInstanceOf(IllegalStateException.class);
  assertThat(queries.projection(admin).generation()).isEqualTo(generation);
  assertThat(reports.operations(admin, FROM, TO).summary()).isEqualTo(before.summary());
}

P6 还用 205 条源订单验证超过一个二百条批次的重建,用并发读者验证旧完整代际可见,并在重建期间提交新订单变化,确认事件后来能追平。

关键验收场景包括:

场景 需要证明什么
付款结果重复或取消后迟到 不重复制造来单提示
同订单连续催单 归属、版本、状态与冷却规则生效
消息投递失败、耗尽和重投 有轨迹,历史仍可补查,不自动确认已读
旧投递尝试迟到 不覆盖新领取状态
票据过期、重放和不允许的 Origin 握手被拒绝
退出或停用已经连接的员工 原会话复验会撤销连接
跨经营日付款、完成、退款 各指标按自己的时间口径计算
重复或旧版本统计事件 不重复计数,不回退投影
首次启动已有 P1—P5 数据 无需旧事件也能补齐报表基线
重建中途失败 删除和重建一起回滚
一日与一年报表 不逐日查询,不加载全部订单后计算
类公式商品名称导出 仍是普通字符串,金额与同一统计规则一致

P6 提交验收记录为全量 144 项测试通过,零失败、零错误、零跳过。新增 25 项验证包括 8 项领域测试、16 项通知/报表/真实 WebSocket 集成测试和 1 项首次启动已有数据初始化测试。

阶段记录中的一日与一年报表查询次数一致且不超过 18;这是一组既定测试场景的结果,不是任何数据量和业务扩展下的永久查询次数保证。

本次文章编写只核对提交中的代码和验收记录,没有重新执行应用测试或验证 Vditor 渲染。

26.4 P6 对 P5 的一个工程调整

可靠事件恢复已从支付轮询配置中移到应用级 EventRecoveryConfiguration。关闭支付轮询,不会顺带关闭通知和统计消费者的事件恢复。

配置开关也分开:

环境变量 当前默认 职责
EVENT_RECOVERY_ENABLED true 应用级未完成事件恢复
NOTIFICATION_SCHEDULING_ENABLED true 实时投递和空闲连接复验
REPORTING_BOOTSTRAP_ENABLED true 首次基线补齐与失败后的周期重试
NOTIFICATION_ALLOWED_ORIGINS 空配置 默认同源;需要时声明明确允许来源

这延续了 P5 的可靠事件语义,但配置归属随真实消费者增加而调整。前一篇文章描述的是 P5 提交,不能把它的配置位置当作后续永远不变的结构。

26.5 当前边界与读者练习

当前 WebSocket 连接注册表只存在于一个进程内。数据库消息持久化与重投,并不自动提供跨实例广播;多实例部署时,需要另外设计连接归属、广播或每实例投递。

当前没有 Flutter 或新管理端界面,也没有短信、邮件、外部消息中间件、完整数仓、永久事件归档或任意消息逐条未读模型。

练习一:游标与时间。 如果改用 occurredAt 作为补查游标,同一时刻的多条事件和迟到事件应当怎样处理?请列出仅使用一个时间字段可能遗漏的情况。

练习二:已投递与已阅读。 一名员工有三台设备,只向其中两台发送成功,投递和个人阅读状态应该如何变化?当前 Outcome 和 Receipt 分别保护了什么?

练习三:票据不是永久会话。 员工签发票据后立刻退出,票据仍未过期。为什么握手还必须查询原会话?连接建立后又为什么需要周期复验?

练习四:完成率的分母。 今天完成十单,但只有两单是今天创建的。请设计同创建群组完成率与当日完成数量两个指标,说明它们分别回答什么问题。

练习五:重建改成逐批提交。 这样会如何影响读者、失败回滚和事件追平?如果需要支持更大的数据量,如何引入独立代际存储和切换,而不是直接丢掉当前原子性?

练习六:新增一个图表。 先写出业务集合、时间字段、金额来源、去重依据和权限,再决定是否需要新的投影字段。不要从 Controller 里临时跨模块查几张表就开始画图。

P7 可以在这些后端契约之上实现真实客户端恢复、工作台展示与部署验证。界面的责任,是把已经定义清楚的业务事实和处理状态准确呈现出来。

延伸阅读

系列:Han Menu 外卖系统实践 · 从第一篇开始

上一篇:第6篇 · 下一篇:第8篇