You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.26 11:00:36