使用node pg(Postgres)客户端处理错误与实现自动恢复相关问题咨询
node-pg 连接容错与查询重试方案优化
连接逻辑问题修复与优化
重复触发重连日志的原因
你遇到的"Reconnecting"日志重复触发,是两个问题共同导致的:
- 旧客户端的error事件监听没有移除,重连创建新客户端后,旧客户端触发的error事件仍会被捕获
- 连接失败时,pg.Client会同时触发底层socket错误和客户端实例错误两个事件,导致回调执行两次
优化方案(替换全局变量+指数退避)
用闭包封装实例状态,避免全局变量污染,同时加入指数退避逻辑,重连成功后重置延迟:
import pg from 'pg' const createPgClient = () => { let client = null let reconnectDelay = 1000 // 初始延迟1s const MAX_DELAY = 30000 // 最大延迟30s const sleep = (ms) => new Promise(resolve => setTimeout(resolve, ms)) const connect = async () => { // 销毁旧客户端,移除所有监听 if (client) { client.removeAllListeners('error') try { await client.end() } catch {} } client = new pg.Client(process.env.CONNECTION_STRING) client.on('error', async (e) => { console.log(`Postgres连接出错,${reconnectDelay}ms后重试:`, e.message) await sleep(reconnectDelay) // 指数退避,不超过最大延迟 reconnectDelay = Math.min(reconnectDelay * 2, MAX_DELAY) await connect() }) try { await client.connect() console.log('Postgres连接成功') reconnectDelay = 1000 // 连接成功重置延迟 } catch (e) { // 主动捕获首次连接错误,走重连逻辑 console.log(`首次连接失败,${reconnectDelay}ms后重试:`, e.message) await sleep(reconnectDelay) reconnectDelay = Math.min(reconnectDelay * 2, MAX_DELAY) await connect() } } // 暴露获取客户端和连接的方法 return { connect, getClient: () => client } } // 初始化实例 export const pgInstance = createPgClient() // 启动时调用连接 pgInstance.connect()
查询逻辑优化
可重试错误判断规则
pg抛出的错误会携带code字段,仅针对以下类别的错误重试即可,其余错误直接抛出避免无效重试:
- 网络错误:
ECONNREFUSED、ENOTFOUND、EAI_AGAIN、ETIMEDOUT - 服务端临时错误:
57P01(管理员关闭连接)、08006(连接异常)、57000(操作被终止)、40001(事务序列化冲突)
通用重试查询封装
封装统一的查询方法,统一处理退避、错误判断、等待连接就绪逻辑,不需要每个查询单独写重试:
const RETRYABLE_ERRORS = new Set([ 'ECONNREFUSED', 'ENOTFOUND', 'EAI_AGAIN', 'ETIMEDOUT', '57P01', '08006', '57000', '40001' ]) const MAX_RETRY = 10 // 最大重试次数,可根据需求调整 const BASE_DELAY = 500 // 基础延迟 const sleep = (ms) => new Promise(resolve => setTimeout(resolve, ms)) export const query = async (sql, params, retryCount = 0) => { const client = pgInstance.getClient() // 客户端未就绪时等待后重试 if (!client) { if (retryCount >= MAX_RETRY) throw new Error('Postgres客户端未就绪,重试次数耗尽') await sleep(BASE_DELAY * Math.pow(2, retryCount)) return query(sql, params, retryCount + 1) } try { return await client.query(sql, params) } catch (err) { const errCode = err.code || err.message // 非可重试错误或者重试次数耗尽直接抛出 if (!RETRYABLE_ERRORS.has(errCode) || retryCount >= MAX_RETRY) { throw err } // 指数退避后重试 const delay = BASE_DELAY * Math.pow(2, retryCount) console.log(`查询出错,${delay}ms后重试,错误码: ${errCode}`) await sleep(delay) return query(sql, params, retryCount + 1) } } // 业务方法直接调用通用query即可 async function getTimestamp() { const res = await query('select current_timestamp from current_timestamp;') return res.rows[0].current_timestamp }
额外建议
优先使用pg.Pool(连接池)替代单个pg.Client,连接池本身内置了连接有效性检测、自动重建连接的能力,并发性能和容错能力都远高于单客户端,上述优化逻辑同样可以适配到连接池的实现上。
内容的提问来源于stack exchange,提问作者Michael
相关产品推荐
相关产品推荐

