如何在ExpressJS中间件之间传递pg-promise事务/任务对象?
Express中间件共享pg-promise事务/任务对象的解决方案
问题背景
在Express中拆分路由的验证逻辑与业务逻辑为独立中间件时,若验证逻辑(如认证)需要访问数据库,会出现每个中间件使用独立数据库连接的问题,随着路由逻辑复杂度提升,连接复用问题会更突出。
需要解决的核心问题:能否将pg-promise的db.tx()/db.task()创建的事务/任务对象,在其回调外部传递,让整个请求链路的中间件和路由处理共享同一个数据库连接?
解决方案:共享任务/事务对象
完全可以实现,核心思路是在请求入口的中间件中初始化任务/事务,将对象挂载到res.locals上供后续链路复用,最后统一处理任务/事务的结束逻辑。
1. 使用db.task()实现共享连接(无事务需求)
如果不需要事务保证,只是想复用数据库连接,db.task()是最优选择——它会在整个任务周期内维持单个数据库连接:
import { Router } from "express"; import { db } from "#db"; const router = Router(); // 初始化共享任务中间件 router.use("/auth", async (req, res, next) => { try { // 启动任务,将task对象挂载到res.locals await db.task(async (task) => { // 认证逻辑:使用task对象执行查询 const authKey = req.headers.authorization?.split(" ")[1]; // 示例提取认证密钥 const session = await task.one("SELECT * FROM sessions WHERE auth_key = $1", [authKey]); if (!session) { throw new Error("无效的会话"); } // 保存会话和任务对象到请求上下文 res.locals.session = session; res.locals.dbTask = task; // 等待后续中间件/路由处理完成 await new Promise(resolve => next(resolve)); }); } catch (error) { next(error); } }); // 路由处理:复用共享的task对象 router.get("/auth/:id/posts", async (req, res) => { const { id } = req.params; const { session, dbTask } = res.locals; // 使用共享连接执行查询 const posts = await dbTask.manyOrNone( "SELECT * FROM posts WHERE user_id = $1 AND session_id = $2", [id, session.id] ); return res.status(200).json(posts); });
2. 使用db.tx()实现事务共享(需原子性保证)
如果整个请求链路需要事务支持(比如认证后的数据插入/更新需要原子性),则用db.tx()——它会在事务内复用连接,且自动处理提交或回滚:
import { Router } from "express"; import { db } from "#db"; const router = Router(); // 初始化事务中间件 router.use("/auth", async (req, res, next) => { try { // 启动事务,将tx对象挂载到res.locals await db.tx(async (tx) => { const authKey = req.headers.authorization?.split(" ")[1]; const session = await tx.one("SELECT * FROM sessions WHERE auth_key = $1", [authKey]); if (!session) { throw new Error("无效的会话"); } res.locals.session = session; res.locals.dbTx = tx; // 等待后续链路完成,错误会触发事务回滚 await new Promise((resolve, reject) => { next((err) => err ? reject(err) : resolve()); }); }); // 无错误则自动提交事务 } catch (error) { next(error); } }); // 路由处理:在事务内执行操作 router.post("/auth/:id/posts", async (req, res) => { const { id } = req.params; const { session, dbTx } = res.locals; const { content } = req.body; // 事务内执行插入操作 await dbTx.none( "INSERT INTO posts(user_id, session_id, content) VALUES($1, $2, $3)", [id, session.id, content] ); // 事务内查询最新数据 const newPosts = await dbTx.manyOrNone( "SELECT * FROM posts WHERE user_id = $1 AND session_id = $2", [id, session.id] ); return res.status(201).json(newPosts); });
关键注意事项
- 任务/事务对象仅在回调周期内有效:不要在
db.task()/db.tx()的回调外部保存或使用这些对象,否则会导致连接泄漏或逻辑错误。 - 错误自动传播:任何中间件或路由中抛出的错误都会终止任务(或触发事务回滚),无需手动处理连接关闭。
- 必须等待后续链路完成:用
new Promise包裹next(),确保任务/事务等待整个请求处理流程结束后再释放连接。
内容的提问来源于stack exchange,提问作者Biller Builder
相关产品推荐
相关产品推荐

