Thalo命令处理深度指南:从Aggregate定义到事件持久化的完整流程
【免费下载链接】thaloAn Event Sourcing runtime with WebAssembly & embedded event store项目地址: https://gitcode.com/gh_mirrors/th/thalo
Thalo是一个基于WebAssembly和嵌入式事件存储的事件溯源运行时(Event Sourcing runtime),它通过Aggregate模式实现命令处理与事件持久化的完整流程。本文将带您深入了解Thalo中命令从定义到最终持久化的全链路实现,帮助您快速掌握事件驱动架构的核心实践。
一、Aggregate定义:业务逻辑的核心载体
在Thalo中,Aggregate是封装业务逻辑和状态的核心单元。每个Aggregate需要实现Aggregatetrait,定义其关联的命令(Command)、事件(Event)类型以及初始状态。
1.1 基础结构定义
以计数器示例为例,Aggregate的基本定义如下:
use thalo::{events, export_aggregate, Aggregate, Apply, Command, Event, Handle}; #[derive(Debug, Default)] pub struct Counter { count: u64, } impl Aggregate for Counter { type Command = CounterCommand; type Event = CounterEvent; type Error = CounterError; fn init(_id: &str) -> Self { Self::default() } }上述代码定义了一个CounterAggregate,指定了它可以处理的CounterCommand命令类型、产生的CounterEvent事件类型,以及初始化方法。
1.2 命令与事件定义
命令(Command)是触发状态变更的输入,通常使用#[derive(Command)]宏来自动实现序列化和反序列化:
#[derive(Command, Deserialize)] pub enum CounterCommand { Increment { amount: u64 }, Decrement { amount: u64 }, }事件(Event)则是状态变更的结果记录,使用#[derive(Event)]宏:
#[derive(Event, Serialize, Deserialize)] pub enum CounterEvent { Incremented { amount: u64 }, Decremented { amount: u64 }, }二、命令处理:从接收请求到生成事件
命令处理是Thalo的核心流程,涉及命令验证、业务逻辑执行和事件生成三个关键步骤。
2.1 实现Handle trait处理命令
通过实现Handletrait来定义命令处理逻辑:
impl Handle<CounterCommand> for Counter { fn handle(&self, cmd: CounterCommand) -> Result<Vec<CounterEvent>, Self::Error> { match cmd { CounterCommand::Increment { amount } => { Ok(events![Incremented { amount }]) } CounterCommand::Decrement { amount } => { if self.count < amount { return Err(CounterError::InsufficientFunds); } Ok(events![Decremented { amount }]) } } } }在handle方法中,Aggregate根据当前状态验证命令合法性,并返回生成的事件列表。例如,当处理Decrement命令时,会检查当前计数是否足够,避免出现负数。
2.2 应用事件更新状态
生成的事件需要通过Applytrait应用到Aggregate状态:
impl Apply<CounterEvent> for Counter { fn apply(&mut self, event: CounterEvent) { match event { CounterEvent::Incremented { amount } => { self.count += amount; } CounterEvent::Decremented { amount } => { self.count -= amount; } } } }apply方法是纯函数,仅根据事件更新状态,不包含业务逻辑判断,确保状态变更的可追溯性。
三、事件持久化:从内存到存储的可靠落地
事件生成后需要持久化到事件存储(MessageStore),这是事件溯源的核心特性。
3.1 MessageStore架构
Thalo的事件存储通过thalo_message_store模块实现,提供了全局事件日志、流(Stream)和投影(Projection)等核心组件:
use thalo_message_store::MessageStore; // 打开事件存储 let message_store = MessageStore::open("path/to/store")?;MessageStore负责事件的持久化、查询和流管理,支持事务性写入和并发读取。
3.2 命令处理流程中的事件存储
在命令处理流程中,事件通过以下步骤完成持久化:
- 命令路由:命令通过
CommandGateway路由到对应的Aggregate处理 - 事件生成:Aggregate处理命令生成事件列表
- 事务写入:事件被原子性地写入事件存储和Outbox
- 状态更新:事件应用到Aggregate状态,完成状态更新
相关实现位于crates/thalo_runtime/src/command/aggregate_command_handler.rs中,核心代码片段如下:
// 从消息存储获取流 let stream = self.message_store.stream(stream_name)?; // 追加事件到流 stream.append(events).await?;四、Thalo命令处理全流程总结
Thalo的命令处理流程可以概括为以下五个关键步骤:
- 定义Aggregate:实现
Aggregatetrait,指定命令、事件类型和初始状态 - 定义命令/事件:使用
#[derive(Command)]和#[derive(Event)]宏定义数据结构 - 实现命令处理:通过
Handletrait实现业务逻辑,生成事件 - 实现状态更新:通过
Applytrait定义事件如何更新状态 - 事件持久化:通过
MessageStore将事件原子性写入存储
4.1 核心模块路径
- Aggregate trait定义:crates/thalo/src/lib.rs
- 命令处理运行时:crates/thalo_runtime/src/command/
- 事件存储实现:crates/thalo_message_store/src/message_store.rs
4.2 实战示例
Thalo提供了多个示例项目帮助理解命令处理流程:
- 计数器示例:examples/counter/src/lib.rs
- 待办事项示例:examples/todos/src/lib.rs
- 银行账户示例:examples/bank_account/src/lib.rs
通过这些示例,您可以快速掌握Thalo命令处理的最佳实践,进而构建自己的事件驱动应用。
五、快速上手Thalo命令处理
要开始使用Thalo进行命令处理,只需按照以下步骤操作:
克隆仓库:
git clone https://gitcode.com/gh_mirrors/th/thalo定义Aggregate:创建您的业务实体并实现
Aggregatetrait实现命令处理:通过
Handle和Applytrait定义业务逻辑运行 runtime:使用Thalo runtime加载模块并处理命令
Thalo的WebAssembly架构使您可以使用多种语言编写Aggregate逻辑,同时保持事件处理的高性能和可靠性。无论是构建微服务还是复杂业务系统,Thalo的命令处理流程都能为您提供清晰的业务逻辑边界和可追溯的状态变更历史。
【免费下载链接】thaloAn Event Sourcing runtime with WebAssembly & embedded event store项目地址: https://gitcode.com/gh_mirrors/th/thalo
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考