@effect/sql-pg 原生 PostgreSQL 协议客户端改造:连接启动、类型编解码与破坏性变更全解读
2026/9/14 15:58:04 网站建设 项目流程

@effect/sql-pg 原生 PostgreSQL 协议客户端改造:连接启动、类型编解码与破坏性变更全解读

【免费下载链接】effectBuild production-ready applications in TypeScript项目地址: https://gitcode.com/GitHub_Trending/ef/effect

导读:@effect/sql-pg是 Effect SQL 生态中面向 PostgreSQL 的驱动,本次 changeset(.changeset/pre/eff-854-pg-connection-startup.md)将其底层从pg(node-postgres)文本协议运行时替换为原生 PostgreSQL 协议客户端,并把PgConnection/PgPool的能力全面升级为自建会话。文章将按变更说明的骨架,逐一拆解新的连接启动流程、make/makeClient构造方式、二进制编解码结果、命名预处理语句与prepare: false的适用场景、LISTEN/NOTIFY的新签名、sql.json的显式包装要求、单语句约束,以及Pool.reserve独占借用与失效替换修复,并给出可复现的代码示例。

1. 变更总览:从pg运行时到原生协议客户端

本次补丁(patch 级别,涉及@effect/sql-pgeffect两个包)的核心事实是:@effect/sql-pgpg运行时被替换为原生 PostgreSQL 客户端。此前依赖 node-postgres 的pg库完成连接管理、查询执行与类型转换;现在PgConnectionPgPool自行处理以下全部职责:

  • 连接启动(connection setup):传输层建立、可选SSLRequest、启动报文(startup message)与认证握手(见 PgConnection.ts);
  • 二进制查询(binary queries):走 PostgreSQL 扩展查询协议,Bind以二进制格式(format = 1)编码参数(见 PgTypes.ts);
  • 预处理语句(prepared statements):命名预处理默认开启,按连接维护一个有界 LRU 缓存;
  • 管道化(pipelining)multiplex开启时,多个 fiber 的语句在同一连接上被合并为单次 socket 写入(见 PgPool.ts 与PgConnectionImpl中的pipelinePending/pipelineInFlight);
  • 流式(streaming):结果不整体收集,会话在流存续期间被钉住(pinned);
  • 通知(notifications)LISTEN/NOTIFY的原生实现;
  • 取消(cancellation):通过旁路连接发送CancelRequest中断活动查询;
  • 自定义编解码器(custom codecs):通过PgTypes.Registry注册自定义类型编解码。

伴随的还有effect包中Pool的修复:新增Pool.reserve用于对池中一个条目进行独占借用,并修复了失效(invalidation)后等待者唤醒与容量替换的问题(该修复在effect包中落地,@effect/sql-pgPgPool.reserve直接建立其上)。

这一改动的直接后果是一组破坏性变更,下面逐一展开。

2. 构造方式:makemakeClient取代fromPool/fromClient/makeWith

2.1 已移除的旧构造器

changeset 明确列出:fromPoolfromClientmakeWith三个构造器被移除。搜索packages/sql/pg/src目录可以确认,旧构造器只剩PgPool.ts内的残留引用,PgClient.ts已不再导出它们。

升级动作:使用make创建基于连接池的客户端,使用makeClient创建基于单连接的客户端。

2.2PgClient.make:连接池客户端

make的签名(见 PgClient.ts):

export const make = ( options: PgPoolConfig ): Effect.Effect<PgClient, SqlError, Scope.Scope | Reactivity.Reactivity>

它内部先PgPool.make(options)建池,再把池的get(普通借用)、use(临时借用)、reserve(独占借用,用于事务)与pool.reserve(用于监听)组装成一个完整的PgClient。返回的 effect 是 scoped 的:作用域关闭时池被关闭、所有会话被释放。

典型用法:

import { Effect, Layer, Redacted } from "effect" import * as PgClient from "@effect/sql-pg/PgClient" const SqlLive = PgClient.layer({ host: "localhost", port: 5432, database: "app", username: "app_user", password: Redacted.make("secret"), maxConnections: 10, minConnections: 0, idleTimeout: "10 seconds" }) const program = Effect.gen(function*() { const sql = yield* PgClient.PgClient const rows = yield* sql<{ value: number }>`SELECT 1 AS value` return rows }) program.pipe(Effect.provide(SqlLive), Effect.runPromise)

2.3PgClient.makeClient:单连接客户端

makeClient(见 PgClient.ts)只建立一个PgConnection,不建池:

export const makeClient = ( options: PgClientConfig & { readonly acquireForStream?: boolean | undefined } ): Effect.Effect<PgClient, SqlError, Scope.Scope | Reactivity.Reactivity>

其语义要点:

  • 查询在主连接上顺序执行;事务通过connection.pin独占主连接,因此单连接上的事务是串行化的;
  • 默认情况下流(stream)也会钉住主连接——测试 Client.integration.test.ts 验证了“流激活期间普通查询会被阻塞”这一行为;
  • 设置acquireForStream: true后,每个流和监听器会额外打开一条独立连接,主连接上的查询在流活跃期间仍可执行(测试见 Client.integration.test.ts)。源码中streamAcquirerlistenAcquirer都根据该开关决定是复用主连接还是新建PgConnection.make(options)
import { Effect, Redacted } from "effect" import * as PgClient from "@effect/sql-pg/PgClient" const program = Effect.gen(function*() { const sql = yield* PgClient.makeClient({ url: Redacted.make("postgres://app_user:secret@localhost:5432/app"), acquireForStream: true }) const rows = yield* sql`SELECT generate_series(1, 3) AS value`.stream.pipe( Stream.runCollect ) return rows })

2.4 Layer 辅助

PgClient同时提供三个 Layer 构造(见 PgClient.ts):

  • layer(config):从PgPoolConfig直接建层;
  • layerConfig(config):从Config.Wrap<PgPoolConfig>(effect 的配置描述子)建层;
  • layerFrom(acquire):从任意Effect<PgClient, E, R>建层,同时提供PgClientSqlClient两个服务标签。

测试工具 utils.ts 正是用PgClient.layerFrom(PgClient.makeClient(...))组装出测试容器,说明layerFrom也是把自定义获取逻辑接入依赖注入的推荐通道。

3. 连接启动细节:URL、Unix Socket、TLS 与认证

PgConnectionConfig支持urlhostportpathssldatabaseusernamepasswordconnectTimeoutapplicationNamestream等字段(见 PgConnection.ts),启动逻辑遵循以下规则:

  • url按 libpq URI(postgres://postgresql://)解析,显式字段(host/port/database等)优先于 URL 中的值;
  • stream工厂优先于hostportpath:传入stream: () => Duplex时完全由调用方提供传输层(测试 PgConnection.in-process.test.ts 正是用内存Duplex完成进程内集成测试);
  • path原样作为 Unix socket 路径;若host/开头,则视为 socket 目录并展开为${host}/.s.PGSQL.${port}
  • sslmode=prefersslmode=allow先尝试 TLS,仅当服务端对SSLRequest回复N时回退明文;与 libpq 不同,allow同样优先 TLS。证书校验默认开启,除非通过ssl选项显式关闭;
  • 传输、可选SSLRequest、启动与认证握手全程受connectTimeout约束,默认 5 秒;超时抛出ConnectionError(reason 为"connect"操作)。make内部以acquireRelease包裹会话,作用域关闭时发送Terminate并结束 socket(见 PgConnection.ts);
  • 使用 Unix socket 或自定义stream时,应显式设置ssl.servername

认证方面,原生客户端实现了MD5 与 SCRAM-SHA-256(见 PgAuth.ts):明文认证直接发送PasswordMessage;SCRAM 交换由客户端计算 nonce、盐与验证消息,服务端迭代次数超过1,000,000会被拒绝;SCRAM-SHA-256-PLUS(需要 TLS 通道绑定的变体)未实现,且密码按 UTF-8 直接使用、不做 SASLprep 归一化。

测试 PgConnection.integration.test.ts 对真实 PostgreSQL 容器验证了连接建立与prepare: false场景下的查询,是了解端到端启动流程的参考。

4. 二进制编解码与结果形状:int8bigint

新客户端结果统一走原生二进制编解码器PgTypes,v1 仅实现二进制线格式,format = 0解码直接报错)。对使用者影响最大的行为变化:

数据库类型解码结果说明
int8(bigint)bigint不再像文本协议那样可能得到字符串或number
date字符串不以Date返回
timestamp/timestamptzUnix 纪元毫秒(number 或超出精确范围时回退biginttimestamp在线上无时区,双向均按 UTC 处理;解码向零截断丢弃亚毫秒精度
bytea或未知 OIDUint8Array字节原样保留
其余常见类型按 OID 编解码,见PgTypes.OID常量表boolint2/int4text/varcharfloat4/float8json/jsonbinet/cidrtime
// int8 列解码为 bigint const rows = yield* sql<{ id: bigint }>`SELECT 1::int8 AS id` // rows[0].id === 1n

另外:

  • executeRaw现在返回原生PgConnection.Result形状{ command, rowCount, oid, rows, fields },见 PgConnection.ts),不再是pg.Result
  • PgClientConfig.types改为接收PgTypes.Registry,取代旧的pg.CustomTypesConfig;自定义编解码器通过 Registry 注册 OID 与编码/解码函数。

4.1 参数推断保持宽松

changeset 特别说明:推断参数仍然宽松。源码 PgConnection.ts 的inferScalar展示了推断规则:

  • null/undefined绑定为未指定 OID(oid = 0);
  • booleanbool
  • bigintint8
  • 整数number:落在int4范围内按int4,超出int4范围但在安全整数范围内则作为int8绑定(转为bigint,非安全整数走float8
  • 字符串按“未指定类型”绑定(文本格式字面量),由后端根据语句上下文推导类型——因此一个字符串既能匹配bigint列也能匹配timestamp列,行为与旧的文本协议驱动一致;
  • Datetimestamptz(毫秒时间戳);Uint8Array/Int8Arraybytea
  • 数组参数要求元素类型一致且不允许嵌套,空数组无法推断元素类型时会抛出CodecError(此时需用PgTypes.array显式指定)。

5. 预处理语句:默认开启、prepare: falseunprepared

5.1 命名预处理默认启用

新客户端命名预处理语句默认开启,每个连接维护一个有界缓存:

  • 缓存默认上限100preparedStatementCacheSize,PgConnection.ts 中defaultPreparedStatements = 100);
  • 缓存以SQL 文本 + 推断出的参数 OID 列表为键(${sql}\u0000${oids.join(",")}),同一文本配不同类型参数在服务端是不同语句;
  • 命中缓存的语句只发送Bind/Execute/Sync,跳过ParseDescribe(列信息首次执行时已随名缓存);被逐出的语句以Close消息搭车,不额外占用往返;
  • 服务端丢失语句或计划与列不匹配(SQLSTATE26000/0A000)时,会从缓存逐出并以未命名语句重试一次,重试跳过缓存因此不会死循环(retryStale,PgConnection.ts)。

5.2 何时必须prepare: false

在使用**无法在查询之间保留预处理语句的连接池中间件(statement-mode pooler)**或会产生大量唯一 SQL 的工作负载时,设置prepare: false

典型场景是 PgBouncer 的 statement 模式、以及需要直连但 SQL 高度动态化的场景。prepare: false会让每个查询走未命名路径(Parse/Bind/Describe/Execute/Sync),同时不再向预处理缓存添加条目。该选项可通过PgClientConfig.preparePgPoolConfig.prepare设置;preparedStatementCacheSize: 0同样等效于关闭缓存(PgConnection.ts)。

5.3Statement.unpreparedStatement.valuesUnprepared

为了在“默认预处理”的世界里仍能对个别语句精确控制,语句对象暴露了两个跳过预处理的执行入口(见 Statement.ts):

  • statement.unprepared:执行该语句但不走命名预处理(返回转换后的对象行);
  • statement.valuesUnprepared:同上,但返回位置值数组。

两者底层走executeUnprepared/executeValuesUnprepared(对应 PgClient.ts 中prepare = false的调用),不使用也未向预处理缓存添加条目

集成测试 Client.integration.test.ts 精确验证了这一点:先对两条语句执行unprepared/valuesUnprepared,查询pg_prepared_statements结果为空;再执行普通版本后,两条语句都出现在服务端预处理列表中。

6.LISTEN/NOTIFY:scoped 队列取代Stream

6.1 新签名

changeset 列出的第二个破坏性变更:

PgClient.listen现在返回 scoped 的Effect<Dequeue<string>, SqlError, Scope>,取代原来的Stream

准确地说,PgClient.listen返回Effect<Queue.Dequeue<PgConnection.Notification>, SqlError, Scope.Scope>,其中Notification = { processId, channel, payload }(见 PgConnection.ts),payload 可从notification.payload读取。

关键语义:获取(acquisition)在 PostgreSQL确认LISTEN之后才完成——即返回的 effect 成功后,通道订阅已经生效,因此在返回之后发送的通知不会被遗漏。监听期间连接被独占钉住,作用域关闭时发送UNLISTEN并关闭队列;注册期的服务端错误会使获取 effect 失败。

6.2 用法示例

import { Effect, Queue, Scope } from "effect" import * as PgClient from "@effect/sql-pg/PgClient" const listenDemo = Effect.gen(function*() { const sql = yield* PgClient.PgClient const channel = "order_events" const payloads = yield* sql.listen(channel) // LISTEN 已确认,队列可取 yield* sql.notify(channel, JSON.stringify({ orderId: 42 })) const notification = yield* Queue.take(payloads) console.log(notification.channel, notification.payload) }).pipe(Effect.scoped)

sql.notify(channel, payload)通过SELECT pg_notify($1, $2)发送(PgClient.ts)。监听与通知都会校验通道名,超过 63 个 UTF-8 字节的通道名会被拒绝(测试见 Client.integration.test.ts)。

6.3 池化场景的行为

在连接池上调用listen时,监听器会通过pool.reserve独占借用一条连接,普通查询仍可使用池中其余连接(测试 Client.integration.test.ts 用两个池连接的容器验证了“监听器占用一条连接时查询依旧可用”)。PgPool.reserve的独占语义因此是监听、事务、多语句工作的推荐入口。

7.sql.json:JSON 参数必须显式包装

破坏性变更第三条:

普通对象参数不再被推断为 JSON;请用sql.json包装。

也就是说,旧驱动中“传普通对象即序列化为 JSON”的隐式行为被移除。现在向json/jsonb列写入对象必须显式调用sql.json(value),它内部以PgTypes.jsonb自定义片段编码(见 PgClient.ts 的onCustom分支)。

const sql = yield* PgClient.PgClient // 正确 yield* sql`INSERT INTO events (payload) VALUES (${sql.json({ kind: "order_created", id: 1 })})` // 旧写法(普通对象)不再被推断为 JSON

测试同样覆盖了sql.jsoninsert片段、updateValues片段以及::jsonb列查询中的用法(Client.integration.test.ts)。PgClientConfig.transformJson(默认 true)控制 JSON 值是否经过transformResultNames/transformQueryNames的转换。

8. 查询字符串:单语句约束

查询字符串必须只包含一条语句——PostgreSQL 扩展协议拒绝多语句字符串。

这是因为每条查询都作为一次扩展查询循环(一个Parse/Bind/Execute/Sync)发送,而服务端在该协议下不支持一次执行多条语句。集成测试 Client.integration.test.ts 验证了含两条语句的 SQL 会以SqlSyntaxError失败。

// 会失败:多语句字符串 yield* sql`CREATE TABLE a (id int); SELECT 1` // 应拆成两条独立语句分别执行

批量 DDL/多语句迁移请交给PgMigrator(见 PgMigrator.ts),它按迁移文件逐条执行。

9. 流式查询与取消

PgConnection.stream以流式返回行而不整体收集结果(PgConnection.ts):

  • 会话在流存续期间被钉住pin);
  • 流在结果完成前被中止时,会通过旁路连接发送CancelRequest取消该语句,并把连接排干回到ReadyForQuery状态(abortDrainTimeoutMillis = 5000cancelRequestTimeoutMillis = 5000);
  • 未钉住的复用连接(multiplex 开启、未独占)上的取消是 no-op——因为活动查询可能属于其他 fiber,interrupt因此在multiplex && !pinned时直接返回(PgConnection.ts)。
import { Stream } from "effect" const rows = yield* sql`SELECT generate_series(1, 1000000) AS n`.stream.pipe( Stream.take(10), Stream.runCollect )

10. 连接池行为:Pool.reserve、TTL 与失效恢复

10.1 池默认值

PgPool.make(见 PgPool.ts)建立的池采用以下默认:

配置项默认值说明
maxConnections10池最大连接数
minConnections0空闲时保留的最小连接数(惰性建连,idleTimeout后释放至下限)
idleTimeout10 秒空闲超时(timeToLiveStrategy: "usage",按使用计)
connectionTTL未设置(不限)超过寿命的连接被替换;每条连接至少使用一次,TTL 为 0 表示禁用复用
multiplexfalse开启后 fiber 可共享池连接做管道化查询,multiplexConcurrency默认 32
multiplexConcurrency32一条连接上最多共享的并发语句数

开启multiplex时,池的并发度被设为multiplexConcurrency(至少 1),同一连接上的多个语句被queueMicrotask合并成一次 socket 写(flushPipeline);但这意味着慢语句会阻塞排在它后面的语句(以吞吐换尾延迟)。事务、流、监听器依然通过reserve独占连接

10.2Pool.reserve:独占借用

changeset 末尾声明:

新增Pool.reserve,用于独占借用池中一个并发条目,并修复失效后的等待者唤醒与容量替换。

PgPool.reserve(PgPool.ts)实现为get之后调用connection.pinpin本身通过pool.reserve预留池条目并取得独占视图(PgConnection.ts),事务、流、监听在multiplex池上必须走这条路径。配套修复确保:致命协议/socket 错误自动失效连接(fatalHooksPool.invalidate),等待中的借用者被及时唤醒,池容量在失效后正确补充(deadConnections集合在下次借用前被逐批清理)。

// 多语句工作(如事务)在 multiplex 池上应使用 reserve const conn = yield* sql.reserve yield* conn.executeRaw("BEGIN", []) try { yield* conn.executeRaw("UPDATE accounts SET balance = balance - 100 WHERE id = $1", [1]) yield* conn.executeRaw("COMMIT", []) } catch { yield* conn.executeRaw("ROLLBACK", []) }

单连接客户端(makeClient)中的事务通过connection.pin独占主连接并串行化;池客户端(make)的事务则通过pool.reserve独占某条连接。

11. 升级检查清单

综合以上变更,从旧版@effect/sql-pg升级时应逐项核对:

  1. 构造器fromPool/fromClient/makeWithmake(池)/makeClient(单连接);
  2. 监听listen返回值由Stream变为 scoped 的Effect<Dequeue<Notification>, SqlError, Scope>,需用Effect.scoped包裹;事件改用Queue.take消费,payload 取notification.payload
  3. types配置pg.CustomTypesConfigPgTypes.Registry
  4. JSON 参数:普通对象必须用sql.json(...)包装;
  5. 单语句:拆分任何多语句 SQL;
  6. 结果类型int8bigintdate→ string、时间戳 → Unix 毫秒、bytea/未知 OID →Uint8ArrayexecuteRawPgConnection.Result
  7. 预处理:默认开启;使用 statement-mode pooler 时设prepare: false;个别语句用.unprepared/.valuesUnprepared
  8. 事务/监听:在 multiplex 池上用reserve保证独占。

适用前提:本变更说明对应@effect/sql-pg(当前仓库版本4.0.0-rc.115)的原生协议客户端实现,代码位于 packages/sql/pg/src,相关行为均有集成测试佐证(见 packages/sql/pg/test)。若仍需使用旧的pg运行时,请参考历史版本或迁移说明(MIGRATION.md、migration/)。

【免费下载链接】effectBuild production-ready applications in TypeScript项目地址: https://gitcode.com/GitHub_Trending/ef/effect

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询