不依赖第三方库,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
相关产品推荐
相关产品推荐

