☰
PgDog 逻辑复制分片实践:用 `data-sync` 将 PostgreSQL 数据分发到分片集群
2026/10/12 3:03:25 网站建设 项目流程
  • 数据库
  • 后端

【免费下载链接】pgdog

PostgreSQL connection pooler, load balancer and database sharder.

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

导读

本文基于 PgDog 仓库中的integration/logical集成示例,讲解如何利用 PostgreSQL 原生的逻辑复制(Logical Replication)能力,通过 PgDog 的data-sync命令行工具把源库中的一张表复制到目标分片集群,并在复制过程中基于分片键(shard key)实时路由 WAL 变更。读完本文,你将掌握逻辑复制分片的环境搭建步骤、data-sync及相关子命令的完整参数用法、配套pgdog.toml/users.toml的配置要点,以及这条链路在 PgDog 源码中的实现原理与关键保障机制。

1. 逻辑复制分片是什么

逻辑复制是 PostgreSQL 内置的复制机制:源端发布(Publication)把表的 INSERT / UPDATE / DELETE 变更写入 WAL,通过pgoutput逻辑解码插件输出,订阅端按行级事件回放。PgDog 把它与自身的分片路由能力结合起来,构成「逻辑复制分片」(Logical replication sharding):

  • 源库(source)中的一张表被完整、持续地分发到多个目标分片(destination shard);
  • 每条 WAL 事件在 PgDog 内经过与线上查询完全相同的分片键评估流程,落到与业务写入一致的 shard 上(见 docs/REPLICATION.md 的StreamContext说明);
  • 整个过程由cargo run -->CREATE DATABASE pgdog; CREATE USER pgdog SUPERUSER PASSWORD 'pgdog' REPLICATION; \c pgdog CREATE SCHEMA pgdog; CREATE TABLE pgdog.books ( id BIGINT PRIMARY KEY, title VARCHAR, content VARCHAR );

    三个要点:

    1. REPLICATION权限必不可少。data-sync需要在源库创建并消费逻辑复制槽(replication slot),而创建复制槽要求用户拥有REPLICATION属性;SUPERUSER则保证data-sync能以 schema owner 身份在目标分片执行 DDL 与数据写入。
    2. 三张表结构必须一致。pgdog.books在源库与每个目标分片都完全同名同构,这是逻辑复制回放的前提。
    3. 主键(PRIMARY KEY)是分片与复制的双重锚点。它在后续既作为默认复制标识(REPLICA IDENTITY DEFAULT)参与 UPDATE / DELETE 定位,也是integration/logical/pgdog.toml中[[sharded_tables]]分片列声明的基础。

    然后在**源库(5432)**上创建发布(Publication),把需要复制的表纳入发布集:

    CREATE PUBLICATION books FOR TABLE pgdog.books;

    data-sync通过--publication books参数定位这张发布;在源码中,发布对象由Publisher管理(pgdog/src/backend/replication/logical/publisher/mod.rs),它会在需要时自动为发布内的表生成建表与复制所需的 SQL。

    3. 运行逻辑复制分片:data-sync命令

    环境就绪后,在仓库根目录执行:

    cargo run -- --from-database source --from-user pgdog --to-database destination --to-user pgdog --publication books

    对照 pgdog/src/cli.rs 中DataSync子命令的定义,该命令的完整可用参数如下:

    参数类型说明
    --from-database <NAME>必填源数据库名,对应pgdog.toml中[[databases]]的名称
    --from-user <USER>必填连接源库使用的用户名
    --to-database <NAME>必填目标数据库名(分片集群),同样对应pgdog.toml配置
    --to-user <USER>必填连接目标库使用的用户名
    --publication <NAME>必填源库上创建的发布名称,如books
    --replicate-only布尔只做增量复制,跳过初始数据拷贝(默认false)
    --sync-only布尔只做数据拷贝与表同步,跳过复制与切换(默认false)
    --replication-slot <NAME>可选指定要创建/使用的复制槽名;不传时自动生成
    --skip-schema-sync布尔跳过前序与后序 schema 同步(默认false)

    说明:README 中的--from-user/--to-user用于指定认证用户;在源码中用户名解析到pgdog.toml中对应数据库条目携带的密码(Databases::passwords,见 pgdog/src/backend/databases.rs),所以两侧数据库条目必须在配置中声明同一用户。

    data-sync的入口在 pgdog/src/cli.rs 的data_sync()函数:它先用--from-database/--to-database/--publication构造ReshardingState(源集群、目标集群、发布名、复制槽名),再以--skip-schema-sync、--replicate-only、--sync-only等开关构建ReshardTask并运行。命令执行期间按 Ctrl-C 会优雅取消任务(复制会在断开前收尾),不会硬杀进程(对应run_to_completion的ctrl_c分支)。

    4. 配套配置解读:pgdog.toml与users.toml

    4.1 数据库与分片声明

    integration/logical/pgdog.toml 完整展示了逻辑复制分片的配置骨架:

    [general] [rewrite] enabled = false shard_key = "ignore" split_inserts = "error" [[databases]] name = "pgdog" host = "127.0.0.1" [[databases]] name = "source" host = "127.0.0.1" port = 5432 database_name = "pgdog" min_pool_size = 0 [[databases]] name = "destination" host = "127.0.0.1" port = 5433 database_name = "pgdog" min_pool_size = 0 shard = 0 # [[databases]] # name = "destination" # host = "127.0.0.1" # port = 5434 # database_name = "pgdog" # min_pool_size = 0 # shard = 1 [[sharded_tables]] database = "destination" name = "books" column = "id" data_type = "bigint"

    配置要点:

    • source与destination都是[[databases]]条目:data-sync的--from-database source、--to-database destination正是按这里的name解析集群连接信息(database_name指向真实库名,port区分实例)。
    • shard = 0/shard = 1声明分片:多个name = "destination"但端口不同、shard不同的条目构成目标分片集群。示例中分片 1 以注释形式给出,取消注释即可启用第二个分片。
    • min_pool_size = 0:避免 PgDog 在启动阶段为这两个仅供迁移使用的库预先建立空闲连接。
    • [[sharded_tables]]声明分片表:database、name指定表所属集群与表名,column = "id"指定分片键,data_type = "bigint"指定键类型。data-sync与复制引擎依据它判断每条 WAL 事件应该路由到哪个分片。
    • [rewrite]中enabled = false、shard_key = "ignore"表明该示例不启用 SQL 改写,分片信息全部交由分片键路由处理。

    4.2 用户与复制模式

    integration/logical/users.toml 中为三种用途各声明了一个用户条目:

    [[users]] database = "pgdog" name = "pgdog" password = "pgdog" replication_mode = true [[users]] database = "source" name = "pgdog" password = "pgdog" [[users]] database = "destination" name = "pgdog" password = "pgdog"
    • replication_mode = true的用户(对应pgdog库)用于与 PgDog 建立复制/管理连接;
    • source与destination两个条目供data-sync分别连接源集群与目标分片集群做数据迁移。

    4.3 可调参数:复制与拷贝的并发、重试

    从源码看,逻辑复制拷贝与复制的行为还可以通过[general]下的参数进一步调优(定义见 pgdog-config/src/general.rs):

    参数默认值作用
    resharding_parallel_copies1并发启动的表拷贝数量(与可用副本数无关)
    resharding_parallel_within_table_copies1单张表内部并发读取的 COPY 源连接数;对 TOAST 重表(大字段多、每次读多块磁盘)建议调高以打满磁盘 IOPS
    resharding_copy_retry_max_attempts5单张表拷贝失败后的最大重试次数,指数退避从resharding_copy_retry_min_delay开始
    resharding_copy_retry_min_delay1000拷贝重试的基础延迟(毫秒),每次尝试翻倍,最高 32 倍
    resharding_replication_retry_max_attempts5复制订阅端连续出错的最大容忍次数;每次失败触发slot.reconnect(),PostgreSQL 会从上次已确认提交处重新流式传输;0表示无限重试
    resharding_replication_retry_min_delay1000复制订阅端重试间隔(毫秒)
    resharding_copy_formatCopyFormat默认值拷贝期间COPY语句使用的格式;注意:主键从INTEGER迁移到BIGINT时必须使用文本格式

    这些参数让大规模表迁移可以在「并行拷贝」与「失败重试」两个维度上按数据特征调优,是integration/logical示例之外的进阶配置。

    5. 大规模数据演练:用 Gutenberg 数据集压测

    仓库提供了配套的数据注入脚本 integration/logical/gutenberg.py,用于演练真实规模的数据同步:

    python3 gutenberg.py /path/to/archive 80000

    脚本会连接127.0.0.1:5432/pgdog,自动建表(CREATE TABLE IF NOT EXISTS pgdog.books),清空旧数据(TRUNCATE TABLE pgdog.books),然后用COPY pgdog.books (id, title, content) FROM STDIN把 Gutenberg 图书元数据与全文批量灌入源库,并实时打印吞吐(KB/s)。README 指出该数据集灌入 PostgreSQL 后约 16GB,适合验证data-sync在大表下的拷贝吞吐与复制稳定性。

    注意:运行脚本前需pip install psycopg tqdm;数据来自公开的 Gutenberg 图书数据集归档目录,需先解压再传入archive文件夹路径。

    6. 源码级实现原理:data-sync 如何完成「拷贝 + 复制 + 分片路由」

    6.1 总览:ReshardTask 的流水线

    data-sync最终运行的是ReshardTask(pgdog/src/api/resharding.rs 与 docs/RESHARDING.md),它把一次数据同步拆成有序阶段:

    Pre-data schema(前序 schema) → Bulk COPY(批量数据拷贝,DataSyncTask) → Post-data schema(后序 schema) → Table synchronization(表同步) → Validation(校验) → Forward replication(前向复制) → Cutover(切换,仅 RESHARD/ReplicateAndCutover)
    • --sync-only停在「表同步」阶段,不做复制与切换;
    • --replicate-only跳过数据拷贝,直接进入前向复制;
    • 手动模式下,COPY_DATA/data-sync之后可用schema-sync --phase post在复制运行期间补跑后序 schema(该阶段默认忽略语句错误,见 docs/RESHARDING.md)。

    6.2 数据拷贝:临时复制槽 + 一致快照 + 并行 COPY

    核心拷贝逻辑在 pgdog/src/backend/replication/logical/data_sync.rs。单表拷贝(copy_table)的流程是:

    1. 创建临时逻辑复制槽:ReplicationSlot::new_temporary调用CREATE_REPLICATION_SLOT ... LOGICAL "pgoutput" (SNAPSHOT 'use')(见 pgdog/src/backend/replication/logical/publisher/replication_slot.rs)。临时槽随连接关闭自动释放,其一致点(consistent point)LSN 成为本次拷贝的水位线。
    2. 导出快照:SELECT pg_export_snapshot()取得与复制槽一致的快照标识;所有并行读连接复用同一快照(connect_reader中SET TRANSACTION SNAPSHOT),保证各读连接看到同一份表数据,避免重复行带来的同步问题。
    3. 并行 COPY:发布端CopyPublisher把表按 ctid 块范围拆分,多个 Tokio 任务并行执行COPY ... TO STDOUT;订阅端CopySubscriber在目标分片执行COPY ... FROM STDIN。并行度受resharding_parallel_within_table_copies控制,且会在表行数很少时自动降级(pgdog/src/backend/replication/logical/data_sync.rs 中依据TableColumnSplit的块数判断)。
    4. 记录水位线并排空:拷贝完成后COMMIT、向复制槽发送StatusUpdate确认已消费到table.lsn,然后排空槽内剩余消息,返回带 LSN 的表信息。该 LSN 是后续表同步与前向复制的衔接点。

    6.3 前向复制:WAL 消息 → 分片路由 → 预处理语句执行

    复制引擎(docs/REPLICATION.md)由两个核心模块构成:

    • Publisher(pgdog/src/backend/replication/logical/publisher/publisher_impl.rs):打开到源库的流式复制连接,消费解码后的XLogPayload,转发给订阅端,并跟踪每张表的复制延迟供切换逻辑使用。
    • StreamSubscriber(pgdog/src/backend/replication/logical/subscriber/stream.rs):有状态的消息处理器。它按表 OID 缓存预处理语句集,维护每张表的 LSN 水位线(滤掉第 6.2 节中已批量拷贝过的行),并为每个目标分片保持一条持久连接。

    每条 INSERT / UPDATE / DELETE 事件在到达目标分片前经过三重闸门:

    事件 → LSN ≤ 表水位线? → 是:跳过(已拷贝) → 否:StreamContext::shard() 评估分片键 → 绑定并执行预处理语句

    其中分片路由由StreamContext(pgdog/src/backend/replication/logical/subscriber/context.rs)完成:它从 WAL 元组中提取分片键,复用与线上查询一致的ContextBuilder → Context::apply()管线,保证 WAL 行落到与应用程序写入相同的分片上。语句形状由Table(pgdog/src/backend/replication/logical/publisher/table.rs)在收到Relation消息时一次性生成并缓存:

    操作语句形状
    INSERT(分片表)普通INSERT
    INSERT(omni 表)INSERT … ON CONFLICT (identity_cols) DO UPDATE SET …
    UPDATEUPDATE … SET 非标识列 = $N WHERE identity_cols = $M
    UPDATE(部分)同上 WHERE,SET 仅限未 TOAST 的列(形状缓存)
    DELETEDELETE … WHERE identity_cols = $N

    6.4 复制标识与 TOAST 边界情况

    实现层面对两种复制标识做了差异化处理:

    • REPLICA IDENTITY DEFAULT/USING INDEX:标识列保证 NOT NULL,UPDATE / DELETE 的 WHERE 用普通=即可;每条事件恰好路由到一个分片,无需广播。
    • REPLICA IDENTITY FULL:无主键/唯一索引的表在 WAL 中携带完整旧行。分片 FULL 表仍按分片键单点路由;omni FULL 表则复制到所有分片,并要求目标分片存在「防止 NULL 键重复」的唯一索引(所有键列NOT NULL,或 PG15+ 的NULLS NOT DISTINCT唯一索引),否则connect()会在开始流式传输前明确拒绝(错误FullIdentityOmniNoUniqueIndex)。

    另一个关键边界是未变化的 TOAST 列:PostgreSQL 对 UPDATE 中未触碰的大字段只写'u'(unchanged)标记而不写数据。复制端若直接把空槽位写进目标库,会把一个有效的大值静默覆盖为空。PgDog 在收到'u'标记时从旧元组补齐该列的值(update_partial路径),从而避免数据损坏(详细说明见 docs/REPLICATION.md)。

    7. 其他 CLI 子命令与使用边界

    围绕逻辑复制分片,pgdog/src/cli.rs 还提供了两个相关子命令:

    • schema-sync:把 schema 从源集群同步到目标集群,按--phase(pre/post/cutover/ 校验)分段执行,--dry-run只打印语句不执行,--ignore-errors忽略错误;post 阶段默认忽略错误。
    • replicate-and-cutover:一次性完成「schema 同步 + 数据同步 + 复制 + 触发切换」,但源码注释明确标注仅供内部测试使用——生产环境的完整切换应通过 admin 数据库的RESHARD <source> <destination> <publication>;命令驱动(多节点部署下的协调切换需要 Enterprise 控制面,见 docs/RESHARDING.md)。

    使用边界提醒:

    • 不要在一个复制任务运行期间启动另一个复制任务;
    • data-sync --skip-schema-sync跳过前后序 schema 同步,但任务仍会在前向复制前完成表同步;
    • post 阶段 schema 同步在目标库写入活跃后不应重跑,重跑会重建已有索引。

    8. 日志与观测

    仓库提供了 integration/logical/log.sh 帮助观察运行过程:

    #!/bin/bash touch log.txt cargo run > log.txt 2>&1 & pid=$! trap shutdown INT function shutdown() { kill -TERM $pid } tail -f log.txt

    它将 PgDog 前台输出重定向到log.txt并后台运行,捕获进程 PID,收到 Ctrl-C 时发送TERM优雅关闭,并实时tail日志。配合data-sync运行过程中打印的data sync for "pgdog"."books" started/finished at lsn ...日志(见 pgdog/src/backend/replication/logical/data_sync.rs),可以直观确认每个分片的拷贝起点与进度。

    总结

    integration/logical示例完整演示了 PgDog「逻辑复制分片」的最小可行路径:三套实例 + 一张发布 + 一条data-sync命令,即可把源表数据按分片键持续分发到目标分片集群。其底层由 PgDog 的复制引擎(Publisher+StreamSubscriber)与分片路由(StreamContext)协作完成:临时复制槽提供一致快照用于并行 COPY,LSN 水位线衔接拷贝与复制,WAL 事件经与线上查询相同的路由管线落到正确的分片。理解这条链路后,无论是为业务搭建读写分离式的多分片集群,还是为后续使用RESHARD做在线分片迁移,都打下了扎实的实践基础。

    • 数据库
    • 后端

    【免费下载链接】pgdog

    PostgreSQL connection pooler, load balancer and database sharder.

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

    相关推荐

    上一篇:【亲测免费】 Pear Admin Flask:Flask后台管理系统的快速开发利器
    下一篇:【亲测免费】 EventOS Nano:轻量级事件驱动嵌入式开发平台

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

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

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

立即咨询