领域实体与事件溯源 (Domain Entity and Event Sourcing)
什么是领域实体事件溯源? (What is Domain Entity Event Sourcing?)
在分布式事务的整个生命周期中,领域实体 (Domain-Entity) 对象充当核心数据载体。 它是一个共享的、强类型的数据容器,使 Saga 工作流中的每个执行器(跨度 / span)能够在业务流程向前推进的过程中消费现有的事务数据,并贡献新产生的业务状态。
StackSaga 并非就地修改单条无历史记录的数据,而是采用*领域实体事件溯源 (Domain-Entity Event Sourcing)* 机制。 每一次重要的原子状态流转都会生成一份不可变的领域实体状态快照 (Snapshot),并立即提交持久化到事件存储库 (Event Store) 中:
-
基准状态 (Baseline State):在事务初始化时,领域实体的初始状态作为基准版本持久化,状态标记为
STARTED。 -
步骤快照 (Step Snapshots):当每个原子执行器完成其操作后,都会持久化一个更新后的领域实体的时间点快照,并盖上唯一的跨度标识符和执行元数据戳记。
领域实体事件溯源提供了两项核心架构能力:
-
事务重新调用与自动恢复(重试与状态还原)
-
如果原子执行因瞬态问题(如临时网络延迟或下游服务不可用)而失败,Saga 引擎可以安全重试该操作。如果事务因系统崩溃或节点重启而中断,则可以从中断的确切位置无缝恢复并重新调用。因为每个里程碑的历史状态都完好保存,引擎能够精确还原故障发生时存在的领域实体快照,彻底杜绝数据损坏和部分状态不一致。
-
-
通过仪表盘实现全局事务可观测性与审计溯源
-
存储每个离散的状态变迁提供了完整的按时间顺序排列的可追溯性。系统管理员和开发人员可以在每个原子执行器运行之前和之后检查领域实体的状态,从而能够通过 Trace-Window(链路追踪窗口)仪表盘轻松审计事务的演进轨迹、诊断极端边缘异常并检查载荷数据变更。
-
领域实体作为 Saga 领域标识符 (Domain Entity as Saga Domain Identifier)
除了作为事务数据载体之外,DomainEntity 类在 StackSaga 框架中还承担着一项基础架构职责:它充当 Saga 领域标识符 (Saga Domain Identifier)。
框架利用 DomainEntity 的类类型来区分业务领域,并将负责处理该特定工作流的所有执行器绑定在一起。换句话说,领域实体类充当了整个 Saga 的强类型锚点:
-
例如,在电子商务平台中,订单履行流程使用
OrderDomainEntity作为其领域标识符,而客户订阅流程使用SubscriptionDomainEntity。 -
每个业务领域定义自己专用的
DomainEntity子类,其中包含与该业务上下文相关的属性字段。
|
每个不同的长事务 (LRT - Long-Running Transaction) 都需要有自己专用的 除了事务识别外,
这种设计确保了编译期类型安全、结构化路由以及跨整个应用程序的清晰领域隔离。 |
Saga 执行中领域实体的生命周期与状态流转
在整个 Saga 事务执行期间,DomainEntity 历经三个清晰的生命周期阶段:
-
阶段 1:状态创建与初始化 (
STARTED) -
阶段 2:渐进式演进与前向变更 (
IN_PROGRESS) -
阶段 3:生命周期终态终止 (
COMPLETED或FAILED)
理解跨这些阶段如何填充、消费、变更和终结数据,对于设计高可靠、幂等的 Saga 工作流至关重要。
阶段 1:状态创建与初始化 (State Creation & Initialization)
该生命周期始于客户端应用程序,随后才将事务移交给 Saga 编排引擎:
-
自定义实例创建:开发者实例化自定义
DomainEntity子类(例如OrderDomainEntity),并使用从客户或上游服务接收到的初始请求载荷填充它(例如user_id、amount、items)。 -
初始字段状态:此时仅填充了初始请求属性。所有下游属性(例如
delivery_details、order_id和payment_id)均保持为null。 -
移交 Saga:通过模板初始化方法将实体传递给编排器:
sagaTemplate.init(orderDomainEntity) .startWith(UserDetailExecutor.class) ... .execute(); -
持久化基准事件:在调用第一个执行器之前,编排引擎序列化实体并将初始快照(版本 0)提交到事件存储库,将事务标记为
STARTED状态。
阶段 2:渐进式演进与前向变更 (Progressive Evolution & Forward Mutation)
一旦 Saga 引擎接管执行控制权,它将协调原子跨度(QueryExecutor 和 CommandExecutor)的顺序执行:
-
选择性状态消费(读流程):
每个执行器并不需要全部数据集;它仅从当前领域实体快照中查询其所需的特定属性。例如:
-
UserDetailExecutor读取user_id以查询客户收货地址详情。 -
OrderInitializeExecutor读取user_id和amount来登记订单。 -
MakePaymentExecutor读取order_id和amount来处理支付扣款。
-
-
前向状态变更 (
doProcess()):当执行器完成外部调用后,它使用 setter 方法丰富领域实体(例如
domainEntity.setDeliveryDetails(…)、domainEntity.setOrderId(…))。 从doProcess()返回时,框架会自动序列化变更后的领域实体,并将不可变快照写入事件存储库,盖上唯一的跨度标识符并标记为IN_PROGRESS状态。 -
非变更跨度中的状态保留 (State Preservation):
并非每个执行器都需要向领域实体添加新属性。 例如, ReserveItemsExecutor读取items和order_id以在外部仓储服务中锁定库存。尽管它执行了关键的事务命令,但它并没有向领域实体追加新字段。在这一步中,既有字段保持完好保留 (preserved),引擎会持久化一个里程碑快照以确认该原子跨度执行成功。 -
补偿不可变规则 (
doRevert()):在补偿执行 (
doRevert()) 期间,领域实体是严格只读的。 补偿操作绝不能修改领域实体的状态,确保对实际已发生事实的历史审计记录保持不可篡改。 补偿所需的任何元数据(例如授权令牌或临时取消标识)均通过 Revert-Hint-Store (补偿提示存储库) 独立管理。
阶段 3:生命周期终态 (Lifecycle Termination)
Saga 生命周期最终会在以下两种确定性状态之一终止:
-
成功完成 (
COMPLETED):-
当工作流中的所有执行器均成功返回
stepManager.next(…)且未遇到关键致命异常(枢轴不可恢复故障)时,Saga 工作流到达终态节点。 -
引擎将事务流转为
COMPLETED状态,最终完全填充的领域实体快照在事件存储库中归档封存,作为已完成业务事务的永久审计凭证。
-
-
关键故障与逆向回滚 (
FAILED):-
如果原子执行遇到不可重试错误(枢轴执行故障),前向演进立即停止。
-
Saga 引擎启动逆向回滚,以相反顺序对先前所有已完成的命令执行器执行
doRevert()。 -
一旦所有补偿全部完成,事务以
FAILED状态终止。领域实体保留故障发生前最后一次有效的正向状态,在 Trace-Window 仪表盘中为开发人员和运维人员提供确切的事后分析取证数据。
-
理解架构图布局
该图沿执行时间线组织为三个架构列:
-
左列(Saga 执行器 / Saga Executors):展示原子执行单元(
QueryExecutor和CommandExecutor),详述具体的方法调用(doProcess()与doRevert())及外部微服务交互。 -
中轴线(执行时间线与数据流):说明从初始化 (
INIT) 经步骤01到04直至终态完成 (END) 的时间流。-
蓝色虚线箭头 (
Reads from State):表示执行器从前一快照读取特定输入属性。 -
绿色实线箭头 (
Updates State):表示从doProcess()流出的状态更新,在事件存储库中生成持久化的新里程碑。
-
-
右列(OrderDomainEntity 快照):展示每一步后提交至事件存储库的领域实体时间点状态:
-
薄荷绿行 (
✓ initial/✓ updated):高亮显示新初始化或变更的属性。 -
灰色行 (
preserved):高亮显示先前捕获且保持完好可访问的属性。 -
白色行 (
null):表示等待下游执行初始化的未赋值属性。
-
逐步状态演进矩阵表
| 步骤 | 执行器与类型 | 状态读取 (输入) | 状态变更与生命周期事件 |
|---|---|---|---|
INIT |
客户端控制器 (Client Controller) |
客户端提交的订单请求载荷。 |
使用初始字段创建 |
01 |
UserDetailExecutor |
从领域实体读取 |
调用外部 |
02 |
OrderInitializeExecutor |
从领域实体读取 |
调用 |
03 |
ReserveItemsExecutor |
从领域实体读取 |
调用 |
04 |
MakePaymentExecutor |
从领域实体读取 |
调用 |
END |
Saga 引擎 (Saga Engine) |
最终事务校验。 |
所有执行器均成功完成且未触发枢轴故障。事务流转至 |
核心架构要点
-
解耦的微服务架构:微服务之间从不直接相互调用来传递上下文。
DomainEntity作为单一事实来源和共享数据载体贯穿整个分布式工作流。 -
单调状态累积 (Monotonic State Accumulation):领域实体渐进式累积状态——每个执行器使用其执行结果丰富共享容器,以便后续执行器消费。
-
状态保留与状态修改:执行器并不强制要求修改领域实体。它可以执行业务逻辑并成功返回,同时保留既有状态不变。
-
事件溯源与重放保证:因为每一步的快照都持久化在事件存储库中,引擎可以恢复确切的历史状态以重试失败操作或补偿先前步骤,避免数据污染。
-
严格变更边界:状态修改严格限制在
doProcess()中。补偿执行 (doRevert()) 保持纯只读,回滚所需的元数据隔离保存在 Revert-Hint-Store 中。
创建自定义领域实体 (Creating a Custom Domain-Entity)
要定义自定义领域实体,需创建一个继承自框架基类 DomainEntity 的类,并为其添加 @SagaDomainEntity 注解。
您可以声明在 Saga 工作流中维护状态所需的任何属性字段(如 username、totalAmount、productItems 等)。启动事务时,将此自定义类的已初始化实例传递给 SagaTemplate.init(…).startWith(…).execute()。
import java.util.List;
import java.util.Map;
import com.fasterxml.jackson.annotation.JsonProperty;
import lombok.Getter;
import lombok.NoArgsConstructor;
import lombok.Setter;
import org.stacksaga.api.DomainEntity;
import org.stacksaga.api.UnknownPropertyArchive;
import org.stacksaga.annotation.SagaDomainEntity;
import org.stacksaga.annotation.SagaDomainEntityVersion;
@Getter
@Setter
(1)
@SagaDomainEntity(
version = @SagaDomainEntityVersion(major = 1, minor = 0, patch = 0),
name = "OrderDomainEntity"
)
public class OrderDomainEntity extends DomainEntity { (2)
(4)
@JsonProperty("username")
private String username;
@JsonProperty("order_id")
private String orderId;
@JsonProperty("total_amount")
private double totalAmount;
@JsonProperty("payment_reference_id")
private String paymentReferenceId;
@JsonProperty("user_validation_data")
private UserValidationData userValidationData;
@JsonProperty("product_items")
private List<ProductItem> productItems;
@JsonProperty("metadata")
private Map<String, String> metadata;
protected OrderDomainEntity() {
(3)
super(OrderDomainEntity.class);
}
(5)
@Getter
@Setter
@NoArgsConstructor
public static class ProductItem extends UnknownPropertyArchive {
@JsonProperty("product_id")
private String productId;
@JsonProperty("quantity")
private int quantity;
@JsonProperty("price")
private double price;
}
@Getter
@Setter
@NoArgsConstructor
public static class UserValidationData extends UnknownPropertyArchive {
@JsonProperty("is_user_validated")
private boolean userValidated;
@JsonProperty("validation_note")
private String validationNote;
}
}
| 1 | @SagaDomainEntity:配置领域实体元数据:
|
| 2 | 继承:自定义类必须继承 DomainEntity,以继承核心 Saga 标识与生命周期能力。 |
| 3 | 构造函数:调用 super(OrderDomainEntity.class) 的 protected 无参构造函数。由于框架在反序列化和状态恢复期间会动态实例化领域实体,因此不应声明有参构造函数。 |
| 4 | 字段映射:声明承载事务状态所需的属性。强烈建议使用 @JsonProperty 注解以确保确定性 JSON 序列化,并避免模式演进期间的命名不一致。 |
| 5 | 嵌套类型与未知属性归档:对于复杂的嵌套对象,声明继承自 org.stacksaga.api.UnknownPropertyArchive 的静态内部类(或独立类)。UnknownPropertyArchive 基类在反序列化期间自动捕获未识别属性,确保架构升级时的无缝前后向兼容性。 |
|
重构安全性:框架通过 |
|
Spring Bean 说明:在 StackSaga 中,自定义领域实体是轻量级的状态载体——它们不是 Spring Bean。
它们无需位于应用程序的组件扫描包路径下。应通过 |
敏感数据保护与追踪窗口防护(数据脱敏)
由于 DomainEntity 代表事务的完整状态,它在每个原子跨度快照处都会被序列化并提交到事件存储库中。
这些序列化快照会直接被检索并显示在 Trace-Window 仪表盘 (/execution/binary) 上,用于取证分析、调试和可观测性监控。
当事务处理敏感业务数据时(例如密码、认证令牌、API 密钥、信用卡号或个人身份信息 PII),开发者必须确保敏感值受到保护,不会以明文形式暴露在事件存储库中或 Trace-Window 仪表盘上。
由于 DomainEntity 是通过 Jackson 序列化的标准 Java 对象,您可以利用 Jackson 注解和自定义 getter/setter 封装来保护敏感属性:
1. 使用 Jackson 注解 (@JsonIgnore) 忽略敏感字段
如果某个字段仅在内存运行时需要,且绝不应持久化到事件存储库或显示在 Trace-Window 中:
-
在 getter 方法或字段上标注
@JsonIgnore。 -
Jackson 将在序列化 JSON 载荷时排除此属性。
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonProperty;
public class OrderDomainEntity extends DomainEntity {
@JsonProperty("card_last_four")
private String cardLastFour;
// 内存中的临时 CVV:绝不序列化或持久化到事件存储库
@JsonIgnore
private String cvv;
// ...
}
2. 加密 / 混淆存储与明文业务访问器
当敏感数据必须跨 Saga 跨度保留(以便后续执行器使用)时,在持久化属性中以加密、哈希或编码形式存储数据,同时提供标注有 @JsonIgnore 的辅助方法,供执行器在内存中解密或解码实际值:
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonProperty;
import org.stacksaga.annotation.SagaDomainEntity;
import org.stacksaga.annotation.SagaDomainEntityVersion;
import org.stacksaga.api.DomainEntity;
@SagaDomainEntity(
version = @SagaDomainEntityVersion(major = 1, minor = 0, patch = 0),
name = "OrderDomainEntity"
)
public class OrderDomainEntity extends DomainEntity {
// 序列化并存储在事件存储库中(对于 Trace-Window 仪表盘安全脱敏)
@JsonProperty("token")
private String encodedToken;
public OrderDomainEntity() {
super(OrderDomainEntity.class);
}
/**
* 业务方法 / 执行器调用此方法获取真实解密数据。
* 标注 @JsonIgnore 以免 Jackson 将其序列化到事件存储库或 Trace-Window 中。
*/
@JsonIgnore
public String getTokenAsString() {
return decodeOrDecrypt(this.encodedToken);
}
public String getEncodedToken() {
return this.encodedToken;
}
public void setEncodedToken(String encodedToken) {
this.encodedToken = encodedToken;
}
}
|
Trace-Window 安全最佳实践:
通过在
|
高级配置 (Further Configurations)
此外,您可以通过 @SagaDomainEntity 注解的属性为自定义领域实体提供一些附加配置:
领域实体自定义映射提供器 (Custom Mapper Provider)
默认情况下,StackSaga 通过 DefaultDomainEntityMapperProvider 使用 Spring Boot 提供的默认 ObjectMapper。
如果您希望针对目标领域实体自定义 ObjectMapper,可以作为 AbstractDomainEntityMapperProvider 的自定义实现为目标领域实体创建并提供专用的 objectMapper 对象。
您可以根据需要为不同的领域实体创建任意数量的自定义映射提供器,如下所示:
@Component (1)
public class OrderDomainEntityMapperProvider extends AbstractDomainEntityMapperProvider { (2)
@Override (3)
protected ObjectMapper provide() {
return new ObjectMapper(); (4)
}
}
//-------------------------------------------------------------------------------
@Getter
@Setter
@SagaDomainEntity(
version = @SagaDomainEntityVersion(major = 1, minor = 0, patch = 0),
name = "OrderDomainEntity",
mapper = OrderDomainEntityMapperProvider.class (5)
)
public class OrderDomainEntity extends DomainEntity {
//...
}
| 1 | @Component:将自定义对象映射器实现声明为 Spring Bean。 | ||
| 2 | 继承 AbstractDomainEntityMapperProvider 抽象类。
|
||
| 3 | 重写该方法以提供自定义的 ObjectMapper 对象。 |
||
| 4 | 返回定制配置后的 ObjectMapper 对象。 |
||
| 5 | mapper:在 DomainEntity 类中指定自定义领域实体映射器提供器类。 |
领域实体自定义键生成器提供器 (Custom Key Generator Provider)
键生成器负责生成事务键前缀以及每个跨度的幂等键。
默认情况下,StackSaga 使用 DefaultDomainEntityKeyGenerator 作为所有领域实体的键生成器。
如果您想为特定领域实体自定义键生成规则,可以通过继承 AbstractDomainEntityKeyGenerator 来创建自定义实现。
您可以按需为不同的领域实体创建独立的自定义键生成器。
AbstractDomainEntityKeyGenerator 提供了两个带有默认实现的方法供您重写:
-
generateTransactionKeyPrefix- 生成事务键前缀。 -
generateIdempotencyKey- 为事务的每个跨度生成幂等键。
事务键生成 (Transaction Key Generation)
每个 Saga 事务都需要一个由 SagaUUID 表示的全局唯一标识符。
SagaUUID 由两部分组成,中间用连字符分隔:
-
<prefix>-<UUIDv7>-
前缀 (Prefix):由
AbstractDomainEntityKeyGenerator中的generateTransactionKeyPrefix(…)生成。 -
后缀 (Suffix):由框架使用 java-uuid-generator 库自动生成的基于时间的唯一 UUIDv7。
-
由于框架会自动追加高熵的基于时间的 UUID 后缀,因此保证了跨事务与分布式节点的绝对唯一性。 前缀则充当简洁、人类可读的修饰符,便于可观测性监控、链路追踪以及业务领域切分。
|
为什么事务标识符采用 UUIDv7? UUIDv7 在其前 48 位嵌入了高精度的纪元时间戳,生成按时间单调递增的有序标识符。 这种设计为事件溯源和事务存储带来了关键的架构优势:
|
默认行为 (first4)
提供自定义键生成器是可选的。
如果您未提供自定义键生成器,StackSaga 将使用 DefaultDomainEntityKeyGenerator。
默认情况下,它提取 @SagaDomainEntity(name = "…") 注解中定义的领域实体名称的前 4 个字符 (first4) 并将其转换为小写。
例如,对于 @SagaDomainEntity(name = "OrderDomainEntity"),前 4 个字符 ("Orde") 被转换为小写 ("orde"),生成的事务 ID 如下:
-
orde-0195655a-350f-786d-96eb-63c1dfc6e9ba -
orde-0195655a-4e20-7a1b-80c2-1249fae53412
前缀约束与小写格式要求
通过 generateTransactionKeyPrefix(…) 自定义前缀时:
-
小写要求:开发者应提供全小写的前缀字符串。如果未提供小写,框架内部也会自动将其转换为小写。
-
允许的字符:仅允许字母、数字和连字符 (
[a-zA-Z0-9-])。 -
禁用字符:空格、空白符以及特殊字符(例如
_、:、#、@、.、/、%)严禁使用,否则会导致框架拒绝该事务 ID。
自定义前缀可用于:
-
融入微服务名称或领域命名空间,
-
嵌入区域 (Region) 或集群标识符以进行多区域路由和日志过滤,
-
满足企业的合规审计或可观测性命名规范。
| 有关如何在代码中提供自定义前缀,请参阅下文的 自定义键生成器示例。 |
为 Saga 跨度生成幂等键 (Generating Idempotency Keys)
| 如果您刚接触幂等性概念,请先阅读 长事务中的原子事务与幂等性 (Atomic Transactions & Idempotency in LRT)。 |
Saga 事务中的每个跨度(即每次原子执行尝试)都需要一个幂等键,以确保安全、无重复副作用的重试。
generateIdempotencyKey 方法接收运行时执行上下文,该上下文划分为两个输入对象:SafeIdempotentInput 与 UnSafeIdempotentInput。
|
切勿使用 务必仅从 |
框架的默认实现通过拼接 SafeIdempotentInput 中的 transactionId、currentExecutionName 和 executionMode,并通过提供的 HashGenerator 对组合字符串进行哈希处理,生成固定长度的 MD5 哈希值。这确保了幂等键紧凑且抗碰撞。
自定义键生成器示例 (Custom Key Generator Example)
以下示例展示了如何重写 generateTransactionKeyPrefix 和 generateIdempotencyKey 以提供自定义逻辑,并在自定义 DomainEntity 类中进行配置:
由于 AbstractDomainEntityKeyGenerator 为所有方法提供了默认实现,您只需重写需要自定义的特定方法(例如,仅需自定义前缀命名时只重写 generateTransactionKeyPrefix,或者仅需自定义幂等逻辑时只重写 generateIdempotencyKey)。
|
@Component (1)
public class OrderDomainEntityKeyGenerator extends AbstractDomainEntityKeyGenerator { (2)
@Override (3)
// 在每个事务初始化时调用该方法以提供前缀
public String generateTransactionKeyPrefix(String serviceName, String applicationVersion, String instanceId, String region, String zone, SagaDomainEntity sagaDomainEntity) {
// 自定义前缀必须为全小写(仅限字母数字和连字符)
return String.format("%s-%s", serviceName.toLowerCase(), region.toLowerCase());
}
@Override (4)
public String generateIdempotencyKey(SafeIdempotentInput safeIdempotentInput, UnSafeIdempotentInput unSafeIdempotentInput) {
final String rowKey = new StringJoiner(":")
.add(safeIdempotentInput.transactionId().toString())
.add(safeIdempotentInput.currentExecutionName())
.add(safeIdempotentInput.executionMode().name().toLowerCase())
.toString();
return this.hashGenerator().generateHash(rowKey, HashGenerator.ALGType.MD5);
}
}
//: 在自定义领域实体上配置自定义 KeyGen
@Getter
@Setter
@SagaDomainEntity(
version = @SagaDomainEntityVersion(major = 1, minor = 0, patch = 0),
name = "OrderDomainEntity",
mapper = OrderDomainEntityMapperProvider.class,
keyGen = OrderDomainEntityKeyGenerator.class (5)
)
public class OrderDomainEntity extends DomainEntity {
//...
}
| 1 | @Component:将自定义键生成器实现声明为 Spring Bean。 |
| 2 | 继承 AbstractDomainEntityKeyGenerator。 |
| 3 | 重写 generateTransactionKeyPrefix 方法,为事务 ID (SagaUUID) 创建自定义前缀。应提供小写前缀;若传入大写字符,框架内部会自动转为小写。该方法在每个事务初始化时被调用。 |
| 4 | 重写 generateIdempotencyKey 方法,为每个跨度生成自定义幂等键。该方法在每个跨度执行被调用之前触发。 |
| 5 | 在 DomainEntity 类上的 @SagaDomainEntity 注解中将自定义类配置为 keyGen。 |
强烈建议使用基类方法 hashGenerator() 生成组合字符串的固定长度哈希值作为幂等键,而不是直接返回输入值的原始拼接字符串。这种方式能保证幂等键紧凑、长度统一,并在输入值较长或内容可变时将碰撞风险降至最低。
|
将实现类注册为 Spring Bean(例如 @Component),并确保其为无状态且线程安全。框架可能会并发调用它。
|
跨版本修改幂等键逻辑与飞行中事务 (In-Flight Transactions)
在生产微服务架构中,应用程序通过零停机滚动部署持续演进。 当发布修改了幂等键构造规则的新服务版本时,必须确保飞行中事务 (In-Flight Transactions) 保持向后兼容性和确定性。
飞行中重试危险场景
考虑以下分布式场景:
-
服务版本
1.0.0已部署并正在执行事务。 -
某个执行器执行原子操作(例如调用外部支付或库存 API),派发了携带幂等键
K1的请求。 -
在响应能够提交持久化到事件存储库之前,发生了临时网络抖动、网关超时或容器重启。该事务在事件存储库中被标记为待重试。
-
紧接着,部署了服务版本
1.1.0,其中在OrderDomainEntityKeyGenerator中更新了幂等键生成公式。 -
当 StackSaga 重试子系统或调度器在版本
1.1.0下重新调用已暂停的事务时,它调用了generateIdempotencyKey(…)。 -
故障模式:如果键生成器对这个被重试的历史事务应用了*新的*
1.1.0公式,它将生成键K2而不是K1。 由于K2 != K1,下游微服务无法将此请求识别为先前尝试的重试。它将K2视为全新事务并二次执行业务副作用(例如重复扣款或双重库存锁定)。
|
幂等键不变性 (Idempotency Key Invariance): 任何给定执行跨度的幂等键在事务生命周期的所有尝试、重试与补偿中必须保持绝对完全一致,即便跨越了应用程序版本部署亦是如此。 |
基于对象级方法的语义版本分支判断
为了防止此类问题,在更改键生成逻辑时应递增 @SagaDomainEntityVersion 中的版本号,并使用 DomainEntity 对象级版本控制 API 进行条件分支:
框架通过 unSafeIdempotentInput.domainEntity() 提供目标 DomainEntity。该实例保留了事务的初始事件版本(事务最初在事件存储库中创建时盖上的版本戳记),使您的键生成器能够将历史飞行中事务与新发起的事务区分开来。
DomainEntity 提供了以下对象级版本比较方法:
-
domainEntity.compare(String versionString):将事务的初始化版本与语义版本字符串进行比较(例如domainEntity.compare("1.1.0"))。如果事务是在较早版本下初始化的,返回负整数;相等返回0;较新则返回正整数。 -
domainEntity.compare(Version versionY):与解析好的 Jacksontools.jackson.core.Version对象进行比较。 -
domainEntity.isInitializedVersionBetween(DomainEntityVersionDetail from, DomainEntityVersionDetail to):判断事务是否在特定历史版本区间内初始化。 -
domainEntity.compareInitializedVersionToCurrentVersion():将事件存储库中记录的事务初始版本与当前运行应用@SagaDomainEntity上声明的版本进行对比。 -
domainEntity.getInitializedVersionAsString():以字符串形式返回初始版本(例如"1.0.0")。
|
为什么使用 |
多版本键生成器实现示例
以下展示了一个支持历史飞行中事务同时对新事务应用新公式的 AbstractDomainEntityKeyGenerator 实现:
@Component
public class OrderDomainEntityKeyGenerator extends AbstractDomainEntityKeyGenerator {
@Override
public String generateIdempotencyKey(
SafeIdempotentInput safeIdempotentInput,
UnSafeIdempotentInput unSafeIdempotentInput
) {
DomainEntity domainEntity = unSafeIdempotentInput.domainEntity();
// 检查事务是否起源于 1.1.0 之前(覆盖 1.0.0、1.0.1、1.0.2 等)
if (domainEntity != null && domainEntity.compare("1.1.0") < 0) {
return generateLegacyIdempotencyKey(safeIdempotentInput);
}
// 1.1.0 及以上版本的标准键生成逻辑
return generateModernIdempotencyKey(safeIdempotentInput);
}
private String generateLegacyIdempotencyKey(SafeIdempotentInput safeInput) {
// 1.1.0 之前版本中使用的旧版键生成公式
final String rawKey = safeInput.transactionId().toString() + ":" + safeInput.currentExecutionName();
return this.hashGenerator().generateHash(rawKey, HashGenerator.ALGType.MD5);
}
private String generateModernIdempotencyKey(SafeInput safeInput) {
// 1.1.0+ 的更新版键生成公式
final String rawKey = new StringJoiner(":")
.add(safeInput.transactionId().toString())
.add(safeInput.currentExecutionName())
.add(safeInput.executionMode().name().toLowerCase())
.toString();
return this.hashGenerator().generateHash(rawKey, HashGenerator.ALGType.MD5);
}
}
领域实体版本控制与版本转换 (Domain-Entity Versioning & Version Casting)
在企业微服务系统中,软件通过功能迭代、重构与模式调整持续演进。 在 StackSaga 编排的分布式 Saga 工作流中,对事务状态或执行序列的任何结构修改都需要更新对应领域实体的版本号。
由于 StackSaga 依赖事件溯源 (Event Sourcing) 和自动化事务重试/恢复机制来保证跨分布式服务的最终一致性,因此确定性地管理模型版本至关重要。
为什么领域实体版本控制至关重要
当原子跨度或执行器遇到瞬态故障(例如下游超时、网络分区)时,事务会被暂停,直到重试引擎或调度器重新调用它。 在此窗口期间,可能已经部署了具有更新业务逻辑或修改了实体字段的新应用版本。
当 Saga 引擎恢复飞行中的事务时,从事件存储库检索的历史快照可能与新部署应用的 Java 类定义不匹配。 引擎需要明确的版本元数据来识别事务发起的确切模式定义,从而安全重构和协调历史状态,而不会发生数据损坏或反序列化失败。
|
什么是事件版本? (What is the Event’s Version?) 启动 Saga 事务时,基准事件(快照)提交持久化到事件存储库,并打上当时自定义 |
何时需要更新领域实体版本
每当所做更改影响到事务的状态模式或执行拓扑时,必须递增在 @SagaDomainEntityVersion(major = …, minor = …, patch = …) 中声明的版本:
1. 领域实体结构变更
如果自定义 DomainEntity 类的属性模式发生变更,必须更新版本。
-
添加新数据(触发向上转换 Up-Casting):
向自定义
DomainEntity添加新属性时,递增版本。 添加字段会在事务重放期间触发*向上转换 (Up-Casting)*,此时事件存储库中的旧快照会被转换为包含新增字段的新模式定义。
-
移除既有数据(触发向下转换 Down-Casting):
从自定义
DomainEntity中删除废弃字段时,递增版本。 删除字段会在事务重放期间触发*向下转换 (Down-Casting)*,此时事件存储库中的历史快照仍然包含当前 Java 类中不再定义的属性。
2. 执行器修改
即使 DomainEntity 类的属性字段未发生改变,对工作流执行器的修改也可能改变事务状态的生成、消费或补偿方式,同样需要递增版本。
-
修改既有执行器内部逻辑:
如果修改了现有执行器的内部业务逻辑(例如更改外部 API 请求载荷或更改下游验证标准),需要确定该修改是否应应用于等待重试的历史挂起事务。 递增版本允许执行器根据
domainEntity.getRealVersionAsString()进行条件分支,对飞行中的历史事务应用旧版处理,同时对新发起的事务执行新逻辑。
-
修改执行器数量:
更改 Saga 工作流中的执行器数量会改变执行图:
-
新增执行器: 当在序列中引入额外步骤(例如
command-executor-4)时,递增领域实体版本。
-
移除既有执行器: 从工作流序列中移除执行器同样需要更新版本。
|
生产环境删除执行器的注意事项:
执行器直接参与自动化重试和补偿 ( |
|
虽然理论上可以将现有的 |
3. 键生成器与幂等逻辑变更
如果您在自定义 AbstractDomainEntityKeyGenerator 实现中更改了幂等键生成规则,必须递增 @SagaDomainEntityVersion。
递增版本确保:
-
在瞬态错误后从事件存储库恢复的飞行中事务可以通过
domainEntity.compare(…)或domainEntity.isInitializedVersionBetween(…)被识别。 -
自定义键生成器可以将历史事务路由至旧版键生成算法,保持跨重试的完全幂等键不变性,防止下游重复执行。
-
详细架构解释和代码示例,请参阅 跨版本修改幂等键逻辑与飞行中事务。
领域实体版本转换(SEC 重放机制)
在滚动更新应用程序期间,运行新版本的实例已上线,而旧实例正被逐步淘汰。 在事件溯源架构中,先前应用版本中挂起或失败的事件仍存储在事件存储库中等待重试或补偿。
当 Saga 执行协调器 (SEC) 恢复这些事务时,必须从旧的序列化 JSON 二进制载荷中重新构造一个实时的 Java DomainEntity 实例:
这种反序列化、映射和转换过程被称为领域实体版本转换 (Domain-Entity Version Casting)。
|
如果跨版本的模式变更未得到正确管理,引擎将无法反序列化历史事件,导致飞行中的事务因不可恢复的反序列化异常而停滞。 开发者必须确保跨滚动版本的向后兼容性。 |
转换分类:向上转换 (Up-Casting) 与向下转换 (Down-Casting)
根据领域实体状态结构的演变方式,版本转换分为两类:
处理向下转换与模式演进
为了安全管理领域实体向下转换,存在两种技术方案:
|
反模式:直接忽略未知属性 ( |
推荐方案:通过 UnknownPropertyArchive 归档未知属性
StackSaga 提供了 UnknownPropertyArchive 基类 (org.stacksaga.api.UnknownPropertyArchive),而不是丢弃未知属性。
当引擎反序列化包含不再声明为 Java 类变量的字段的旧 JSON 载荷时,UnknownPropertyArchive 会拦截并将这些未识别的键值对存储在内部的 missingProperties 映射中。
-
根领域实体自动归档: StackSaga 中的基类
DomainEntity本身就继承了UnknownPropertyArchive。因此,从自定义DomainEntity类中删除的任何顶层属性都会自动捕获到missingProperties中,无需任何自定义注解或方法。 -
嵌套类归档: 对于复杂的嵌套类型或内部类(例如明细行项、客户详情或元数据),开发者必须在每个嵌套静态类上显式继承
UnknownPropertyArchive,以捕获已删除的子属性。
代码演练:处理向下转换
考虑一个订单处理工作流,其中版本 1.0.0 演进为版本 1.1.0。
旧模式 (OrderDomainEntity 版本 1.0.0):
@SagaDomainEntity(
version = @SagaDomainEntityVersion(major = 1, minor = 0, patch = 0),
name = "OrderDomainEntity"
)
@Getter
@Setter
public class OrderDomainEntity extends DomainEntity {
public OrderDomainEntity() {
super(OrderDomainEntity.class);
}
@JsonProperty("order_id")
private String orderId;
@JsonProperty("username")
private String username;
@JsonProperty("total")
private Double total;
@JsonProperty("is_active")
private Integer isActive;
@JsonProperty("comment")
private String comment; (1)
@JsonProperty("item_details")
private List<ItemDetail> itemDetails = new ArrayList<>();
@Getter
@Setter
@AllArgsConstructor
@NoArgsConstructor
public static class ItemDetail implements Serializable { (2)
@JsonProperty("item_name")
private String itemName;
@JsonProperty("qty")
private int qty;
@JsonProperty("price")
private double price;
@JsonProperty("note")
private String note; (3)
}
}
| 1 | comment 定义在根实体上。 |
| 2 | ItemDetail 是一个嵌套静态类。 |
| 3 | note 定义在嵌套的 ItemDetail 类上。 |
新模式 (OrderDomainEntity 版本 1.1.0):
在全新发布中,业务需求发生变化:根属性 comment 和行项属性 note 已被废弃,并从活动的 Java 模型中移除:
@SagaDomainEntity(
version = @SagaDomainEntityVersion(major = 1, minor = 1, patch = 0),
name = "OrderDomainEntity"
)
@Getter
@Setter
public class OrderDomainEntity extends DomainEntity {
public OrderDomainEntity() {
super(OrderDomainEntity.class);
}
@JsonProperty("order_id")
private String orderId;
@JsonProperty("username")
private String username;
@JsonProperty("total")
private Double total;
@JsonProperty("is_active")
private Integer isActive;
// 从根类中移除了字段 'comment' (1)
@JsonProperty("item_details")
private List<ItemDetail> itemDetails = new ArrayList<>();
@Getter
@Setter
@AllArgsConstructor
@NoArgsConstructor
public static class ItemDetail extends UnknownPropertyArchive { (2)
@JsonProperty("item_name")
private String itemName;
@JsonProperty("qty")
private int qty;
@JsonProperty("price")
private double price;
// 从 ItemDetail 中移除了字段 'note' (3)
}
}
| 1 | 从根类中移除了 comment 属性。因为 OrderDomainEntity 继承了 DomainEntity(其自身继承自 UnknownPropertyArchive),历史事件中任何旧的 comment 值都会自动路由至 domainEntity.getMissingProperties().get("comment")。 |
| 2 | ItemDetail 现在继承 UnknownPropertyArchive,以捕获已移除的嵌套字段。 |
| 3 | 从 ItemDetail 中移除了 note 属性。因为 ItemDetail 继承了 UnknownPropertyArchive,旧的 note 值会被捕获在 itemDetail.getMissingProperties().get("note") 中。 |
|
如果内部嵌套类移除了属性却未继承 |
在执行器中访问已保留的历史属性
当执行器运行时,它可以检查事务是否源自历史旧版本,并从 getMissingProperties() 中检索保留的属性:
@SagaExecutor(executeFor = "order-service", value = "OrderSaveExecutor")
@AllArgsConstructor
public class OrderSaveExecutor implements CommandExecutor<OrderDomainEntity> {
@Override
public ProcessStepManager<OrderDomainEntity> doProcess(
OrderDomainEntity currentDomainEntity,
ProcessStepManagerUtil<OrderDomainEntity> stepManager,
String idempotencyKey
) throws RetryableExecutorException, NonRetryableExecutorException {
// 检查事务是否起源于 1.1.0 之前(覆盖 1.0.0、1.0.1 等)
if (currentDomainEntity.compare("1.1.0") < 0) { (1)
// 检索根级缺失属性
Object comment = currentDomainEntity.getMissingProperties().get("comment"); (2)
if (comment != null) {
log.info("Processing legacy transaction with comment: {}", comment);
}
// 从内部行项中检索嵌套缺失属性
for (OrderDomainEntity.ItemDetail itemDetail : currentDomainEntity.getItemDetails()) { (3)
Object note = itemDetail.getMissingProperties().get("note");
if (note != null) {
log.info("Legacy item note: {}", note);
}
}
} else {
// 1.1.0+ 的标准执行路径
log.info("Processing standard order: {}", currentDomainEntity.getOrderId());
}
return stepManager.next(UpdateStockExecutor.class, "SAVED_ORDER");
}
@Override
public SagaExecutionEventName doRevert(
NonRetryableExecutorException primaryExecutionException,
OrderDomainEntity finalDomainEntityState,
RevertHintStore revertHintStore,
String idempotencyKey
) throws RetryableExecutorException {
// 回滚逻辑同样可以按需检查 finalDomainEntityState.compare(...) 或 getMissingProperties()
return SagaExecutionEventName.of("REVERTED_ORDER_SAVE");
}
}
| 1 | 使用 currentDomainEntity.compare("1.1.0") < 0(或 currentDomainEntity.getRealVersionAsString())检测事务的基准事件版本。与阈值进行比较涵盖了所有历史发布版本(如 1.0.0、1.0.1),无需为每个版本手动编写字符串匹配。 |
| 2 | 通过 currentDomainEntity.getMissingProperties().get(…) 检索根级被移除的属性。 |
| 3 | 通过 itemDetail.getMissingProperties().get(…) 检索嵌套行项中被移除的属性。 |
|
|