Kafka MySQL源连接器如何同步MySQL函数至Oracle Sink连接器
当前基于Kafka连接器搭建MySQL到Oracle的CDC(变更数据捕获)同步链路,环境规格如下:
- MySQL 8.0,部署在CentOS 7.8(64-bit)系统
- Oracle 18c XE,部署在Oracle Linux 7.6(64-bit)系统
- 源端连接器:debezium-connector-mysql-1.9.5
- 下沉端连接器:confluentinc-kafka-connect-jdbc-10.5.0
目前链路已正常支持MySQL侧建表、insert操作的变更同步到Oracle,但创建MySQL函数时,虽然连接器可以读取binlog生成对应topic,变更始终无法下沉到Oracle。经初步排查:Confluent JDBC Sink连接器按表维度的topic生成并执行对应SQL,MySQL函数生成时不会产生独立的表维度专属topic,导致下沉逻辑无法触发。
首先纠正一个排查误区:函数事件没有独立专属topic不是bug,是连接器的设计逻辑。
Debezium MySQL连接器默认只针对表对象创建独立的表级topic,所有非表类的DDL事件(包括CREATE FUNCTION、存储过程、触发器创建/修改/删除)都会统一写入服务级别的全局schema change topic,不会拆分独立topic。而Confluent JDBC Sink从设计上就只消费表级topic,内置的SQL生成逻辑完全围绕表数据变更、表结构变更实现,既不会消费全局schema change topic里的函数事件,也没有内置MySQL函数语法到Oracle PL/SQL语法的转换能力,这才是函数无法同步的核心原因。
方案1:轻量改造现有CDC链路
适合函数数量少、逻辑简单的场景,不需要重构现有链路:
- 确认源端Debezium连接器配置开启
include.schema.changes=true,保证所有DDL事件(含函数操作)都正常写入全局schema change topic。 - 新增一层消息处理逻辑,可以用Kafka Connect单消息转换(SMT)实现,也可以部署轻量的Kafka消费端程序:
- 消费全局schema change topic的消息,过滤出和函数操作相关的DDL事件(CREATE/ALTER/DROP FUNCTION)
- 做语法转换,把MySQL的函数定义改成Oracle可直接执行的PL/SQL语法:去掉MySQL特有的
DELIMITER标记、把字段类型映射为Oracle对应类型、替换MySQL独有的内置函数为Oracle等价写法、调整函数块结构为Oracle要求的CREATE OR REPLACE FUNCTION ... IS BEGIN ... END;格式 - 把转换完成的合法Oracle DDL语句发送到专门的自定义topic(比如命名为
oracle_ddl_execute)
- 配置支持原生DDL执行的Sink连接器消费这个自定义topic,跳过自动生成SQL的逻辑,直接执行消息体里携带的DDL语句即可完成函数同步。
方案2:独立元数据同步链路(生产环境推荐)
CDC链路本身定位是同步表数据和表结构,函数、存储过程这类程序化逻辑不适合走CDC自动同步,稳定性差:
- 单独搭建非表对象的同步流程:定时(或者每次代码发版前)用
mysqldump --routines --no-create-info --no-data --no-create-db --skip-opt <数据库名>命令导出MySQL侧全量函数、存储过程定义 - 提前做语法适配,把MySQL函数语法批量转换为Oracle兼容的PL/SQL脚本,和业务代码一起纳入版本库管理,做人工校验
- 定时调度或者发版时直接在Oracle侧执行校验通过的脚本,完成函数同步,原有CDC链路只负责表和数据的同步。
这个方案是异构数据库同步的通用实践,稳定性远高于全自动CDC同步,可以提前拦截语法差异、逻辑不兼容的问题,避免线上执行异常。
方案3:二次开发Sink连接器(不推荐)
如果有极强的全链路自动化需求,可以基于Confluent JDBC Sink源码做二次开发:新增全局schema change topic消费逻辑、内置MySQL到Oracle的函数语法转换引擎、新增函数幂等执行逻辑。但这个方案开发维护成本极高,MySQL和Oracle的函数语法差异覆盖范围极广(变量定义、异常处理、内置函数、游标逻辑、返回值规则都存在大量差异),长期维护成本很高,非特殊需求不建议采用。
内容的提问来源于stack exchange,提问作者coffeethesharing

