@effect/sql-pg 全解析:基于 Effect 的 PostgreSQL 客户端架构与演进史
【免费下载链接】t3code项目地址: https://gitcode.com/GitHub_Trending/t3/t3code
本文以@effect/sql-pg的官方变更日志(.repos/effect-smol/packages/sql/pg/CHANGELOG.md)为骨架,结合其源码(PgClient.ts、PgProtocol.ts、PgTypes.ts、PgAuth.ts、PgMigrator.ts),系统梳理这个把 PostgreSQL 接入 Effect 类型化错误体系、依赖注入与资源管理模型中的官方 SQL 客户端:你将掌握它的连接池与单客户端构造方式、完整配置项、LISTEN/NOTIFY 与 JSON 片段等 PostgreSQL 专属能力、SQLSTATE 到类型化错误的分类映射,以及 v4 新增的低层协议(protocol 3.0)、二进制类型编解码与 MD5/SCRAM-SHA-256 认证实现,并能从源码级理解每一次关键变更背后的设计动机。
一、包定位:一个「Effect SQL 驱动的 PostgreSQL 客户端」
@effect/sql-pg是 Effect 生态中面向 PostgreSQL 的官方 SQL 客户端包。它的定位在 README 中一句话说得很清楚:built on thepglibrary,即运行时底层复用 Node.js 生态最成熟的pg驱动,而上层完全以 Effect 的模型(Effect、Layer、Scope、Stream、Config、Redacted)重新组织连接、查询、事务与流式读取。对应地,package.json把pg、pg-pool、pg-cursor、pg-connection-string、pg-types列为运行时依赖,而effect是唯一的 peerDependency。
安装方式(来自 README):
npm install effect@rc @effect/sql-pg@rc包名从 changelog 中可以读出清晰的版本脉络:0.1.0随 Effect 3.0 一同发布,随后沿0.x推进;在 4.0.0 大版本上以4.0.0-beta.N(beta)与4.0.0-rc.N(候选)两条线演进,最新的4.0.0-rc.112与effect@4.0.0-rc.112同步发布。模块导出采用扁平结构(见 index.ts),依次暴露PgClient、PgProtocol、PgTypes、PgAuth、PgMigrator五个命名空间。
版本历史中值得注意的一个节点:
0.2.0起@effect/sql变得「dialect agnostic」,所有客户端实现共享同一个Context.Tag,因此你可以编写同时支持多种 SQL 方言的服务,仅在需要实现特有功能时(如 PostgreSQL 的LISTEN/NOTIFY)才从本包取专用 Tag。
二、核心服务PgClient:接口、配置与构造方式
2.1 服务接口
源码中PgClient接口(PgClient.ts)在通用SqlClient之上追加了 PostgreSQL 专属成员:
export interface PgClient extends Client.SqlClient { readonly [TypeId]: TypeId readonly config: PgClientConfig readonly json: (_: unknown) => Fragment // 构造 JSON 参数片段 readonly listen: (channel: string) => Stream.Stream<string, SqlError> readonly notify: (channel: string, payload: string) => Effect.Effect<void, SqlError> }json:把任意值包装成 JSON 参数片段,由编译器生成$N占位符并序列化;listen/notify:基于 PostgreSQL 的 LISTEN/NOTIFY 机制实现通道订阅与消息推送(见第五节)。
2.2 配置项全表
PgClientConfig(PgClient.ts)与PgPoolConfig(同文件 L135-L141)是两个核心配置模型,其字段如下:
| 配置项 | 所属 | 类型 | 说明 |
|---|---|---|---|
url | Client | Redacted.Redacted | 连接串(connection string),运行时经Redacted.value取出使用 |
host/port/path | Client | string/number/string | 连接地址;path对应 unix socket 场景 |
ssl | Client | boolean \| ConnectionOptions | TLS 开关或完整tls.ConnectionOptions(0.14.1 起支持传入 TLS 选项对象) |
database/username/password | Client | string/string/Redacted.Redacted | 认证凭据,密码用Redacted包装避免意外泄漏 |
connectTimeout | Client | Duration.Input | 连接超时(默认 5 秒,见 0.24.3 变更) |
stream | Client | () => Duplex | 自定义网络流(0.6.3 新增) |
applicationName | Client | string | 映射为connection.application_name(0.6.3 新增,默认"@effect/sql-pg") |
spanAttributes | Client | Record<string, unknown> | 附加到观测 span 的属性(0.3.1 起可传入) |
transformResultNames/transformQueryNames | Client | (str: string) => string | 结果列名 / SQL 标识符的命名变换函数 |
transformJson | Client | boolean | 是否对 JSON 参数应用命名变换(默认true) |
types | Client | Pg.CustomTypesConfig | pg自定义类型解析器配置(0.1.17 新增) |
idleTimeout | Pool | Duration.Input | 池中空闲连接回收时间 |
maxConnections/minConnections | Pool | number | 连接池上/下限 |
connectionTTL | Pool | Duration.Input | 连接最大存活时长(maxLifetimeSeconds) |
这些字段在make(连接池路径)与makeClient(单客户端路径)中会逐一映射为pg.Pool/pg.Client的原生配置,例如connectTimeout经Duration.toMillis转成connectionTimeoutMillis,connectionTTL经Duration.toSeconds转成maxLifetimeSeconds。
2.3 五种构造方式与 Layer
PgClient提供了一组从易到难、可组合的构造器(PgClient.ts):
| 构造器 | 底层形态 | 说明 |
|---|---|---|
make(options: PgPoolConfig) | pg.Pool | 托管连接池,最常用;构造时先执行SELECT 1探活 |
makeClient(options) | pg.Client | 单个托管客户端,可选acquireForStream让流式/订阅操作使用独立连接 |
fromPool({ acquire }) | 外部pg.Pool | 由你提供的池构建客户端,派生事务、流式与 LISTEN/NOTIFY 支持 |
fromClient({ acquire, acquireForStream }) | 外部pg.Client | 由你提供的客户端构建,用信号量串行化共享访问 |
makeWith(...) | 自定义连接获取器 | 完全自定义的 acquirer / transactionAcquirer / listenAcquirer |
对应地有三个 Layer:
layer(config):接收裸配置对象(0.19.0 起采用layer/layerConfig命名约定,layer直接收裸对象);layerConfig(config: Config.Wrap<PgPoolConfig>):接收Config.Config,便于从环境变量等来源读取配置;layerFrom(acquire):从任意的PgClient获取 Effect 构建 Layer。
三者都会同时提供PgClient与通用SqlClient两个服务标签。layerConfig的失败通道是Config.ConfigError | SqlError,layer则是SqlError。
2.4 连接生命周期与健壮性细节
从 changelog 与源码对照,连接管理经历了多轮打磨:
- 0.15.2「把连接测试纳入客户端构造」:
make在 acquire 阶段执行SELECT 1,若失败则分类为SqlError(reason 为"connect"),并用Effect.timeoutOrElse兜底超时(默认 5 秒,0.24.3 起connectTimeout可配置); - beta.99「修复
PgClient.makeClient的连接时序」:单客户端路径在资源获取阶段即调用client.connect(),确保拿到手的就是已连接对象; - beta.106「防止
makeClient连接期间的未处理 error 事件」:源码中通过client.on("error", onError)提前挂接空处理器(onError() {}),避免pg客户端在无人监听 error 时把异常抛向全局;释放时再off("error", onError); - 关闭兜底:pool/client 的关闭操作都套了
Effect.timeoutOption(1000),防止资源释放被卡死。
三、语句编译、类型化错误与 SQLSTATE 分类
3.1 编译器:$N占位符与 PostgreSQL 方言
makeCompiler(PgClient.ts)构建方言编译器:占位符统一为 PostgreSQL 的$N形式,标识符用双引号转义(defaultEscape("\"")),INSERT ... ON CONFLICT与RETURNING子句按 PostgreSQL 语法生成;onCustom分支处理PgJson自定义片段——这正是sql.json(value)的底层实现。
3.2 reason-based 错误模型
0.37.0 将SqlError重构为「reason-based」结构:所有 SQL 驱动把数据库原生失败归类为结构化 reason,拿不到原生错误码时回退Unknown。在@effect/sql-pg中,这个映射由classifyError(PgClient.ts)基于 PostgreSQL 的 SQLSTATE 代码完成:
| SQLSTATE 前缀/值 | 分类结果(reason) |
|---|---|
08* | ConnectionError(连接失败) |
28* | AuthenticationError(认证失败) |
42501 | AuthorizationError(权限不足) |
42* | SqlSyntaxError(语法错误) |
23505 | UniqueViolation(唯一约束违反,见下) |
23* | ConstraintError(其余约束违反) |
40P01 | DeadlockError(死锁) |
40001 | SerializationError(序列化失败) |
55P03 | LockTimeoutError(锁等待超时) |
57014 | StatementTimeoutError(语句超时) |
| 其他/无码 | UnknownError |
3.3UniqueViolation的细化
4.0.0-beta.65 新增UniqueViolation作为独立错误 reason:受支持的唯一约束冲突从宽泛的ConstraintError中拆出,适用于 PostgreSQL、PGlite、MySQL、MSSQL 以及 SQLite 家族共享的分类。UniqueViolation.constraint字段承载「能拿到的最佳约束/索引/键标识」,拿不到可靠标识时精确回退为字符串"unknown"。源码中pgConstraintFromCause(PgClient.ts)负责从pg错误对象读取constraint字段并做 trim 归一化。
四、LISTEN/NOTIFY:订阅式消息通道
listen/notify是PgClient最典型的 PostgreSQL 专属能力(0.8.2 首次加入),其实现历经三次关键改造:
- 0.8.2:新增
listen/notify; - 4.0.0-beta.38:
notify改用pg_notify($1, $2)函数,通道与负载都作为参数传递,而非拼接成NOTIFY语句字符串,从根上消除了通道名/负载注入 SQL 的隐患。源码印证见 makeWith 的 notify 实现(SELECT pg_notify($1, $2)); - 4.0.0-beta.37:
LISTEN / UNLISTEN订阅改为使用专用 PostgreSQL 客户端,而不是占用一个池化连接在整个监听生命周期内——这样订阅不再挤占连接池配额,也不会因池连接被回收而中断监听。源码中fromPool通过RcRef懒加载一个独立new Pg.Client(pool.options)作为listenAcquirer(PgClient.ts),订阅结束的 finalizer 中执行UNLISTEN并移除notification监听器。
实际使用形态:listen(channel)返回Stream.Stream<string, SqlError>,你可以在 Effect 程序中订阅并消费;notify(channel, payload)是Effect<void, SqlError>。
五、事务、流式查询与资源获取
5.1 事务期间持有连接许可
4.0.0-beta.104 明确「为整个事务生命周期持有共享的 PostgreSQL 客户端许可」。在fromClient路径中可以看到对应实现:普通查询经semaphore.withPermit串行化,而事务获取使用Effect.uninterruptibleMask先semaphore.take(1),再把semaphore.release(1)注册为 Scope finalizer(PgClient.ts)——信号量许可伴随整个事务作用域,而不是查询结束即释放。
5.2 事务连接获取失败保持在错误通道内
4.0.0-beta.44 修复PgClient.fromPool的事务连接获取:pool.connect回调中的失败统一转换为SqlError(reason 为acquireConnection),保证失败始终落在类型化错误通道,而不会以未处理异常的形式泄漏。reserveRaw(PgClient.ts)即此路径的实现,它还处理了「回调返回空客户端」「已结束但仍回调」等边界。
5.3 流式结果与取消
executeStream基于pg-cursor实现:每条流式查询先reserve一个连接,创建Cursor,以 128 行为批次cursor.read推送数据,结束时cursor.close()(PgClient.ts)。与之配套的makeCancel使用pg_cancel_backend(processId)做尽力而为的查询取消(带 5 秒超时),取消失败不报错。
六、v4 新增的低层协议栈:PgProtocol / PgTypes / PgAuth
4.0.0-rc.112是本包最重要的里程碑之一:新增了低层 PostgreSQL 协议、二进制类型编解码与认证 codec,对应的三个新模块组成一条完整、自洽的「裸协议栈」:
6.1PgProtocol:protocol 3.0 的消息编解码
PgProtocol.ts负责 protocol 3.0 的编码与增量解析,全部为纯函数:「bytes in, bytes or plain data out」。关键设计(源码顶部注释可印证):
- 消息格式:1 字节类型 + 4 字节
int32长度(长度字段包含自身但不包含类型字节)+ 负载;整数大端序,字符串为 NUL 结尾的 UTF-8; makeParser提供有状态增量解析器:push(chunk)返回当前完整的所有消息,未凑整的尾部消息保留到后续字节到达;默认maxMessageSize为 16 MiB(defaultMaxMessageSize常量);解析或字段读取错误是终态错误——parser 抛错后不可复用,且当次 push 已解码的消息会被丢弃;- 缓冲池机制:已交给调用方的字节永不重写,因此
DataRow字段可以零拷贝地以view 形式返回(每缓冲一次分配而非每列一次分配);池上限 64 KiB,翻倍增长;代价是持有 view 会连带持有整个池缓冲,凡是要活得比消息更久的数据必须拷贝; - 启动握手前特殊应答(无类型字节的 SSL 响应)由
decodeSslResponse单独处理,不进入通用 parser。
6.2PgTypes:按 OID 的二进制编解码
PgTypes.ts实现二进制线格式(format = 1)的标量与一维数组 codec,布局对齐 rust-postgres 的postgres-types,包含 infinity 哨兵;假设服务端以integer_datetimes构建(PostgreSQL 10 起唯一受支持的配置)。要点:
- 不做类型推断:OID 必须显式给出(直接传 OID,或通过
int4这类自带 OID 的构造器); - 错误模型:公开 codec 返回类型化
Result失败(PgTypesCodecError),而 parser 的字段读取器走内部抛异常的快路径;timestamp在线路上无时区,双向均按 UTC 处理,解码向零截断亚毫秒精度; - 性能取向:用预分配 scratch buffer(
scratch4/scratch8)代替逐字节 DataView 包装,用两位数字表避免padStart,ASCII 字符串编码按 48 字符区分「逐字符循环」与encodeInto两条路径(V8 实测交叉点约 50 字符)。
6.3PgAuth:MD5 与 SCRAM-SHA-256
PgAuth.ts实现协议认证交换中的 MD5 与 SCRAM-SHA-256,均返回类型化Result(PgAuthError)。SASL 帧结构属于PgProtocol,本模块只负责帧内部的计算(createHash/createHmac/pbkdf2Sync+ XOR)。两个明确边界(源码注释):SCRAM-SHA-256-PLUS(通道绑定)未实现,因为通道绑定需要 TLS socket 而 codec 不持有它;密码按 UTF-8 直接使用、不做 SASLprep 归一化,因此需要归一化的非 ASCII 密码不受支持。明文认证无需本模块——直接发送PasswordMessage即可。
重要澄清:rc.112 中PgClient本身保持原样,运行时依然走pg;新协议栈是独立交付的低层能力,供需要自行实现客户端、或希望脱离pg的定制场景使用。
七、迁移工具PgMigrator与语句辅助能力
- 4.0.0-beta.8:实现
PgMigrator(同时修复ChildProcess选项类型)。PgMigrator.ts复用effect/unstable/sql/Migrator,暴露run与layer:run用当前 SQL 客户端执行待应用的迁移文件;schema dump 走pg_dump子进程(附带--no-owner --no-privileges,并通过PGHOST/PGPORT/PGUSER/PGPASSWORD环境变量传连接信息),因此run依赖ChildProcessSpawner、FileSystem、Path等服务。 - 4.0.0-beta.86:新增
Statement.valuesUnprepared——把未预处理的 SQL 语句行以数组形式返回(与executeValues的rowMode: "array"对应,见 PgClient.ts)。 - 0.1.11:新增
sql\...`.unprepared,执行不尝试PREPARE的查询,并加入 SQL 事务 tracing span、把 span 属性对齐语义约定(0.2.9 起改用@opentelemetry/semantic-conventions常量;0.43.0 又随上游把属性名升级为db.system.name、db.namespace等新一代命名——PgClient.ts中的 span 属性常量ATTR_DB_SYSTEM_NAME = "db.system.name"、ATTR_DB_NAMESPACE = "db.namespace"` 正是这一演进的结果)。 - 0.1.17:
PgClientConfig增加prepare与types。
八、API 与模块演进的代表性节点
从 changelog 可以梳理出几条贯穿始终的设计主线:
- 方言无关化(0.2.0):
@effect/sql改为方言无关,客户端共享同一Context.Tag;需要方言特定能力时再用实现包专属 Tag(如PgClient),或用sql.onDialect({ pg: ... })按方言分支。 - 扁平导入与命名约定(0.4.0 扁平化;0.19.0
layer/layerConfig约定):导入路径与构造器命名走向统一规范。 - 模块重命名(4.0.0-beta.44):
ServiceMap模块重命名为Context,贯穿导出、文档与测试。 - 入口点精简(4.0.0-beta.103):移除显式
./index入口点。 - v4 大版本(4.0.0-beta.0):以
v4 beta标记的 Major Changes 开启 4.0 系列;随后 beta(功能推进)→ rc(候选发布)两个阶段,直至4.0.0-rc.112。
九、测试与基准验证
仓库为上述能力提供了完备的测试与基准支撑(见 test/ 与 benchmark/):
PgProtocol.test.ts/PgTypes.test.ts/PgAuth.test.ts:覆盖协议解析、二进制编解码与认证交换,fixtures 目录(goldens.ts)保存黄金样本;SqlErrorClassification.test.ts/TransactionAcquire.test.ts:验证 SQLSTATE 分类映射与事务连接获取的错误通道;Client.integration.test.ts、KeyValueStore.integration.test.ts、Persistence.integration.test.ts等集成测试基于@testcontainers/postgresql起真实 PostgreSQL 容器(见 package.json 的 devDependencies);- 基准:
pnpm benchmark:codec运行 benchmark/PgCodec.ts(基于 tinybench),量化 codec 性能。
十、总结
@effect/sql-pg的演进史本质上是一部「如何在保证类型安全与资源安全的前提下,把成熟驱动(pg)全面纳入 Effect 生态」的工程实践记录:从0.1.0的初版连接管理,到0.19.0的统一layer约定,再到 v4 时代原生协议栈(PgProtocol/PgTypes/PgAuth)的独立交付。理解这条脉络,不仅能帮你正确配置PgClientConfig/PgPoolConfig并合理选用make/makeClient/fromPool等构造路径,也能让你在遇到连接超时、唯一约束冲突、LISTEN/NOTIFY 通道行为异常等问题时,快速定位到对应的 SQLSTATE 分类逻辑与资源获取实现,真正做到「知其然,更知其所以然」。
【免费下载链接】t3code项目地址: https://gitcode.com/GitHub_Trending/t3/t3code
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考