Flink SQL CDC:无需代码,实现 PostgreSQL 到 ClickHouse 的实时数据同步
·
Flink SQL CDC:PostgreSQL 到 ClickHouse 实时同步方案
核心原理
通过 Flink SQL 的 CDC(Change Data Capture) 功能,直接捕获 PostgreSQL 的实时数据变更(插入/更新/删除),并写入 ClickHouse。整个过程无需编写代码,仅需 SQL 配置。
实现步骤
1. 环境准备
- Flink 集群:1.13+ 版本(需支持 SQL CDC)
- 依赖 JAR 包:
flink-sql-connector-postgres-cdc-2.3.0.jar(PostgreSQL 连接器)flink-connector-clickhouse-1.14.5.jar(ClickHouse 连接器)
2. 创建 PostgreSQL CDC 源表
CREATE TABLE pg_source (
id INT PRIMARY KEY,
name STRING,
update_time TIMESTAMP(3)
) WITH (
'connector' = 'postgres-cdc',
'hostname' = 'localhost',
'port' = '5432',
'username' = 'postgres',
'password' = 'password',
'database-name' = 'test_db',
'schema-name' = 'public',
'table-name' = 'source_table',
'decoding.plugin.name' = 'pgoutput'
);
3. 创建 ClickHouse 目标表
CREATE TABLE ch_sink (
id INT,
name STRING,
update_time TIMESTAMP(3),
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'clickhouse',
'url' = 'clickhouse://localhost:8123',
'database-name' = 'default',
'table-name' = 'sink_table',
'username' = 'default',
'password' = '',
'sink.batch-size' = '1000' -- 批量写入提升性能
);
4. 启动实时同步任务
INSERT INTO ch_sink
SELECT id, name, update_time
FROM pg_source;
关键技术点
-
变更数据捕获:
- PostgreSQL 通过逻辑解码(
pgoutput)捕获INSERT/UPDATE/DELETE事件。 - Flink 自动将数据转换为
INSERT流(删除操作转为-D消息)。
- PostgreSQL 通过逻辑解码(
-
数据一致性保证:
- Exactly-Once 语义:Flink Checkpoint 机制确保数据不丢失不重复。
- 自动 Schema 同步:源表结构变更时自动适配(需重启任务)。
-
性能优化:
- 批量写入:通过
sink.batch-size控制 ClickHouse 写入批次。 - 并行处理:根据数据量调整 Flink 并行度。
- 批量写入:通过
验证与监控
- 数据验证:
-- 在 ClickHouse 中查询数据 SELECT count(*) FROM sink_table; - 任务监控:
- Flink Web UI 查看
INSERT任务状态 - 日志检查
pg_cdc事件消费进度
- Flink Web UI 查看
常见问题处理
| 问题现象 | 解决方案 |
|---|---|
| 同步延迟高 | 增加 Flink 并行度或调大 sink.batch-size |
| 主键冲突 | 确认 ClickHouse 表引擎支持更新(如 ReplacingMergeTree) |
| 字段类型不匹配 | 在 SELECT 子句中使用 CAST 转换类型 |
注:此方案适用于 分钟级延迟 场景。如需亚秒级延迟,可启用 Flink 的
MiniBatch优化。
更多推荐
所有评论(0)