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;


关键技术点
  1. 变更数据捕获

    • PostgreSQL 通过逻辑解码(pgoutput)捕获 INSERT/UPDATE/DELETE 事件。
    • Flink 自动将数据转换为 INSERT 流(删除操作转为 -D 消息)。
  2. 数据一致性保证

    • Exactly-Once 语义:Flink Checkpoint 机制确保数据不丢失不重复。
    • 自动 Schema 同步:源表结构变更时自动适配(需重启任务)。
  3. 性能优化

    • 批量写入:通过 sink.batch-size 控制 ClickHouse 写入批次。
    • 并行处理:根据数据量调整 Flink 并行度。

验证与监控
  1. 数据验证
    -- 在 ClickHouse 中查询数据
    SELECT count(*) FROM sink_table;
    

  2. 任务监控
    • Flink Web UI 查看 INSERT 任务状态
    • 日志检查 pg_cdc 事件消费进度

常见问题处理
问题现象解决方案
同步延迟高增加 Flink 并行度或调大 sink.batch-size
主键冲突确认 ClickHouse 表引擎支持更新(如 ReplacingMergeTree
字段类型不匹配SELECT 子句中使用 CAST 转换类型

:此方案适用于 分钟级延迟 场景。如需亚秒级延迟,可启用 Flink 的 MiniBatch 优化。

更多推荐