最近接了公司一个比较硬核的任务:把核心业务库 Oracle 里的订单、用户、库存这些表,实时同步到数据仓库里。说实话,听到这个需求我第一反应不是兴奋,是头疼。MySQL 生态里玩 CDC(Change Data Capture)已经很顺手了,但 Oracle 这套东西,归档模式、补充日志、LogMiner、权限模型,每一项都有讲究,稍不注意就是一堆 ORA- 报错等你填坑。
这一轮折腾下来,踩了不少坑,也总结了一套可以复用的流程。今天这篇文章就把我用 Flink CDC 实时同步 Oracle 的完整过程拆开来讲:为什么选这套方案、Oracle 侧到底要做什么准备、核心原理是怎么工作的、具体怎么用 Flink SQL 和最新的 YAML Pipeline 跑起来,以及那些文档里不会写但实战必踩的问题。
1. 先想清楚:Flink CDC 在同步链路里到底扮演什么角色
1.1 一句话讲透 CDC 与 Flink CDC
CDC 说白了就是数据库变更数据捕获,核心思路不是用定时任务去SELECT轮询哪些数据变了,而是直接读取数据库日志文件里的变更记录,然后把每一条 insert、update、delete 都解析出来,形成一条持续的数据流。
Flink CDC 是 Apache Flink 社区维护的一套连接器,它把 CDC 能力和 Flink 的实时计算、状态管理、Checkpoint 机制结合在一起。你要理解 Flink CDC 到底解决了什么问题,可以这么类比:普通 JDBC 同步相当于你每 5 分钟拿个本子去仓库清点一遍货物,记下变化;Flink CDC 则是在仓库门口装了一个摄像头,每一件货进去出来都被实时记录下来,而且记录绝对精确、不会漏也不会重复。
这套方案最核心的价值有几个:实时性高,秒级延迟;消费的是日志,对源库压力远比轮询查询小;同时天然具备断点续传能力,任务挂了可以恢复,不会丢数据。整个同步过程不需要在 Oracle 上装任何 Agent,只需要建一个账号、给权限、读日志就行。
1.2 为什么最终选择了 Flink CDC
在做方案选型的时候,我对比了市面上常见的几条路线,这也是你在立项时最容易纠结的地方。
| 方案 | Oracle 支持 | 实时性 | 精确一次 | 上手门槛 | 典型场景 |
|---|---|---|---|---|---|
| Flink CDC | 原生支持 | 秒级 | 支持(依赖 Checkpoint) | 中等,SQL 即可 | 实时入仓、入湖、单表/整库同步 |
| Debezium | 支持,但仅提供数据流 | 秒级 | 需要配合 Kafka 事务 | 高,需自建整个管道 | 流平台集成,接入 Kafka 生态 |
| Canal | 开源版只支持 MySQL | 秒级 | 支持 | 中等 | MySQL 生态为主,Oracle 需商业版 |
| Maxwell | 不支持 Oracle | 秒级 | 不支持 | 低 | 只适合 MySQL |
| DataX / 定时任务 | 支持 | 分钟级以上 | 不需要 | 低 | 离线/准实时批量 |
我最终的结论是:如果同步目标是数据仓库、数据湖或者消息队列,并且团队已经用了 Flink 技术栈,那么 Flink CDC 是综合成本最低的选择。不需要额外维护 Kafka 集群作为中转,不需要自己解析 Debezium 消息,用 Flink SQL 写几个建表语句就能把同步链路搭起来。如果你要做的是流式 ETL,比如同步过来之后马上做清洗、打宽表、聚合,那 Flink CDC 更是直接把计算和同步合并到了一套引擎里,省掉了中间环节。
1.3 一个真实场景:订单数据从 Oracle 到数仓的实时链路
我说一个完整的业务背景,方便你对照自己的场景。业务库是 Oracle 19c,里面有几张核心表:订单表 orders、订单明细表 order_items、用户表 users。原本数仓是通过每天凌晨的批量任务拉数据,但运营需要看到当天实时的销售数据,特别是大促期间,隔天报表完全不够用。
于是同步链路就变成了这样:Flink 集群通过 Oracle LogMiner 读取 redo log 和 archive log 中的变更,经过 Flink 的 Checkpoint 机制保证精确一次,再把数据写入数据仓库的明细表。同时,因为 Flink 本身就是计算引擎,我可以在同步过程中顺便完成数据清洗、字段映射、类型转换,不需要先把原始数据落一遍再做二次加工。
这套链路用到今天,运行了几个月,整体非常稳。剩下的文章,我会把你需要知道的所有细节全部讲清楚。
2. 动手前的关键准备:Oracle 这一侧千万别省事
很多朋友上来就写 Flink SQL,结果一跑就报权限错误,回头一看 Oracle 侧啥都没配。Oracle 和 MySQL 有个很大区别:MySQL 的 binlog 默认在大多数云数据库上都是开着的,但 Oracle 默认不开归档模式,补充日志也没开,日志挖掘更不可能让你随便读。所以第一步必须把源库准备工作做扎实。
2.1 开启归档模式与最关键的补充日志
先确定你的 Oracle 是否已经开启归档模式。用有 DBA 权限的账号执行:
SELECT log_mode, supplemental_log_data_min FROM v$database;如果log_mode返回的是NOARCHIVELOG,就必须开启归档模式。这个过程需要重启数据库,所以在生产环境操作一定要走变更审批流程。单实例环境下的标准步骤如下:
SHUTDOWN IMMEDIATE; STARTUP MOUNT; ALTER DATABASE ARCHIVELOG; ALTER DATABASE OPEN; ALTER DATABASE FORCE LOGGING;如果是 RAC 环境,过程会更复杂一些,需要把所有实例都停下来,然后在一个节点上执行ALTER DATABASE ARCHIVELOG。这里我不展开 RAC 的完整步骤,但一定要提醒你:RAC 下操作不当会影响整个集群,务必在维护窗口、按照官方文档逐条执行。
归档模式开了之后,还有更关键的一步:开启补充日志(Supplemental Log)。补充日志的作用是确保 LogMiner 在解析日志时能拿到足够的信息去还原变更前后的完整数据。Oracle 默认的日志记录在 UPDATE 时可能不记录修改前所有列的值,补充日志就是解决这个问题的。
最简单的方式是开数据库级最小补充日志加所有列级补充日志:
ALTER DATABASE ADD SUPPLEMENTAL LOG DATA; ALTER DATABASE ADD SUPPLEMENTAL LOG DATA (ALL) COLUMNS;第一行开的是最小补充日志,第二行让所有被修改行的所有列都写入日志。从 CDC 的角度,我建议直接两条都执行。验证结果:
SELECT supplemental_log_data_min, supplemental_log_data_all FROM v$database;两个字段都返回YES就说明配置没问题。只开最小补充日志在某些更新场景下可能拿不到旧值,尤其是没有主键的表,到时排查问题很痛苦。
2.2 给同步账号放权限,一个都不能少
这一步特别容易遗漏。Flink CDC 连接 Oracle 时,需要读取数据字典、访问日志、执行 LogMiner 相关 API。我给生产环境建同步账号的通用脚本如下:
CREATE USER flinkuser IDENTIFIED BY "YourStrongPassword"; GRANT CREATE SESSION TO flinkuser; GRANT SET CONTAINER TO flinkuser; -- 19c 多租户环境需要 GRANT LOGMINING TO flinkuser; GRANT SELECT ON V_$DATABASE TO flinkuser; GRANT SELECT ON V_$ARCHIVED_LOG TO flinkuser; GRANT SELECT ON V_$LOGMNR_CONTENTS TO flinkuser; GRANT SELECT ON V_$LOGFILE TO flinkuser; GRANT SELECT ON V_$LOG TO flinkuser; GRANT SELECT ON V_$LOGMNR_PARAMETERS TO flinkuser; GRANT SELECT ANY TABLE TO flinkuser; GRANT SELECT ANY DICTIONARY TO flinkuser; GRANT EXECUTE ON DBMS_LOGMNR TO flinkuser; GRANT EXECUTE ON DBMS_LOGMNR_D TO flinkuser; GRANT EXECUTE ON DBMS_LOGMNR_LOGREP_DICT TO flinkuser; GRANT FLASHBACK ANY TABLE TO flinkuser;其中SELECT ANY TABLE是给全量快照阶段读数据用的,SELECT ANY DICTIONARY是读数据字典用的,LOGMINING和DBMS_LOGMNR相关权限是给 LogMiner 解析日志用的。最后一个FLASHBACK ANY TABLE容易被忽略,群里有朋友遇到过全量快照中途报ORA-08181或者 flashback 相关错误,就是因为没有这个权限。
如果你的数据库是 19c 多租户架构(CDB/PDB),尤其要注意连接串里填的 service_name 是 PDB 的,账号也要建在 PDB 里。不要在 CDB 里建同步账号,否则后面扫表和数据字典都会有问题。
2.3 监听与网络配置:很多同步失败其实是连不上
同步任务要从 Flink 端主动发起 JDBC 连接,所以 Oracle 监听必须正常,网络要通。我这里说的不只是tnsping通不通,还包括几个很隐蔽的坑。
第一个坑:Oracle 19c 默认用的认证协议比较新,如果 Flink 端用的 JDBC 驱动版本偏老,连接时可能报ORA-28040: No matching authentication protocol。解决方法是升级 JDBC 驱动到ojdbc8的新版本,或者为了避免麻烦,在 Oracle 服务端sqlnet.ora里临时放宽允许的认证版本:
SQLNET.ALLOWED_LOGON_VERSION_CLIENT=8这个配置改完不用重启监听,新连接就会生效。但从安全角度,我建议优先升级驱动而不是长期降低服务端认证等级。
第二个坑:Oracle 监听服务没有启动。这个在测试环境特别常见,你明明觉得数据库是好的,但 Flink 端就是连不上。先到 Oracle 服务器上执行lsnrctl status,如果监听没起来就执行lsnrctl start。另外检查实例是否注册到了监听器,Oracle 19c 一般会自动注册,但如果你改了端口或 hostname,可能需要重启监听。
第三个坑是防火墙。Flink 机器到 Oracle 机器的 1521 端口要放通。别笑,我见过好几次在云服务器上排了半天,最后发现是安全组没放行。
2.4 源表结构检查:主键和大字段提前摸底
Flink CDC 对同步表的主键是有要求的。Oracle CDC 在做增量阶段时,LogMiner 返回的变更数据需要有一个 key 去标记行,如果没有主键,update 和 delete 事件就无法正确映射到具体行。所以你要同步的每张表尽量都要有主键;如果确实没有,起码要有唯一索引,并在配置中指定scan.incremental.snapshot.chunk.key-column。
我在一个灰度表上吃过亏:测试表没有主键,结果同步 select 全量数据正常,但一执行 update,下游就报主键冲突。最后只能回 Oracle 补主键,重新初始化同步任务。
另外,如果表里有 CLOB/BLOB 这类大字段,快照阶段读取会比较慢。建议在同步之前就明确哪些大字段是必须同步的,能裁剪就在同步 SQL 里裁剪,不要让大字段拖慢整个链路。
3. 核心原理:Flink CDC 同步 Oracle 是怎么工作的
3.1 全量快照阶段:增量快照框架的分片思想
Flink CDC 同步一张大表时,不是简单地把历史数据全量捞一遍,再切换到增量日志。传统的 CDC 方案是全量阶段不能同时消费增量,必须等全量跑完才能开始读日志,这会导致同步完成前积压大量变更数据。
Flink CDC 的增量快照框架(Incremental Snapshot)解决了这个问题。它会把一张表的主键范围分成多个 chunk,每个 chunk 是一个区间,比如[1, 10000]、[10001, 20000]。多个 chunk 由多个并行子任务同时读取,每个 chunk 读取完成时,都会记录当前数据库日志位点(SCN)。这样在全量阶段,Oracle 上持续发生的增量变更也会被记录,当所有 chunk 都读取完成后,再统一从最早的位点开始消费增量,保证数据不丢、不重。
这个设计你可以理解成大扫除:把整个屋子分成几个区域,A 区域扫完后,在门口贴一个时间记号,B 区域扫完再贴一个,这样扫完所有区域后,只要从最早的那个记号开始继续清理新产生的垃圾,就不会有任何遗漏。
增量快照框架有两个关键收益:一是全量阶段和增量阶段可以并行,二是大表初始化速度快很多,因为并行度可以横向扩展。
3.2 增量阶段:LogMiner 解析 Oracle 日志
全量快照完成之后,Flink CDC 会进入增量阶段。Oracle 这边,Flink CDC 底层依赖的是 Debezium 的 Oracle Connector,而 Debezium 的 Oracle 实现用的是 LogMiner。
Oracle 每次数据变更都会写 redo log,归档模式下还会产生 archive log。LogMiner 就是 Oracle 官方提供的日志解析工具,它可以把日志里的变更记录还原成 SQL 级别的操作描述,比如某一行在某一个 SCN 下被 update 了,旧值是什么,新值是什么。
Flink CDC 的 Oracle Connector 在配置里有一组debezium.log.mining.*参数,其中最重要的是debezium.log.mining.strategy。它有几种取值:
online_catalog:从在线数据字典获取元数据,启动快,资源消耗相对小,但长事务场景下元数据可能不完整。redo_log_catalog:从 redo log 中记录的数据字典获取元数据,完整但解析慢。hybrid:混合模式,优先用在线字典,不够时回退到 redo 字典,也是我推荐使用的默认策略。
还有一个重要参数是debezium.log.mining.continuous.mine,默认是 true,表示持续从在线日志和归档日志中挖掘变更。如果你把它设为 false,LogMiner 就只在每个查询窗口内挖掘,消费完就停止,这种模式适用于测试,不适合生产实时同步。
3.3 精确一次与断点续传的底层逻辑
Flink CDC 的完整链路之所以可靠,核心在 Flink 的 Checkpoint 机制。Flink 会周期性对作业状态做快照,其中就包括当前消费日志的位置,也就是 offset。当作业失败重启时,从最近一次成功的 Checkpoint 恢复,继续从那个位点消费日志,配合下游 Sink 的两阶段提交,就能做到端到端的精确一次。
这里有个前提:你必须开启了 Checkpoint,并且在任务重启时指定从 Checkpoint 或 Savepoint 恢复。很多人任务一挂就直接重新提交,结果又开始全量扫描,时间全浪费了。SQL 任务里最少要这样设置:
SET 'execution.checkpointing.interval' = '3s'; SET 'execution.checkpointing.mode' = 'EXACTLY_ONCE'; SET 'state.backend.type' = 'rocksdb'; SET 'state.checkpoint-storage' = 'filesystem'; SET 'state.checkpoints.dir' = 'hdfs:///flink/checkpoints';上面配置最好写进 Flink 集群的配置文件里,而不是每次在 SQL Client 里手工执行。
4. 实操:从零搭起一套 Oracle 实时同步任务
4.1 环境准备:Flink 2.2.1 + Flink CDC 3.5.0 的 Docker 快速部署
我这次使用的组合是 Flink 2.2.1 和 Flink CDC 3.5.0,用 Docker 来跑,环境隔离和版本管理都省心。官方会发布对应的连接器 jar,文件名一般是flink-sql-connector-oracle-cdc-3.5.0.jar,把它和 Flink 镜像合在一起即可。
推荐用 Dockerfile 来构建一个带连接器的镜像:
FROM flink:2.2.1 COPY flink-sql-connector-oracle-cdc-3.5.0.jar /opt/flink/lib/构建并启动:
docker build -t flink-oracle-cdc:2.2.1 . docker run -d --name flink-jobmanager \ -p 8081:8081 \ -p 6123:6123 \ --network flink-net \ flink-oracle-cdc:2.2.1 jobmanager docker run -d --name flink-taskmanager \ --network flink-net \ flink-oracle-cdc:2.2.1 taskmanager进入容器操作 SQL Client:
docker exec -it flink-jobmanager /opt/flink/bin/sql-client.sh如果你打算跑 Flink CDC 3.x 的 YAML Pipeline,还需要下载flink-cdc-pipeline-connector-oracle相关 jar 放到 lib 目录。具体的目录结构在所有版本迭代中会有调整,最稳的做法是到当时的 Flink CDC 文档页面找到对应版本的 binary 下载清单。
4.2 用 Flink SQL 建 Source 表,直接读取 Oracle
进入 SQL Client 之后,第一步是把 Oracle 数据源注册成一张 Flink 表。这里我以订单表ORDERS为例,完整 DDL 如下:
CREATE TABLE orders_source ( order_id BIGINT PRIMARY KEY NOT ENFORCED, customer_id STRING, amount DECIMAL(10, 2), order_status INT, create_time TIMESTAMP(3) ) WITH ( 'connector' = 'oracle-cdc', 'hostname' = '192.168.10.20', 'port' = '1521', 'username' = 'flinkuser', 'password' = 'YourStrongPassword', 'database-name' = 'ORCLPDB1', 'schema-name' = 'SCOTT', 'table-name' = 'ORDERS', 'scan.incremental.snapshot.enabled' = 'true', 'scan.incremental.snapshot.chunk.size' = '8096', 'debezium.log.mining.strategy' = 'hybrid', 'debezium.log.mining.continuous.mine' = 'true' );这里几个参数我展开说一下,因为它们直接影响同步正确性和性能。
database-name填的是 Oracle 的 service name。单实例环境是 ORCL,19c 多租户环境必须填 PDB 的 service name,这个前面已经强调过。schema-name填的是 Oracle Schema,通常和用户名同名,比如 SCOTT。table-name就是你要同步的表名,如果表名是大写,这里也要大写。
scan.incremental.snapshot.enabled我建议保持 true,除非你明确知道不要增量快照。开启后大表全量同步会快很多。scan.incremental.snapshot.chunk.size默认是 8096,控制每个 chunk 的行数。对特别大的表可以调大,比如 20000 到 50000,但也要考虑源库压力,别一次把数据库搞挂了。
debezium.log.mining.strategy我用 hybrid,适应大部分场景。如果你的 Oracle 有大量长事务,建议先试 hybrid,再根据日志和监控调整。
类型映射上需要注意,Oracle 的NUMBER一般映射成 Flink 的DECIMAL或BIGINT,DATE映射成TIMESTAMP(3),VARCHAR2映射成STRING。如果你在 Oracle 端用了NUMBER(1)表示布尔,Flink 这边默认会映射成数字,不会自动变成布尔类型,需要你在查询语句里手动转换。
4.3 跑一个最简单的 Print 任务验证链路
注册完 Source 表之后,为了快速验证整个链路是通的,我建议先建一个print类型的 Sink 表,把数据打到控制台。这样不需要依赖任何外部存储,一条 SQL 就能验证 Oracle 日志解析是否正常。
CREATE TABLE orders_print ( order_id BIGINT, customer_id STRING, amount DECIMAL(10, 2), order_status INT, create_time TIMESTAMP(3), PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( 'connector' = 'print' ); INSERT INTO orders_print SELECT order_id, customer_id, amount, order_status, create_time FROM orders_source;提交任务后,打开 Flink Web UI(默认 8081 端口),可以看到作业已经进入 RUNNING 状态。这时候回 Oracle 执行几条 DML:
INSERT INTO SCOTT.ORDERS (order_id, customer_id, amount, order_status, create_time) VALUES (1001, 'C001', 199.90, 1, SYSDATE); UPDATE SCOTT.ORDERS SET order_status = 2 WHERE order_id = 1001; DELETE FROM SCOTT.ORDERS WHERE order_id = 1001;正常情况下,Print Sink 对应的 TaskManager 日志里会依次出现插入、更新、删除三条记录。看到这个,就说明 Oracle 侧配置没问题,LogMiner 解析正常,Flink 链路完整。
验证没问题之后,再把printSink 替换成真实的数仓 Sink,比如 Doris、StarRocks、Kafka 或者 JDBC。Flink CDC 3.x 对多种 Sink 都有原生支持,配置方式和普通 Flink SQL Sink 没有区别。
4.4 进阶:Flink CDC 3.x YAML Pipeline 整库同步
如果你要同步的不是一两张表,而是整个 Schema 甚至整个库,用 Flink SQL 一张一张建表太累了。Flink CDC 3.x 引入了 Pipeline 模式,用一份 YAML 配置就能完成整库同步,不需要写 SQL DDL。
下面是一个 Oracle 整库同步到 Doris 的简化配置示例:
source: type: oracle hostname: 192.168.10.20 port: 1521 username: flinkuser password: YourStrongPassword database-name: ORCLPDB1 schema-name: SCOTT tables: SCOTT.ORDERS, SCOTT.USERS, SCOTT.ORDER_ITEMS sink: type: doris username: doris_user password: doris_password fenodes: 192.168.10.100:8030 properties: format: json pipeline: parallelism: 4 schema.change: true这份配置提交方式不是用 SQL Client,而是用 Flink CDC 自带的提交脚本:
bin/flink-cdc.sh /path/to/pipeline.yamlschema.change设为 true 后,源端的表结构变更有机会自动同步到下游,这对于整库同步场景非常实用。比如 Oracle 加了一个字段,下游 Doris 表也会自动加列。不过这依赖下游 Sink 对 Schema 演变的支持程度,Doris 和 StarRocks 目前支持得比较好,其他 Sink 需要单独验证。
YAML Pipeline 的方式特别适合整库迁移、分库分表汇聚,但它的灵活性没有 Flink SQL 高。如果你需要在同步过程中做复杂的清洗和关联计算,建议还是用 Flink SQL,Pipeline 适合“原样同步”类需求。
5. 常见问题与排查技巧实录
5.1 Oracle 权限和日志挖掘相关的报错
这一块是踩坑重灾区。我按错误现象、原因、解决方式整理了一个速查表,方便你直接对照。
| 错误现象 | 根本原因 | 处理方式 |
|---|---|---|
ORA-01031 insufficient privileges | 同步账号缺少 LOGMINING 或 V$ 视图权限 | 确认补齐GRANT LOGMINING和所有 V$ 视图授权 |
ORA-00942 table or view does not exist,报错涉及 V$ 或 DBMS_LOGMNR | SELECT ANY DICTIONARY权限缺失 | 执行GRANT SELECT ANY DICTIONARY TO flinkuser; |
ORA-00604/ORA-08181出现在快照阶段 | 缺少 FLASHBACK ANY TABLE 或 undo 数据不足 | 先补权限;如果补了权限还报错,检查undo_retention |
| 启动任务后过几分钟报日志找不到 | Oracle 归档日志被清理,LogMiner 需要的日志段已不存在 | 调大归档日志保留时间,或在任务运行中保持在线日志挖掘 |
连接阶段报ORA-28040 | 客户端 JDBC 与服务端认证协议不一致 | 升级 ojdbc8 驱动;或临时调整sqlnet.ora的SQLNET.ALLOWED_LOGON_VERSION_CLIENT |
特别说一下ORA-01555: snapshot too old。这个错误在全量快照阶段容易出现,尤其是几亿行的大表。原因是全量快照在跑的时候,如果某个 chunk 读取耗时太长,Oracle 的 undo 数据被后续事务覆盖,flashback query 就取不到那个时间点的数据了。解决办法第一是调大undo_retention,第二是把scan.incremental.snapshot.chunk.size调小,让每个 chunk 的执行时间变短。
5.2 时区问题:时间字段少了 8 小时
同步 Oracle 里的DATE或TIMESTAMP字段时,如果你发现下游时间比源库时间少 8 小时,大概率是时区处理的问题。Debezium 在处理带时区的时间类型时,默认会转成 UTC,而 Flink 消费时如果没有设置正确的会话时区,就会出现偏移。
处理方式是在提交 SQL 前显式指定 Flink 本地时区:
SET 'table.local-time-zone' = 'Asia/Shanghai';然后在建表时对时间字段使用TIMESTAMP_LTZ类型,这样 Flink 会在查询结果输出时转换回目标时区。如果你直接同步到下游数据库的时间字段,下游也保持 Asia/Shanghai,就不会有偏差。测试环境强烈建议同步前先用printSink 验证一下时间字段,别等数据进了数仓才发现偏移。
5.3 大表初始化同步太慢怎么办
全量同步一张 5 亿行的表,如果默认参数跑,可能要好几个小时。加速的思路有三个方向。
第一个是调整增量快照参数。把scan.incremental.snapshot.chunk.size从默认 8096 调到 20000 到 50000,减少分片数量,降低每片的调度开销。但不要调得太大,否则单个 chunk 内 flashback 查询时间变长,反而容易触发ORA-01555。
第二个是增加并行度。如果是 Flink SQL 模式,可以在作业级别的 SET 语句里提高 source 并行度:
SET 'parallelism.default' = '8';如果是 YAML Pipeline 模式,就设置pipeline.parallelism。并行度增加后,多个 chunk 同时读取,对 Oracle 的查询压力也会增大,生产环境建议结合源库负载逐步调整。
第三个是确认瓶颈到底在哪。很多人并行度调上去了,发现全量阶段还是慢,一看监控发现 Oracle 服务器的 CPU 或者磁盘 IO 已经打满了。这时并行度再高也没用,瓶颈在源库。反过来,如果 Oracle 负载不高但速度上不去,就看 Flink TaskManager 的堆内存和 GC 情况,snapshot 读取阶段对内存消耗不小,堆内存紧张会频繁 GC,严重影响吞吐。
5.4 任务挂掉之后断点续传失败
正常情况下,Flink 作业挂了之后你只需要从最近一次 Checkpoint 恢复。但很多人的写法不对,导致恢复后重新全量扫描。
先确认 Checkpoint 是开着的。我之前见过有人只在代码里配置了 Checkpoint,但 SQL 作业是从 SQL Client 提交的,根本没加载那份代码,导致 No Checkpoint。SQL 任务要在会话里执行SET 'execution.checkpointing.interval' = '3s';并确认 Web UI 里 Checkpoint 数量在增长。
任务重启时,用flink run -s指定 Checkpoint 或 Savepoint 路径:
bin/flink run -s hdfs:///flink/checkpoints/xxx/chk-123 \ -c org.apache.flink.table.gateway.rest.util.SqlSubmit \ /path/to/your-sql-job.jar如果你用的是 SQL Gateway 或者 YAML Pipeline,也都有对应的--from-savepoint参数。
还有一种情况是任务恢复成功了,但增量阶段比较慢,这往往是因为任务停止期间积压了大量日志,LogMiner 需要从头挖掘。这时可以临时调大debezium.log.mining.batch.size.default,让每个批次处理更多日志行,等追平之后再把参数调回来。我自己遇到过积压了 3 个多小时日志的情况,调整 batch size 后十几分钟就追平了。
5.5 Schema 变更带来的连锁问题
整库同步场景下,上游经常加字段、删字段,如果同步链路处理不当,任务可能直接失败。Flink CDC 3.x 的 YAML Pipeline 对 Schema 变更支持已经不错,但使用 Flink SQL 手动建表时,源表结构变了,Flink 里的表定义不会自动变。
所以如果你用 Flink SQL 做生产同步,我建议和开发流程绑定:上游表结构变更时,同步修改 Flink 的 DDL,并从 Savepoint 重启任务。这个过程要提前演练,尤其是字段顺序、类型变化对下游的影响。
如果是临时需要新增同步一张表,不需要重启已有的同步任务,直接在同一个 SQL 作业里再创建一个新的 source 表和新的 sink 表,提交即可。但要注意新表的全量同步会占用一定资源,最好评估一下对已有任务的影响。
最后分享一点个人心得
整套 Flink CDC 同步 Oracle 的方案,我实际跑了几个月,最大的体会是:核心链路非常可靠,问题大多出在前置准备和参数调优上。Oracle 的归档、权限、补充日志这三件事没做好,后面所有精力都会消耗在报错排查上;但只要前期准备扎实,Flink CDC 的稳定性是能让人放心的。
再给大家一个实用建议:任何新表接入,先用printSink 小流量验证,在测试库里完整跑一遍 DML 覆盖,包括 insert、update、delete 以及跨日期的数据变更,确认数据都正确了再接生产。别嫌麻烦,这一套验证流程能帮你避开 90% 的线上事故。
另外,Flink 版本和 Flink CDC 版本一直在迭代,新版本对 Oracle 的支持和性能优化都更完善。如果你有条件,尽量用较新的稳定版本,同时升级前一定要看官方文档中版本的兼容性矩阵,别盲目升级导致连接器不兼容。我这次用的 Flink 2.2.1 搭配 Flink CDC 3.5.0 的整体组合,实测下来很稳,可以作为你选型的参考之一。