Quarkus中基于Postgres LISTEN/NOTIFY实现每条通知的事务封装
基于Quarkus + Postgres JDBC的事务化MQ监听实现
核心依赖
确保你的项目中包含以下Quarkus依赖(以Maven为例):
<dependency> <groupId>io.quarkus</groupId> <artifactId>quarkus-jdbc-postgresql</artifactId> </dependency> <dependency> <groupId>io.quarkus</groupId> <artifactId>quarkus-reactive-pg-client</artifactId> </dependency> <dependency> <groupId>io.quarkus</groupId> <artifactId>quarkus-narayana-jta</artifactId> </dependency>
最小实现代码
以下代码实现了从PgChannel监听通知,到事务化处理的完整逻辑,同时提供两种API风格,并包含异常回滚演示:
import io.quarkus.reactive.pg.client.PgChannel import io.quarkus.reactive.pg.client.PgSubscriber import io.smallrye.mutiny.Multi import jakarta.enterprise.context.ApplicationScoped import jakarta.enterprise.event.Observes import jakarta.inject.Inject import org.eclipse.microprofile.context.ManagedExecutor import java.sql.Connection import javax.sql.DataSource @ApplicationScoped class PostgresMQListener { @Inject lateinit var pgSubscriber: PgSubscriber @Inject lateinit var dataSource: DataSource @Inject lateinit var managedExecutor: ManagedExecutor private lateinit var messageChannel: PgChannel // 应用启动时初始化通知监听 fun init(@Observes startup: Any) { val topic = "message_notifications" // 替换为你的触发器NOTIFY主题 messageChannel = pgSubscriber.channel(topic) } // 风格1:暴露Multi供调用方订阅,每条通知在独立事务中处理 fun messageStream(): Multi<String> { return messageChannel.getMessages() .onItem().transformToMulti { notification -> // 将同步JDBC操作提交到阻塞线程池,避免阻塞Reactor线程 Multi.createFrom().completionStage( managedExecutor.runAsync { dataSource.connection.use { conn -> try { conn.autoCommit = false // 1. 从通知中获取消息ID,查询原始消息 val messageId = notification.payload() val messageContent = conn.prepareStatement("SELECT content FROM messages WHERE id = ?").use { stmt -> stmt.setString(1, messageId) stmt.executeQuery().use { rs -> if (rs.next()) rs.getString("content") else null } } ?: throw IllegalArgumentException("消息ID $messageId 不存在") // 2. 执行下游业务操作:插入处理日志 conn.prepareStatement("INSERT INTO message_process_logs (message_id, content) VALUES (?, ?)").use { stmt -> stmt.setString(1, messageId) stmt.setString(2, messageContent) stmt.executeUpdate() } // 【回滚演示】取消注释下面的代码,插入后抛出异常,日志记录会被回滚 // throw RuntimeException("模拟处理失败,触发事务回滚") conn.commit() messageId } catch (e: Exception) { conn.rollback() throw e // 将异常传递给Multi的错误处理流程 } } } ) } .concurrency(5) // 控制并发处理的事务数量,根据系统性能调整 } // 风格2:接受自定义Handler,处理每条通知 fun subscribeToMessages(handler: (String, String) -> Unit) { messageChannel.getMessages() .onItem().invoke { notification -> managedExecutor.runAsync { dataSource.connection.use { conn -> try { conn.autoCommit = false val messageId = notification.payload() val messageContent = conn.prepareStatement("SELECT content FROM messages WHERE id = ?").use { stmt -> stmt.setString(1, messageId) stmt.executeQuery().use { rs -> if (rs.next()) rs.getString("content") else null } } ?: throw IllegalArgumentException("消息ID $messageId 不存在") // 执行下游业务操作 conn.prepareStatement("INSERT INTO message_process_logs (message_id, content) VALUES (?, ?)").use { stmt -> stmt.setString(1, messageId) stmt.setString(2, messageContent) stmt.executeUpdate() } // 【回滚演示】取消注释下面的代码,触发回滚 // throw RuntimeException("模拟处理失败") conn.commit() handler(messageId, messageContent) } catch (e: Exception) { conn.rollback() // 这里可以添加错误日志、重试逻辑等 e.printStackTrace() } } } } .subscribe().with { } } }
关键说明
- 事务控制逻辑:通过手动关闭
autoCommit,在操作完成后commit,异常时rollback,确保每条通知的处理逻辑完全包裹在独立事务中。 - 线程适配:由于JDBC是同步阻塞API,而
PgSubscriber的回调运行在Reactor非阻塞线程池,因此使用ManagedExecutor将JDBC操作提交到专用阻塞线程池,避免影响Reactive性能。 - 回滚验证:取消代码中
throw RuntimeException的注释后,执行插入操作后抛出异常,数据库中的message_process_logs记录会被回滚,验证事务生效。
内容的提问来源于stack exchange,提问作者Gabriel Shanahan
相关产品推荐
相关产品推荐

