You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

能否在Quarkus中结合Reactive SQL Client使用事务观察者?

关于Reactive SQL Client结合事务事件的实现方案

可行性说明

直接使用JBoss Weld文档中描述的事务观察者机制与Reactive SQL Client结合是不可行的。原因在于Weld的事务观察者依赖JTA同步事务模型,而Reactive SQL Client采用异步非阻塞的事务处理模式,两者的事务生命周期回调机制不兼容,无法直接对接。

手动实现基于事件机制的事务生命周期回调

我们可以通过Reactive SQL Client的事务API,结合自定义事件发布/订阅机制,实现事务开始前、提交前、提交后的业务逻辑。下面以Quarkus环境下的Reactive PostgreSQL Client为例给出实现示例:

1. 定义事务生命周期事件类

public class TransactionLifecycleEvent {
    // 事务开始前事件
    public static class BeforeStart {}
    // 提交前事件
    public static class BeforeCommit {}
    // 提交后事件
    public static class AfterCommit {}
    // 事务回滚事件(可选)
    public static class AfterRollback {}
}

2. 实现事件发布者与事务处理封装

import io.quarkus.vertx.ReactiveSqlClient;
import io.smallrye.mutiny.Uni;
import io.vertx.mutiny.sqlclient.SqlConnection;
import jakarta.enterprise.context.ApplicationScoped;
import jakarta.enterprise.event.Event;

@ApplicationScoped
public class TransactionalEventManager {

    private final ReactiveSqlClient sqlClient;
    private final Event<TransactionLifecycleEvent.BeforeStart> beforeStartEvent;
    private final Event<TransactionLifecycleEvent.BeforeCommit> beforeCommitEvent;
    private final Event<TransactionLifecycleEvent.AfterCommit> afterCommitEvent;
    private final Event<TransactionLifecycleEvent.AfterRollback> afterRollbackEvent;

    // 构造注入依赖
    public TransactionalEventManager(ReactiveSqlClient sqlClient,
                                     Event<TransactionLifecycleEvent.BeforeStart> beforeStartEvent,
                                     Event<TransactionLifecycleEvent.BeforeCommit> beforeCommitEvent,
                                     Event<TransactionLifecycleEvent.AfterCommit> afterCommitEvent,
                                     Event<TransactionLifecycleEvent.AfterRollback> afterRollbackEvent) {
        this.sqlClient = sqlClient;
        this.beforeStartEvent = beforeStartEvent;
        this.beforeCommitEvent = beforeCommitEvent;
        this.afterCommitEvent = afterCommitEvent;
        this.afterRollbackEvent = afterRollbackEvent;
    }

    // 封装事务执行逻辑,触发对应事件
    public <T> Uni<T> runInTransaction(Uni<T> transactionalLogic) {
        // 触发事务开始前事件
        beforeStartEvent.fire(new TransactionLifecycleEvent.BeforeStart());

        return sqlClient.begin()
                .flatMap(connection -> {
                    return transactionalLogic
                            // 触发提交前事件
                            .invoke(() -> beforeCommitEvent.fire(new TransactionLifecycleEvent.BeforeCommit()))
                            .flatMap(result -> connection.commit().replaceWith(result))
                            // 触发提交后事件
                            .invoke(() -> afterCommitEvent.fire(new TransactionLifecycleEvent.AfterCommit()))
                            .onFailure().invoke(throwable -> {
                                // 触发回滚事件
                                afterRollbackEvent.fire(new TransactionLifecycleEvent.AfterRollback());
                                // 确保事务回滚
                                connection.rollback().subscribe().with(ignore -> {});
                            });
                });
    }
}

3. 实现事件观察者(业务逻辑)

import jakarta.enterprise.event.Observes;
import jakarta.enterprise.context.ApplicationScoped;

@ApplicationScoped
public class TransactionBusinessLogic {

    // 事务开始前执行的逻辑
    public void onTransactionBeforeStart(@Observes TransactionLifecycleEvent.BeforeStart event) {
        System.out.println("事务开始前:执行初始化检查、日志记录等逻辑");
        // 可添加参数校验、缓存预热等业务代码
    }

    // 提交前执行的逻辑
    public void onTransactionBeforeCommit(@Observes TransactionLifecycleEvent.BeforeCommit event) {
        System.out.println("提交前:执行数据一致性校验、预提交通知等逻辑");
        // 可添加数据完整性检查、待提交通知发送等业务代码
    }

    // 提交后执行的逻辑
    public void onTransactionAfterCommit(@Observes TransactionLifecycleEvent.AfterCommit event) {
        System.out.println("提交后:执行缓存更新、消息发送、统计上报等逻辑");
        // 可添加缓存更新、MQ消息发送等业务代码
    }

    // 回滚后执行的逻辑(可选)
    public void onTransactionAfterRollback(@Observes TransactionLifecycleEvent.AfterRollback event) {
        System.out.println("事务回滚后:执行告警通知、数据恢复等逻辑");
    }
}

4. 使用示例

import jakarta.inject.Inject;
import jakarta.ws.rs.GET;
import jakarta.ws.rs.Path;
import io.smallrye.mutiny.Uni;

@Path("/demo")
public class TransactionDemoResource {

    @Inject
    TransactionalEventManager transactionalEventManager;

    @Inject
    ReactiveSqlClient sqlClient;

    @GET
    public Uni<String> executeTransactionalLogic() {
        // 定义事务内的数据库操作逻辑
        Uni<String> dbLogic = sqlClient.withConnection(conn ->
                conn.query("INSERT INTO users(name) VALUES ('test')")
                        .execute()
                        .map(result -> "数据插入成功")
        );

        // 通过封装的事务管理器执行,自动触发生命周期事件
        return transactionalEventManager.runInTransaction(dbLogic);
    }
}

关键说明

  • 示例基于Quarkus的ReactiveSqlClient和CDI事件机制,若使用其他Reactive SQL Client实现,可替换对应的事务API,核心思路是在事务各阶段手动触发自定义事件。
  • 事件发布与订阅是异步的,符合Reactive编程的非阻塞特性,不会阻塞事务执行流程。
  • 可根据业务需求扩展事件类,比如携带事务ID、操作数据等上下文信息,方便观察者逻辑获取更多事务相关数据。

内容的提问来源于stack exchange,提问作者Alexey Poletaev

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.13 16:35:50