- 消息队列
- 后端
- 通信
【免费下载链接】centrifugo
Scalable real-time messaging server in a language-agnostic way. Self-hosted alternative to Pubnub, Pusher, Ably, socket.io, Phoenix.PubSub, SignalR. Set up once and forever.
导读
本文围绕 Centrifugo 项目中 PostgreSQL MapBroker 的 schema 迁移体系展开,系统讲解迁移文件的命名与编写规范、整数版本号驱动的自动升级机制、以及从新增一列到函数签名变更等各类变更的落地路径。读完本文,你将掌握如何为internal/pgmapbroker/的存储层新增一次安全、幂等、兼容滚动部署的数据库迁移,并理解EnsureSchema()在启动时自动建表、升级与自愈的完整原理。
本文的权威依据来自 migrations/README.md 与同目录的 postgres_schema.md,并以 pgmapbroker.go、schema.sql、pgschema.go 及 schema_template_test.go 等源码与测试作为实现佐证。
一、背景:MapBroker 的 Schema 从何而来
PostgreSQL MapBroker 是 Centrifugo 中基于 PostgreSQL 实现的 MapBroker(映射型实时数据存储),提供持久化订阅、CAS 操作与事务化发布能力,典型场景包括协作白板、文档协同、库存/预订系统和带持久状态的游戏大厅。
其数据库对象(表、索引、函数)并非由外部 SQL 脚本手动维护,而是由PostgresMapBroker.EnsureSchema()在启动时自动管理:
- 一次调用同时创建JSONB 与 BYTEA 两套 schema 变体(与运行时
BinaryData配置无关,两套都会被创建); - 使用整数版本号记录 schema 版本;
- 支持正向迁移(forward migrations)。
核心实现位于 pgmapbroker.go 的EnsureSchema方法,它内部复用internal/pgschema包提供的共享原语(版本读取、迁移锁、事务化迁移应用、降级拒绝等),保证所有 PostgreSQL 后端(MapBroker、StreamBroker、控制器)的 schema 管理行为一致。
二、迁移文件约定(README 核心规范)
migrations/README.md定义了迁移文件的硬性约定,这是所有迁移作者必须遵守的底线:
- 文件命名:文件名为
NNN.sql,其中NNN是目标版本号。例如002.sql表示将 schema 升级到版本 2。迁移从版本 2 开始,版本 1 是基线(由完整 DDL 应用)。 - 幂等性:每个迁移必须幂等——使用
IF NOT EXISTS、IF EXISTS等构造,允许重复执行而不会报错或产生副作用。 - 双前缀覆盖:每个迁移必须同时作用于两套前缀的表:
cf_map_*(JSONB 变体)与cf_binary_map_*(BYTEA 变体)。注意cf是默认table_prefix,用户可配置为其他前缀。 - 无模板占位符:迁移文件内不得使用模板占位符,必须写显式表名——生产环境迁移中显式更安全。
标准迁移示例(摘自 README):
-- Migration 002: Add foo column ALTER TABLE cf_map_state ADD COLUMN IF NOT EXISTS foo TEXT; ALTER TABLE cf_binary_map_state ADD COLUMN IF NOT EXISTS foo TEXT;完整开发者工作流见 postgres_schema.md。
三、EnsureSchema 自动处理哪些变更
EnsureSchema()通过幂等 DDL 与版本化迁移的组合,把 schema 变更分为两大类:无需迁移与必须迁移。
| 变更类型 | 处理机制 | 需要迁移吗 |
|---|---|---|
| 新增表 | CREATE TABLE IF NOT EXISTS | 否 |
| 新增索引 | CREATE INDEX IF NOT EXISTS | 否 |
| 函数体改动 | CREATE OR REPLACE FUNCTION | 否 |
| 新增带 DEFAULT 的函数参数 | CREATE OR REPLACE FUNCTION | 否 |
| 现有表新增列 | — | 是(ALTER TABLE ADD COLUMN IF NOT EXISTS) |
| 列类型变更 | — | 是(ALTER TABLE ALTER COLUMN TYPE) |
| 函数签名变更(参数/返回类型) | — | 是(DROP FUNCTION+CREATE) |
| 删除列 | — | 是(两阶段:先停止使用,再删除) |
这张表的关键启示:只要 DDL 模板 schema.sql 里用幂等语句能表达且不破坏已有数据结构的变更,都交给EnsureSchema自动完成;只有那些幂等 DDL 无法安全表达的变更(新增列、改列类型、改函数签名、删列)才需要手写迁移。
四、Schema 版本管理机制
4.1 版本存储
schema 版本存放在<prefix>schema_version表中(cf_map_schema_version与cf_binary_map_schema_version两张),该表由 schema.sql 中的幂等 DDL 创建并初始化:
CREATE TABLE IF NOT EXISTS __PREFIX__schema_version ( id INTEGER PRIMARY KEY, schema_version INTEGER NOT NULL ); INSERT INTO __PREFIX__schema_version (id, schema_version) VALUES (1, 1) ON CONFLICT (id) DO NOTHING;即全新安装时版本先初始化为 1,随后由EnsureSchema的收尾步骤更新到当前schemaVersion(见 pgmapbroker.go 的SetSchemaVersion调用)。
4.2 启动时的三条路径
EnsureSchema在每次进程启动时执行,具体行为分三种情况:
- 快速路径(fast path):若
schema_version与当前schemaVersion一致,且两张主表(cf_map_stream、cf_binary_map_stream)的探测查询成功,则跳过全部 DDL,只做shard_lock对齐与分区 lookahead 刷新后直接返回。源码见 pgmapbroker.go,其中bothVariantsPresent会探测两套变体的主表,防止"只有一套变体建成"的部分安装误走快速路径。 - 全新安装:DDL 直接创建最新形态的 schema,版本置为当前值,跳过迁移链(没有旧版本可升级)。
- 升级:幂等 DDL 被重新应用,然后按顺序执行从
dbVersion+1到schemaVersion之间的所有迁移。
4.3 迁移先于 DDL 的原因
EnsureSchema中迁移在 DDL之前执行(pgmapbroker.go)。原因在源码注释中明确:DDL 模板反映的是最新形态,可能包含CREATE INDEX IF NOT EXISTS idx ON tbl(newcol)这类语句——解析它要求newcol已经存在于表中。先跑迁移,保证 DDL 引用的新列已就位。
五、迁移机制的源码级实现
5.1 迁移注册表:schemaMigrations
迁移通过 Go 侧的map[int]string注册,定义于 pgmapbroker.go:
var schemaVersion = 1 var schemaMigrations = map[int]string{}schemaVersion是当前 schema 版本,新增迁移时需手动递增;schemaMigrations将目标版本号映射到迁移 SQL模板;- 版本 1 是基线(由完整 DDL 应用),迁移从 2 开始。
5.2 启动即校验:ValidateMigrationMap
pgmapbroker.go的包级init()调用 pgschema.go 的ValidateMigrationMap:
func init() { pgschema.ValidateMigrationMap("pgmapbroker", schemaVersion, schemaMigrations) }该校验函数要求[2..schemaVersion]区间内每个版本都恰好注册了一个迁移、区间外不能有任何迁移,否则进程启动即 panic。设计意图在源码注释中说明:只递增schemaVersion却漏注册迁移,会让schema_version通过收尾 UPDATE 直接跳到新版本,而数据库形态还停留在旧版——静默损坏。用启动期 panic 代替运行时损坏。
5.3 占位符渲染:一处编写,两个变体自动生成
虽然 README 要求迁移文件写显式表名,但schemaMigrations中的模板仍支持两个占位符(这与"无模板占位符"的迁移文件约定并不冲突——仓库源码当前的做法是:迁移在注册到 map 时仍可使用与 schema.sql 相同的占位符机制,由renderSchemaTemplate渲染;而落盘的迁移文件本身不依赖占位符渲染,见下文 §7 讨论):
__PREFIX__→ 例如cf_map_或用户自定义前缀__DATA_TYPE__→JSONB或BYTEA(按变体)
渲染逻辑在 pgmapbroker.go:
func renderSchemaTemplate(template, prefix string, binary bool) string { dataType := "JSONB" if binary { dataType = "BYTEA" } return strings.NewReplacer( "__PREFIX__", prefix, "__DATA_TYPE__", dataType, ).Replace(template) }migrationVariants(pgmapbroker.go)把同一个迁移模板渲染为两个变体,并分别关联各自的版本表:
func (e *PostgresMapBroker) migrationVariants(template string) []pgschema.MigrationVariant { return []pgschema.MigrationVariant{ {SQL: renderSchemaTemplate(template, e.names.jsonbPrefix, false), VersionTable: e.names.jsonbPrefix + "schema_version"}, {SQL: renderSchemaTemplate(template, e.names.binaryPrefix, true), VersionTable: e.names.binaryPrefix + "schema_version"}, } }由此,无论用户把table_prefix配成什么(如prod_us_cf)、无论变体是 JSONB 还是 BYTEA,迁移作者都只需写一次 SQL。
5.4 事务化迁移应用:ApplyMigrationInTx
迁移执行在 pgschema.go 的ApplyMigrationInTx中完成:
- 每个迁移的所有变体在同一个事务内执行,任意一个失败则全部回滚(包括部分版本自增);
- 每个变体执行后立即
UPDATE <prefix>schema_version SET schema_version = $1 WHERE id = 1; - 若
id=1行不存在,返回错误提示手动修复版本表,而不是静默通过; - 对死锁(40P01)与
tuple concurrently updated(XX000,滚动部署中并发CREATE OR REPLACE的典型冲突)做最多 3 次、间隔 200ms 的重试; - 刻意不把
duplicate_object类错误纳入重试——迁移运行在迁移锁之下,重复对象是真实 bug,必须浮出水面。
5.5 集群级迁移串行化:AcquireMigrationLock
滚动部署时多个节点会同时跑EnsureSchema。迁移链由 pgschema.go 的AcquireMigrationLock以 PostgreSQL advisory lock 串行化:
- 锁 ID 由
FNV-64("pgschema/migration:" + label)派生,pgmapbroker与pgstreambroker互不阻塞; - 锁持有期间会重新读取
schema_version:若等待期间其他节点已完成升级,当前节点的迁移循环自动变成空操作; - 释放函数显式执行
pg_advisory_unlock后才归还连接,避免下一个池用户继承锁;进程崩溃时锁随会话结束自动释放。
5.6 快速路径的两项配套维护
即使走快速路径,EnsureSchema仍执行两项必要维护(pgmapbroker.go):
reconcileShardLock:确保shard_lock表恰好包含[0, NumShards)每行一个分片。缺失行不只是性能问题——会破坏 per-shard 发布串行化,导致 stream ID 乱序、outbox 游标跳行;ensurePartitionedStream:刷新分区 lookahead。分区是 schema 中唯一会"过期"的部分,集群停机超过PartitionLookaheadDays后今天的分区可能缺失,此调用即时补齐,避免发布一直报 "no partition for value"。
六、添加一次迁移的完整开发流程
postgres_schema.md给出了开发者新增迁移的六步工作流,结合源码可细化为:
- 提升
schemaVersion:在 pgmapbroker.go 中递增var schemaVersion。漏改会导致迁移永远不执行;只改不改注册表则会在启动时被ValidateMigrationMappanic 拦截。 - 创建迁移文件:在 internal/pgmapbroker/internal/sql/migrations/ 下新建
NNN.sql,写显式 SQL、同时覆盖两套前缀、保持幂等。 - 注册到
schemaMigrations:在 pgmapbroker.go 的 map 中新增条目并嵌入该文件。 - 同步更新
schema.sql模板:让模板反映最新形态——全新安装只跑 DDL、不跑迁移链,因此模板必须已经包含迁移的最终结果。 - 重新生成 schema 产物:文档工作流提到运行
make pg-schemas重新生成。需要说明的是,当前仓库根目录 Makefile 中未包含该目标,且 schema 模板通过//go:embed internal/sql/schema.sql直接嵌入(pgmapbroker.go),模板不变式由 schema_template_test.go 在测试期自动校验。因此在当前仓库中,修改模板后执行go test ./internal/pgmapbroker/...即可验证。 - 维持不变式:全新安装(DDL)与升级(迁移链)必须收敛到完全一致的 schema 形态。
源码注释中还强调了一个关键设计:EnsureSchema收尾处对两个变体版本表统一执行SetSchemaVersion(pgmapbroker.go)。对升级而言迁移事务已写版本号,该 UPDATE 是空操作;对全新安装而言,它把 DDL 初始化的1提升到当前schemaVersion——若此步失败是致命错误,否则下次启动会误判为"待升级"而重跑整个迁移链。
七、迁移作者硬性要求(Requirements on migration authors)
postgres_schema.md对迁移作者提出三条硬性要求:
- 必须幂等:使用
ADD COLUMN IF NOT EXISTS等构造。即便迁移在迁移锁 + 事务保护下运行,幂等仍然必要——部分部署或运维手动重跑必须安全。 - 必须向后兼容:迁移执行后旧版本代码必须能继续正常工作。这直接服务于滚动部署。
- 必须同时覆盖两套前缀:
cf_map_*与cf_binary_map_*都要被迁移语句命中。
README 中"迁移文件写显式表名、无模板占位符"与源码注册表用模板渲染看似存在张力,实际是一致的:落盘的NNN.sql文件是给运维与审计看的生产级显式 SQL;运行时通过renderSchemaTemplate渲染注册表中的模板,则是为了在一个地方覆盖任意用户前缀与两种变体。两者服务于同一目标——生产迁移的显式、安全与可重复执行。
八、滚动部署规则(Rolling deploy rules)
多节点滚动升级时,EnsureSchema的并发行为由以下规则约束:
- 加法变更安全:带 DEFAULT 的新列、带 DEFAULT 的新函数参数——旧节点忽略、新节点使用,天然兼容。
- 破坏性变更必须两阶段:第一阶段部署"不再使用旧列"的新代码;第二阶段再部署删除该列的迁移。
- 函数签名变更需协调:两种选择——接受
DROP+CREATE期间的短暂报错(通常 <1s),或改用新函数名实现零停机。 EnsureSchema每次启动运行一次:并发运行安全(幂等 DDL + 死锁重试)。DDL 的并发冲突由RetrySchemaExec处理:识别deadlock_detected、internal_error、duplicate_table、duplicate_object、unique_violation五类冲突,做最多 5 次、100ms 起步指数退避(带抖动)的重试,见 pgschema.go。NumShards变更必须全量重启:不能滚动。不同NumShards的并发节点会在shard_lock填充上互相竞争。- 回滚(降级)安全:
EnsureSchema会把函数覆盖回旧版本;新版本迁移多出的列仍保留但被旧代码忽略,无数据丢失。注意该规则针对的是同一 schema 版本区间内的函数体差异;真正的版本降级(DB 版本 > 二进制支持的版本)会被CheckDowngrade显式拒绝——见下文。
其中降级保护的实现位于 pgschema.go:当 DB 的schema_version高于当前二进制支持的上限时,EnsureSchema直接返回明确错误(提示运行新版本二进制或恢复旧快照),绝不静默回写版本号。
九、手动迁移路径(Manual migration path)
对于不使用自动EnsureSchema或需要人工审计的场景,文档给出四条手动路径:
- 全新安装:应用基线 DDL。文档中提到的
schema_all.sql在当前仓库中对应的实际文件为 internal/pgmapbroker/internal/sql/schema.sql(模板,含__PREFIX__与__DATA_TYPE__占位符,渲染后即完整 DDL;当前仓库未提供预渲染的全量 SQL 文件)。 - 检查版本:
SELECT schema_version FROM cf_map_schema_version WHERE id = 1; - 按序应用迁移:按版本号顺序执行 internal/pgmapbroker/internal/sql/migrations/ 下的
NNN.sql文件。当前schemaVersion为 1,该目录暂未包含实际的NNN.sql迁移文件——README 定义的是首个迁移(如002.sql)及后续迁移的编写规范。 EnsureSchema可选:保持开启是安全的(版本一致时为空操作),也可以按需禁用。
十、测试与不变式保障
schema 管理有一套自动化防护网:
- 启动校验:
ValidateMigrationMap保证迁移注册表无缺口、无越界(pgschema.go)。 - 模板不变式测试:schema_template_test.go 对 JSONB 与 BYTEA 两个变体分别渲染模板并断言:
- 不出现
CONCURRENTLY(CREATE INDEX CONCURRENTLY无法在隐式事务内执行,会破坏RetrySchemaExec的重试不变式); - DDL 半段不出现显式事务控制语句(
BEGIN/COMMIT/ROLLBACK),保证整个批次作为服务端单一隐式事务运行,失败时整体回滚。
- 不出现
- 迁移并发语义:迁移的重试面刻意窄于 DDL(pgschema.go 注释),确保重复对象这类真实 bug 不被重试掩盖。
- 版本读取错误判别:
ReadSchemaVersion把"表缺失 / 行缺失"视为全新安装、把瞬时连接失败等错误向上传播而不是当作全新安装——避免瞬时故障导致跳过迁移并强制版本前进的静默损坏(pgschema.go)。
十一、常见疑问与最佳实践小结
- 为什么版本 1 不是
001.sql?版本 1 是基线,由完整 DDL(schema.sql)一次性创建,不需要迁移文件;NNN.sql从002.sql开始。 - 为什么必须写
IF NOT EXISTS?迁移在滚动部署的多个节点间由 advisory lock 串行化,但部分部署、运维手动重跑、失败后重试都要求迁移可安全重复执行。 - 为什么每个迁移都要碰
cf_binary_map_*?EnsureSchema无条件创建两套变体(与运行时BinaryData无关),版本号也是两套各记一份,因此迁移必须同步两套,否则二进制变体的 schema 会与版本号脱节。 - 删列为什么必须两阶段?旧节点代码仍在引用该列时直接删除会导致线上 SQL 报错;先发布"停止使用"的代码、再发布删列迁移,才能与滚动部署兼容。
总而言之,MapBroker 的迁移体系把"能自动化的交给幂等 DDL、不能自动化的交给版本化迁移",配合集群级迁移锁、事务化应用、启动期校验与降级拒绝,形成了一套可审计、可回滚、面向滚动部署的 PostgreSQL schema 演进方案。后续新增迁移时,只要遵守 README 的命名与幂等约定、按六步工作流执行并保持"DDL 与迁移收敛一致"的不变式,即可安全地把 schema 演进融入日常发布流程。
- 消息队列
- 后端
- 通信
【免费下载链接】centrifugo
Scalable real-time messaging server in a language-agnostic way. Self-hosted alternative to Pubnub, Pusher, Ably, socket.io, Phoenix.PubSub, SignalR. Set up once and forever.
相关推荐
Centrifugo PostgreSQL MapBroker 架构解析:EnsureSchema 自动建表与整数版本化迁移机制
Centrifugo PostgreSQL MapBroker 架构解析:EnsureSchema 自动建表与整数版本化迁移机制 Centrifugo 的 Po
消息队列后端通信WeKnora 数据库 schema 与迁移机制全解:PostgreSQL / ParadeDB / SQLite 版本化迁移实战
WeKnora 数据库 schema 与迁移机制全解:PostgreSQL / ParadeDB / SQLite 版本化迁移实战 WeKnora 是一个开源的
人工智能大模型RAGAI Agent后端前端MCP 服务知识库dsh-plugin工具调用7个终极Node.js数据库迁移工具:Schema管理与版本控制完整指南
7个终极Node.js数据库迁移工具:Schema管理与版本控制完整指南 在Node.js开发中,数据库迁移和Schema管理是确保应用数据结构一致性的关键环节
文档
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考