如何实现支持事务的事件总线?
嘿,我太懂被Rust生命周期问题卡得头疼的感觉了!你要做的这个带事务传递的事件总线,核心是让事务能跟着事件流走,所有订阅者都在同一个事务上下文里干活,咱们一步步来拆解怎么实现,避开那些坑。
第一步:抽象事务Trait
首先得把事务的行为抽象出来,这样不管是用sqlx、Diesel还是其他数据库的事务,都能无缝接入。咱们定义一个基础的Transaction Trait,包含事务最核心的提交、回滚方法:
pub trait Transaction: Send + Sync { // 提交事务 fn commit(&mut self) -> Result<(), Box<dyn std::error::Error>>; // 回滚事务 fn rollback(&mut self) -> Result<(), Box<dyn std::error::Error>>; } // 给sqlx的Postgres事务实现这个Trait(其他数据库同理) impl<'c> Transaction for sqlx::Transaction<'c, sqlx::Postgres> { fn commit(&mut self) -> Result<(), Box<dyn std::error::Error>> { self.commit().map_err(|e| e.into()) } fn rollback(&mut self) -> Result<(), Box<dyn std::error::Error>> { self.rollback().map_err(|e| e.into()) } }
这里加Send + Sync是为了让事务能在多线程或异步场景下安全传递,如果你的总线是单线程同步的,也可以去掉,但留着更通用。
第二步:定义事件监听器Trait
接下来是订阅者要实现的EventListener,这里关键是让监听器接收事务的可变引用,而不是所有权——这样所有订阅者都能共享同一个事务,不会出现生命周期冲突:
pub trait EventListener<E, T: Transaction> { // 处理事件,同时拿到事务的可变引用 fn handle(&self, event: &E, tx: &mut T) -> Result<(), Box<dyn std::error::Error>>; }
这里E是事件类型,T是事务类型,泛型设计让监听器可以绑定特定的事件和事务实现。
第三步:实现事件总线
事件总线需要管理所有订阅者,并且在发布事件时把事务传递给每个监听器。咱们先做一个同步版本的总线,用Arc<Mutex>来保证并发安全(如果是异步场景,可以换成tokio::sync::Mutex):
use std::sync::{Arc, Mutex}; pub struct EventBus<E, T: Transaction> { // 用Arc+Mutex管理监听器列表,支持多线程订阅/发布 listeners: Arc<Mutex<Vec<Box<dyn EventListener<E, T>>>>>, } impl<E, T: Transaction + 'static> EventBus<E, T> { // 创建新的事件总线 pub fn new() -> Self { Self { listeners: Arc::new(Mutex::new(Vec::new())), } } // 订阅事件:添加监听器到总线 pub fn subscribe(&self, listener: Box<dyn EventListener<E, T>>) { self.listeners.lock().unwrap().push(listener); } // 发布事件:遍历所有监听器,传递事件和事务引用 pub fn publish(&self, event: E, tx: &mut T) -> Result<(), Box<dyn std::error::Error>> { let listeners = self.listeners.lock().unwrap(); for listener in listeners.iter() { // 只要有一个监听器出错,整个流程终止,后续可以回滚事务 listener.handle(&event, tx)?; } Ok(()) } }
这里T: 'static是因为我们把监听器存在Box<dyn ...>里,需要事务类型的生命周期至少和总线一样长——这在大多数业务场景下都是合理的,因为事务一般是在业务逻辑里临时创建,用完就提交/回滚。
第四步:实际使用示例
咱们用一个用户创建的场景来演示怎么用:
// 定义一个事件:用户创建完成 pub struct UserCreated { pub user_id: u64, pub email: String, } // 实现一个监听器:在事务中插入用户日志 pub struct UserCreatedLogger; impl EventListener<UserCreated, sqlx::Transaction<'_, sqlx::Postgres>> for UserCreatedLogger { fn handle(&self, event: &UserCreated, tx: &mut sqlx::Transaction<'_, sqlx::Postgres>) -> Result<(), Box<dyn std::error::Error>> { // 在同一个事务里插入日志 sqlx::query!( "INSERT INTO user_logs (user_id, action) VALUES ($1, 'created')", event.user_id ) .execute(tx) .map_err(|e| e.into())?; Ok(()) } } // 业务逻辑中的使用流程 async fn create_user(db: &sqlx::PgPool) -> Result<(), Box<dyn std::error::Error>> { // 1. 开启数据库事务 let mut tx = db.begin().await?; // 2. 执行核心业务操作:插入用户数据 let user_id = sqlx::query!( "INSERT INTO users (email) VALUES ($1) RETURNING id", "test@example.com" ) .fetch_one(&mut tx) .await? .id; // 3. 创建事件总线,订阅监听器 let bus = EventBus::new(); bus.subscribe(Box::new(UserCreatedLogger)); // 4. 发布事件,把事务传递给所有监听器 bus.publish(UserCreated { user_id, email: "test@example.com".into() }, &mut tx)?; // 5. 所有操作完成,提交事务 tx.commit().await?; Ok(()) }
避开生命周期坑的关键技巧
- 用引用传递事务:绝对不要让监听器拿走事务的所有权,否则总线没法把事务传给其他订阅者,还会触发生命周期错误。
- 控制事务的生命周期:事务要在业务逻辑中创建,整个发布流程都要在事务的生命周期范围内,之后再统一提交/回滚。
- 异步场景的适配:如果你的总线是异步的,需要用
async-traitcrate来定义异步监听器Trait,同时要给事务和事件加上Send + Sync约束,确保异步任务能安全跨线程调度:use async_trait::async_trait; #[async_trait] pub trait EventListener<E, T: Transaction + Send + Sync> { async fn handle(&self, event: &E, tx: &mut T) -> Result<(), Box<dyn std::error::Error + Send + Sync>>; }
扩展:支持多事件类型
如果要让总线支持多种事件,可以把事件也抽象成Trait:
pub trait Event: Send + Sync + 'static {} // 给所有事件实现这个Trait impl Event for UserCreated {}
然后把EventBus的泛型改成E: Event,这样就能接收任意实现了Event的事件类型。
备注:内容来源于stack exchange,提问作者Lyle

