命令处理管线
本页解释一条非 Void 命令如何穿过 Wow 运行时。如何构造和发送命令见发送命令,如何选择等待阶段见完成语义;这里仅讨论实现顺序和边界。
组件地图
CommandBus 只负责投递和接收信封;CommandDispatcher 按具名聚合建立处理器,并把同一聚合 ID 映射到稳定的调度组;DefaultCommandHandler 按固定顺序执行命令管道;聚合执行、事件持久化、transport ack、领域事件发布与状态事件发布是不同步骤。
发送前管道
DefaultCommandGateway 是门面:在一个不归它所有的 CommandBus 前面加上一条准入链。所有发送入口都按以下顺序执行同一条链:
- 命令体实现
CommandValidator时先执行自校验,再交给 JakartaValidator。 RequestIdChecker.check(aggregateId, requestId)做 request-ID 预检;返回false时以DuplicateRequestIdException终止。校验在前,校验失败的命令不占用 request ID(自 9.2.3 起)。- 只对
sendAndWait与sendAndWaitStream:等待计划必须支持Void命令,注册等待句柄,并把要发送的消息构建为调用方消息的副本,副本的 Header 带上等待键。调用方的消息不被修改(自 9.3.0 起;此前网关在发送前把等待键写进调用方的 Header)。 CommandBus.send;发送失败时调用RequestIdChecker.release释放这次预留。
SENT 信号只在一处产生:CommandBus.send 成功或失败之后,交给等待它的一方——已登记的句柄、Saga 为等待链发出的命令的上游等待,或 sendAndWaitForSent 的结果。sendAndWaitForSent 不登记句柄、不写等待 Header。每个等待都只有一个端到端截止时间,在调用被订阅时设定一次,由网关自己的定时器调度(自 9.3.0 起;此前流式等待每收到一个元素就在共享调度器上重新设定超时)。关闭网关只释放这个定时器,不关闭 CommandBus,总线归它的创建者所有(自 9.3.0 起)。
预检不是持久并发裁决。处理节点在聚合执行之前还会再查一次 request ID(见下文),最终的 request-ID 和版本冲突仍由 EventStore.append 的原子边界负责,详见失败与幂等。
Bus 到 Dispatcher
CommandBus.receiver(运行时的 CommandDispatcher 使用 runtime-owned 订阅)产生 ServerCommandExchange。CommandDispatcher 先过滤 isVoid 消息:这些消息会被确认但不会进入聚合命令链;普通命令继续按 NamedAggregate 分派。
自 9.3.0 起每个 AggregateCommandDispatcher 通过一个接收器服务一个限界上下文的全部聚合。它持有这些聚合的 metadata,把每条命令所属聚合的 metadata 传给 CommandHandler,并在运行时的 KeyedExecutor 上按 aggregate ID 使用邮箱。同一 ID 的命令按收到的顺序逐条执行,不同 ID 在共享工作线程上并行;这避免同一聚合在本进程内并发执行,但不替代 EventStore 的持久版本约束。
DefaultCommandHandler 执行一条固定顺序的管道,命令侧不再有过滤器链(自 9.3.0 起;此前是按 @Order 排序的 CommandFilter Bean):
CommandInstrumentation (each, the first outermost)
-> PROCESSED report
-> request-ID check, aggregate processing, then acknowledgement
-> DomainEventBus.send
-> StateEventBus.send attemptrequest-ID 检查在处理命令的节点上、聚合的处理函数执行之前进行(自 9.3.0 起)。它使用自己的布隆过滤器,不与网关共用,只有过滤器见过这个 request ID 时才查询 EventStore;聚合已经提交过的 request ID 会让命令以 DuplicateRequestIdException 失败,处理函数不会执行。wow.command.idempotency.enabled=false 会同时关闭它和网关的检查。
外层步骤包住内层步骤,因此观察的是内部整条管线的完成或错误,而不是只观察聚合函数返回。原来需要命令过滤器的场景,对应到类型化的扩展点:
| 需求 | 自 9.3.0 起 |
|---|---|
| 每条命令的追踪、指标、日志 | 注册 CommandInstrumentation Bean:around(exchange, handling) 包住整条管道,不得改变结果。OpenTelemetry 模块的 TraceCommandInstrumentation 取代 TraceAggregateFilter。多个 instrumentation 按 @Order 依次包裹。 |
| 命令执行前的检查或拒绝 | 在网关处校验命令(CommandValidator、Jakarta 校验),或在命令函数中检查;它抛出的异常让命令失败。 |
| 响应已提交的事件 | 在发布的事件上编写事件处理器、Saga 或投影。 |
| 改变命令函数收到的参数 | 注入参数(Spring Bean,或 exchange 提供的值)。 |
聚合恢复与调用
DefaultCommandHandler 为 exchange 放入 ServiceProvider,再按聚合身份与 metadata 创建 AggregateProcessor。默认 RetryableAggregateProcessor:
- 创建命令直接构造空的 StateAggregate;
- 其他命令从
StateAggregateRepository恢复状态; - 用恢复后的状态构造
SimpleCommandAggregate;聚合模式下命令根接收状态对象,非聚合模式直接复用状态对象; - 只对标记为 recoverable 的失败按内置退避策略重建状态并重试;
- 每次尝试都从第一次尝试之前的 exchange 开始:失败尝试留下的错误、事件流、聚合版本、命令调用结果、命令结果和命令聚合都不带到下一次,等待信号不会报告一个没有持久化的版本(自 9.2.3 起);
@OnError函数只在最终失败后(重试耗尽时取其原因)在最近一次加载的聚合上执行一次:通常是最后一次尝试的聚合,最后一次尝试在加载前失败时用更早一次尝试的聚合;没有任何尝试加载到聚合时不执行,自定义的非SimpleCommandAggregate的CommandAggregate由它自己的process处理错误(自 9.2.3 起)。
SimpleCommandAggregate.process 随后检查期望版本、创建许可、owner、space、删除/恢复状态和命令函数是否存在。检查通过后,它在聚合的 AggregateModel 中查找命令。该模型在解析聚合元数据时编译一次:命令条目(含匹配的 after-command 函数以及内置的删除、恢复、资源标签处理)、错误函数,以及该类型所有状态聚合共享的溯源表。处理函数把命令根或状态根作为参数接收,不再按聚合实例或按命令绑定(自 9.3.0 起)。命令条目调用匹配函数及有序的 after-command 函数,把返回值展平为一条 DomainEventStream 并放入 exchange。每个函数在模型编译时选定唯一的结果适配器,把各种返回形态(普通值、Mono、Flux、其他 Publisher、Flow、suspend 结果)转成这条事件流,并使用同一条异常规则:函数抛出的异常不经包装直接传出,返回 Flow 的函数在返回之前抛出的异常也一样(自 9.3.0 起)。
决定、应用、再追加
命令函数只读状态。它产出的事件流先应用到状态,再追加,三者是一个原子单元(自 9.3.0 起):已持久化的事件总能加载,内存状态绝不偏离存储。
invoke command(只读状态)
-> build DomainEventStream
-> 将事件溯源到状态
-> EventStore.append- 应用之前失败(守卫、命令函数):什么都不改变。
- 溯源失败:命令失败,什么都不存储、不发布,因此无法加载的事件永远不会被持久化。失败以 ERROR 记录日志。
- 追加失败(版本冲突、重复请求 ID、存储错误):命令失败,不发布任何内容,也不发送
StateEvent。 - 上述任一失败后,该状态实例可能持有存储中没有的事件,因此被丢弃:它不再接受任何命令。每次尝试(包括重试)都加载自己的聚合,测试 DSL 每一步之后都从存储重新加载状态,因此后续命令和读取都看不到它。
@OnError总是看到已提交的状态。失败的尝试应用了未存储的事件时,@OnError在重新加载的聚合上执行(只在这种情况下,且命令有@OnError函数时才加载;创建命令不读存储,而是由状态工厂新建聚合)。exchange 不会保留被丢弃的聚合,因此命令错误处理器、CommandInstrumentation与测试 DSL 也看不到它的状态。若重新加载也失败,则跳过@OnError,报告原始错误并把加载失败作为 suppressed 附上,同时以 ERROR 记录跳过。直接调用CommandAggregate.process的手工处理器则在被丢弃的实例上执行@OnError。- exchange 上的聚合版本(以及等待信号和
CommandResult)只在追加成功后才变为事件流的版本;失败的命令报告它做决定时的版本(N)。@OnError看到的是重新加载的已提交状态,版本冲突时它是存储中更新的版本,而不是 N。 StateAggregate.onSourcing在事件流的每个溯源函数都执行完之后,才推进版本、事件 ID、操作人、事件时间和系统元数据(拥有者、空间、删除标记、标签);VersionAware状态也在此时获得新版本。溯源函数抛出异常时,它们都保持在之前的版本。
事件历史与状态恢复的完整合同见事件溯源。
ack/事件发送顺序
DefaultCommandHandler 对聚合处理结果使用 finallyAck。因此无论聚合处理成功还是报错,都会先执行 exchange 的 transport ack;只有成功路径才发布。它发送处理器返回的事件流,并在继续之前等待 DomainEventBus.send 完成;随后在状态已初始化且已应用这条事件流(版本等于事件流的版本;这只是防御性检查,已存储的事件流总已应用)时复制事件流与当前状态,转换成 StateEvent 并尝试 StateEventBus.send。
实际顺序是:
EventStore.append
-> command exchange ack
-> DomainEventBus.send
-> StateEventBus.send attempt
-> PROCESSED signal如果聚合在形成事件流前失败,仍会 ack,但不会发布任何内容。若事件已经追加,而 DomainEventBus.send 失败,transport ack 已经发生,错误会继续传播,StateEventBus.send 不会执行,PROCESSED 会观察到失败;因此不能把领域事件发布失败解释为“事件未保存”,也不能假定 command transport 会重投它。
StateEventBus.send 的失败边界不同:它的错误被记录日志并恢复为空完成。于是成功的 PROCESSED 只证明 StateEvent 发布已经被尝试并返回,不证明 StateEvent 已经发布;依赖该输入的快照与投影可能没有收到消息。事件侧消费过程见事件分发管线。
PROCESSED 错误边界
PROCESSED 报告用 MonoCommandWaitNotifier 包住内部管道:
- 内部链正常完成时,从 exchange 的函数、版本、结果和可能的业务错误生成
PROCESSED信号; - 内部链抛错时,先生成失败信号,再把原异常继续传给处理器的 error handler(记录到 exchange 并打日志);retry-exhausted 包装会先还原其 cause;
- 没有等待 Header,或目标阶段不需要
PROCESSED时,不生成信号; - 通知采用 fire-and-forget,通知失败只记录日志,不改写命令处理结果。
所以 PROCESSED 成功表示聚合执行、事件追加、command ack 和 DomainEventBus.send 已经完成;状态已应用这条事件流时,StateEventBus.send 尝试已经返回。它不保证 StateEvent 发布成功,也不表示快照、投影、事件处理器或 Saga 已完成。失败信号也不能单独证明事件未追加,必须按失败与幂等检查权威历史。
API 分层
本页的类型是实现,不是应用 API。自 9.3.0 起,wow-core 在代码中标明这一点:
CommandAggregate、它的父接口AggregateProcessor、CommandAggregateFactory与SimpleCommandAggregateFactory标注@WowSpi。自行提供命令聚合的代码用@OptIn(WowSpi::class)选择加入,不加入时编译器给出警告。它们在同一个次版本线内保持二进制签名不变,次版本可以修改它们,并写进发布说明。AggregateProcessorFactory、RetryableAggregateProcessorFactory、DefaultCommandHandler、SimpleStateAggregate,函数元数据类型(FunctionAccessorMetadata、InjectParameter、FirstParameterKind、AfterCommandFunctionMetadata、MessageFunctionRegistrar、SimpleMessageFunctionRegistrar),事件分发器基类(CompositeEventDispatcher、AbstractEventFunctionRegistrar、EventHandler),COMMAND_GATEWAY_FUNCTION,以及 exchange 上调用结果的存取方法和事件流、版本的设置方法标注@InternalWowApi:由 Wow 自己的模块装配,任何版本都可能修改。RetryableAggregateProcessor、SimpleCommandAggregate、编译后的聚合模型(AggregateModel及其命令条目和编译后的函数)、exchange 属性键、函数访问器以及聚合与状态事件分发器是internal。
应用通过 CommandGateway 发送命令,用 @OnCommand 函数处理命令,读取 ServerCommandExchange.getEventStream(),这些都不需要选择加入。命令函数需要当前状态时,声明 ReadOnlyStateAggregate<S> 参数(例如读取 initialized),而不是 CommandAggregate。