
做了几年的数据接入和实时计算我越来越觉得Flink CDC已经成了数据同步场景里绕不开的一个方案。尤其是PostgreSQL配合Flink CDC几行SQL就能把业务库的变更实时流式地送到下游不需要额外部署一堆Agent也不用手写DataX任务去轮询。这段时间我把整套链路从环境准备到生产踩坑完整梳理了一遍今天直接把这套实战流程和避坑记录分享出来想快速上手的朋友可以直接照着抄作业。这篇内容适合正在做数据库实时同步、准备把PostgreSQL数据接入数据仓库或消息队列的开发者。会用到Docker、Flink SQL和PostgreSQL的基础操作如果你已经对这几个组件有基本认知那整个上手过程会非常流畅。我会从方案选型讲到环境搭建再到核心配置和最终验证最后把最常见的坑全部列出来争取让你一次跑通。1. 为什么选择Flink CDC做PostgreSQL实时同步1.1 CDC的本质把数据库日志变成数据流CDCChange Data Capture变更数据捕获的思路其实很简单就是让下游系统能感知到数据库里的数据变化而不需要每次全量拉一遍数据。实现方式有很多种有的基于触发器有的基于时间戳轮询但最优雅、对源库影响最小的方式是解析数据库的事务日志。对PostgreSQL来说就是基于WALWrite-Ahead Log预写日志机制。你可以把WAL理解为数据库的“日记”——每次增删改操作发生之前PostgreSQL都会先把变更记录写到WAL里用来保证崩溃恢复时的数据一致性。既然这些日志已经记录了每一次数据变化那CDC工具直接去订阅和解析这份日志就顺理成章了。这样做的好处是实时性极高数据从提交到被下游感知通常只有毫秒级延迟而且完全不需要修改业务表结构也不会给业务数据库带来明显的查询压力。Flink CDC本质上就是把这套日志解析能力封装成了Flink的连接器。对Flink来说一个CDC数据源就是一个流每条记录代表一次数据库变更事件包括INSERT、UPDATE、DELETE操作以及变更前后的数据快照。这种统一抽象让后续的数据处理变得非常简单你不需要关心底层日志格式直接用SQL就能操作这些变更流。1.2 Flink CDC相比其他同步方案的优势早些年做PostgreSQL实时同步常见的选择是Canal主要面向MySQL、Debezium或者自己写程序消费WAL。Debezium其实非常强大它基于Kafka Connect构建也是Red Hat在维护的顶级开源项目Flink CDC底层有一部分能力就是借鉴了Debezium的格式设计。但直接用Debezium的问题在于你需要额外维护一套Kafka Connect集群还要处理Topic管理和Schema演进链路变长之后运维成本直线上升。Flink CDC最核心的优势是“一体化”。它把数据捕获、数据处理、数据投递整合到了同一个Flink作业里。比如你想把PostgreSQL的表同步到Doris如果用Debezium你需要部署Kafka Connect、配置Connector、再写一个Flink作业去消费Kafka里的变更数据最后写入Doris。而用Flink CDC一个作业就搞定了而且可以直接用Flink SQL定义同步逻辑开发效率完全不在一个量级。另外Flink CDC天然继承了Flink的流处理能力。实时同步不仅仅是简单的“搬数据”往往还需要过滤字段、类型转换、维表关联、多流合并。这些在纯同步工具里实现很麻烦在Flink里就是几个SQL算子的事。我之前做过一个场景需要把PostgreSQL的订单表和MySQL的用户表实时关联后写入数仓用Flink CDC加一条JOIN就实现了换做其他方案简直是灾难。1.3 一个典型的实时链路长什么样我要演示的场景是比较常见的一种PostgreSQL业务库中的一张订单表实时同步到另一个PostgreSQL分析库或者Doris、StarRocks等OLAP引擎。链路非常简单业务应用写入PostgreSQL主库PostgreSQL通过逻辑复制将WAL变更发送给Flink CDC连接器Flink作业解析并处理这些变更然后通过JDBC连接器写入目标库。整个过程不需要业务侧做任何改造不需要双写不需要中间表。Flink会先做一次全量快照同步把表里已有的历史数据搬过去然后自动无缝切换到增量日志消费继续同步新产生的变更。这个“全量增量”的自动衔接是Flink CDC非常受欢迎的原因实际使用中我基本不需要关心切换的衔接点框架内部处理得很干净。2. 环境准备与配套版本选型2.1 软件版本矩阵才是第一个坑很多人上手Flink CDC直接被版本搞懵。Flink CDC不是一个独立运行的框架它是以连接器JAR包的形式存在需要嵌入到Flink环境里。所以你必须保证Flink CDC连接器版本、Flink版本、PostgreSQL驱动版本、甚至PostgreSQL数据库版本之间是兼容的否则各种诡异报错都会冒出来。我这次实战使用的版本组合如下经过验证可以稳定运行组件版本说明PostgreSQL14.x源库和目标库都用14版本Flink1.17.x稳定版兼容CDC 2.4Flink CDC连接器flink-sql-connector-postgres-cdc-2.4.x直接下载带依赖的包PostgreSQL JDBC驱动42.5.x写入端需要用到Doris Flink Connector1.x按需如果目标是Doris才需要这里我要特别强调一下建议直接使用flink-sql-connector开头的那种带依赖的JAR包因为它把Debezium相关的依赖、PostgreSQL驱动等都已经打包进去了你不需要自己去找一堆依赖往Flink的lib目录里塞。我第一次用的时候贪图方便下载了不带依赖的包结果启动作业时报了各种ClassNotFoundException折腾半天。用带依赖的包一个文件就够。PostgreSQL的版本也需要注意。Flink CDC的PostgreSQL连接器基于Debezium的PostgreSQL Connector实现逻辑复制功能在PostgreSQL 9.4之后就支持了但我在实际使用中发现10以上的版本兼容性最好14和15都很稳。如果你的库还是9.x甚至更老的版本建议先升级再考虑实时同步。2.2 用Docker快速搭建PostgreSQL测试环境如果你的本机还没有现成的PostgreSQL环境用Docker起一个是最省事的方式。我这次实战用Docker Compose管理源库和目标库配置很简单实际项目中我用的是docker-compose:postgresql这个镜像模板里面关于版本号和端口映射的部分我按自己的环境做了调整version: 3 services: postgres-source: image: postgres:14 container_name: pg-source environment: POSTGRES_USER: postgres POSTGRES_PASSWORD: postgres POSTGRES_DB: source_db ports: - 5432:5432 command: - postgres - -c - wal_levellogical - -c - max_replication_slots10 - -c - max_wal_senders10 postgres-sink: image: postgres:14 container_name: pg-sink environment: POSTGRES_USER: postgres POSTGRES_PASSWORD: postgres POSTGRES_DB: sink_db ports: - 5433:5432注意看我在源库的启动命令里显式指定了wal_levellogical这一项非常关键。CDC要读取逻辑变更必须把WAL级别设置为logical默认的replica级别是不够的。如果漏掉这个配置Flink CDC作业启动时会直接报错提示wal_level配置不正确。启动容器很简单在docker-compose.yml所在目录执行docker-compose up -d然后等待镜像拉取和容器启动完成。用docker ps确认两个容器都处于Up状态然后用docker exec -it pg-source psql -U postgres -d source_db进入源库建一张测试表并写入一些初始数据。2.3 准备Flink运行环境Flink本地跑不需要复杂的集群部署用Flink的Standalone模式就够了。到Apache Flink官网下载flink-1.17.x的二进制包解压后就能用。所谓“flink一定要hdfs”的说法其实要分场景——如果你的同步作业是全量加增量的CDC任务状态后端用RocksDB并配置本地文件系统或者HDFS都可以如果是纯本地验证用默认的配置就能跑不需要强依赖HDFS。我这里直接采用本地模式把Flink的StateBackend设为filesystem检查点目录指向本机路径即可。解压之后把下载好的Flink CDC连接器JAR包放到FLINK_HOME/lib目录下。如果你用的是SQL客户端方式还需要把PostgreSQL JDBC驱动也放到lib目录因为写入端创建目标表时需要加载驱动。我用的是postgresql-42.5.4.jar放到lib目录后重启Flink SQL客户端让依赖生效。Flink下载安装的具体步骤在官方文档里有很详细的说明跟着做基本不会出问题。有一个小细节是如果本机内存不是特别充裕建议在conf/flink-conf.yaml里把jobmanager.memory.process.size和taskmanager.memory.process.size调小一点比如都设置成1024m否则默认值在开发机上可能撑不住。之前有朋友用2G内存的云主机跑Flink跑着跑着就OOM把内存调小后稳定多了。3. 5分钟跑通核心同步流程3.1 第一步PostgreSQL源库开启逻辑复制前置条件在正式创建同步任务之前源库还需要配置几个前置条件。直接用psql执行下面几条SQL-- 创建用于数据同步的账号并授予复制权限 CREATE USER flink_user WITH PASSWORD flink_pass; ALTER USER flink_user WITH REPLICATION; GRANT CONNECT ON DATABASE source_db TO flink_user; -- 授予表级权限 GRANT SELECT ON ALL TABLES IN SCHEMA public TO flink_user; ALTER DEFAULT PRIVILEGES IN SCHEMA public GRANT SELECT ON TABLES TO flink_user; -- 创建逻辑复制发布Publication CREATE PUBLICATION flink_pub FOR TABLE orders;这里有两件事容易被忽略。第一CDC连接器读取的是PostgreSQL的逻辑复制流所以同步账号必须具备REPLICATION权限否则作业连上去之后拿不到日志数据。第二在PostgreSQL 10及以上版本逻辑复制是基于Publication机制的你必须为需要同步的表创建PublicationFlink CDC才能识别并订阅这些表的变更。如果漏掉这一步任务启动时会报错说找不到表的复制标识。权限和Publication配置完成之后可以先用一条命令验证一下复制槽是否正常后面排查问题时也常用到SELECT * FROM pg_replication_slots;此时应该没有复制槽因为Flink作业还没启动。等Flink作业启动后再看这里会多出一条记录。3.2 第二步在Flink SQL中创建源表和目标表映射Flink SQL连接器的使用方式非常直观先创建一张映射到PostgreSQL源表的虚拟表再创建一张映射到目标库的表最后一行INSERT INTO就能启动同步。抛开命令行的准备工作5分钟跑通核心流程是完全可行的。先启动Flink SQL客户端$FLINK_HOME/bin/sql-client.sh embedded -l $FLINK_HOME/lib然后在SQL客户端里执行建表语句CREATE TABLE orders_source ( id INT, user_id INT, product_name STRING, amount DECIMAL(10, 2), order_time TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector postgres-cdc, hostname localhost, port 5432, username flink_user, password flink_pass, database-name source_db, schema-name public, table-name orders, decoding.plugin.name decoderbufs, slot.name flink_slot );有几个参数我需要解释一下。decoding.plugin.name指定了解码WAL日志的插件方式可选值有decoderbufs、wal2json和pgoutput。其中pgoutput是PostgreSQL 10原生支持的输出插件稳定性和兼容性最好我实际用下来推荐直接选pgoutput。有些老教程会推荐decoderbufs那是因为当时的历史原因现在用pgoutput不会有问题。slot.name是逻辑复制槽的名称需要唯一如果冲突可以换一个名字。目标端表定义如下我用了PostgreSQL的JDBC连接器CREATE TABLE orders_sink ( id INT, user_id INT, product_name STRING, amount DECIMAL(10, 2), order_time TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector jdbc, url jdbc:postgresql://localhost:5433/sink_db, table-name orders, username postgres, password postgres, sink.buffer-flush.max-rows 1, sink.buffer-flush.interval 1s );如果你希望源表和目标表的字段名不一致可以在建表SQL里用别名方式重新映射。需要注意的是目标表在主键或唯一键冲突时的行为取决于JDBC连接器的实现默认是INSERT如果同步的场景里UPDATE操作非常多建议把sink.buffer-flush.max-rows调小一些减少批量写入时的延迟。3.3 第三步提交同步作业并验证结果核心语句就一行直接提交INSERT INTO orders_sink SELECT * FROM orders_source;Flink SQL会把这段逻辑编译成一个流式计算作业启动后先做全量数据同步把orders_source表里已有的数据读出来写入orders_sink然后自动切换到增量模式持续监听PostgreSQL的WAL变更。同步作业跑起来后可以做一个简单的验证。先在源库插入几条新数据INSERT INTO orders (user_id, product_name, amount, order_time) VALUES (1001, iPhone 15, 6999.00, NOW()), (1002, MacBook Pro, 12999.00, NOW()), (1001, AirPods Pro, 1899.00, NOW());然后去目标库查询SELECT * FROM orders ORDER BY id;正常情况下这几条数据会在一秒内出现在目标库中。再试试UPDATE和DELETEUPDATE orders SET amount 6500.00 WHERE id 1; DELETE FROM orders WHERE id 3;回到目标库查询你会发现这些变更也完整地同步过来了。UPDATE会变成目标表里的新值DELETE会把对应主键的数据删掉。这就是CDC和普通批量同步的本质区别它同步的是“操作”而不只是“数据结果”。整个核心流程如果准备工作做得充分5分钟确实够了。我自己第一次操作的时候卡在版本匹配和权限配置上花了很久但一旦这两个前置问题解决后面的SQL配置和启动作业其实就是几分钟的事这也是为什么我坚持把版本选型和权限配置放到前面重点讲的原因。3.4 全量和增量衔接时的时间戳与时区问题很多人在看到全量数据同步完成后对增量同步是否就绪会有疑虑。Flink CDC会自动管理这个衔接过程它会先读取当前WAL日志的位点然后执行全量快照快照完成后从记录的位置继续消费增量变更。这期间产生的数据变更不会丢也不会重复消费。实际效果就是全量和增量是平滑衔接的中间不需要人工干预。这里有一个很容易被忽略的细节时间戳类型。PostgreSQL的TIMESTAMP不带时区信息而TIMESTAMPTZ带时区信息。Flink SQL里TIMESTAMP(3)映射到PostgreSQL是TIMESTAMP WITHOUT TIME ZONETIMESTAMP_LTZ(3)映射带时区类型。如果你在同步过程中发现时间字段差了几个小时大概率是类型映射没配对用TIMESTAMP_LTZ来代替TIMESTAMP通常能解决这个问题。我朋友在生产环境做过一个案例订单时间字段在源库是TIMESTAMPTZ建Flink表时用了TIMESTAMP结果同步到下游数仓的时间全部变成了UTC时间而不是北京时间凌晨的业务数据全部对不上。排查了大半天才发现是类型映射的问题。所以建表时一定要确认时间字段的映射类型。4. 常见的坑与排查技巧实录4.1 PostgreSQL侧配置引发的各种问题WAL级别未修改。这是遇到最多的报错任务启动时直接抛异常提示wal_level不满足要求。解决办法很简单修改postgresql.conf或通过启动参数设置wal_levellogical然后重启PostgreSQL实例。需要注意的是修改wal_level需要重启数据库才能生效ALTER SYSTEM SET只能改配置不能热加载。如果你在云上使用托管的PostgreSQL比如阿里云RDS管理控制台里切换到逻辑复制模式通常也需要实例重启。生产环境操作前务必评估重启影响窗口。复制槽不清理导致WAL膨胀。当Flink作业停止或异常退出时PostgreSQL上的复制槽还会保留。此时WAL日志无法被清理随着数据变更不断产生磁盘空间会持续增长。这种情况在业务高峰期尤其危险我见过一个案例一个复制槽一周时间吃掉了上百GB的磁盘空间。排查方式也比较简单检查pg_replication_slots视图找到active为false的复制槽手动删除SELECT * FROM pg_replication_slots; SELECT pg_drop_replication_slot(flink_slot);在Flink作业里合理配置checkpoint和重启策略也能减少复制槽不清理的情况但最稳妥的方式还是定期巡检。如果你用Flink CDC做长周期任务建议监控PostgreSQL的WAL目录大小超过阈值及时告警。权限不足导致连接失败。如果同步账号没有授予REPLICATION权限作业启动时会报权限相关的错误。确保账号有REPLICATION属性并且对要同步的表有SELECT权限。如果表所在的Schema不是public还需要单独授权GRANT USAGE ON SCHEMA your_schema TO flink_user; GRANT SELECT ON ALL TABLES IN SCHEMA your_schema TO flink_user;4.2 Flink SQL运行阶段的经典报错找不到合适的驱动或连接器类。这类错误多半是JAR包没放对位置或者版本冲突。检查$FLINK_HOME/lib目录下是否有Flink CDC连接器包和JDBC驱动包确认版本是兼容的。有时候Flink SQL客户端缓存了旧的依赖执行:quit退出重启客户端也能解决。“flink type is datev2, but arrow type is dateday”这类类型不匹配错误。这个报错我见过多次通常出现在Flink CDC同步PostgreSQL到Doris的场景。PostgreSQL的DATE类型在Flink中被解析为DATE但Doris的Flink Connector内部映射成了DATEV2而Arrow的日期类型是DATE32三者之间存在一个映射断层。解决方法是在写入目标表前显式进行类型转换CREATE TABLE doris_sink ( dt STRING, ... ) WITH (...); INSERT INTO doris_sink SELECT CAST(order_time AS STRING), ... FROM orders_source;或者在Flink表定义中直接把时间字段声明为STRING从源头避免类型转换。“flink的jdbc连接器异常”一般源于写入端连接串配置错误。比如目标端url写错了端口或者目标表不存在或者用户没有对应权限。这类异常日志比较明显排查起来相对容易。还有一种情况是目标表的主键和源表不一致导致UPDATE语义不明确JDBC连接器会报主键冲突。建目标表时尽量保持主键一致。4.3 常见问题速查表现象可能原因解决办法作业启动报wal_level错误PostgreSQL未开启逻辑复制设置wal_levellogical并重启实例连接被拒绝端口错误或pg_hba.conf限制检查端口和pg_hba.conf中是否允许该IP连接逻辑复制插件冲突replication slot已存在删除旧复制槽或更换slot.name表数据不同步未创建Publication或表不在Publication里为表创建Publication全量同步后增量无效表没有主键或复制标识不完整为表添加主键或设置REPLICA IDENTITY FULL时间字段差8小时时区类型映射错误使用TIMESTAMP_LTZ而不是TIMESTAMP同步作业内存持续增长checkpoint配置不合理调整checkpoint间隔使用RocksDB状态后端写入目标库乱码字符集不一致确保源库、目标库、连接串都使用UTF-84.4 我的一些排查习惯我做CDC任务排查时有个固定思路可以分享。遇到问题先确定是源库连接问题、数据读取问题还是写入目标问题用排除法把范围缩小。我会先看Flink Web UI里的TaskManager日志很多异常在日志里其实已经写得很直白了。如果是全量同步阶段报错大概率是类型映射或者权限问题如果是增量阶段报错优先检查复制槽状态和WAL日志是否能正常读取如果同步任务运行一段时间后出现延迟而Flink本身CPU和内存都不高那就要看看目标端的写入瓶颈比如Doris的导入频率限制或者PostgreSQL的锁竞争。另外有一个小习惯每次新建同步任务前我都会用数据库自带的命令手动确认源库的连接信息。比如先测试PostgreSQL能否正常连通再确认测试表的权限PGPASSWORDflink_pass psql -h localhost -p 5432 -U flink_user -d source_db -c SELECT * FROM orders LIMIT 1;如果这一步能通过Flink连接器大概率也不会出问题。很多连接报错其实在数据库侧验证一下就能提前发现。5. 让同步作业更稳定的一些进阶操作5.1 关于主键和REPLICA IDENTITY的注意事项Flink CDC要求同步的表必须有主键否则无法正确处理UPDATE和DELETE操作。如果源表没有主键全量同步没问题但增量阶段对已有数据的UPDATE事件会识别不了。PostgreSQL的逻辑复制机制需要知道每一行的唯一标识。解决的办法有两个。第一个是给源表补上主键这在业务库中通常比较难操作。第二个是设置REPLICA IDENTITY FULL让WAL日志记录整行所有字段的旧值ALTER TABLE orders REPLICA IDENTITY FULL;这样设置之后即使表没有主键Debezium也能通过对比整行数据来识别变更。缺点是WAL日志的体积会变大因为每次UPDATE都会记录整行数据。如果不是磁盘特别紧张这个方案是可用的。但需要说明的是Flink CDC的PostgreSQL连接器在语义上更推荐有主键的表所以我建议源表设计时尽量规范主键。5.2 检查点配置决定同步语义Flink CDC的“精确一次”语义依赖检查点Checkpoint机制。如果checkpoint没配置好极端情况下可能出现数据重复或丢失。Flink默认的checkpoint间隔可能不适合实时同步场景建议在提交作业时显式设置SET execution.checkpointing.interval 10s; SET execution.checkpointing.mode EXACTLY_ONCE; SET state.backend.type rocksdb; SET state.checkpoint-storage filesystem; SET state.checkpoints.dir file:///tmp/flink-checkpoints;我建议把checkpoint间隔设置在10秒到30秒之间太频繁会增加存储压力太稀疏会导致故障恢复时间变长。状态后端用RocksDB可以避免大量状态数据占用堆内存导致OOM。5.3 当单表同步变成整库同步如果想把整个数据库的所有表都同步过来不需要为每张表单独建一个Flink作业。Flink PostgreSQL CDC连接器支持通过正则表达式匹配表名一张源表映射就能覆盖多张物理表。上游的table-name参数可以用正则表达式CREATE TABLE all_tables_source ( ... ) WITH ( connector postgres-cdc, database-name source_db, schema-name public, table-name orders|users|products, ... );但这里要提醒一句多表同步到一个Flink作业时所有表的Schema必须完全一致否则下游处理会非常麻烦。如果表结构差异较大还是建议拆成多个作业单独同步。5.4 从PostgreSQL到Doris的实战补充热词里反复出现的“flink type is datev2, but arrow type is dateday”这个问题我再展开说一下。这个问题的根源在于Doris的新版本引入了DATEV2类型来替代DATE而Flink CDC在向Doris写入时内部用Arrow格式传输数据Arrow的日期类型是DATE32。三者映射关系不统一就会出现日期字段无法写入的报错。解决方法有三种。最直接的是在导入前用SQL做一次CAST把日期字段转成STRING再写入Doris目标表。还有一种方案是升级Doris Connector到较新的版本新版本已经适配了DATEV2和Arrow的映射。第三种是调整Doris表的字段类型把DATEV2改回DATE但这需要Doris版本支持。实践中我通常直接在Flink SQL里CAST为STRING简单有效对下游数仓的日期处理也没有影响。6. 写在最后的一点个人体验做了这么多次Flink CDC同步我的整体感受是这个方案的上手门槛确实不高但要把它用稳、用好需要对底层机制有一定理解。版本匹配、权限配置、复制槽管理、checkpoint设置这些都是决定作业能否长期稳定运行的关键细节。千万不要等到作业上线后才开始关注这些问题前置配置省掉的每一步后面都可能变成半夜爬起来救火的代价。我个人在刚开始接触Flink CDC时最受益的一个做法是先在本地用Docker把整套环境完整跑一遍包括模拟故障场景比如杀掉Flink作业再重启、在同步过程中对源表执行大事务、观察checkpoint失败时作业的表现。把这些场景都在本地验证过之后上生产环境心里就有底了。这个经验后面也一直沿用每次换新的连接器版本都会先做一轮本地演练再发布。最后再分享一个小技巧如果你不确定某个Flink CDC连接器的版本和你的Flink版本是否兼容最直接的办法是去Flink CDC的GitHub Release页面看对应的文档说明或者在本地用最小案例验证一下不要盲目相信网上的教程版本兼容这种事实测过了才算数。