从零重构外卖系统(七):可靠通知、鉴权 WebSocket 与经营统计,怎样让消息可补查、报表可重建
系列:Han Menu 外卖系统实践 · 从第一篇开始
系列: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 催单。未支付、取消处理中、退款中、已取消和已完成都不允许。
冷却规则为两次成功催单至少间隔六十秒:
它与前端按钮防抖不同。即使两台设备同时发请求,后端仍然要保护该规则。
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 中同一条控制行:
- 检查稳定事件 ID 是否已经生成通知。
- 在托管控制行上递增 sequence。
- 保存该 sequence 对应的消息。
- 两者共同提交,或共同回滚。
第二个追加事务必须等第一个提交或回滚后再分配。因此客户端不会先越过一个尚未提交的较小消息游标。
这里的 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=0/version=0;第一次真正前移后持久化版本推进为 1,思路与 P3 首次加购相似。
一个标量游标还隐含了产品约定:确认序号 20,表示此前应处理的消息已经形成连续前缀。它不适合直接表达“只读了 20,但 18、19 明确未读”的任意集合。
如果未来需要逐消息未读状态,就要重新设计阅读模型,不能继续用一个最高序号掩盖中间空洞。
11. 浏览器 WebSocket 为什么先换一次性票据
普通浏览器 WebSocket 构造方式不能像 fetch 一样随意设置业务 Authorization 请求头。把长期员工 Bearer 放进 URL 查询参数又容易进入访问日志和复制链接。
P6 采用两步:
- 先通过带员工 Bearer 的普通 HTTP 接口申请票据。
- 用固定业务子协议和一次性票据建立连接。
示意调用如下,它是接口用法示例,不是 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 的已确认退款。
当前口径为:
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 外卖系统实践 · 从第一篇开始