Centrifugo PostgreSQL MapBroker 数据库迁移指南:Schema 版本管理与迁移文件规范
2026/9/24 2:04:56 网站建设 项目流程
  • 消息队列
  • 后端
  • 通信

【免费下载链接】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.

项目地址:https://gitcode.com/gh_mirrors/ce/centrifugo
点击查看免费下载

导读

本文围绕 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定义了迁移文件的硬性约定,这是所有迁移作者必须遵守的底线:

  1. 文件命名:文件名为NNN.sql,其中NNN是目标版本号。例如002.sql表示将 schema 升级到版本 2。迁移从版本 2 开始,版本 1 是基线(由完整 DDL 应用)。
  2. 幂等性:每个迁移必须幂等——使用IF NOT EXISTSIF EXISTS等构造,允许重复执行而不会报错或产生副作用。
  3. 双前缀覆盖:每个迁移必须同时作用于两套前缀的表:cf_map_*(JSONB 变体)与cf_binary_map_*(BYTEA 变体)。注意cf是默认table_prefix,用户可配置为其他前缀。
  4. 无模板占位符:迁移文件内不得使用模板占位符,必须写显式表名——生产环境迁移中显式更安全。

标准迁移示例(摘自 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_versioncf_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_streamcf_binary_map_stream)的探测查询成功,则跳过全部 DDL,只做shard_lock对齐与分区 lookahead 刷新后直接返回。源码见 pgmapbroker.go,其中bothVariantsPresent会探测两套变体的主表,防止"只有一套变体建成"的部分安装误走快速路径。
  • 全新安装:DDL 直接创建最新形态的 schema,版本置为当前值,跳过迁移链(没有旧版本可升级)。
  • 升级:幂等 DDL 被重新应用,然后按顺序执行从dbVersion+1schemaVersion之间的所有迁移。

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__JSONBBYTEA(按变体)

渲染逻辑在 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)派生,pgmapbrokerpgstreambroker互不阻塞;
  • 锁持有期间会重新读取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给出了开发者新增迁移的六步工作流,结合源码可细化为:

  1. 提升schemaVersion:在 pgmapbroker.go 中递增var schemaVersion。漏改会导致迁移永远不执行;只改不改注册表则会在启动时被ValidateMigrationMappanic 拦截。
  2. 创建迁移文件:在 internal/pgmapbroker/internal/sql/migrations/ 下新建NNN.sql,写显式 SQL、同时覆盖两套前缀、保持幂等。
  3. 注册到schemaMigrations:在 pgmapbroker.go 的 map 中新增条目并嵌入该文件。
  4. 同步更新schema.sql模板:让模板反映最新形态——全新安装只跑 DDL、不跑迁移链,因此模板必须已经包含迁移的最终结果。
  5. 重新生成 schema 产物:文档工作流提到运行make pg-schemas重新生成。需要说明的是,当前仓库根目录 Makefile 中未包含该目标,且 schema 模板通过//go:embed internal/sql/schema.sql直接嵌入(pgmapbroker.go),模板不变式由 schema_template_test.go 在测试期自动校验。因此在当前仓库中,修改模板后执行go test ./internal/pgmapbroker/...即可验证。
  6. 维持不变式:全新安装(DDL)与升级(迁移链)必须收敛到完全一致的 schema 形态。

源码注释中还强调了一个关键设计:EnsureSchema收尾处对两个变体版本表统一执行SetSchemaVersion(pgmapbroker.go)。对升级而言迁移事务已写版本号,该 UPDATE 是空操作;对全新安装而言,它把 DDL 初始化的1提升到当前schemaVersion——若此步失败是致命错误,否则下次启动会误判为"待升级"而重跑整个迁移链。

七、迁移作者硬性要求(Requirements on migration authors)

postgres_schema.md对迁移作者提出三条硬性要求:

  1. 必须幂等:使用ADD COLUMN IF NOT EXISTS等构造。即便迁移在迁移锁 + 事务保护下运行,幂等仍然必要——部分部署或运维手动重跑必须安全。
  2. 必须向后兼容:迁移执行后旧版本代码必须能继续正常工作。这直接服务于滚动部署。
  3. 必须同时覆盖两套前缀cf_map_*cf_binary_map_*都要被迁移语句命中。

README 中"迁移文件写显式表名、无模板占位符"与源码注册表用模板渲染看似存在张力,实际是一致的:落盘的NNN.sql文件是给运维与审计看的生产级显式 SQL;运行时通过renderSchemaTemplate渲染注册表中的模板,则是为了在一个地方覆盖任意用户前缀与两种变体。两者服务于同一目标——生产迁移的显式、安全与可重复执行。

八、滚动部署规则(Rolling deploy rules)

多节点滚动升级时,EnsureSchema的并发行为由以下规则约束:

  1. 加法变更安全:带 DEFAULT 的新列、带 DEFAULT 的新函数参数——旧节点忽略、新节点使用,天然兼容。
  2. 破坏性变更必须两阶段:第一阶段部署"不再使用旧列"的新代码;第二阶段再部署删除该列的迁移。
  3. 函数签名变更需协调:两种选择——接受DROP+CREATE期间的短暂报错(通常 <1s),或改用新函数名实现零停机。
  4. EnsureSchema每次启动运行一次:并发运行安全(幂等 DDL + 死锁重试)。DDL 的并发冲突由RetrySchemaExec处理:识别deadlock_detectedinternal_errorduplicate_tableduplicate_objectunique_violation五类冲突,做最多 5 次、100ms 起步指数退避(带抖动)的重试,见 pgschema.go。
  5. NumShards变更必须全量重启:不能滚动。不同NumShards的并发节点会在shard_lock填充上互相竞争。
  6. 回滚(降级)安全EnsureSchema会把函数覆盖回旧版本;新版本迁移多出的列仍保留但被旧代码忽略,无数据丢失。注意该规则针对的是同一 schema 版本区间内的函数体差异;真正的版本降级(DB 版本 > 二进制支持的版本)会被CheckDowngrade显式拒绝——见下文。

其中降级保护的实现位于 pgschema.go:当 DB 的schema_version高于当前二进制支持的上限时,EnsureSchema直接返回明确错误(提示运行新版本二进制或恢复旧快照),绝不静默回写版本号。

九、手动迁移路径(Manual migration path)

对于不使用自动EnsureSchema或需要人工审计的场景,文档给出四条手动路径:

  1. 全新安装:应用基线 DDL。文档中提到的schema_all.sql在当前仓库中对应的实际文件为 internal/pgmapbroker/internal/sql/schema.sql(模板,含__PREFIX____DATA_TYPE__占位符,渲染后即完整 DDL;当前仓库未提供预渲染的全量 SQL 文件)。
  2. 检查版本
    SELECT schema_version FROM cf_map_schema_version WHERE id = 1;
  3. 按序应用迁移:按版本号顺序执行 internal/pgmapbroker/internal/sql/migrations/ 下的NNN.sql文件。当前schemaVersion为 1,该目录暂未包含实际的NNN.sql迁移文件——README 定义的是首个迁移(如002.sql)及后续迁移的编写规范。
  4. EnsureSchema可选:保持开启是安全的(版本一致时为空操作),也可以按需禁用。

十、测试与不变式保障

schema 管理有一套自动化防护网:

  • 启动校验ValidateMigrationMap保证迁移注册表无缺口、无越界(pgschema.go)。
  • 模板不变式测试:schema_template_test.go 对 JSONB 与 BYTEA 两个变体分别渲染模板并断言:
    • 不出现CONCURRENTLYCREATE INDEX CONCURRENTLY无法在隐式事务内执行,会破坏RetrySchemaExec的重试不变式);
    • DDL 半段不出现显式事务控制语句(BEGIN/COMMIT/ROLLBACK),保证整个批次作为服务端单一隐式事务运行,失败时整体回滚。
  • 迁移并发语义:迁移的重试面刻意窄于 DDL(pgschema.go 注释),确保重复对象这类真实 bug 不被重试掩盖。
  • 版本读取错误判别ReadSchemaVersion把"表缺失 / 行缺失"视为全新安装、把瞬时连接失败等错误向上传播而不是当作全新安装——避免瞬时故障导致跳过迁移并强制版本前进的静默损坏(pgschema.go)。

十一、常见疑问与最佳实践小结

  • 为什么版本 1 不是001.sql版本 1 是基线,由完整 DDL(schema.sql)一次性创建,不需要迁移文件;NNN.sql002.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.

项目地址:https://gitcode.com/gh_mirrors/ce/centrifugo
点击查看免费下载

相关推荐

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

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

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

立即咨询