SeaTunnel Oracle-CDC 连接器实战指南:从 LogMiner 环境配置到 Exactly-Once 全增量同步
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
本文围绕 SeaTunnel 的 Oracle CDC 源连接器展开:它基于 Oracle LogMiner 实现「快照 + 增量」两阶段数据捕获,覆盖 CDC 任务在 Oracle 端的全部前置准备(归档日志、补充日志、LogMiner 账户授权)、完整的 Source 配置参数、数据类型映射,以及自定义主键、Exactly-Once、按时间戳/SCN 启动、Schema 变更过滤等实战场景,并结合 connector-cdc-oracle 模块源码解析其偏移量结构、启动模式校验与分块切分策略的底层实现。
一、连接器定位与能力概览
Oracle CDC 连接器(插件名Oracle-CDC)用于从 Oracle 数据库读取快照数据与增量数据,通过 Debezium 嵌入式引擎驱动 LogMiner 挖掘 redo log 实现变更捕获。它同时支持 SeaTunnel Zeta 与 Flink 引擎,从源码声明来看(OracleIncrementalSource)实现了SupportParallelism与SupportSchemaEvolution两个能力接口。
| 能力 | 是否支持 |
|---|---|
| 批量(batch) | 否 |
| 流式(stream) | 是 |
| exactly-once | 是 |
| 列投影(column projection) | 否 |
| 并行(parallelism) | 是 |
| 用户自定义 split | 是 |
一个重要的前提性说明(文档原文 Notice):Debezium Oracle 连接器不依赖log.mining.continuous.mine连续挖掘选项——连接器自身负责检测日志切换并自动调整被挖掘的日志,因此不能在debezium块中设置该属性。
数据源驱动信息
| Datasource | 支持的版本 | Driver | Url 示例 |
|---|---|---|---|
| Oracle | 不同依赖版本对应不同驱动类 | oracle.jdbc.OracleDriver | jdbc:oracle:thin:@datasource01:1523:xe |
根据文档 FAQ 部分说明,Oracle CDC 支持 Oracle Database11g、12c、19c、21c;12c 及之后的多租户环境需使用 CDB root 连接配合 common user(C##前缀)。
二、JDBC 驱动依赖安装
Oracle CDC 的 JDBC 驱动不在连接器 jar 中,需要按运行引擎分别放置:
Flink / Spark 引擎
- 将 Oracle JDBC 驱动(ojdbc8) 放置到
${SEATUNNEL_HOME}/plugins/目录; - 如需支持 i18n 字符集,将
orai18n.jar一并复制到$SEATUNNEL_HOME/plugins/。
SeaTunnel Zeta 引擎
- 将 Oracle JDBC 驱动放置到
${SEATUNNEL_HOME}/lib/目录; - 同理,
orai18n.jar复制到$SEATUNNEL_HOME/lib/。
从源码可以印证驱动加载机制:OracleIncrementalSourceFactory 的restoreSource与 OracleSourceConfigFactory 的create方法中都有Class.forName("oracle.jdbc.OracleDriver")(配置工厂内为oracle.jdbc.driver.OracleDriver全限定名)的加载动作,加载失败仅打印 warn——这也是驱动 jar 缺失时任务难以排查的典型根因。
三、Oracle 数据库侧准备:启用 LogMiner
SeaTunnel 使用 Oracle 内置的 LogMiner 工具进行 CDC。数据库端需完成两类工作:开启归档日志 + 补充日志,以及创建具备 LogMiner 权限的采集账户。以下按部署形态分述(继承自文档中的完整操作步骤)。
3.1 非 CDB(容器数据库)模式
第 1 步:在操作系统层面创建归档日志与用户表空间目录:
mkdir -p /opt/oracle/oradata/recovery_area mkdir -p /opt/oracle/oradata/ORCLCDB chown -R oracle /opt/oracle/***第 2 步:以管理员登录并启用归档日志:
sqlplus /nolog; connect sys as sysdba; alter system set db_recovery_file_dest_size = 10G; alter system set db_recovery_file_dest = '/opt/oracle/oradata/recovery_area' scope=spfile; shutdown immediate; startup mount; alter database archivelog; alter database open; ALTER DATABASE ADD SUPPLEMENTAL LOG DATA (ALL) COLUMNS; archive log list;第 3 步:创建logminer_user采集账户(密码 oracle)并授权读取表与日志:
CREATE TABLESPACE logminer_tbs DATAFILE '/opt/oracle/oradata/ORCLCDB/logminer_tbs.dbf' SIZE 25M REUSE AUTOEXTEND ON MAXSIZE UNLIMITED; CREATE USER logminer_user IDENTIFIED BY oracle DEFAULT TABLESPACE logminer_tbs QUOTA UNLIMITED ON logminer_tbs; GRANT CREATE SESSION TO logminer_user; GRANT SELECT ON V_$DATABASE to logminer_user; GRANT SELECT ON V_$LOG TO logminer_user; GRANT SELECT ON V_$LOGFILE TO logminer_user; GRANT SELECT ON V_$LOGMNR_LOGS TO logminer_user; GRANT SELECT ON V_$LOGMNR_CONTENTS TO logminer_user; GRANT SELECT ON V_$ARCHIVED_LOG TO logminer_user; GRANT SELECT ON V_$ARCHIVE_DEST_STATUS TO logminer_user; GRANT EXECUTE ON DBMS_LOGMNR TO logminer_user; GRANT EXECUTE ON DBMS_LOGMNR_D TO logminer_user; GRANT SELECT ANY TRANSACTION TO logminer_user; GRANT SELECT ON V_$TRANSACTION TO logminer_user;注:Oracle 11g 不支持
GRANT LOGMINING语句;仅需对采集表授权时,使用:
GRANT SELECT ANY TABLE TO logminer_user; GRANT ANALYZE ANY TO logminer_user;3.2 CDB(容器数据库)+ PDB(可插拔数据库)模式
第 1 步:创建目录:
mkdir -p /opt/oracle/oradata/recovery_area mkdir -p /opt/oracle/oradata/ORCLCDB mkdir -p /opt/oracle/oradata/ORCLCDB/ORCLPDB1 chown -R oracle /opt/oracle/***第 2 步:以管理员启用日志:
sqlplus /nolog connect sys as sysdba; # Password: oracle alter system set db_recovery_file_dest_size = 10G; alter system set db_recovery_file_dest = '/opt/oracle/oradata/recovery_area' scope=spfile; shutdown immediate startup mount alter database archivelog; alter database open; archive log list;第 3 步:在 CDB 中为目标表启用补充日志:
ALTER TABLE TEST.* ADD SUPPLEMENTAL LOG DATA (ALL) COLUMNS; ALTER TABLE TEST.T2 ADD SUPPLEMENTAL LOG DATA (ALL) COLUMNS;第 4 步:分别创建采集账户使用的表空间。在 CDB 中:
sqlplus sys/top_secret@//localhost:1521/ORCLCDB as sysdba CREATE TABLESPACE logminer_tbs DATAFILE '/opt/oracle/oradata/ORCLCDB/logminer_tbs.dbf' SIZE 25M REUSE AUTOEXTEND ON MAXSIZE UNLIMITED; exit;在 PDB 中:
sqlplus sys/top_secret@//localhost:1521/ORCLPDB1 as sysdba CREATE TABLESPACE logminer_tbs DATAFILE '/opt/oracle/oradata/ORCLCDB/ORCLPDB1/logminer_tbs.dbf' SIZE 25M REUSE AUTOEXTEND ON MAXSIZE UNLIMITED; exit;第 5 步:在 CDB 中创建 common user(c##dbzuser)并授予跨容器权限:
sqlplus sys/top_secret@//localhost:1521/ORCLCDB as sysdba CREATE USER c##dbzuser IDENTIFIED BY dbz DEFAULT TABLESPACE logminer_tbs QUOTA UNLIMITED ON logminer_tbs CONTAINER=ALL; GRANT CREATE SESSION TO c##dbzuser CONTAINER=ALL; GRANT SET CONTAINER TO c##dbzuser CONTAINER=ALL; GRANT SELECT ON V_$DATABASE to c##dbzuser CONTAINER=ALL; GRANT FLASHBACK ANY TABLE TO c##dbzuser CONTAINER=ALL; GRANT SELECT ANY TABLE TO c##dbzuser CONTAINER=ALL; GRANT SELECT_CATALOG_ROLE TO c##dbzuser CONTAINER=ALL; GRANT EXECUTE_CATALOG_ROLE TO c##dbzuser CONTAINER=ALL; GRANT SELECT ANY TRANSACTION TO c##dbzuser CONTAINER=ALL; GRANT LOGMINING TO c##dbzuser CONTAINER=ALL; GRANT CREATE TABLE TO c##dbzuser CONTAINER=ALL; GRANT LOCK ANY TABLE TO c##dbzuser CONTAINER=ALL; GRANT CREATE SEQUENCE TO c##dbzuser CONTAINER=ALL; GRANT EXECUTE ON DBMS_LOGMNR TO c##dbzuser CONTAINER=ALL; GRANT EXECUTE ON DBMS_LOGMNR_D TO c##dbzuser CONTAINER=ALL; GRANT SELECT ON V_$LOG TO c##dbzuser CONTAINER=ALL; GRANT SELECT ON V_$LOG_HISTORY TO c##dbzuser CONTAINER=ALL; GRANT SELECT ON V_$LOGMNR_LOGS TO c##dbzuser CONTAINER=ALL; GRANT SELECT ON V_$LOGMNR_CONTENTS TO c##dbzuser CONTAINER=ALL; GRANT SELECT ON V_$LOGMNR_PARAMETERS TO c##dbzuser CONTAINER=ALL; GRANT SELECT ON V_$LOGFILE TO c##dbzuser CONTAINER=ALL; GRANT SELECT ON V_$ARCHIVED_LOG TO c##dbzuser CONTAINER=ALL; GRANT SELECT ON V_$ARCHIVE_DEST_STATUS TO c##dbzuser CONTAINER=ALL; exit;四、Source 配置参数详解
以下为 Oracle-CDC 源插件的完整参数表(含默认值与说明):
| 参数名 | 类型 | 必填 | 默认值 | 说明 |
|---|---|---|---|---|
| url | String | 是 | - | JDBC 连接 URL,例如jdbc:oracle:thin:@//oracle-host:1521/ORCLCDB |
| username | String | 是 | - | 连接 Oracle 数据库的用户名 |
| password | String | 是 | - | 连接数据库服务器的密码 |
| database-names | List | 否 | - | 要监控的数据库名 |
| schema-names | List | 否 | - | 要监控的 schema 名 |
| table-names | List | 条件必填 | - | 要监控的表,格式为database.schema.table,例如ORCLCDB.DEBEZIUM.FULL_TYPES。与table-pattern二选一 |
| table-pattern | String | 条件必填 | - | 匹配表名的正则表达式。与table-names二选一 |
| table-names-config | List | 否 | - | 按表配置,例如[{"table": "ORCLCDB.DEBEZIUM.FULL_TYPES","primaryKeys": ["ID"],"snapshotSplitColumn": "ID"}]。适用于表无主键、需要自定义主键或显式指定快照切分列的场景 |
| startup.mode | Enum | 否 | INITIAL | 启动模式,可选initial、latest、timestamp、specific。initial:先同步历史快照再同步增量;latest:从最新位点开始、跳过初始快照;timestamp:从startup.timestamp换算的 SCN 开始;specific:从用户指定的 SCN 开始 |
| startup.timestamp | Long | 否 | - | 从指定时间戳(Unix 纪元毫秒)开始,startup.mode = timestamp时该时间戳会结合server-time-zone换算为 SCN。startup.mode为timestamp时必填 |
| startup.specific-offset.scn | Long | 否 | - | 从指定 Oracle SCN 开始。startup.mode为specific时必填,且该 SCN 必须仍可被所选 Oracle 日志挖掘后端访问 |
| stop.mode | Enum | 否 | NEVER | 停止模式,唯一合法值为never,流式 Oracle CDC 源会持续运行直到任务被停止 |
| snapshot.split.size | Integer | 否 | 8096 | 快照切分大小(行数),读表快照时表被切分为多个 split |
| snapshot.fetch.size | Integer | 否 | 1024 | 读表快照时每次 poll 的最大 fetch size |
| server-time-zone | String | 否 | UTC | 数据库服务器的会话时区,未设置时使用ZoneId.systemDefault();同时用于将startup.timestamp换算为 SCN。数据库时区与 JVM 时区不一致时应显式设置 |
| connect.timeout.ms | Duration | 否 | 30000 | 连接数据库超时时间上限 |
| connect.max-retries | Integer | 否 | 3 | 建立数据库连接的最大重试次数 |
| connection.pool.size | Integer | 否 | 20 | JDBC 连接池大小 |
| incremental.parallelism | Integer | 否 | 1 | 快照阶段结束后进入增量日志读取时使用的并行 reader 数 |
| chunk-key.even-distribution.factor.upper-bound | Double | 否 | 100 | chunk 键分布因子上限。因子 (MAX(id) - MIN(id) + 1) / 行数 小于等于该上界时按均匀分布优化切分;否则视为不均匀分布,若估算分片数超过sample-sharding.threshold则启用采样分片策略 |
| chunk-key.even-distribution.factor.lower-bound | Double | 否 | 0.05 | chunk 键分布因子下限,语义同上:因子落在上下界之内按均匀分布切分,之外触发采样策略 |
| sample-sharding.threshold | Integer | 否 | 1000 | 触发采样分片策略的估算分片数阈值(估算行数 / chunk size 超过该值时启用采样分片),用于更高效地处理超大表 |
| inverse-sampling.rate | Integer | 否 | 1000 | 采样分片策略中的采样率倒数,例如 1000 表示 1/1000 采样率,控制最终分片数粒度 |
| split.allow-sampling | Boolean | 否 | true | 设为 false 时无论分片数多少都回退到不等大 chunk 切分(迭代查询方式) |
| enable_concurrent_read | Boolean | 否 | true | 快照阶段是否启用基于 split 的并发读取;设为 false 时源跳过 split 分析、将整表作为单个 split 读取,适用于无索引表 |
| exactly_once | Boolean | 否 | false | 启用 Exactly-Once 语义 |
| use_select_count | Boolean | 否 | false | 全量阶段用select count直接统计行数,而非通过 analysis 表估算(当 statistics 更新慢时直接 count 更快) |
| skip_analyze | Boolean | 否 | false | 全量阶段跳过表行数分析(适用于已周期性调度 analysis 更新统计信息或数据不频繁变化的场景) |
| format | Enum | 否 | DEFAULT | 输出格式,可选DEFAULT、COMPATIBLE_DEBEZIUM_JSON |
| schema-changes.enabled | Boolean | 否 | false | 默认禁用 Schema Evolution;当前仅支持加列、删列、改列名、改列类型 |
| schema-changes.include | List | 否 | - | 仅让列出的 Schema 变更事件类型向下游传递(需schema-changes.enabled = true),空表示全部放行 |
| schema-changes.exclude | List | 否 | - | 列出的 Schema 变更事件类型不向下游传递,在 include 之后生效,冲突时exclude 优先 |
| debezium | Config | 否 | - | 透传给 Debezium Embedded Engine 的属性,用于捕获 Oracle 服务器数据变更 |
| common-options | - | 否 | - | 源插件通用参数,参见 Source Common Options |
| decimal_type_narrowing | Boolean | 否 | true | 十进制类型窄化:为 true 时,在不损失精度的前提下将 decimal 窄化为 int 或 long 类型,当前仅 Oracle 支持,见下文详解 |
4.1 启动模式的源码校验逻辑
OracleIncrementalSourceFactory.optionRule 中定义了条件校验规则:startup.specific-offset.scn仅在startup.mode = specific时必需且必须大于 0;startup.timestamp仅在timestamp模式时必需;exactly_once仅在initial模式下可选。此外table-names与table-pattern被声明为互斥项(exclusive)。
进一步地,OracleIncrementalSource.getOracleStartupConfig 对specific模式做了显式处理:SCN 会被封装成RedoLogOffset结构的scn、commit_scn=0、lcr_position=null三键位点,而通用 CDC 基类中的 file/position 位点在 Oracle 下会被直接拒绝——因为 Oracle 的位点语义是 SCN(System Change Number)而非日志文件偏移。
4.2 命名约束:库名与表名必须大写
从 OracleSourceConfigFactory.validateConfig 源码可见三条硬校验:
- Oracle仅支持单个 database(
database.names列表只能有一个元素); - 数据库名、表名中出现的字母必须全部为大写;
- 表名格式必须是
${database}.${schema}.${table}或${schema}.${table}。
配置示例中统一使用ORCLCDB.DEBEZIUM.FULL_TYPES这种三段式写法即是此原因的体现。
4.3 decimal_type_narrowing 类型窄化
decimal_type_narrowing = true(默认)时,不损失精度则窄化为整型;false时保留 DECIMAL:
| Oracle | true 时映射 | false 时映射 |
|---|---|---|
| NUMBER(1, 0) | Boolean | Decimal(1, 0) |
| NUMBER(6, 0) | INT | Decimal(6, 0) |
| NUMBER(10, 0) | BIGINT | Decimal(10, 0) |
4.4 快照行数统计的三种策略
全量阶段需要先估算表行数来决定切分策略,文档提供了三种可组合的手段:
- 默认:通过 Oracle analysis 表(统计信息)估算行数;
use_select_count = true:直接执行select count(*)计数,适用于统计信息过期、直接 count 更快的场景;skip_analyze = true:跳过行数分析,直接读取all_tables的NUM_ROWS,适用于已周期性调度 analyze 更新统计信息、或数据不频繁变化的表。
五、数据类型映射
| Oracle 类型 | SeaTunnel 类型 |
|---|---|
| INTEGER | INT |
| FLOAT | DECIMAL(38, 18) |
| NUMBER(precision <= 9, scale == 0) | INT |
| NUMBER(9 < precision <= 18, scale == 0) | BIGINT |
| NUMBER(18 < precision, scale == 0) | DECIMAL(38, 0) |
| NUMBER(precision == 0, scale == 0) | DECIMAL(38, 18) |
| NUMBER(scale != 0) | DECIMAL(38, 18) |
| BINARY_DOUBLE | DOUBLE |
| BINARY_FLOAT / REAL | FLOAT |
| CHAR / NCHAR / NVARCHAR2 / VARCHAR2 / LONG / ROWID / NCLOB / CLOB | STRING |
| DATE | DATE |
| TIMESTAMP / TIMESTAMP WITH LOCAL TIME ZONE | TIMESTAMP |
| BLOB / RAW / LONG RAW / BFILE | BYTES |
类型转换的实现位于 OracleTypeUtils,它将 Debezium 的Column元信息(类型名、长度、精度、小数位、默认值)桥接为 SeaTunnel 的BasicTypeDefine,再委托给OracleTypeConverter完成到 SeaTunnel 类型系统的转换;其中对TIMESTAMP开头的类型会把length作为 scale 处理,以保留时间精度信息。
六、任务配置示例
以下示例完整继承自官方文档,可直接复制修改后使用。
6.1 基础多表读取
source { Oracle-CDC { plugin_output = "customers" username = "system" password = "top_secret" database-names = ["ORCLCDB"] schema-names = ["DEBEZIUM"] table-names = ["ORCLCDB.DEBEZIUM.FULL_TYPES", "ORCLCDB.DEBEZIUM.FULL_TYPES2"] url = "jdbc:oracle:thin:@//oracle-host:1521/ORCLCDB" source.reader.close.timeout = 120000 connection.pool.size = 1 debezium { database.oracle.jdbc.timezoneAsRegion = "false" } } }6.2 用 select count 直接统计行数(use_select_count)
source { Oracle-CDC { plugin_output = "customers" use_select_count = true username = "system" password = "top_secret" database-names = ["ORCLCDB"] schema-names = ["DEBEZIUM"] table-names = ["ORCLCDB.DEBEZIUM.FULL_TYPES"] url = "jdbc:oracle:thin:@//oracle-host:1521/ORCLCDB" source.reader.close.timeout = 120000 } }6.3 读取 all_tables 的 NUM_ROWS 并跳过 analyze(skip_analyze)
source { Oracle-CDC { plugin_output = "customers" skip_analyze = true username = "system" password = "top_secret" database-names = ["ORCLCDB"] schema-names = ["DEBEZIUM"] table-names = ["ORCLCDB.DEBEZIUM.FULL_TYPES"] url = "jdbc:oracle:thin:@//oracle-host:1521/ORCLCDB" source.reader.close.timeout = 120000 } }6.4 自定义主键(table-names-config)
source { Oracle-CDC { plugin_output = "customers" url = "jdbc:oracle:thin:@//oracle-host:1521/ORCLCDB" source.reader.close.timeout = 120000 username = "system" password = "top_secret" database-names = ["ORCLCDB"] schema-names = ["DEBEZIUM"] table-names = ["ORCLCDB.DEBEZIUM.FULL_TYPES"] table-names-config = [ { table = "ORCLCDB.DEBEZIUM.FULL_TYPES" primaryKeys = ["ID"] } ] } }6.5 启用 Exactly-Once CDC
exactly_once = true与默认startup.mode = "initial"路径配合使用;当下游 sink 也配置为 exactly-once 投递(例如开启 XA 的 JDBC sink)时启用:
env { parallelism = 1 job.mode = "STREAMING" checkpoint.interval = 5000 } source { Oracle-CDC { plugin_output = "customers" username = "system" password = "top_secret" database-names = ["ORCLCDB"] schema-names = ["DEBEZIUM"] table-names = ["ORCLCDB.DEBEZIUM.FULL_TYPES"] url = "jdbc:oracle:thin:@//oracle-host:1521/ORCLCDB" exactly_once = true connection.pool.size = 1 debezium { database.oracle.jdbc.timezoneAsRegion = "false" } } }6.6 从指定时间戳启动
使用startup.mode = "timestamp"从 Unix 毫秒时间戳换算出的 Oracle SCN 开始同步:
source { Oracle-CDC { plugin_output = "customers" username = "system" password = "top_secret" database-names = ["ORCLCDB"] schema-names = ["DEBEZIUM"] table-names = ["ORCLCDB.DEBEZIUM.FULL_TYPES"] url = "jdbc:oracle:thin:@//oracle-host:1521/ORCLCDB" startup.mode = "timestamp" startup.timestamp = 1700000000000 server-time-zone = "UTC" debezium { database.oracle.jdbc.timezoneAsRegion = "false" } } }6.7 配置 Debezium 心跳(Heartbeat)
低流量场景下,Oracle LogMiner 的 SCN 只有在 redo log 发生变化时才会前进。配置 Debezium 心跳可以让 SCN 持续移动,使 checkpoint 位点被定期记录、复制延迟保持可观测。心跳表必须在任务启动前已在 Oracle 服务端存在:
source { Oracle-CDC { username = "system" password = "top_secret" database-names = ["ORCLCDB"] schema-names = ["DEBEZIUM"] table-names = ["ORCLCDB.DEBEZIUM.FULL_TYPES"] url = "jdbc:oracle:thin:@//oracle-host:1521/ORCLCDB" debezium { database.oracle.jdbc.timezoneAsRegion = "false" heartbeat.interval.ms = 100 heartbeat.action.query = "INSERT INTO DEBEZIUM.heartbeat (ts) VALUES (SYSTIMESTAMP)" } } }6.8 读取无主键表
选择与源表保证程度匹配的路径:
- 仅追加(Append-only)负载(下游永不产生 UPDATE/DELETE):保持
exactly_once = false且不声明主键。源会退化为尽力而为的行标识;没有可用键时,连接器无法安全地应用 UPDATE/DELETE 事件。 - 存在唯一非主键列:通过
table-names-config.primaryKeys声明该列,并设置exactly_once = true,使快照阶段与 redo log 阶段使用同一把键保持行标识一致。
source { Oracle-CDC { username = "system" password = "top_secret" database-names = ["ORCLCDB"] schema-names = ["DEBEZIUM"] url = "jdbc:oracle:thin:@//oracle-host:1521/ORCLCDB" table-names = ["ORCLCDB.DEBEZIUM.FULL_TYPES_NO_PRIMARY_KEY"] table-names-config = [ { table = "ORCLCDB.DEBEZIUM.FULL_TYPES_NO_PRIMARY_KEY" primaryKeys = ["ID"] } ] exactly_once = true } }没有可用主键时,连接器无法安全地应用 UPDATE/DELETE 事件,该模式只应用于仅追加负载。
6.9 Schema 变更事件过滤
当schema-changes.enabled = true时,可用schema-changes.include/schema-changes.exclude进一步控制哪些变更事件向下游传播。使用 SeaTunnel 的规范事件名:
| 规范名 | 操作 |
|---|---|
add.column | 加列 |
drop.column | 删列 |
modify.column | 修改列类型/属性,列名不变 |
change.column | 重命名列,可选地改类型 |
update.columns | 上述四种列级变更的组别名 |
优先级是确定性的:① 若设置了schema-changes.include,只有被包含的事件类型有资格下发;② 随后应用schema-changes.exclude;③ 同一类型同时出现在两个列表时exclude 胜出。
source { Oracle-CDC { # ... schema-changes.enabled = true schema-changes.include = ["add.column", "drop.column"] schema-changes.exclude = ["change.column"] } }关于排除drop.column的数据处理注意:对于被保留的NOT NULL列,NULL写入会被 sink 拒绝——如果对某个源端已停止供数的 NOT NULL 列排除drop.column,会在 sink 端失败。
6.10 使用 Debezium 兼容格式发送到 Kafka
format = COMPATIBLE_DEBEZIUM_JSON必须与 Kafka sink 连接器配合使用,细节参见 cdc-compatible-debezium-json 格式文档。
七、源码级实现解析
结合 connector-cdc-oracle 模块,可以进一步理解文档中各项行为的底层支撑:
1. LogMiner 挖掘策略与只读模式。OracleSourceConfigFactory.create 会向 Debezium 属性集中写入若干 SeaTunnel 内部约定:log.mining.read.only = true(只读挖掘模式,对应 changelog 中的 ReadOnlyLogWriterFlushStrategy 支持);tombstones.on.delete = false(禁用墓碑消息);默认log.mining.strategy = online_catalog,而当schema-changes.enabled = true时切换为redo_log_catalog——因为online_catalog无法正确捕获 DDL。源码中还有一处显式防御:若用户在debezium块中把include.schema.changes设为 true 但同时指定了online_catalog策略,会直接抛出 IllegalArgumentException。
2. RedoLogOffset 位点结构。RedoLogOffset 用scn、commit_scn、lcr_position三个键描述 redo log 事件位点,compareTo通过 Debezium 的Scn值做顺序比较;NO_STOPPING_OFFSET(Long.MIN_VALUE)表示永不停止,对应stop.mode = never的流式语义。
3. 快照切分(ChunkSplitter)。OracleChunkSplitter 实现了queryMinMax、sampleDataFromColumn(调用OracleUtils.skipReadAndSortSampleData做跳读采样)、queryNextChunkMax(迭代式不等大分片)与queryApproximateRowCnt等方法,正是第四节中chunk-key.even-distribution.factor.*、inverse-sampling.rate、split.allow-sampling等参数的执行落点;对ROWID类型切分键还使用了ROWID.compareBytes做字节级比较。enable_concurrent_read = false时,该切分分析整体被跳过、整表作为单 split 读取。
4. Schema 变更解析。Oracle 的 DDL 变更由 OracleSchemaChangeResolver 委托给基于 ANTLR 的 CustomOracleAntlrDdlParser 解析 redo log 中的 DDL 语句,产出AlterTableColumnEvent事件序列——这解释了为什么当前仅支持加列、删列、改名、改类型这四类列级变更(与 OracleIncrementalSource.supports 中声明的ADD_COLUMN / DROP_COLUMN / RENAME_COLUMN / UPDATE_COLUMN一致)。
5. 任务分发。OracleDialect.createFetchTask 依据 split 类型分发到OracleSnapshotFetchTask(快照阶段)或OracleRedoLogFetchTask(增量阶段,内部复用 Debezium 的 LogMiner streaming event source 与只读日志写入冲刷策略 ReadOnlyLogWriterFlushStrategy)。
八、FAQ
Q1:CDC 需要哪些 Oracle 权限?
LogMiner 用户需要以下权限:
GRANT CREATE SESSION TO logminer_user; GRANT SET CONTAINER TO logminer_user; GRANT SELECT ON V_$DATABASE TO logminer_user; GRANT FLASHBACK ANY TABLE TO logminer_user; GRANT SELECT ANY TABLE TO logminer_user; GRANT SELECT_CATALOG_ROLE TO logminer_user; GRANT EXECUTE_CATALOG_ROLE TO logminer_user; GRANT SELECT ANY TRANSACTION TO logminer_user; GRANT LOGMINING TO logminer_user; GRANT CREATE TABLE TO logminer_user; GRANT LOCK ANY TABLE TO logminer_user; GRANT CREATE SEQUENCE TO logminer_user; GRANT EXECUTE ON DBMS_LOGMNR TO logminer_user; GRANT EXECUTE ON DBMS_LOGMNR_D TO logminer_user; GRANT SELECT ON V_$LOG TO logminer_user; GRANT SELECT ON V_$LOG_HISTORY TO logminer_user; GRANT SELECT ON V_$LOGMNR_LOGS TO logminer_user; GRANT SELECT ON V_$LOGMNR_CONTENTS TO logminer_user; GRANT SELECT ON V_$LOGMNR_PARAMETERS TO logminer_user; GRANT SELECT ON V_$LOGFILE TO logminer_user; GRANT SELECT ON V_$ARCHIVED_LOG TO logminer_user; GRANT SELECT ON V_$ARCHIVE_DEST_STATUS TO logminer_user; GRANT SELECT ON V_$TRANSACTION TO logminer_user;同时在数据库级和表级启用补充日志:
ALTER DATABASE ADD SUPPLEMENTAL LOG DATA; ALTER TABLE schema_name.table_name ADD SUPPLEMENTAL LOG DATA (ALL) COLUMNS;Q2:支持多租户(CDB/PDB)数据库吗?
支持。将database-names设置为 CDB 名,JDBC URL 指向 CDB root;用户必须是 common user(C##前缀),且上述权限以CONTAINER = ALL授予所有容器。
Q3:支持无主键表吗?
默认要求主键。若表中存在合适的唯一列,可通过table-names-config的primaryKeys字段指定自定义主键列;无可用主键时只能用于仅追加负载。
Q4:如何使用自定义快照查询?
通过debezium块内的 Debezium 属性snapshot.select.statement.overrides配置。该查询在 SeaTunnel 添加快照切分边界之前应用,因此必须包含配置表 schema 及其切分键所需的全部列:
debezium { snapshot.select.statement.overrides = "DEBEZIUM.FULL_TYPES" snapshot.select.statement.overrides.DEBEZIUM.FULL_TYPES = "SELECT * FROM DEBEZIUM.FULL_TYPES WHERE ACTIVE = 1" }Q5:如何提升 LogMiner 性能?
应将其首先当作数据库与 redo log 调优问题处理:优先复用本文的 LogMiner 配置与补充日志小节,仅为所需表启用日志;在此基础上再考虑通过debezium透传属性调优,且须先验证这些属性在你实际部署的 Oracle CDC 运行时版本中受支持。
九、延伸阅读
- 覆盖「全量 + 增量」同步完整生命周期、2PC sink 配置、Schema Evolution 与排障的生产级指南:CDC Production Cookbook
- 兼容 Debezium JSON 格式(用于 Kafka sink):cdc-compatible-debezium-json
- Oracle CDC 连接器变更历史:connector-cdc-oracle changelog
- 端到端测试:connector-cdc-oracle-e2e
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考