开篇先从一个实际场景切入:数据库从自建机房迁到云数据库,或者从 PostgreSQL 老版本升级到新版本时,通常要考虑停窗口够不够、数据量太大怎么同步、表结构差异如何处理。Pgcopydb 是 PostgreSQL 生态里很好用的迁移工具,但它是 C 语言实现,编译部署在某些内网环境里并不轻松。最近社区有人用 Go 重写这套思路,项目名为 Pgmigrate。本文会围绕 Pgmigrate 的设计思路、核心模块、Go 实现要点和工程落地展开,帮你理解这类迁移工具的原理,也方便你在项目里按需改造或直接使用。
1. Pgcopydb 是什么,为什么值得用 Go 重写
1.1 Pgcopydb 的核心能力
Pgcopydb 是 PostgreSQL 官方生态里一个非常实用的迁移工具。它做的事情可以简单概括为:把一个 PostgreSQL 实例中的数据库结构、数据、序列、索引、约束等全部复制到另一个 PostgreSQL 实例中。
它和普通的pg_dump+psql这种逻辑备份恢复方式相比,主要优势在于:
- 支持网络直连的在线迁移方式,不要求先落地成文件。
- 自动并行处理多张表的数据拷贝,充分利用源库和目标库的 IO 能力。
- 对无主键的大表采用分批主键范围扫描,避免一次性加载全部数据导致内存溢出。
- 迁移过程中会自动处理序列值的同步,避免迁移后自增主键冲突。
- 支持先同步结构,再同步数据,最后加约束和索引的策略,减少锁持有时间。
在实际生产环境中,Pgcopydb 常用于:
- 把 PostgreSQL 从旧版本升级到新版本。
- 把本地数据库迁移到云数据库 RDS。
- 在两个自建 PostgreSQL 集群之间做数据搬迁。
- 作为大版本升级前的演练工具。
1.2 为什么有人想用 Go 重写
Pgcopydb 本身是 C 语言实现的,能力很强,但对不少团队来说存在几个现实问题:
- C 语言项目编译依赖较多,需要 autoconf、libpq、zstd 等开发库,在一些离线内网环境里配置很麻烦。
- 二进制的分发和跨平台编译不如 Go 方便。
- C 语言手动管理内存,遇到复杂错误处理时对开发者要求高,二次开发成本不低。
- 社区中不少数据库工具链已经转向 Go 生态,团队内部有 Go 语言技术积累。
所以 Pgmigrate 这个项目思路很自然:保留 Pgcopydb 的迁移策略和流程设计,用 Go 重写一套可分发、易部署、并发模型更友好的实现。这样既能利用 Go 的 goroutine 做高并发数据拷贝,又能非常方便地交叉编译出 Linux、macOS、Windows 三个平台的可执行文件。
1.3 Pgmigrate 的定位
Pgmigrate 不是 Pgcopydb 的简单复制,而是用 Go 语言把 Pgcopydb 的核心流程重新实现。它面向的场景是:
- 开发人员希望用 Go 工具链管理 PostgreSQL 迁移。
- 运维团队需要在不同发布环境里快速部署迁移工具。
- 项目需要定制化开发,比如在迁移过程中对接内部监控、消息队列或审核系统。
从技术形态上看,Pgmigrate 的核心是一个命令行工具,通过参数指定源库连接串、目标库连接串、迁移模式(结构、数据、全量)等,然后完成迁移任务。
2. Pgmigrate 的整体架构与核心流程
理解 Pgmigrate 之前,先看一张简单的流程示意:
源库 PostgreSQL │ ├─ 读取源库元数据(表、索引、约束、序列) │ ├─ 在目标库中创建 Schema │ ├─ 创建辅助迁移表(可选) │ ├─ 并行读取源表数据 │ ├── 按主键分片 │ ├── goroutine 并发读取 │ └── COPY 协议写入目标库 │ ├─ 同步序列值 │ └─ 校验数据一致性(行数、校验和)整体流程可以拆成以下几个核心阶段。
2.1 元数据发现阶段
迁移工具首先需要知道“要迁移什么”。Pgmigrate 通过查询 PostgreSQL 的系统目录表(pg_catalog)来获取所有用户表、视图、序列、索引、约束、触发器的定义。这一步通常会用以下系统视图:
pg_class:表、索引等关系对象。pg_attribute:表的列信息。pg_index:索引定义。pg_constraint:约束定义。pg_sequences:序列信息。information_schema.tables:用户表的基础信息。
工具需要判断哪些表需要迁移数据,哪些表只需要结构。比如分区表的父表通常不存储数据,子表才需要迁移;视图和物化视图的处理策略也不一样。
2.2 Schema 同步阶段
拿到元数据后,Pgmigrate 先在目标库中创建对应的 Schema、表、序列等对象。这个阶段只需要执行 DDL 语句。很多迁移工具会把“结构同步”和“数据同步”分开,原因是:
- 先建表再导数据,结构问题可以提前暴露。
- 先不创建索引和约束,可以大幅提升数据导入速度。
- 所有表都完成数据导入后,再统一加索引和约束,减少锁等待时间。
在 Go 实现中,Schema 同步通常使用database/sql标准库,直接执行 DDL 字符串。由于 DDL 不能参数化,所以构造 SQL 时要注意标识符的转义,防止特殊表名或列名导致语法错误。
2.3 数据拷贝阶段
这是整个迁移工具最核心、也最复杂的部分。Pgcopydb 的经验告诉我们,数据拷贝要关注下面几件事。
并发模型
不是表越多并发线程越多就越好。并发数过高会把源库的 IO 打满,影响线上业务。一般建议先按表大小排序,大表单独处理,小表用较高并发批量处理。Go 的 goroutine 非常适合做这种任务池模型。
分片策略
对于大表,不能一次性SELECT * FROM table,因为如果这张表有上亿行,查询结果集会占用巨量内存,而且网络传输中断后需要整体重来。更稳妥的方式是:
- 如果表有主键,用主键范围分片。
- 每次取一个主键区间,例如
id > x AND id <= y。 - 每个分片独立读取、独立写入。
- 某个分片失败时只重试该分片。
写入方式
Go 的lib/pq和pgx都支持 PostgreSQL 的 COPY 协议。COPY 协议比逐条 INSERT 快非常多,因为它减少了大量网络往返和 SQL 解析开销。Pgmigrate 在 Go 中的写入路径通常是:
- 从源库按分片查询数据。
- 将行数据编码成 COPY 格式。
- 在目标库执行
COPY table FROM STDIN。 - 将数据流式写入。
2.4 序列同步阶段
PostgreSQL 的序列(Sequence)默认是独立对象,不会随着表数据的 INSERT 自动更新。如果直接把表数据从源库拷贝到目标库,但序列起始值没有同步,后续应用插入新记录时就会撞上主键冲突。
Pgmigrate 在数据拷贝完成后,会读取源库每个序列的last_value,然后在目标库执行:
SELECT setval('序列名', last_value, true);这样可以保证迁移后,目标库的新增数据主键从正确的位置继续递增。
2.5 数据校验阶段
数据迁移完成不等于任务结束。生产环境通常要求行数一致和抽样内容一致。Pgmigrate 会在迁移完成后,对每张表执行计数对比:
SELECT count(*) FROM 源表; SELECT count(*) FROM 目标表;对于核心业务表,还可以对某些关键列做sum(checksum)或哈希校验。这个阶段的目标是发现数据丢失或重复等问题。
3. 环境准备与项目结构
3.1 环境说明
本文示例以 Go 1.21 以上版本为例,操作系统为 Linux,数据库为 PostgreSQL 14 或更高版本。实际操作中版本可以按你的环境调整,核心思路一致。
需要准备的工具:
- Go 1.21+
- PostgreSQL 源库和目标库各一个
- 能访问两个数据库的测试网络环境
psql命令行工具(用于查看迁移结果)
3.2 项目结构
为了演示 Pgmigrate 的核心原理,我们创建一个教学版小工具,目录结构如下:
pgmigrate-demo/ ├── go.mod ├── main.go ├── config.go ├── meta/ │ └── metadata.go ├── migrate/ │ ├── schema.go │ ├── data.go │ └── sequence.go └── util/ └── db.go这个结构对应了前面讲的几个核心阶段:
config.go:处理命令行参数,解析源库和目标库连接串。meta/metadata.go:读取源库的表、序列等元数据。migrate/schema.go:在目标库创建 Schema 和表结构。migrate/data.go:并发拷贝表数据。migrate/sequence.go:同步序列。util/db.go:统一的数据库连接管理。
4. 使用 Go 实现 Pgmigrate 核心逻辑
接下来我们逐步实现一个简化版 Pgmigrate。它不能完全替代 Pgcopydb,但可以帮助大家理解迁移工具的核心套路,也为二次开发打基础。
4.1 初始化 Go 模块
首先创建项目目录并初始化模块。
mkdir pgmigrate-demo cd pgmigrate-demo go mod init pgmigrate-demo然后安装 PostgreSQL 驱动。这里我们使用pgx标准库,它支持COPY协议,性能比lib/pq好。
go get github.com/jackc/pgx/v5最终go.mod大致如下:
module pgmigrate-demo go 1.21 require github.com/jackc/pgx/v5 v5.5.04.2 配置文件与参数解析
创建一个config.go,用于解析命令行参数。一个迁移工具常用的参数包括:
--source:源库连接串。--target:目标库连接串。--schema:指定要迁移的 Schema,默认public。--concurrency:数据拷贝并发数。--with-data:是否拷贝数据,默认 true。
// config.go package main import ( "flag" "fmt" ) type Config struct { SourceDSN string TargetDSN string Schema string Concurrency int WithData bool OnlySchema bool } func ParseConfig() (*Config, error) { cfg := &Config{} flag.StringVar(&cfg.SourceDSN, "source", "", "source PostgreSQL DSN") flag.StringVar(&cfg.TargetDSN, "target", "", "target PostgreSQL DSN") flag.StringVar(&cfg.Schema, "schema", "public", "schema to migrate") flag.IntVar(&cfg.Concurrency, "concurrency", 4, "number of concurrent table copy workers") flag.BoolVar(&cfg.WithData, "with-data", true, "copy table data") flag.BoolVar(&cfg.OnlySchema, "schema-only", false, "only migrate schema") flag.Parse() if cfg.SourceDSN == "" || cfg.TargetDSN == "" { return nil, fmt.Errorf("source and target DSN are required") } return cfg, nil }这里使用了 Go 标准库flag,方便在命令行里直接传参。实际生产工具可以用cobra做更复杂的子命令,但教学版用标准库就够了。
4.3 数据库连接管理
创建一个util/db.go,统一管理源库和目标库的连接。这里使用pgxpool,它是pgx提供的连接池实现,适合并发场景。
// util/db.go package util import ( "context" "fmt" "github.com/jackc/pgx/v5/pgxpool" ) func NewPool(ctx context.Context, dsn string) (*pgxpool.Pool, error) { pool, err := pgxpool.New(ctx, dsn) if err != nil { return nil, fmt.Errorf("connect to database failed: %w", err) } if err := pool.Ping(ctx); err != nil { pool.Close() return nil, fmt.Errorf("ping database failed: %w", err) } return pool, nil }连接串的格式是标准 PostgreSQL DSN 格式,例如:
postgres://user:password@localhost:5432/dbname?sslmode=disable4.4 元数据读取
创建一个meta/metadata.go,用来读取源库中的表信息和序列信息。
// meta/metadata.go package meta import ( "context" "fmt" "github.com/jackc/pgx/v5" ) type TableMeta struct { Schema string Name string HasPK bool PKCols []string } type SequenceMeta struct { Schema string Name string } // LoadTables 读取指定 schema 下的所有普通表 func LoadTables(ctx context.Context, conn *pgx.Conn, schema string) ([]TableMeta, error) { rows, err := conn.Query(ctx, ` SELECT c.relname, EXISTS ( SELECT 1 FROM pg_index i WHERE i.indrelid = c.oid AND i.indisprimary ) AS has_pk FROM pg_class c JOIN pg_namespace n ON n.oid = c.relnamespace WHERE c.relkind = 'r' AND n.nspname = $1 ORDER BY c.relname `, schema) if err != nil { return nil, err } defer rows.Close() var tables []TableMeta for rows.Next() { var t TableMeta if err := rows.Scan(&t.Name, &t.HasPK); err != nil { return nil, err } t.Schema = schema if t.HasPK { pkCols, err := loadPKColumns(ctx, conn, schema, t.Name) if err != nil { return nil, err } t.PKCols = pkCols } tables = append(tables, t) } return tables, rows.Err() } func loadPKColumns(ctx context.Context, conn *pgx.Conn, schema, table string) ([]string, error) { rows, err := conn.Query(ctx, ` SELECT a.attname FROM pg_index i JOIN pg_attribute a ON a.attrelid = i.indrelid AND a.attnum = ANY(i.indkey) WHERE i.indrelid = ($1 || '.' || $2)::regclass AND i.indisprimary `, schema, table) if err != nil { return nil, fmt.Errorf("load pk columns for %s.%s failed: %w", schema, table, err) } defer rows.Close() var cols []string for rows.Next() { var col string if err := rows.Scan(&col); err != nil { return nil, err } cols = append(cols, col) } return cols, rows.Err() } // LoadSequences 读取指定 schema 下的所有序列 func LoadSequences(ctx context.Context, conn *pgx.Conn, schema string) ([]SequenceMeta, error) { rows, err := conn.Query(ctx, ` SELECT sequencename FROM pg_sequences WHERE schemaname = $1 ORDER BY sequencename `, schema) if err != nil { return nil, err } defer rows.Close() var seqs []SequenceMeta for rows.Next() { var s SequenceMeta if err := rows.Scan(&s.Name); err != nil { return nil, err } s.Schema = schema seqs = append(seqs, s) } return seqs, rows.Err() }需要注意的是,查询主键列时用了$1 || '.' || $2)::regclass这种方式把 schema 和表名拼成一个regclass类型。这样可以直接定位到系统目录里的表对象。实际项目中建议再对表名做合法性校验,避免特殊字符拼接出非预期对象。
4.5 Schema 同步实现
接下来实现migrate/schema.go。在目标库中,我们只需要读取源库中每张表的建表语句,然后在目标库执行。
获取建表语句最简单的方式是用pg_dump的--schema-only模式,然后过滤出需要的部分。但作为 Go 程序,我们可以直接查询pg_get_tabledef类型的函数。PostgreSQL 提供了一些实用函数:
pg_get_tabledef(oid):获取表的 CREATE TABLE 语句。pg_get_indexdef(index_oid):获取 CREATE INDEX 语句。pg_get_constraintdef(constraint_oid):获取约束定义。
为了简化,这里用一条 SQL 直接获取建表语句。
// migrate/schema.go package migrate import ( "context" "fmt" "github.com/jackc/pgx/v5" ) // SyncSchema 在目标库中创建源库中的表结构 func SyncSchema(ctx context.Context, src, dst *pgx.Conn, schema string) error { rows, err := src.Query(ctx, ` SELECT c.relname, pg_get_tabledef(c.oid) FROM pg_class c JOIN pg_namespace n ON n.oid = c.relnamespace WHERE c.relkind = 'r' AND n.nspname = $1 ORDER BY c.relname `, schema) if err != nil { return fmt.Errorf("query table definitions failed: %w", err) } defer rows.Close() for rows.Next() { var tableName, ddl string if err := rows.Scan(&tableName, &ddl); err != nil { return err } // 先尝试删除目标库中的同名表,避免冲突 // 注意:生产环境需要更谨慎的策略,这里仅用于演示 _, err := dst.Exec(ctx, fmt.Sprintf(`DROP TABLE IF EXISTS %s.%s CASCADE`, schema, tableName)) if err != nil { return fmt.Errorf("drop target table %s failed: %w", tableName, err) } // 执行源库的建表语句 if _, err := dst.Exec(ctx, ddl); err != nil { return fmt.Errorf("create table %s in target failed: %w", tableName, err) } fmt.Printf("schema synced: %s.%s\n", schema, tableName) } return rows.Err() }这里有一个非常重要的工程问题:pg_get_tabledef在不同 PostgreSQL 版本里返回的语句格式可能不同,而且不一定包含SET子句、注释、权限等信息。如果你要做一个严肃的迁移工具,不能只靠这一个函数。更稳妥的做法是:
- 用
pg_dump --schema-only生成完整结构文件。 - 内部调用
psql执行。 - 或者通过解析系统目录自己重建 DDL。
教学示例中用pg_get_tabledef是为了展示思路,实际生产一定要在目标版本数据库上提前验证。
4.6 数据拷贝实现
数据拷贝是整个工具的重头戏。我们分别来实现无主键表的全量拷贝和有主键表的分片并发拷贝。
4.6.1 使用 COPY 协议写入目标库
在 Go 中,pgx的CopyFrom方法可以直接把行数据批量导入目标表。它的签名是:
func (c *Conn) CopyFrom(ctx context.Context, tableName pgx.Identifier, columnNames []string, rowSrc pgx.CopyFromSource) (int64, error)pgx.Identifier是一个字符串切片,比如[]string{"public", "users"},它会自动处理好标识符转义。
所以数据拷贝的总体思路是:从源库查询数据,然后通过一个实现了pgx.CopyFromSource接口的对象,把数据流式传给目标库。
4.6.2 无主键表全量拷贝
无主键表无法分片,只能一次性读取并写入。对于中小表这完全够用。
// migrate/data.go package migrate import ( "context" "fmt" "github.com/jackc/pgx/v5" ) type sliceRows [][]any func (r sliceRows) Next() bool { return false } func (r sliceRows) Values() ([]any, error) { return nil, nil } func (r sliceRows) Err() error { return nil } // CopyTableFull 拷贝整张表数据 func CopyTableFull(ctx context.Context, src, dst *pgx.Conn, schema, table string) (int64, error) { // 读取列信息 rows, err := src.Query(ctx, fmt.Sprintf(` SELECT column_name FROM information_schema.columns WHERE table_schema = $1 AND table_name = $2 ORDER BY ordinal_position `, ), schema, table) if err != nil { return 0, err } var columns []string for rows.Next() { var col string if err := rows.Scan(&col); err != nil { rows.Close() return 0, err } columns = append(columns, col) } rows.Close() if len(columns) == 0 { return 0, fmt.Errorf("table %s.%s has no columns", schema, table) } // 从源库查询全表数据 srcRows, err := src.Query(ctx, fmt.Sprintf(`SELECT %s FROM %s.%s`, quoteIdentifiers(columns), schema, table)) if err != nil { return 0, err } defer srcRows.Close() // 使用 CopyFrom 写入目标库 count, err := dst.CopyFrom(ctx, pgx.Identifier{schema, table}, columns, srcRows, // pgx.Rows 实现了 CopyFromSource 接口 ) if err != nil { return 0, fmt.Errorf("copy data to %s.%s failed: %w", schema, table, err) } return count, nil } func quoteIdentifiers(cols []string) string { s := "" for i, c := range cols { if i > 0 { s += ", " } s += pgx.Identifier{c}.Sanitize() } return s }这里pgx.Rows本身就实现了CopyFromSource接口,所以可以直接把源库查询结果传给目标库的CopyFrom,省去了手动数据转换的步骤。这是一个非常优雅的设计。
4.6.3 有主键表的分片并发拷贝
对于大表,我们需要按主键范围分片。下面用一个简化版本演示分片逻辑:每次从源库取一个主键区间,然后拷贝到目标库。
// migrate/data.go 追加内容 // CopyTableSharded 按主键分片拷贝大表 func CopyTableSharded(ctx context.Context, src, dst *pgx.Conn, schema, table string, pkCols []string, shardSize int64) error { // 获取主键范围 var minVal, maxVal any err := src.QueryRow(ctx, fmt.Sprintf( `SELECT min(%s), max(%s) FROM %s.%s`, quoteIdentifiers(pkCols), quoteIdentifiers(pkCols), schema, table, )).Scan(&minVal, &maxVal) if err != nil { return fmt.Errorf("get pk range failed: %w", err) } if minVal == nil || maxVal == nil { fmt.Printf("table %s.%s is empty, skip\n", schema, table) return nil } // 这里简化处理:只支持整数类型主键的分片 // 实际工具需要根据列类型选择不同的分片策略 minInt, okMin := minVal.(int64) maxInt, okMax := maxVal.(int64) if !okMin || !okMax { return fmt.Errorf("unsupported pk type, fallback to full copy") } for start := minInt; start <= maxInt; start += shardSize { end := start + shardSize - 1 if end > maxInt { end = maxInt } // 查询分片数据 srcRows, err := src.Query(ctx, fmt.Sprintf( `SELECT * FROM %s.%s WHERE %s >= $1 AND %s <= $2`, schema, table, quoteIdentifiers(pkCols), quoteIdentifiers(pkCols), ), start, end) if err != nil { return fmt.Errorf("query shard failed: %w", err) } // 获取列名 fields := srcRows.FieldDescriptions() columns := make([]string, len(fields)) for i, f := range fields { columns[i] = string(f.Name) } // 写入目标库 count, err := dst.CopyFrom(ctx, pgx.Identifier{schema, table}, columns, srcRows, ) srcRows.Close() if err != nil { return fmt.Errorf("copy shard [%d, %d] failed: %w", start, end, err) } fmt.Printf("shard [%d, %d] copied, rows=%d\n", start, end, count) } return nil }需要注意,这个分片逻辑只适用于整数类型主键,实际迁移工具还需要支持字符串、UUID、时间戳等类型。Pgcopydb 本身也支持多种分片策略,比如按ctid分片,因为ctid是 PostgreSQL 表内部的物理行标识。
4.6.4 并发调度
有了分片能力,还需要一个调度器来并发处理多张表。在 Go 中,可以用errgroup或原生 goroutine + WaitGroup。
// migrate/data.go 追加内容 import ( "sync" "golang.org/x/sync/errgroup" ) // CopyAllTables 并发拷贝所有表 func CopyAllTables(ctx context.Context, src, dst *pgx.Conn, tables []meta.TableMeta, concurrency int, shardSize int64) error { sem := make(chan struct{}, concurrency) var mu sync.Mutex var totalRows int64 g, ctx := errgroup.WithContext(ctx) for _, t := range tables { t := t g.Go(func() error { select { case sem <- struct{}{}: defer func() { <-sem }() case <-ctx.Done(): return ctx.Err() } var copied int64 var err error if t.HasPK && len(t.PKCols) == 1 { // 有单列主键,尝试分片 // 这里简化判断,实际需要确认类型 err = CopyTableSharded(ctx, src, dst, t.Schema, t.Name, t.PKCols, shardSize) } else { copied, err = CopyTableFull(ctx, src, dst, t.Schema, t.Name) } if err != nil { return fmt.Errorf("copy table %s.%s failed: %w", t.Schema, t.Name, err) } if copied > 0 { mu.Lock() totalRows += copied mu.Unlock() } return nil }) } if err := g.Wait(); err != nil { return err } fmt.Printf("total rows copied: %d\n", totalRows) return nil }errgroup是 Go 官方扩展包提供的并发错误聚合工具,非常适合这种任务并发场景。这里用带缓冲的 channel 做信号量,限制同时执行的表数量。
4.7 序列同步
同步序列时,只需要读取源库序列的值,然后在目标库执行setval。
// migrate/sequence.go package migrate import ( "context" "fmt" "github.com/jackc/pgx/v5" ) // SyncSequences 同步源库序列到目标库 func SyncSequences(ctx context.Context, src, dst *pgx.Conn, schema string) error { rows, err := src.Query(ctx, ` SELECT sequencename, last_value, is_called FROM pg_sequences WHERE schemaname = $1 ORDER BY sequencename `, schema) if err != nil { return err } defer rows.Close() for rows.Next() { var name string var lastValue int64 var isCalled bool if err := rows.Scan(&name, &lastValue, &isCalled); err != nil { return err } // 在目标库设置序列值 _, err := dst.Exec(ctx, fmt.Sprintf( `SELECT setval('%s.%s', $1, $2)`, schema, name, ), lastValue, isCalled) if err != nil { return fmt.Errorf("set sequence %s.%s failed: %w", schema, name, err) } fmt.Printf("sequence synced: %s.%s = %d\n", schema, name, lastValue) } return rows.Err() }注意setval的第三个参数:is_called为true时,下次调用nextval返回last_value + 1;为false时,下次调用返回last_value。直接从pg_sequences读取的is_called就是源库的真实状态,直接传过去最准确。
4.8 main 函数串联全流程
最后在main.go中把整个流程串起来。
// main.go package main import ( "context" "fmt" "log" "time" "pgmigrate-demo/meta" "pgmigrate-demo/migrate" "pgmigrate-demo/util" ) func main() { cfg, err := ParseConfig() if err != nil { log.Fatalf("parse config failed: %v", err) } ctx, cancel := context.WithTimeout(context.Background(), 30*time.Minute) defer cancel() // 连接源库和目标库 srcPool, err := util.NewPool(ctx, cfg.SourceDSN) if err != nil { log.Fatalf("connect source failed: %v", err) } defer srcPool.Close() dstPool, err := util.NewPool(ctx, cfg.TargetDSN) if err != nil { log.Fatalf("connect target failed: %v", err) } defer dstPool.Close() // 获取单连接,便于方法复用 srcConn, err := srcPool.Acquire(ctx) if err != nil { log.Fatalf("acquire source conn failed: %v", err) } defer srcConn.Release() dstConn, err := dstPool.Acquire(ctx) if err != nil { log.Fatalf("acquire target conn failed: %v", err) } defer dstConn.Release() // 1. 读取源库元数据 tables, err := meta.LoadTables(ctx, srcConn.Conn(), cfg.Schema) if err != nil { log.Fatalf("load tables failed: %v", err) } fmt.Printf("found %d tables in schema %s\n", len(tables), cfg.Schema) // 2. 同步 Schema if err := migrate.SyncSchema(ctx, srcConn.Conn(), dstConn.Conn(), cfg.Schema); err != nil { log.Fatalf("sync schema failed: %v", err) } // 3. 同步数据(可选) if cfg.WithData && !cfg.OnlySchema { if err := migrate.CopyAllTables(ctx, srcConn.Conn(), dstConn.Conn(), tables, cfg.Concurrency, 10000); err != nil { log.Fatalf("copy data failed: %v", err) } } // 4. 同步序列 if err := migrate.SyncSequences(ctx, srcConn.Conn(), dstConn.Conn(), cfg.Schema); err != nil { log.Fatalf("sync sequences failed: %v", err) } fmt.Println("migration completed successfully") }到这一步,一个具备基础能力的 Pgmigrate 简化版就完成了。它能够完成 schema 创建、表数据并发拷贝、序列同步这几个核心环节。
5. 编译运行与验证
5.1 准备测试数据库
先在本地创建两个 PostgreSQL 数据库,分别模拟源库和目标库。
# 创建源库 createdb source_db # 创建目标库 createdb target_db在源库中创建测试表和序列。
-- 连接 source_db CREATE TABLE users ( id BIGSERIAL PRIMARY KEY, name TEXT NOT NULL, email TEXT UNIQUE, created_at TIMESTAMP DEFAULT now() ); INSERT INTO users (name, email) SELECT 'user_' || g, 'user_' || g || '@example.com' FROM generate_series(1, 10000) g;5.2 运行迁移命令
编译并运行工具。
go build -o pgmigrate-demo . ./pgmigrate-demo \ --source "postgres://postgres:postgres@localhost:5432/source_db?sslmode=disable" \ --target "postgres://postgres:postgres@localhost:5432/target_db?sslmode=disable" \ --schema public \ --concurrency 4预期输出类似:
found 1 tables in schema public schema synced: public.users shard [1, 10000] copied, rows=10000 total rows copied: 10000 sequence synced: public.users_id_seq = 10000 migration completed successfully5.3 验证数据
连接目标库,查询行数和序列值。
psql target_db -c "SELECT count(*) FROM users;" psql target_db -c "SELECT last_value, is_called FROM users_id_seq;"如果输出显示10000,且序列is_called为t,说明迁移成功。
6. 常见问题与排查思路
在开发和测试 Pgmigrate 类工具时,下面几个问题出现频率最高。
| 问题现象 | 常见原因 | 解决思路 |
|---|---|---|
| 连接源库超时 | 网络不通或防火墙限制 | 检查端口连通性:nc -vz host 5432,确认pg_hba.conf允许远程连接 |
pg_get_tabledef返回空 | PostgreSQL 版本过旧或函数不支持 | 改为使用pg_dump --schema-only,或通过系统目录自行组装 DDL |
COPY 数据时报invalid input syntax | 数据类型或格式不兼容 | 检查源库和目标库的 PostgreSQL 版本差异,确认没有使用不兼容的扩展类型 |
| 序列重复导致主键冲突 | 只同步了表数据,没有同步序列 | 数据拷贝后必须执行 sequence 同步,且is_called要带上 |
| 并发太高,源库 CPU 100% | 并发数设置过大 | 适当降低--concurrency,可以用 2 到 8 之间的值,配合监控逐步调整 |
| 目标表已存在导致创建失败 | 迁移多次运行,目标库有残留对象 | 迁移前执行清理策略,例如DROP TABLE IF EXISTS ... CASCADE,或者做增量迁移设计 |
| 字符串类型主键无法分片 | 分片逻辑只支持整数 | 改用ctid分片或全表拷贝,不要强行按字符串范围分片 |
再补充一个比较隐蔽的问题:使用information_schema.columns获取列名时,返回的顺序依赖列定义顺序,但是当表有 dropped 列时,information_schema的结果可能与pg_attribute的物理顺序不一致。遇到这种情况,建议直接读取pg_attribute.attnum来排序。
排查这类问题,可以按下面的顺序走:
- 先单独在源库执行元数据查询,确认返回结果是否符合预期。
- 再单独执行建表 DDL,确认目标库能否成功建表。
- 然后用
psql手动执行 COPY 导入,验证列格式是否对齐。 - 最后再用工具做全流程迁移,缩小问题范围。
7. 生产环境使用 Pgmigrate 类工具的工程建议
7.1 连接串与密钥管理
生产环境不要把数据库连接串直接写在命令行参数里,否则容易通过进程列表泄露密码。建议通过环境变量读取。
export PG_SOURCE_DSN="postgres://...?..." export PG_TARGET_DSN="postgres://...?..."7.2 先演练再迁移
每次正式迁移前,建议用生产环境的备份恢复一个测试实例,先跑一遍 Pgmigrate 全流程,记录耗时、并发数、源库性能指标,然后根据结果调整参数。
7.3 关注锁与在线业务的影响
数据拷贝过程中,源库需要执行SELECT查询,普通SELECT不会阻塞读写,但如果有VACUUM、DDL操作,可能会有短暂锁竞争。建议:
- 在业务低峰期执行迁移。
- 使用
row_exclusive_lock级别,避免长时间持锁。 - 大表迁移前先提前收集统计信息,减少执行计划偏差。
7.4 目标库关闭归档和复制
迁移期间目标库如果开启了 WAL 归档和流复制,写入放大倍数会很高。对于一次性迁移任务,可以临时关闭归档,迁移完成后再开启。
7.5 失败重试与断点续传
生产环境网络不会永远稳定。一个成熟迁移工具必须具备断点续传能力。最简单的方式是:为每张表记录当前拷贝到的分片位置,下次启动时跳过已完成分片。这需要在目标库创建一张迁移状态表,记录表名、分片范围、完成状态。
CREATE TABLE migration_state ( schema_name text NOT NULL, table_name text NOT NULL, shard_start bigint, shard_end bigint, status text NOT NULL DEFAULT 'pending', updated_at timestamptz DEFAULT now(), PRIMARY KEY (schema_name, table_name, shard_start) );这个表和业务数据放在同一个 Schema 或独立 Schema 中,迁移完成后删除或保留均可。
7.6 数据一致性校验
数据拷贝完成并不意味着绝对正确。生产环境建议做两层校验:
- 行数校验:每张表
count(*)一致。 - 抽样校验:对关键表用
ORDER BY 主键 LIMIT 1000随机抽几段,对比列值是否一致。
如果要更严格,可以对整表做md5(string_agg(...)),但这种校验在超大数据量下非常耗资源,通常只用于核心配置表。
8. 总结与下一步学习思路
这篇文章围绕Pgmigrate这个技术方向,完整拆解了 Pgcopydb 的核心迁移思路,并用 Go 实现了一个简化但可运行的版本。通过这个教学项目,可以掌握以下内容:
- Pgcopydb 的核心流程:元数据发现、Schema 同步、数据拷贝、序列同步、一致性校验。
- 为什么用 Go 重写这类工具会带来部署和并发模型上的优势。
pgx驱动中CopyFrom的使用方法,以及pgx.Rows如何直接作为CopyFromSource。- 大表分片拷贝的思路,以及并发调度中信号量和
errgroup的应用。 - 迁移工具生产化时需要考虑的断点续传、状态记录、连接串安全、数据校验等问题。
如果还想继续深入,可以从下面几个方向入手:
- 完整阅读 Pgcopydb 官方文档,理解它对无主键表、分区表、大对象、外键依赖顺序的处理方式。
- 研究 PostgreSQL 的逻辑复制协议,尝试把迁移工具改为增量同步模式。
- 学习
pgx底层COPY协议实现,理解流式写入和内存控制细节。 - 调研
FerretDB、CockroachDB等数据库工具链中 Go 语言的使用方式,学习它们如何组织数据库连接和任务调度。
数据库迁移是高风险操作,每一次生产迁移都要当成一次正式发布来做:先备份、再演练、做好回退预案。希望这篇文章能让你在理解 Pgmigrate 原理的同时,也能在自己的团队里从容落地一次 PostgreSQL 迁移任务。