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

不依赖第三方库,Node.js中如何创建PostgreSQL连接池并查询?

不依赖第三方库在Node.js中实现PostgreSQL连接池与查询

我想在不依赖pg、socket.io等第三方库的情况下,在Node.js中创建PostgreSQL连接池并执行数据库查询。目前尝试用net模块创建Socket连接PostgreSQL的5432端口,但返回连接重置错误。我知道连接需要指定数据库名、用户名、密码等参数,但不清楚具体的协议格式,希望得到指导。

我的测试代码:

import * as net from "net";

let client = new net.Socket();

client.connect(5432, 'localhost', () => {
  console.log('Connected');
  client.write('Hello');
});

client.on('error', (err) => {
  console.log('error', err);
});

执行后错误信息:

Connected
error Error: read ECONNRESET                                        
  errno: -4077,
  code: 'ECONNRESET',
  syscall: 'read'
}

问题根源

PostgreSQL有自定义的TCP连接协议,不是直接发送文本字符串就能通信。你的代码发送Hello后,服务器无法识别不符合协议的报文,直接断开了连接,导致ECONNRESET错误。要实现通信,必须严格按照PostgreSQL的协议格式完成握手、认证、查询等步骤。

实现步骤

1. 实现单个PostgreSQL连接

首先要完成PostgreSQL的连接握手流程,核心步骤:

  • 发送启动报文:包含协议版本、用户名、目标数据库名等参数
  • 处理服务器的认证响应:如果是密码认证,需要发送加密后的密码
  • 握手完成后,发送查询报文并处理结果

以下是单个连接的示例代码:

import * as net from "net";
import { createHash } from "crypto";

function createPgConnection(config) {
  return new Promise((resolve, reject) => {
    const client = new net.Socket();
    let buffer = Buffer.alloc(0);

    client.connect(5432, config.host, () => {
      // 构造启动报文:协议版本+参数键值对
      const startupMsg = Buffer.concat([
        // 报文长度(含自身4字节)
        Buffer.from([0, 0, 0, 35]),
        // PostgreSQL协议版本号(3.0)
        Buffer.from([0, 3, 0, 0]),
        // 参数:user
        Buffer.from('user\0' + config.user + '\0'),
        // 参数:database
        Buffer.from('database\0' + config.db + '\0'),
        // 结束标记
        Buffer.from([0])
      ]);
      client.write(startupMsg);
    });

    client.on('data', (data) => {
      buffer = Buffer.concat([buffer, data]);
      // 解析报文类型
      const msgType = buffer.readUInt8(0);
      
      // 认证请求(R表示需要密码认证)
      if (msgType === 82) {
        // 提取salt
        const salt = buffer.slice(8, 12);
        // 加密密码(PostgreSQL的md5认证格式)
        const md5Hash = createHash('md5')
          .update(config.password + config.user)
          .digest('hex');
        const passwordHash = createHash('md5')
          .update(md5Hash + salt.toString('hex'))
          .digest('hex');
        const passwordMsg = Buffer.concat([
          Buffer.from([112]), // 'p'表示密码报文
          Buffer.from([0, 0, 0, 3 + passwordHash.length + 1]),
          Buffer.from('md5' + passwordHash + '\0')
        ]);
        client.write(passwordMsg);
        buffer = Buffer.alloc(0);
      }
      // 连接成功(K表示后端密钥交换,可忽略;Z表示就绪)
      else if (msgType === 90) {
        resolve({
          client,
          query: (sql) => {
            return new Promise((qResolve) => {
              let queryBuffer = Buffer.alloc(0);
              // 构造查询报文
              const queryMsg = Buffer.concat([
                Buffer.from([81]), // 'Q'表示查询报文
                Buffer.from([0, 0, 0, sql.length + 5]),
                Buffer.from(sql + '\0')
              ]);
              client.write(queryMsg);
              
              client.on('data', (resData) => {
                queryBuffer = Buffer.concat([queryBuffer, resData]);
                // 处理结果(简化版,只提取行数据)
                if (queryBuffer.readUInt8(0) === 68) { // 'D'表示数据行
                  const rowData = queryBuffer.slice(8).toString().split('\0');
                  qResolve(rowData.filter(item => item));
                } else if (queryBuffer.readUInt8(0) === 90) { // 'Z'表示查询结束
                  queryBuffer = Buffer.alloc(0);
                }
              });
            });
          }
        });
      }
    });

    client.on('error', reject);
  });
}

// 使用示例
(async () => {
  const conn = await createPgConnection({
    host: 'localhost',
    user: 'your_username',
    db: 'your_database',
    password: 'your_password'
  });
  const result = await conn.query('SELECT * FROM your_table LIMIT 1');
  console.log('查询结果:', result);
})();

2. 实现连接池

连接池的核心是维护一个可用连接队列,当请求连接时:

  • 如果队列有空闲连接,直接返回
  • 如果队列已满且未达到最大连接数,创建新连接
  • 如果已达最大连接数,等待队列中有连接释放

以下是简化版连接池实现:

class PgPool {
  constructor(config, maxConnections = 5) {
    this.config = config;
    this.maxConnections = maxConnections;
    this.availableConnections = [];
    this.pendingRequests = [];
    this.activeConnections = 0;
  }

  async getConnection() {
    // 优先使用空闲连接
    if (this.availableConnections.length > 0) {
      return this.availableConnections.shift();
    }

    // 未达最大连接数,创建新连接
    if (this.activeConnections < this.maxConnections) {
      this.activeConnections++;
      try {
        const conn = await createPgConnection(this.config);
        // 包装连接,释放时放回队列
        const wrappedConn = {
          ...conn,
          release: () => {
            this.availableConnections.push(wrappedConn);
            // 处理等待的请求
            if (this.pendingRequests.length > 0) {
              const resolve = this.pendingRequests.shift();
              resolve(wrappedConn);
            }
          }
        };
        return wrappedConn;
      } catch (err) {
        this.activeConnections--;
        throw err;
      }
    }

    // 达到最大连接数,等待空闲连接
    return new Promise((resolve) => {
      this.pendingRequests.push(resolve);
    });
  }

  async query(sql) {
    const conn = await this.getConnection();
    try {
      const result = await conn.query(sql);
      return result;
    } finally {
      conn.release();
    }
  }
}

// 使用连接池示例
(async () => {
  const pool = new PgPool({
    host: 'localhost',
    user: 'your_username',
    db: 'your_database',
    password: 'your_password'
  }, 5);

  const result = await pool.query('SELECT version()');
  console.log('查询结果:', result);
})();

注意事项

  • 以上代码是简化实现,仅处理了基本的密码认证和简单查询,实际生产环境需要处理更多协议细节(比如错误报文、不同类型的认证、结果集的完整解析等)
  • 直接实现PostgreSQL协议复杂度很高,除非有特殊需求,否则推荐使用成熟的第三方库(如pg)

内容的提问来源于stack exchange,提问作者j.ss

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 06:37:02