如何在Debezium Oracle源连接器中把RAW(16)转换为GUID字符串?
Debezium Oracle RAW(16)转GUID字符串的连接器转换方案
Debezium默认会将Oracle的RAW(16)类型序列化为Base64编码字符串,要在源连接器层面将其转换为十六进制格式的GUID字符串,可通过自定义Single Message Transform(SMT)实现,避免下游消费者重复处理。
步骤1:编写自定义SMT类
创建Java类实现Kafka Connect的Transformation接口,完成Base64解码、字节数组转十六进制字符串的逻辑,同时处理Debezium CDC消息中的before和after字段:
package com.example.transforms; import org.apache.kafka.common.config.ConfigDef; import org.apache.kafka.connect.connector.ConnectRecord; import org.apache.kafka.connect.data.Struct; import org.apache.kafka.connect.transforms.Transformation; import org.apache.kafka.connect.transforms.util.SimpleConfig; import java.util.Base64; import java.util.Map; public class RawToGuidTransform<R extends ConnectRecord<R>> implements Transformation<R> { private static final String FIELD_NAME_CONFIG = "field.name"; private static final ConfigDef CONFIG_DEF = new ConfigDef() .define(FIELD_NAME_CONFIG, ConfigDef.Type.STRING, ConfigDef.Importance.HIGH, "目标转换字段名"); private String targetField; @Override public void configure(Map<String, ?> configs) { SimpleConfig config = new SimpleConfig(CONFIG_DEF, configs); targetField = config.getString(FIELD_NAME_CONFIG); } @Override public R apply(R record) { if (!(record.value() instanceof Struct)) { return record; } Struct rootStruct = (Struct) record.value(); Struct updatedRoot = new Struct(rootStruct.schema()).putAll(rootStruct); // 处理after字段 handleStructField(updatedRoot, "after"); // 处理before字段 handleStructField(updatedRoot, "before"); return record.newRecord( record.topic(), record.kafkaPartition(), record.keySchema(), record.key(), updatedRoot.schema(), updatedRoot, record.timestamp() ); } private void handleStructField(Struct rootStruct, String structName) { Struct dataStruct = rootStruct.getStruct(structName); if (dataStruct == null || dataStruct.schema().field(targetField) == null) { return; } String base64Value = dataStruct.getString(targetField); byte[] rawBytes = Base64.getDecoder().decode(base64Value); StringBuilder guidStr = new StringBuilder(); for (byte b : rawBytes) { guidStr.append(String.format("%02x", b)); } Struct updatedDataStruct = new Struct(dataStruct.schema()) .putAll(dataStruct) .put(targetField, guidStr.toString()); rootStruct.put(structName, updatedDataStruct); } @Override public ConfigDef config() { return CONFIG_DEF; } @Override public void close() {} }
步骤2:编译打包SMT
将上述代码编译为JAR文件,放入Kafka Connect的类路径中,例如:
- Debezium Oracle连接器的lib目录:
kafka-connect/plugins/debezium-connector-oracle/lib/ - 或者Kafka Connect的全局插件目录
步骤3:配置Debezium源连接器
在连接器配置中添加转换规则,指定自定义SMT和目标字段:
name=oracle-table-connector connector.class=io.debezium.connector.oracle.OracleConnector tasks.max=1 database.hostname=your-oracle-host database.port=1521 database.user=your-db-user database.password=your-db-pass database.dbname=your-db-name database.server.name=oracle-cdc-server table.include.list=YOUR_SCHEMA.TABLE_NAME # 添加转换配置 transforms=convertRowno transforms.convertRowno.type=com.example.transforms.RawToGuidTransform transforms.convertRowno.field.name=rowno
步骤4:生效配置
重启Kafka Connect服务,或重新加载连接器配置。此后,Kafka中传输的CDC消息里,rowno字段将直接是十六进制格式的GUID字符串(如2b09ccae06aa45aaa2f1dd0e4c047249),无需下游消费者再做转换。
内容的提问来源于stack exchange,提问作者user1913559
相关产品推荐
相关产品推荐

