ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

PgDog 逻辑复制分片实践:用 `data-sync` 将 PostgreSQL 数据分发到分片集群

PgDog 逻辑复制分片实践:用 `data-sync` 将 PostgreSQL 数据分发到分片集群 数据库后端【免费下载链接】pgdogPostgreSQL 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 );三个要点REPLICATION权限必不可少。data-sync需要在源库创建并消费逻辑复制槽replication slot而创建复制槽要求用户拥有REPLICATION属性SUPERUSER则保证data-sync能以 schema owner 身份在目标分片执行 DDL 与数据写入。三张表结构必须一致。pgdog.books在源库与每个目标分片都完全同名同构这是逻辑复制回放的前提。主键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.toml4.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 pgdogreplication_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 重表大字段多、每次读多块磁盘建议调高以打满磁盘 IOPSresharding_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最终运行的是ReshardTaskpgdog/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的流程是创建临时逻辑复制槽ReplicationSlot::new_temporary调用CREATE_REPLICATION_SLOT ... LOGICAL pgoutput (SNAPSHOT use)见 pgdog/src/backend/replication/logical/publisher/replication_slot.rs。临时槽随连接关闭自动释放其一致点consistent pointLSN 成为本次拷贝的水位线。导出快照SELECT pg_export_snapshot()取得与复制槽一致的快照标识所有并行读连接复用同一快照connect_reader中SET TRANSACTION SNAPSHOT保证各读连接看到同一份表数据避免重复行带来的同步问题。并行 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的块数判断。记录水位线并排空拷贝完成后COMMIT、向复制槽发送StatusUpdate确认已消费到table.lsn然后排空槽内剩余消息返回带 LSN 的表信息。该 LSN 是后续表同步与前向复制的衔接点。6.3 前向复制WAL 消息 → 分片路由 → 预处理语句执行复制引擎docs/REPLICATION.md由两个核心模块构成Publisherpgdog/src/backend/replication/logical/publisher/publisher_impl.rs打开到源库的流式复制连接消费解码后的XLogPayload转发给订阅端并跟踪每张表的复制延迟供切换逻辑使用。StreamSubscriberpgdog/src/backend/replication/logical/subscriber/stream.rs有状态的消息处理器。它按表 OID 缓存预处理语句集维护每张表的 LSN 水位线滤掉第 6.2 节中已批量拷贝过的行并为每个目标分片保持一条持久连接。每条 INSERT / UPDATE / DELETE 事件在到达目标分片前经过三重闸门事件 → LSN ≤ 表水位线 → 是跳过已拷贝 → 否StreamContext::shard() 评估分片键 → 绑定并执行预处理语句其中分片路由由StreamContextpgdog/src/backend/replication/logical/subscriber/context.rs完成它从 WAL 元组中提取分片键复用与线上查询一致的ContextBuilder → Context::apply()管线保证 WAL 行落到与应用程序写入相同的分片上。语句形状由Tablepgdog/src/backend/replication/logical/publisher/table.rs在收到Relation消息时一次性生成并缓存操作语句形状INSERT分片表普通INSERTINSERTomni 表INSERT … ON CONFLICT (identity_cols) DO UPDATE SET …UPDATEUPDATE … SET 非标识列 $N WHERE identity_cols $MUPDATE部分同上 WHERESET 仅限未 TOAST 的列形状缓存DELETEDELETE … WHERE identity_cols $N6.4 复制标识与 TOAST 边界情况实现层面对两种复制标识做了差异化处理REPLICA IDENTITY DEFAULT/USING INDEX标识列保证 NOT NULLUPDATE / DELETE 的 WHERE 用普通即可每条事件恰好路由到一个分片无需广播。REPLICA IDENTITY FULL无主键/唯一索引的表在 WAL 中携带完整旧行。分片 FULL 表仍按分片键单点路由omni FULL 表则复制到所有分片并要求目标分片存在「防止 NULL 键重复」的唯一索引所有键列NOT NULL或 PG15 的NULLS NOT DISTINCT唯一索引否则connect()会在开始流式传输前明确拒绝错误FullIdentityOmniNoUniqueIndex。另一个关键边界是未变化的 TOAST 列PostgreSQL 对 UPDATE 中未触碰的大字段只写uunchanged标记而不写数据。复制端若直接把空槽位写进目标库会把一个有效的大值静默覆盖为空。PgDog 在收到u标记时从旧元组补齐该列的值update_partial路径从而避免数据损坏详细说明见 docs/REPLICATION.md。7. 其他 CLI 子命令与使用边界围绕逻辑复制分片pgdog/src/cli.rs 还提供了两个相关子命令schema-sync把 schema 从源集群同步到目标集群按--phasepre/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 21 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 的复制引擎PublisherStreamSubscriber与分片路由StreamContext协作完成临时复制槽提供一致快照用于并行 COPYLSN 水位线衔接拷贝与复制WAL 事件经与线上查询相同的路由管线落到正确的分片。理解这条链路后无论是为业务搭建读写分离式的多分片集群还是为后续使用RESHARD做在线分片迁移都打下了扎实的实践基础。赞分享数据库后端【免费下载链接】pgdogPostgreSQL connection pooler, load balancer and database sharder.项目地址https://gitcode.com/gh_mirrors/pg/pgdog点击查看免费下载相关推荐终极指南如何使用PgDog实现COPY命令自动分片到多个PostgreSQL数据库终极指南如何使用PgDog实现COPY命令自动分片到多个PostgreSQL数据库 PgDog是一款功能强大的PostgreSQL连接池、负载均衡器和数据库分数据库后端如何实现PostgreSQL跨分片数据同步PgDog逻辑复制的终极指南如何实现PostgreSQL跨分片数据同步PgDog逻辑复制的终极指南 PgDog作为一款功能强大的PostgreSQL连接池、负载均衡器和数据库分片工具其数据库后端终极指南如何使用PgDog实现PostgreSQL分片集群的自动查询路由终极指南如何使用PgDog实现PostgreSQL分片集群的自动查询路由 PgDog是一款功能强大的PostgreSQL连接池、负载均衡器和数据库分片工具能数据库后端上一篇【亲测免费】 Pear Admin FlaskFlask后台管理系统的快速开发利器下一篇【亲测免费】 EventOS Nano轻量级事件驱动嵌入式开发平台创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进