企业应用架构16:拆单路由与交付结果聚合
一份租赁合同包含相机与投影仪,两种设备由不同端点处理。相机完成消息到了两次,投影仪的结果一直没有到。如果聚合器只计算收到了几条消息,合同就可能在缺少一件设备时被标为全部完成。
本篇用真实 ActiveMQ 通道传递拆分任务和完成结果,H2 保存批次与明细。完成条件按合同、交付批次和设备任务判断,实验依次观察乱序、重复、缺项、超时及上一批次的迟到结果。
拆分之前固定完成条件
Splitter 把复合消息拆成可以独立处理的部分;Content-Based Router 根据内容选择通道;Aggregator 再按关联标识与完成条件收集结果。这里需要先定义“全部”,再决定如何拆分,否则最后一条到达的消息没有资格宣布整份合同完成。Splitter;Content-Based Router
实验在批次创建时保存预期任务数。合同 R1 的 B1 批次包含相机 E1 和投影仪 P1,预期为两项;合同 R2 的 B1 批次只有相机 E2,预期为一项。批次号可以在不同合同内重复,因此数据库主键由合同号和批次号共同组成。
每项任务也有独立标识,例如 R1/B1/E1。它既出现在持久化明细表,也作为业务消息编号发送。这个样本假定同一批次内一台设备只有一项交付任务;若一台设备需要分阶段交付,就应另外引入任务序号,而不能继续使用设备编号代替任务身份。
flowchart LR
C[合同R1 批次B1 两项设备] --> S[Splitter 保存预期明细]
S -->|camera E1| Q1[delivery.camera]
S -->|projector P1| Q2[delivery.projector]
Q1 --> W1[相机端点]
Q2 --> W2[投影端点]
W1 --> R[delivery.results]
W2 --> R
R --> A[Aggregator 合同加批次]
A --> D[H2明细与批次状态]
路由规则位于拆分发送处,设备类型直接决定通道名称。相机端点实际收到 E1、E2,投影端点收到 P1,并逐项断言。该实现没有增加独立路由服务器;模式中的职责可以由一个小函数实现,部署数量不是判断模式是否成立的依据。
明细集合比消息计数可靠
聚合器收到完成消息时,先查合同和批次,再用任务编号、合同号、批次号、设备号共同查找已登记的明细。只有任务存在且当前批次仍在等待,才允许记录完成。未知任务不能通过伪造一个较大的计数提前完成批次。
明细的 done 从假变为真后,聚合器重新统计该批次已完成的明细,并与创建时保存的 expected 比较。重复完成消息只命中已完成明细,不再改变集合大小。完成条件实际是预期集合的覆盖关系,计数只是建立在身份校验之上的实现方式。
Aggregator 的关联规则、完成条件和聚合算法是三个需要分别决定的部分。本实验的关联是合同加批次,完成条件是全部已登记任务完成,聚合结果只是一个批次状态,并没有合并设备的详细物流轨迹。Aggregator
固定预期集合还有一个实际后果:执行中增加一台设备,不能只新增任务而不更新批次定义。修改预期数、已经收到的结果以及截止时间之间需要一致的规则。样本通过新建批次表达重新发起,没有测试运行中修改旧批次,因此不存在“动态加项也已验证”的结论。
用交错结果检查合同隔离
工作端点先取出三条真实任务消息,再按 R2 的 E2、R1 的 E1、R1 的 E1 重复件发送结果。接收顺序刻意不同于原任务发送顺序。结果队列中确实存在三条消息,消费者逐条接收并记录业务编号,没有在本地数组上直接调用聚合函数替代消息交付。
R2 先完成,并不会帮助 R1 完成。第一阶段读回显示 R2 的状态为 COMPLETE,R1 仍为 WAITING,R1 完成明细只有一项。随后到来的 E1 重复件输出 DEDUP,完成数仍是一。缺少 P1 时,收到三条结果这个总量没有业务意义。
投影任务已经被端点接收,但端点主动不发送完成结果。这个受控故障代表结果缺项,不能据此断定实物一定没有交付。聚合状态只能表明证据尚未齐全,真实设备位置需要向交付系统查询,或者交给后续人工处理流程。
独立负例在这一步断言 R1 已有两项完成,实际只有一项,因此进程以一退出。它检查的是错误验收条件会被测试阻止,正常场景则另外验证 WAITING 状态。负例不会把业务等待本身当成程序异常,也不会覆盖正常运行证据。
超时关闭一批结果的接受窗口
批次创建时记录截止毫秒值。第一阶段聚合完成后,运行器等待一千一百毫秒,再执行带有 deadline 条件的 SQL 更新。只有状态仍为 WAITING 且期限已到的 R1/B1 才转为 TIMED_OUT,受影响行数必须为一。
这次实验使用真实时钟与数据库条件,没有把睡眠本身当作超时成功。即使进程调度比预期更慢,最终判断仍来自持久化截止时间。它不测试系统时钟倒退、多主机时差或分布式计时,生产系统需要另外规定期限来源和时钟异常策略。
stateDiagram-v2
[*] --> WAITING
WAITING --> WAITING: 合法新明细或重复明细
WAITING --> COMPLETE: 所有预期任务完成
WAITING --> TIMED_OUT: 期限已到且仍缺项
COMPLETE --> COMPLETE: 迟到结果隔离
TIMED_OUT --> TIMED_OUT: 旧批次结果隔离
超时以后创建 R1/B2,新批次只预期 E1 一项。接着发送旧批次的 P1 结果和新批次的 E1 结果。两个报文使用不同批次标识,即使属于同一合同,接收方也不会把旧 P1 填入新的预期集合。
最终 SQL 读回中,R1/B1 保持 TIMED_OUT,R1/B2 为 COMPLETE。旧 P1 被发送到 delivery.late 并实际从隔离通道读出。迟到结果仍可供审计和人工核对,但不能自动改写已关闭批次的决定。
终态、重复与后续核对
当前实现将非等待批次的结果统一隔离,因此已经完成的批次再收到重复件,也会进入迟到通道。这个选择保持终态稳定,却意味着隔离量中同时包含无害重复和需要人工核对的旧结果。若用于生产,应在诊断记录中保留更细的原因分类。
缺项时可以选择一直等待、部分完成、到期失败或人工审核,每一种都对应不同业务承诺。样本选择到期关闭,并不意味着所有设备租赁都适用相同规则。允许部分交付的合同还要保存已经交付的设备和应付金额,不能只新增一个 PARTIAL 名字而省略后续义务。
明细状态与批次状态在同一 H2 事务更新,随后确认 JMS 消息。二者没有使用 XA,因此数据库提交和消息确认之间仍存在重复窗口。按任务身份记录 done 可以抵御本场景的重复结果,但并未覆盖并发聚合器、冲突载荷或多个进程同时修改同一批次的竞态。
拆分阶段同样先写数据库,再提交消息会话。当前场景没有在这两个提交之间终止进程。若要求批次创建以后每个任务都最终发出,应沿用 Outbox 的持久化发布安排,而不能因为 broker 已使用持久消息就忽略发送前的窗口。
从实验包检查实际结果
批次历史不应在重试时删除。R1/B1 的超时记录解释了为何创建 B2,也记录了哪些设备结果当时已经到达。若将旧批次清空再复用同一个编号,迟到消息就失去了可区分的上下文;即使业务消息编号保持唯一,也无法判断它属于哪一次交付承诺。
截止时间和消息有效期同样属于不同层次。broker 可以因为消息过期而不再交付,但应用仍需要保存“这项结果尚未收到”的事实。只依靠队列有效期清理消息,会让聚合器永久等待已经被删除的结果。本文没有设置 JMS 消息有效期,所有超时决定来自批次表,便于单独观察这个规则。
聚合结果也不是所有后续动作的统一触发器。批次 COMPLETE 可以允许生成交付确认,却不一定意味着应立即扣费;合同可能要求客户签收或核对损伤。将这些步骤都塞进聚合回调,会使重投一次结果时难以确定哪些外部副作用已经执行。更清晰的边界是先持久化聚合事实,再由后续流程决定业务动作。
读日志时应同时检查任务身份和状态更新。只有通道接收记录,无法证明数据库完成条件正确;只有最终 COMPLETE,也无法证明两个端点真的收到了各自设备。实验把这两类记录放在同一证据目录中,仍分别保留来源,便于定位错误发生在路由、身份校验还是完成判定。
下载 实验包,保留 examples 目录结构,用 Java 21 执行入口。它复用第15篇的固定依赖与真实 broker 支持代码,业务表和运行目录各自独立。
1 | |
正常运行退出零,日志依次包含路由接收、交错结果、重复处理、超时条件更新及新旧批次终态。命令末尾添加 negative,会重新建立独立场景,在缺项验收处退出一。数据库 SQL 快照可以直接核对两张表,不必只相信控制台的 PASS 标签。
依赖清单固定 ActiveMQ Classic 5.18.6 和 H2 2.1.214 等坐标与摘要。可写缓存由 ENTERPRISE_INTEGRATION_CACHE 指定,证据输出位置由 ENTERPRISE_EVIDENCE_DIR 指定。运行器记录每个独立 Java 进程的命令和退出码;本篇的聚合验收与系列核心合同测试属于两组证据,不能相互替代。






