混合云数据库同步:基于 Debezium 实现私有云 PostgreSQL→AWS RDS 的 CDC 同步
·
混合云数据库同步:基于 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. 关键配置项
| 组件 | 参数 | 说明 |
|---|---|---|
| PostgreSQL | wal_level=logical | 必须启用逻辑 WAL 解码 |
| Debezium | plugin.name=pgoutput | PostgreSQL 10+ 原生逻辑复制插件 |
slot.name | 防止数据丢失的持久化复制槽 | |
| AWS RDS | insert.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)
- Kafka 启用 SSL (配置
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 仅处理增量变更。
更多推荐
所有评论(0)