尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Flink CDC 实战:PostgreSQL 实时同步到数据仓库的完整指南与避坑手册

Flink CDC 实战:PostgreSQL 实时同步到数据仓库的完整指南与避坑手册 先交代一个背景。我在2023年维护过一条订单数据同步链路用的是 Debezium Kafka Kafka Connect 的组合从 PostgreSQL 同步到数据仓库。听起来挺正规但实际维护起来相当烦一段同步任务要写 Java UDF配置要拆成三份线上任务出了延迟排查链路要跨三个系统来回跳每次看日志都感觉自己是个接线员。后来 Flink CDC 3.0 发布我用它重构了这条链路才意识到实时同步这件事——尤其是从 PostgreSQL 到下游数仓的 CDC 同步——完全可以做到只写一份 YAML 配置跑起来前后不到五分钟。这篇文章不打算重复官方文档我想按照我实际踩完一圈之后的经验把 Flink CDC 同步 PostgreSQL 的完整步骤、核心原理、以及那些文档里不会写的坑一次性讲清楚。阅读对象是两类人一是刚接触 CDC、想快速把 PostgreSQL 数据实时同步到 StarRocks / Doris / Kafka 的工程师二是已经在用 Flink CDC 但被时区、WAL 暴涨、类型映射这类问题纠缠过的人。内容会有点长因为我不想只给结论排查过程和配置背后的理由更重要。1. 为什么是 Flink CDC摆脱 DebeziumKafka 这条老路的理由1.1 CDC 本质PostgreSQL 的逻辑复制是怎么工作的在聊 Flink CDC 之前必须先搞懂一个底层概念CDC 到底是靠什么机制捕捉数据变化的。PostgreSQL 要实现 CDC依赖的是逻辑复制Logical Replication能力。它的核心工作流程可以这样理解PostgreSQL 会把所有数据变更写进一个叫 WALWrite-Ahead Log的日志文件里你可以把它想象成数据库的“流水账本”——每一条 INSERT、UPDATE、DELETE 都会被追加记录进去。普通的流复制物理复制是把整本账本原样拷贝给备库而逻辑复制则是借助一个叫“逻辑解码”Logical Decoding的机制把 WAL 里的二进制变更记录解析成结构化的、人可以读懂的操作流比如“第几行数据被更新成了什么值”。这个解析过程需要一个解码插件PostgreSQL 自带的叫pgoutput它也是 Flink CDC 默认使用的插件。要让 PG 开启逻辑复制能力必须把wal_level参数设置为logical。注意这个参数修改后需要重启 PostgreSQL 才能生效。很多第一次做 CDC 的人卡在第一步就是改了参数忘了重启然后 Flink 任务一直报“WAL 不可用”之类的错误。还有一个概念必须理解复制槽Replication Slot。逻辑复制靠复制槽来记录下游消费到了 WAL 的哪个位置。Flink CDC 在启动任务时会在 PG 里自动创建一个复制槽每消费一条变更就把复制槽的 LSN日志序列号往前推进一点。这个机制的好处是即使 Flink 任务崩溃重启它也能从上次消费的位置继续保证不丢数据。坏处是如果任务停了或者复制槽没人消费PostgreSQL 为保留未消费的 WAL 文件日志会一直堆积直到把磁盘撑爆。这个坑我在第 4 节会重点讲因为它是我踩过最深的坑。1.2 对比旧方案为什么一份 YAML 比三套系统省心在 Flink CDC 3.0 之前业界主流的 CDC 方案是 Debezium Kafka Kafka Connect。这套架构本身没问题它很成熟生产环境里有大量验证。但从工程维护的视角看它有几个让人头疼的地方。Debezium 负责从数据库抓取变更Kafka Connect 负责运行 Debezium 连接器Kafka 负责把变更事件在中间传一轮。如果想对数据做一些清洗加工还得自己写 Kafka Streams 或者起一个 Flink 任务去消费 Kafka再把结果写进目标系统。也就是说一条简单的“PG 到数仓”的链路至少需要维护三个组件每个组件都有自己独立的配置、监控、日志和故障恢复方式。Flink CDC 3.0 引入的 YAML Pipeline 模式把这件事大幅简化了。它把“捕获变更”“传输变更”“写入目标”三个阶段全部封装进一个 Flink 任务里用户只需要用 YAML 声明“从哪里读、写到哪里、表怎么映射”启动时用一个flink-cdc.sh脚本一份配置文件就能拉起整个同步任务。对比维度Debezium Kafka Kafka ConnectFlink CDC YAML Pipeline组件数量至少 3 个1 个 Flink 任务同步逻辑开发需要写 Java/SQL UDF声明式 YAML 配置全量增量衔接需额外配置 Snapshot 模式内置自动完成Schema 变更同步需额外开发或依赖特定 Sink3.0 起支持 Schema Evolution运维监控分散在多个系统Flink Web UI 统一查看上手成本高涉及多种概念低会写配置即可从这张表能看出来Flink CDC 3.0 之后PG 实时同步这件事的门槛确实被拉低了一个数量级。但门槛低不代表没有坑接下来我就从环境准备开始把每一步需要注意的东西完整过一遍。2. 环境搭建除了 Flink 本体这三个细节最容易暗算人2.1 版本兼容性不要无脑选最新版很多教程会直接让你下载最新版的 Flink 和 Flink CDC但版本组合不对任务启动了也起不来或者起来之后行为诡异。我自己在 3.0.1 版本上配合 Flink 1.18.0 跑过完整的 PG 同步链路稳定运行了很久这个组合可以作为默认参考。版本兼容的基本逻辑是这样的Flink CDC 是一个独立于 Flink 的发行包但不同的 Flink CDC 版本支持的 Flink 版本范围不同PostgreSQL Connector 又依赖 Flink CDC 核心版本。建议遵循两个原则第一Flink 和 Flink CDC 的版本尽量参考官方 Release Notes 里明确列出的组合不要凭感觉混搭第二PostgreSQL 的版本不要太老9.6 以上的 PG 都能支持但老版本 PG 的pgoutput插件在某些功能上有限制比如旧版本不支持流式变更的某些元数据字段。我整理了一个常用组合参考表以我实操过的为准Flink CDC 版本匹配 Flink 版本适配 PostgreSQL 版本3.0.x1.15.x - 1.18.x9.6 - 163.1.x1.17.x - 1.18.x9.6 - 163.2.x1.18.x - 1.19.x9.6 - 17这里给一个真实教训不要为了用最新功能就去追 Flink CDC 3.2 Flink 1.19 的组合如果下游 Sink 连接器比如 StarRocks Connector还没有适配这个组合你会花大量时间在排查连接器冲突上。生产环境优先选已经发布半年以上、社区验证充分的版本组合。2.2 PostgreSQL 侧配置改两个地方否则任务起不来Flink CDC 要能从 PG 读取变更流PG 侧必须满足三个前置条件开启逻辑复制、创建有复制权限的账号、创建发布Publication。第一步修改wal_level为logical。我建议直接修改postgresql.conf文件或者用ALTER SYSTEM SET命令两者都会在重启后生效ALTER SYSTEM SET wal_level logical;重启 PG 后确认配置生效SHOW wal_level;如果看到的还是replica说明没改成功检查一下 PG 版本和配置文件路径老版本 PG9.6 以下对ALTER SYSTEM的支持有差异这种情况直接改postgresql.conf更保险。第二步创建专用同步账号。我不建议直接用超级管理员跑同步任务风险太大。最小权限原则下CDC 账号只需要逻辑复制权限和读取目标表的权限CREATE USER flinkuser WITH REPLICATION LOGIN PASSWORD your_password; GRANT CONNECT ON DATABASE your_db TO flinkuser; GRANT SELECT ON ALL TABLES IN SCHEMA public TO flinkuser;第三步创建发布CREATE PUBLICATION flink_pub FOR ALL TABLES;这里有个容易被忽略的点如果后续新建了表默认不会自动加入发布需要手动执行ALTER PUBLICATION flink_pub ADD TABLE new_table;或者在建表时就加进去。Flink CDC 在启动任务时如果检测到表不在发布里会尝试自动创建发布但需要账号有相应权限所以上面创建账号时如果业务允许可以顺带授权CREATE PUBLICATION。2.3 JDBC 驱动与依赖一个 2MB 的 jar 能卡你半小时Flink CDC 的 PostgreSQL Connector 本身是内置在发行包里的但它运行时需要 PostgreSQL 的 JDBC 驱动来做初始的全量快照读取。这个驱动不会自动下载需要你手动把postgresql-42.x.x.jar放到 Flink 的lib目录下。我刚上手时就被这个细节坑过任务启动后没有任何报错但全量阶段一直卡住不动日志里反复出现“Loading classcom.mysql.jdbc.Driver ...”后来才意识到是驱动缺失。放好驱动后还要确认驱动版本和 PG 服务端版本匹配一般 PG 12 以上用postgresql-42.6.0 之后的版本不会有问题。依赖方面再提醒一点如果下游是 StarRocks 或 Doris记得把对应的 Sink Connector 也一起放进 Flink 的lib目录。很多人任务启动时报“ClassNotFoundException: com.starrocks.connector.flink.StarRocksDynamicTableSinkFactory”就是因为只放了 CDC 核心包没放 Sink 包。3. 五步跑通实时同步从建表到 YAML 配置的完整实操3.1 准备一张源表用什么数据验证最直观为了能充分验证同步效果我建议源表里包含几种典型的数据类型主键字段、字符串字段、数值字段、时间字段。下面这张订单表比较有代表性CREATE TABLE orders ( id BIGINT PRIMARY KEY, order_no VARCHAR(50) NOT NULL UNIQUE, amount NUMERIC(10, 2) NOT NULL, status VARCHAR(20) DEFAULT CREATED, created_at TIMESTAMPTZ DEFAULT NOW() );这张表故意设计了几个“坑点”amount用了NUMERIC类型、created_at用了TIMESTAMPTZ。后面你会看到这两个类型在同步过程中各自有坑。先插入三行数据做全量验证INSERT INTO orders (id, order_no, amount, status) VALUES (1, ORD-2024001, 199.90, CREATED), (2, ORD-2024002, 59.00, PAID), (3, ORD-2024003, 1299.00, SHIPPED);3.2 编写 YAML Pipeline只写一份配置搞定全量和增量Flink CDC 3.0 的 YAML Pipeline 语法非常直观。核心就四段source数据源、sink目标端、route表映射、pipeline任务属性。如果你的目标只是先验证 CDC 链路是否打通用内置的 Print Sink 是最快的它会把变更事件打印到日志里source: type: postgres hostname: localhost port: 5432 username: flinkuser password: your_password database-name: your_db schema-name: public table-name: orders slot.name: flink_pg_slot server-time-zone: Asia/Shanghai sink: type: print pipeline: name: PG-to-Print parallelism: 1这份配置里需要注意的参数是slot.name。如果你不指定Flink CDC 会自动生成一个但指定固定值的好处是任务重启时能快速定位到 PG 侧对应的复制槽排查问题方便很多。确认链路没问题后再把 Sink 换成真实的 StarRocks 表。我这里以 StarRocks 为例因为它在 3.0 之后和 Flink CDC 的集成度最高支持 Schema Change 自动同步sink: type: starrocks username: root password: 123456 jdbc-url: jdbc:mysql://127.0.0.1:9030 load-url: 127.0.0.1:8030 database-name: ods table-name: orders sink.properties.format: json sink.properties.strip_outer_array: true route: - source-table: your_db.public.orders sink-table: ods.orders注意route段的作用它负责把源端“库.模式.表”的完整标识映射到目标端的“库.表”。如果不写这段Flink 默认会用源表的全限定名去目标端找表通常会导致“表不存在”。3.3 启动任务一次启动完成全量快照 增量监听启动命令很简单在 Flink CDC 的安装目录下执行bin/flink-cdc.sh /path/to/pipeline.yaml如果是提交到已有的 Flink 集群记得带上--target参数指定 job manager 地址或者用bin/flink-cdc.sh -t yarn-per-job /path/to/pipeline.yaml这类形式取决于你的部署方式。任务启动后理想情况下你会看到两层日志先是全量快照阶段的状态Flink CDC 会把源表已有的三行数据分批地读出来通过 Sink 写入目标端全量阶段完成后任务自动切换到增量阶段开始实时监听并解析 WAL 里出现的变更事件。验证增量同步最直接的办法在 PG 里执行几类典型的 DML 操作然后观察目标端数据的变化。INSERT INTO orders (id, order_no, amount, status) VALUES (4, ORD-2024004, 88.00, CREATED); UPDATE orders SET status PAID WHERE id 1; DELETE FROM orders WHERE id 2;如果用的是 Print Sink日志里会清楚看到I[4, ORD-2024004, 88.00, CREATED, ...]这样的插入事件、-U[1, ...]和U[1, ...]组合出现的更新事件先删旧值再加新值、以及-D[2, ...]的删除事件。看到这三种事件格式说明 CDC 链路已经完全打通。如果你对 Flink 的 Changelog 事件格式不熟悉这里简单解释一下CDC 同步底层是流式变更事件INSERT 对应IUPDATE 在大多数情况下会拆成-U旧值和U新值两个事件DELETE 对应-D。所以流里的数据是“带操作类型的”下游 Sink 会依据这个操作类型去执行对应的写入或删除。4. 避坑指南WAL 撑爆磁盘、时区错乱、类型映射失败的排查实录4.1 坑一复制槽不释放WAL 像滚雪球一样涨这是 CDC 同步里最典型、破坏力最大的坑没有之一。我自己第一次遇到时是某天夜里收到 PG 所在服务器的磁盘告警登录上去一看df -h显示磁盘使用率已经到了 97%检查pg_wal目录发现这个目录的占用空间比正常情况大了几十个 G。排査链路按照下面几步走第一步确认 WAL 目录的大小du -sh $PGDATA/pg_wal第二步查看复制槽状态SELECT slot_name, slot_type, active, wal_status, restart_lsn FROM pg_replication_slots;如果看到某一行active f但wal_status reserved或者restart_lsn一直不变基本可以锁定问题这个复制槽已经没有消费者了但 PostgreSQL 仍为它保留着从restart_lsn开始的全部 WAL 日志。根因通常是 Flink 任务停止不管是异常退出还是手动取消时没有正确调用 PG 侧的清除复制槽的逻辑尤其是任务被kill -9强杀、或者集群崩溃时PG 那边残留的复制槽就成了“只进不出”的黑洞。处理方法确认这个复制槽对应的 Flink 任务确实已经不需要恢复之后手动删除它SELECT pg_drop_replication_slot(flink_pg_slot);删除后PG 会立刻释放保留的 WAL 文件磁盘占用快速下降。长期预防方案有两个。一是给 PG 设置合理的max_slot_wal_keep_size参数限制复制槽最多保留多少 WAL超过限制后 PG 会主动丢弃那些 WAL代价是复制槽不再可用下游必须重新做全量二是在 Flink 任务侧配置自动清理比如在任务停止时通过脚本调用 PG 管理函数删除复制槽。即使有这些预防措施我仍然建议把 PG 服务器的磁盘监控、复制槽监控纳入日常运维因为这类问题一旦发生故障恢复速度直接取决于发现速度。4.2 坑二TIMESTAMPTZ 同步过去比实际时间慢了 8 小时第二个高频坑是时区问题。TIMESTAMPTZ在 PostgreSQL 里存储的是 UTC 时间显示时会根据会话时区做转换。Flink CDC 在读取这个类型时会基于任务自身配置的时区把它转成字符串或TIMESTAMP如果时区设置不一致就会闹出“数据没丢但时间全偏了”的诡异问题。我当时排查某个报表任务时发现PG 里created_at明明显示的是2024-05-01 14:30:0008同步到 StarRocks 后却变成了2024-05-01 06:30:00正好差了 8 小时。第一反应是“数据是不是丢了”后来才发现是时区换算的问题。解决方式有两个层面。第一层在建 Sink 表时明确把目标字段类型定义成DATETIME同时在 Flink 侧的 YAML 配置中加入source: server-time-zone: Asia/Shanghai第二层确保 Flink 任务的 JVM 时区是东八区可以在flink-conf.yaml里设置taskmanager-env.java-opts: -Duser.timezoneAsia/Shanghai jobmanager-env.java-opts: -Duser.timezoneAsia/Shanghai这两步做完后重新同步的数据时间就是正确的了。需要注意的是如果链路里用了 Kafka 做中转Kafka 里的消息体保存的是 Debezium 格式的时间字符串消费端解析时同样要统一时区不然问题会在更下游重复出现。4.3 坑三NUMERIC 精度丢失与枚举类型直接报错PostgreSQL 的NUMERIC(p, s)是任意精度数值类型但大部分下游系统包括 StarRocks、Doris对数值精度的支持是有上限的一般用DECIMAL(p, s)表示。同步过程中如果两边精度定义不一致轻则丢失小数位重则任务直接失败。我建议处理方式源表建表时尽量用明确的精度和标度比如NUMERIC(10, 2)避免裸写NUMERIC。PostgreSQL 的裸NUMERIC类型意味着任意精度Flink CDC 无法知道它最大到多少位只能按照一个很高的默认精度去映射下游如果建的是DECIMAL(10, 2)数据就可能溢出报错。另外PostgreSQL 的枚举类型CREATE TYPE order_status AS ENUM (CREATED,PAID,SHIPPED)在同步时也很容易报“unsupported type”错误。Flink CDC 3.0 对枚举类型的支持不完整不同 Sink 的表现也不一样。最简单粗暴的办法是源表不要直接用枚举类型改用VARCHAR加 CHECK 约束来约束合法值或者建一个视图把枚举字段用CAST转成textFlink 任务读取这个视图而不是原表。4.4 坑四Schema 变更引发任务中断生产环境的表结构不是一成不变的经常会有ALTER TABLE ADD COLUMN之类的操作。Flink CDC 3.0 支持 Schema Evolution但这个能力高度依赖下游 Sink 的实现。以 StarRocks 为例如果在源端执行ALTER TABLE orders ADD COLUMN remark VARCHAR(100);Flink CDC 会捕获这个 DDL 事件并尝试把它透传给 StarRocks 执行ALTER TABLE ods.orders ADD COLUMN remark VARCHAR(100)。多数情况下这个流程能自动完成但有两个限制你要知道第一破坏性操作如DROP COLUMN、ALTER COLUMN TYPE这类有很大概率导致任务失败官方建议不要在生产上做这种操作第二新增列如果带有非确定性默认值比如DEFAULT random()PG 侧产生的 DDL 事件里包含了函数调用Sink 可能无法处理。面对 Schema 变更有两套策略。如果下游能做重建简单的做法是源表变更后直接停任务、改下游表、再重启任务靠 Flink CDC 的全量快照机制重新同步。如果下游不能随便重建就严格约束源端的 DDL 操作只允许新增可空列。说到底CDC 链路里的 Schema 变更从来不只是“数据同步”问题它本质上是数据治理问题需要研发规范来兜底。5. 跑通之后从演示环境到生产环境还差这几步5.1 并行度与 Source Chunk 拆分别让全量快照拖垮源库YAML 里我前面写的parallelism: 1在演示环境没问题但生产环境的数据量动辄千万行、上亿行单并行度做全量快照会非常慢。Flink CDC 的 PostgreSQL Connector 内置了增量快照Incremental Snapshot机制它会根据主键把表的数据拆成多个 chunk每个 chunk 可以交给不同的 subtask 并行读取。调整并行度的方式pipeline: name: PG-to-StarRocks parallelism: 4这里有一个权衡并行度越高全量阶段拉取数据越快但对源库的压力也越大。尤其是主键不是均匀分布的情况比如自增主键前 1000 万行是历史数据、后面只有几万行新数据chunk 默认按 8096 行一个切分可能产生大量小 chunk调度开销反而变高。可以调整chunk.key.even-distribution.factor.upper-bound这类参数但第一优先级是观察源库的 QPS 和磁盘 IO不要让同步任务影响在线业务。5.2 监控与告警不盯复制槽就是在给磁盘埋雷生产环境跑 Flink CDC有两类监控我强烈建议配置。第一类是 Flink 任务自身的监控包括任务是否 Running、Checkpoint 是否成功、消费延迟是多少。Flink Web UI 能看一部分但要实现告警建议把指标接入 Prometheus Grafana。重点盯两个指标Source Operator 的 currentFetchEventTimeLag代表了从数据库产生变更到 Flink 捕获到变更的时间差和 Checkpoint 失败的次数。第二类是 PostgreSQL 侧的监控尤其是复制槽的 WAL 保留量。查询语句可以定时执行并暴露成指标SELECT slot_name, pg_size_pretty(pg_current_wal_lsn() - restart_lsn) AS retained_wal_size FROM pg_replication_slots;一旦retained_wal_size超过阈值比如 20GB立刻告警。这条监控能帮你避免第 4.1 节那种磁盘被撑爆的事故成本极低收益极高属于做过一次就再也不敢省的监控项。5.3 消费位点恢复任务重启后如何不重不漏Flink CDC 的位点恢复依赖 Flink 的 Checkpoint / Savepoint 机制。默认情况下Flink 会定期把当前消费到的 LSN 位置保存下来任务重启后会从最近一次 Checkpoint 的位置继续消费。但这里有个容易理解错的地方Checkpoint 恢复位置只是“不会丢数据”不一定是“刚刚好的位置”。如果任务崩溃前最后一批数据已经写入 PG但对应的变更事件还没被 Flink 处理到 Checkpoint 里重启后那批数据会被重新消费一遍。对于大部分数仓场景这个“至少一次”At Least Once的语义是可以接受的写入目标端时靠主键覆盖就能去重。如果业务要求严格精确一次需要 Sink 端支持幂等写入 Flink 的两阶段提交协议StarRocks Sink 对这块支持得比较好Kafka Sink 则需要额外配置。Savepoint 还有一个用途是任务升级。发布新版本代码后用bin/flink savepoint jobId /savepoint/path触发一次快照然后用-s参数从快照恢复新任务整个过程不需要重新做全量同步。这个操作在模块升级或者修改 Pipeline 配置时很实用建议上线前演练几遍。从我个人的实测感受来说Flink CDC 这套 YAML Pipeline 的体验是真的“爽”五分钟左右跑通 PG 同步完全不是宣传话术前提是你绕开了我上面写的这些坑。特别是复制槽监控真的提前配好能帮你躲过一次灾难级的磁盘告警。这套体系现在已经是我的标准同步方案了后面如果有机会我再单独写写从 MySQL 同步到 Kafka 与 Doris 的配置对比那又是另一套完全不同的注意事项。
返回列表