领域实体与事件溯源 (Domain Entity and Event Sourcing)

什么是领域实体事件溯源? (What is Domain Entity Event Sourcing?)

在分布式事务的整个生命周期中,领域实体 (Domain-Entity) 对象充当核心数据载体。 它是一个共享的、强类型的数据容器,使 Saga 工作流中的每个执行器(跨度 / span)能够在业务流程向前推进的过程中消费现有的事务数据,并贡献新产生的业务状态。

StackSaga 并非就地修改单条无历史记录的数据,而是采用*领域实体事件溯源 (Domain-Entity Event Sourcing)* 机制。 每一次重要的原子状态流转都会生成一份不可变的领域实体状态快照 (Snapshot),并立即提交持久化到事件存储库 (Event Store) 中:

  1. 基准状态 (Baseline State):在事务初始化时,领域实体的初始状态作为基准版本持久化,状态标记为 STARTED。

  2. 步骤快照 (Step Snapshots):当每个原子执行器完成其操作后,都会持久化一个更新后的领域实体的时间点快照,并盖上唯一的跨度标识符和执行元数据戳记。

领域实体事件溯源提供了两项核心架构能力:

  1. 事务重新调用与自动恢复(重试与状态还原)

    • 如果原子执行因瞬态问题(如临时网络延迟或下游服务不可用)而失败,Saga 引擎可以安全重试该操作。如果事务因系统崩溃或节点重启而中断,则可以从中断的确切位置无缝恢复并重新调用。因为每个里程碑的历史状态都完好保存,引擎能够精确还原故障发生时存在的领域实体快照,彻底杜绝数据损坏和部分状态不一致。

  2. 通过仪表盘实现全局事务可观测性与审计溯源

    • 存储每个离散的状态变迁提供了完整的按时间顺序排列的可追溯性。系统管理员和开发人员可以在每个原子执行器运行之前和之后检查领域实体的状态,从而能够通过 Trace-Window(链路追踪窗口)仪表盘轻松审计事务的演进轨迹、诊断极端边缘异常并检查载荷数据变更。

领域实体作为 Saga 领域标识符 (Domain Entity as Saga Domain Identifier)

除了作为事务数据载体之外,DomainEntity 类在 StackSaga 框架中还承担着一项基础架构职责:它充当 Saga 领域标识符 (Saga Domain Identifier)。

框架利用 DomainEntity 的类类型来区分业务领域,并将负责处理该特定工作流的所有执行器绑定在一起。换句话说,领域实体类充当了整个 Saga 的强类型锚点:

  • 例如,在电子商务平台中,订单履行流程使用 OrderDomainEntity 作为其领域标识符,而客户订阅流程使用 SubscriptionDomainEntity。

  • 每个业务领域定义自己专用的 DomainEntity 子类,其中包含与该业务上下文相关的属性字段。

每个不同的长事务 (LRT - Long-Running Transaction) 都需要有自己专用的 DomainEntity 子类。 StackSaga 使用类类型本身(而非运行时的任意字符串或属性值)作为判别器来标识 Saga 领域。 例如,PlaceOrder(下单)和 CancelSubscription(取消订阅)工作流必须实现为独立的类:OrderDomainEntity 和 CancelSubscriptionDomainEntity。

除了事务识别外,DomainEntity 类还作为泛型类型锚点,统一了该 Saga 领域的所有相关框架组件:

  • SagaTemplate<OrderDomainEntity, ?>

  • AbstractEventManager<OrderDomainEntity, ?>

  • TransactionEventListener<OrderDomainEntity>

  • CommandExecutor<OrderDomainEntity> / QueryExecutor<OrderDomainEntity>

这种设计确保了编译期类型安全、结构化路由以及跨整个应用程序的清晰领域隔离。

Saga 执行中领域实体的生命周期与状态流转

在整个 Saga 事务执行期间,DomainEntity 历经三个清晰的生命周期阶段:

  1. 阶段 1:状态创建与初始化 (STARTED)

  2. 阶段 2:渐进式演进与前向变更 (IN_PROGRESS)

  3. 阶段 3:生命周期终态终止 (COMPLETED 或 FAILED)

理解跨这些阶段如何填充、消费、变更和终结数据,对于设计高可靠、幂等的 Saga 工作流至关重要。

阶段 1:状态创建与初始化 (State Creation & Initialization)

该生命周期始于客户端应用程序,随后才将事务移交给 Saga 编排引擎:

  1. 自定义实例创建:开发者实例化自定义 DomainEntity 子类(例如 OrderDomainEntity),并使用从客户或上游服务接收到的初始请求载荷填充它(例如 user_id、amount、items)。

  2. 初始字段状态:此时仅填充了初始请求属性。所有下游属性(例如 delivery_details、order_id 和 payment_id)均保持为 null。

  3. 移交 Saga:通过模板初始化方法将实体传递给编排器:

    sagaTemplate.init(orderDomainEntity)
        .startWith(UserDetailExecutor.class)
        ...
        .execute();
  4. 持久化基准事件:在调用第一个执行器之前,编排引擎序列化实体并将初始快照(版本 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 生命周期最终会在以下两种确定性状态之一终止:

  1. 成功完成 (COMPLETED):

    • 当工作流中的所有执行器均成功返回 stepManager.next(…​) 且未遇到关键致命异常(枢轴不可恢复故障)时,Saga 工作流到达终态节点。

    • 引擎将事务流转为 COMPLETED 状态,最终完全填充的领域实体快照在事件存储库中归档封存,作为已完成业务事务的永久审计凭证。

  2. 关键故障与逆向回滚 (FAILED):

    • 如果原子执行遇到不可重试错误(枢轴执行故障),前向演进立即停止。

    • Saga 引擎启动逆向回滚,以相反顺序对先前所有已完成的命令执行器执行 doRevert()。

    • 一旦所有补偿全部完成,事务以 FAILED 状态终止。领域实体保留故障发生前最后一次有效的正向状态,在 Trace-Window 仪表盘中为开发人员和运维人员提供确切的事后分析取证数据。

状态流转可视化:下单架构示例 (Place-Order Architecture)

以下架构图展示了在下单事务执行期间 OrderDomainEntity 完整的分步生命周期与状态流转过程:

StackSaga 领域实体状态流转架构图

理解架构图布局

该图沿执行时间线组织为三个架构列:

  1. 左列(Saga 执行器 / Saga Executors):展示原子执行单元(QueryExecutor 和 CommandExecutor),详述具体的方法调用(doProcess() 与 doRevert())及外部微服务交互。

  2. 中轴线(执行时间线与数据流):说明从初始化 (INIT) 经步骤 01 到 04 直至终态完成 (END) 的时间流。

    • 蓝色虚线箭头 (Reads from State):表示执行器从前一快照读取特定输入属性。

    • 绿色实线箭头 (Updates State):表示从 doProcess() 流出的状态更新,在事件存储库中生成持久化的新里程碑。

  3. 右列(OrderDomainEntity 快照):展示每一步后提交至事件存储库的领域实体时间点状态:

    • 薄荷绿行 (✓ initial / ✓ updated):高亮显示新初始化或变更的属性。

    • 灰色行 (preserved):高亮显示先前捕获且保持完好可访问的属性。

    • 白色行 (null):表示等待下游执行初始化的未赋值属性。

逐步状态演进矩阵表

步骤 执行器与类型 状态读取 (输入) 状态变更与生命周期事件

INIT

客户端控制器 (Client Controller)

客户端提交的订单请求载荷。

使用初始字段创建 OrderDomainEntity(user_id = "mafei"、amount = 100.0、items = […​])。下游字段均为 null。快照持久化到事件存储库,状态为 STARTED。

01

UserDetailExecutor
(QueryExecutor)

从领域实体读取 user_id。

调用外部 User-Service 查询送货地址。变更属性 delivery_details = "No344 New York."。快照持久化,状态为 IN_PROGRESS。(作为查询执行器,只读操作无需实现 doRevert())。

02

OrderInitializeExecutor
(CommandExecutor)

从领域实体读取 user_id 和 amount。

调用 Order-Service 记录订单条目。变更属性 order_id = "1245"。实现 doRevert() 以便在后续步骤失败时作废订单。快照持久化,状态为 IN_PROGRESS。

03

ReserveItemsExecutor
(CommandExecutor)

从领域实体读取 items 和 order_id。

调用 Stock-Service 锁定库存。注意: 没有追加新字段;既有字段完好保留 (preserved)。实现 doRevert() 以便在回滚时释放库存锁定。持久化快照以记录跨度里程碑。

04

MakePaymentExecutor
(CommandExecutor)

从领域实体读取 order_id 和 amount。

调用 Payment-Service 向客户扣款。变更属性 payment_id = "tx01343743131083017"。实现 doRevert() 以便在需要时发起退款。持久化快照,此时所有必需属性均已填充完毕。

END

Saga 引擎 (Saga Engine)

最终事务校验。

所有执行器均成功完成且未触发枢轴故障。事务流转至 COMPLETED 状态;最终领域实体状态在事件存储库中封存归档。

核心架构要点

  • 解耦的微服务架构:微服务之间从不直接相互调用来传递上下文。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:配置领域实体元数据:
  • name:领域实体在框架内的全局唯一标识名称。用于绑定事务与执行器。

  • version:通过 @SagaDomainEntityVersion 声明的模型结构版本号。对于滚动发布过程中的领域实体向上转换 (Up-Casting) 和向下转换 (Down-Casting) 至关重要。
    参见自定义映射器与键生成器的高级配置。

2 继承:自定义类必须继承 DomainEntity,以继承核心 Saga 标识与生命周期能力。
3 构造函数:调用 super(OrderDomainEntity.class) 的 protected 无参构造函数。由于框架在反序列化和状态恢复期间会动态实例化领域实体,因此不应声明有参构造函数。
4 字段映射:声明承载事务状态所需的属性。强烈建议使用 @JsonProperty 注解以确保确定性 JSON 序列化,并避免模式演进期间的命名不一致。
5 嵌套类型与未知属性归档:对于复杂的嵌套对象,声明继承自 org.stacksaga.api.UnknownPropertyArchive 的静态内部类(或独立类)。UnknownPropertyArchive 基类在反序列化期间自动捕获未识别属性,确保架构升级时的无缝前后向兼容性。

重构安全性:框架通过 @SagaDomainEntity(name = "…​") 中声明的逻辑 name 属性标识领域实体,而不是通过 Java 类名或包名。只要 @SagaDomainEntity 的 name 保持一致,您就可以安全地重构、重命名或移动 Java 类文件,而不会中断活动事务或历史事件存储库记录。 两个不同的自定义领域实体类不能声明相同的 name。

Spring Bean 说明:在 StackSaga 中,自定义领域实体是轻量级的状态载体——它们不是 Spring Bean。 它们无需位于应用程序的组件扫描包路径下。应通过 stacksaga.domain-entity-scan 配置属性注册包含领域实体的包。

敏感数据保护与追踪窗口防护(数据脱敏)

由于 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 安全最佳实践: 通过在 @JsonProperty 字段中仅存储脱敏或加密后的值(例如 encodedToken、maskedPan),并仅通过 @JsonIgnore 访问器方法暴露解密值:

  • 事务状态在节点故障与重试期间保持完全持久化和可恢复。

  • Trace-Window 仪表盘可以安全地向开发人员和审计人员呈现非敏感或脱敏后的载荷,而不会泄漏机密客户凭证或 PII 数据。

高级配置 (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 抽象类。
provide() 方法仅在框架初始化阶段调用一次,以获取目标领域实体的自定义 ObjectMapper 对象。
3 重写该方法以提供自定义的 ObjectMapper 对象。
4 返回定制配置后的 ObjectMapper 对象。
5 mapper:在 DomainEntity 类中指定自定义领域实体映射器提供器类。

领域实体自定义键生成器提供器 (Custom Key Generator Provider)

键生成器负责生成事务键前缀以及每个跨度的幂等键。 默认情况下,StackSaga 使用 DefaultDomainEntityKeyGenerator 作为所有领域实体的键生成器。 如果您想为特定领域实体自定义键生成规则,可以通过继承 AbstractDomainEntityKeyGenerator 来创建自定义实现。 您可以按需为不同的领域实体创建独立的自定义键生成器。

AbstractDomainEntityKeyGenerator 提供了两个带有默认实现的方法供您重写:

事务键生成 (Transaction Key Generation)

每个 Saga 事务都需要一个由 SagaUUID 表示的全局唯一标识符。 SagaUUID 由两部分组成,中间用连字符分隔:

  • <prefix>-<UUIDv7>

    1. 前缀 (Prefix):由 AbstractDomainEntityKeyGenerator 中的 generateTransactionKeyPrefix(…​) 生成。

    2. 后缀 (Suffix):由框架使用 java-uuid-generator 库自动生成的基于时间的唯一 UUIDv7。

由于框架会自动追加高熵的基于时间的 UUID 后缀,因此保证了跨事务与分布式节点的绝对唯一性。 前缀则充当简洁、人类可读的修饰符,便于可观测性监控、链路追踪以及业务领域切分。

为什么事务标识符采用 UUIDv7?

UUIDv7 在其前 48 位嵌入了高精度的纪元时间戳,生成按时间单调递增的有序标识符。 这种设计为事件溯源和事务存储带来了关键的架构优势:

  • 数据库分区裁剪 (Partition Pruning):在配置了时间范围分区的 SQL 和分布式事件存储实现中,可以直接从事务 ID 中提取嵌入的时间戳,将查询直接路由至对应分区,消除开销巨大的全分区扫描。

  • 索引局部性与写入吞吐量:与导致频繁 B-Tree 页分裂和随机 I/O 的随机 UUIDv4 不同,按顺序递增的 UUIDv7 键能自然追加到数据库索引页末尾,保持页面密度并最大化写入吞吐量。

  • 按时间顺序自然可追溯:事务直接通过标识符即可按创建时间自然排序,极大简化了时间窗口查询、日志关联分析和审计核对。

默认行为 (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)

Saga 事务中的每个跨度(即每次原子执行尝试)都需要一个幂等键,以确保安全、无重复副作用的重试。 generateIdempotencyKey 方法接收运行时执行上下文,该上下文划分为两个输入对象:SafeIdempotentInput 与 UnSafeIdempotentInput。

切勿使用 UnSafeIdempotentInput 中的属性来构造幂等键。UnSafeIdempotentInput 中的值在重试尝试或集群节点之间可能发生变化,这会导致生成非确定性键并破坏重试安全性。UnSafeIdempotentInput 严格用于诊断日志记录和调试目的。

务必仅从 SafeIdempotentInput 派生幂等键,它提供了确定、稳定的属性,保证在事务生命周期内的重试期间始终保持完全一致。

框架的默认实现通过拼接 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. 服务版本 1.0.0 已部署并正在执行事务。

  2. 某个执行器执行原子操作(例如调用外部支付或库存 API),派发了携带幂等键 K1 的请求。

  3. 在响应能够提交持久化到事件存储库之前,发生了临时网络抖动、网关超时或容器重启。该事务在事件存储库中被标记为待重试。

  4. 紧接着,部署了服务版本 1.1.0,其中在 OrderDomainEntityKeyGenerator 中更新了幂等键生成公式。

  5. 当 StackSaga 重试子系统或调度器在版本 1.1.0 下重新调用已暂停的事务时,它调用了 generateIdempotencyKey(…​)。

  6. 故障模式:如果键生成器对这个被重试的历史事务应用了*新的* 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):与解析好的 Jackson tools.jackson.core.Version 对象进行比较。

  • domainEntity.isInitializedVersionBetween(DomainEntityVersionDetail from, DomainEntityVersionDetail to):判断事务是否在特定历史版本区间内初始化。

  • domainEntity.compareInitializedVersionToCurrentVersion():将事件存储库中记录的事务初始版本与当前运行应用 @SagaDomainEntity 上声明的版本进行对比。

  • domainEntity.getInitializedVersionAsString():以字符串形式返回初始版本(例如 "1.0.0")。

为什么使用 compare(…​) 而不是原始字符串比较?
在活跃生产系统中,部署 1.1.0 时,事件存储库中的飞行中事务可能跨越多个旧版本(例如 1.0.0、1.0.1、1.0.2)。硬编码类似于 "1.0.0".equals(…​) 的字符串比较非常脆弱,需要手动列出每个历史版本。而通过 domainEntity.compare("1.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 事务时,基准事件(快照)提交持久化到事件存储库,并打上当时自定义 DomainEntity 类上声明的 @SagaDomainEntityVersion 戳记。 该初始版本称为事件版本 (Event Version)(或事务真实版本 / Transaction Real Version)。 即使随后部署了新的应用版本,该历史事务仍会保留其原始事件版本,执行器可以在运行时通过 domainEntity.getRealVersionAsString() 查询该版本。

何时需要更新领域实体版本

每当所做更改影响到事务的状态模式或执行拓扑时,必须递增在 @SagaDomainEntityVersion(major = …​, minor = …​, patch = …​) 中声明的版本:

1. 领域实体结构变更

如果自定义 DomainEntity 类的属性模式发生变更,必须更新版本。

  • 添加新数据(触发向上转换 Up-Casting):

    向自定义 DomainEntity 添加新属性时,递增版本。 添加字段会在事务重放期间触发*向上转换 (Up-Casting)*,此时事件存储库中的旧快照会被转换为包含新增字段的新模式定义。

    StackSaga 领域实体变更 - 添加新数据

  • 移除既有数据(触发向下转换 Down-Casting):

    从自定义 DomainEntity 中删除废弃字段时,递增版本。 删除字段会在事务重放期间触发*向下转换 (Down-Casting)*,此时事件存储库中的历史快照仍然包含当前 Java 类中不再定义的属性。

    StackSaga 领域实体变更 - 移除数据

2. 执行器修改

即使 DomainEntity 类的属性字段未发生改变,对工作流执行器的修改也可能改变事务状态的生成、消费或补偿方式,同样需要递增版本。

  • 修改既有执行器内部逻辑:

    如果修改了现有执行器的内部业务逻辑(例如更改外部 API 请求载荷或更改下游验证标准),需要确定该修改是否应应用于等待重试的历史挂起事务。 递增版本允许执行器根据 domainEntity.getRealVersionAsString() 进行条件分支,对飞行中的历史事务应用旧版处理,同时对新发起的事务执行新逻辑。

    StackSaga 执行器修改

  • 修改执行器数量:

    更改 Saga 工作流中的执行器数量会改变执行图:

  • 新增执行器: 当在序列中引入额外步骤(例如 command-executor-4)时,递增领域实体版本。

    StackSaga 新增执行器

  • 移除既有执行器: 从工作流序列中移除执行器同样需要更新版本。

    StackSaga 移除执行器

生产环境删除执行器的注意事项: 执行器直接参与自动化重试和补偿 (doRevert())。 如果事件存储库中仍有历史事务挂起在引用了您打算删除的执行器的步骤处,若从代码库中彻底删除了该 Java 执行器类,Saga 引擎将无法重放或回滚这些事务。 最佳实践:在应用程序中保留已弃用的执行器类,直到旧版本的所有飞行中历史事务全部完成或终止。

虽然理论上可以将现有的 QueryExecutor 升级为 CommandExecutor,但严禁将 CommandExecutor 改为 QueryExecutor,因为补偿历史命令必须具备回滚能力。StackSaga 官方建议尽量避免更改既有执行器类型。

3. 键生成器与幂等逻辑变更

如果您在自定义 AbstractDomainEntityKeyGenerator 实现中更改了幂等键生成规则,必须递增 @SagaDomainEntityVersion。

递增版本确保:

  • 在瞬态错误后从事件存储库恢复的飞行中事务可以通过 domainEntity.compare(…​) 或 domainEntity.isInitializedVersionBetween(…​) 被识别。

  • 自定义键生成器可以将历史事务路由至旧版键生成算法,保持跨重试的完全幂等键不变性,防止下游重复执行。

  • 详细架构解释和代码示例,请参阅 跨版本修改幂等键逻辑与飞行中事务。

领域实体版本转换(SEC 重放机制)

在滚动更新应用程序期间,运行新版本的实例已上线,而旧实例正被逐步淘汰。 在事件溯源架构中,先前应用版本中挂起或失败的事件仍存储在事件存储库中等待重试或补偿。

当 Saga 执行协调器 (SEC) 恢复这些事务时,必须从旧的序列化 JSON 二进制载荷中重新构造一个实时的 Java DomainEntity 实例:

引擎在事务重试时的事件重构过程

StackSaga 事件重构图

这种反序列化、映射和转换过程被称为领域实体版本转换 (Domain-Entity Version Casting)。

如果跨版本的模式变更未得到正确管理,引擎将无法反序列化历史事件,导致飞行中的事务因不可恢复的反序列化异常而停滞。 开发者必须确保跨滚动版本的向后兼容性。

转换分类:向上转换 (Up-Casting) 与向下转换 (Down-Casting)

根据领域实体状态结构的演变方式,版本转换分为两类:

1. 领域实体向上转换 (Up-Casting)

当当前应用程序版本的属性多于存储在事件存储库中的历史事件快照时,就会发生向上转换。 例如,历史事件持久化时包含 3 个字段,而当前的领域实体类定义了 4 个或更多字段:

StackSaga 领域实体版本向上转换

StackSaga 中的向上转换处理: StackSaga 会自动处理向上转换。 框架默认的 JSON 序列化器/反序列化器 (Jackson ObjectMapper) 可以无缝将较旧的 JSON 载荷反序列化为较新的类,而不会报错,并将任何新增字段初始化为其默认值(null、0 或默认集合实例)。 开发者无需为向上转换编写自定义转换逻辑。

2. 领域实体向下转换 (Down-Casting)

当当前应用程序版本的属性少于存储在事件存储库中的历史事件快照时(例如,从 Java 类中删除了过时的废弃属性),就会发生向下转换:

StackSaga 领域实体版本向下转换

当包含旧字段的历史事件反序列化为新类时,标准反序列化器要么在遇到未知属性时报错失败,要么默默丢弃未识别的数据。 在 Saga 编排中,静默丢弃这些数据是极其危险的,因为挂起的执行器或补偿回滚可能仍需要这些旧属性值来完成任务。

处理向下转换与模式演进

为了安全管理领域实体向下转换,存在两种技术方案:

反模式:直接忽略未知属性 (@JsonIgnoreProperties(ignoreUnknown = true)): 禁用 FAIL_ON_UNKNOWN_PROPERTIES 或在类上标注 @JsonIgnoreProperties(ignoreUnknown = true) 虽然可以避免反序列化报错,但会从内存中永久丢弃历史属性。 如果正在重试的执行器或补偿回滚需要那些已删除的值,事务将因数据缺失而彻底失败。强烈不建议采用此方法。

推荐方案:通过 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") 中。

如果内部嵌套类移除了属性却未继承 UnknownPropertyArchive,反序列化历史事件将会失败。 StackSaga 会在应用启动时对照样本载荷(通过 SagaSerializable 配置)验证领域实体模式,若检测到类型转换不兼容则快速失败报错。

在执行器中访问已保留的历史属性

当执行器运行时,它可以检查事务是否源自历史旧版本,并从 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(…​) 检索嵌套行项中被移除的属性。

getMissingProperties() 映射以及对象级版本控制方法(compare(…​)、isInitializedVersionBetween(…​)、getRealVersionAsString())在所有执行器类型中均可访问,包括 CommandExecutor、QueryExecutor 和 RevertAfterExecutor。