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

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 { }
    }
}

关键说明

  1. 事务控制逻辑:通过手动关闭autoCommit,在操作完成后commit,异常时rollback,确保每条通知的处理逻辑完全包裹在独立事务中。
  2. 线程适配:由于JDBC是同步阻塞API,而PgSubscriber的回调运行在Reactor非阻塞线程池,因此使用ManagedExecutor将JDBC操作提交到专用阻塞线程池,避免影响Reactive性能。
  3. 回滚验证:取消代码中throw RuntimeException的注释后,执行插入操作后抛出异常,数据库中的message_process_logs记录会被回滚,验证事务生效。

内容的提问来源于stack exchange,提问作者Gabriel Shanahan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 12:55:57