混合云数据库同步:基于 Debezium 实现私有云 PostgreSQL → AWS RDS 的 CDC 同步

1. 核心原理
  • 变更数据捕获 (CDC):Debezium 通过 PostgreSQL 的逻辑解码(Logical Decoding)捕获数据变更事件(INSERT/UPDATE/DELETE),生成事件流。
  • 数据流架构
    PostgreSQL (私有云) → Debezium → Kafka → Sink Connector → AWS RDS
    

  • 关键公式
    • 事件流处理延迟:$$ \Delta t = t_{\text{process}} + t_{\text{network}} $$ 其中 $t_{\text{process}}$ 为处理时间,$t_{\text{network}}$ 为跨云网络延迟。
2. 实现步骤
(1) 源端配置(私有云 PostgreSQL)
-- 启用逻辑复制
ALTER SYSTEM SET wal_level = logical;
ALTER SYSTEM SET max_replication_slots = 10;
-- 创建复制用户
CREATE ROLE debezium WITH REPLICATION LOGIN PASSWORD 'secure_pass';
-- 授予表权限
GRANT SELECT ON ALL TABLES IN SCHEMA public TO debezium;

(2) Debezium 连接器部署
// debezium-source.json
{
  "name": "pg-cdc-connector",
  "config": {
    "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
    "database.hostname": "private-cloud-pg-host",
    "database.port": "5432",
    "database.user": "debezium",
    "database.password": "secure_pass",
    "database.dbname": "source_db",
    "plugin.name": "pgoutput",
    "slot.name": "debezium_slot",
    "publication.name": "dbz_publication",
    "table.include.list": "public.*",
    "topic.prefix": "cdc_pg"
  }
}

(3) 目标端同步(AWS RDS)
// jdbc-sink.json (使用 Kafka JDBC Sink Connector)
{
  "name": "rds-sink-connector",
  "config": {
    "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
    "connection.url": "jdbc:postgresql://aws-rds-endpoint:5432/target_db",
    "connection.user": "rds_user",
    "connection.password": "aws_secure_pass",
    "topics": "cdc_pg.public.*",
    "insert.mode": "upsert",
    "pk.mode": "record_key",
    "pk.fields": "id",
    "auto.create": "false",
    "auto.evolve": "false"
  }
}

3. 关键配置项
组件参数说明
PostgreSQLwal_level=logical必须启用逻辑 WAL 解码
Debeziumplugin.name=pgoutputPostgreSQL 10+ 原生逻辑复制插件
slot.name防止数据丢失的持久化复制槽
AWS RDSinsert.mode=upsert处理主键冲突时自动转为 UPDATE
pk.fields=id指定目标表主键字段
4. 网络与安全
  • 网络打通
    • 私有云 ⇄ AWS VPC 通过 VPN 或 Direct Connect 建立加密通道
    • 安全组开放端口:PostgreSQL(5432), Kafka(9092)
  • 加密传输
    • Kafka 启用 SSL (配置 security.protocol=SSL)
    • PostgreSQL 连接使用 SSL 模式 (sslmode=require)
5. 验证与监控
# 检查 Debezium 事件流
kafka-console-consumer --bootstrap-server kafka-host:9092 --topic cdc_pg.public.users

# 监控指标
curl -s http://debezium-host:8083/connectors/pg-cdc-connector/status | jq .

  • 关键指标
    • 延迟:source_record_lag_ms (目标:< 1000ms)
    • 积压:max_offset_diff (目标:≈0)
6. 故障处理
  • 数据冲突:在 Sink Connector 中配置 delete.enabled=true 级联删除
  • 网络中断:Kafka 保留策略设置 retention.ms=604800000(保留7天)
  • Schema 变更:启用 auto.evolve=true(测试环境),生产环境需手动同步 DDL

:全量初始化需使用 pg_dump 或 AWS DMS,Debezium 仅处理增量变更。

更多推荐