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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 12:05:42