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

Apache Kafka自定义SMT咨询:处理CDC数据中的Unicode Null值

Kafka Connector 处理Unicode Null垃圾数据的SMT方案

现有SMT是否支持全局替换

目前官方及主流社区维护的Kafka Connector SMT(如org.apache.kafka.connect.transforms下的工具类)没有直接支持全记录级的replaceAll操作。大部分SMT聚焦于单字段的增删改查、类型转换等场景,无法自动遍历所有字段并替换Unicode Null(\u0000)。

自定义SMT实现方案

1. 依赖准备

引入Kafka Connector核心依赖(Maven示例):

<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>connect-api</artifactId>
    <version>你的Kafka版本</version>
    <scope>provided</scope>
</dependency>

2. 核心实现代码

自定义SMT需实现Transformation接口,支持处理Struct(CDC常见格式)和原始字符串类型的记录,递归清理所有字段中的Unicode Null:

import org.apache.kafka.common.config.ConfigDef;
import org.apache.kafka.connect.connector.ConnectRecord;
import org.apache.kafka.connect.data.Field;
import org.apache.kafka.connect.data.Schema;
import org.apache.kafka.connect.data.SchemaBuilder;
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.Map;

public class RemoveUnicodeNull implements Transformation<ConnectRecord<?>> {

    @Override
    public void configure(Map<String, ?> configs) {
        new SimpleConfig(conf(), configs);
    }

    @Override
    public ConnectRecord<?> apply(ConnectRecord<?> record) {
        if (record.value() == null) {
            return record;
        }

        // 处理结构化的CDC记录(Struct类型)
        if (record.value() instanceof Struct) {
            Struct cleanedStruct = cleanStruct((Struct) record.value());
            return record.newRecord(
                    record.topic(),
                    record.kafkaPartition(),
                    record.keySchema(),
                    record.key(),
                    cleanedStruct.schema(),
                    cleanedStruct,
                    record.timestamp()
            );
        }

        // 处理原始JSON字符串格式的记录
        if (record.value() instanceof String) {
            String cleanedValue = ((String) record.value()).replaceAll("\\u0000", "");
            return record.newRecord(
                    record.topic(),
                    record.kafkaPartition(),
                    record.keySchema(),
                    record.key(),
                    record.valueSchema(),
                    cleanedValue,
                    record.timestamp()
            );
        }

        return record;
    }

    // 递归清理嵌套Struct中的所有字段
    private Struct cleanStruct(Struct original) {
        Schema cleanedSchema = SchemaBuilder.struct().name(original.schema().name())
                .fields(original.schema().fields())
                .build();
        Struct cleanedStruct = new Struct(cleanedSchema);

        for (Field field : original.schema().fields()) {
            Object value = original.get(field);
            if (value instanceof String) {
                cleanedStruct.put(field.name(), ((String) value).replaceAll("\\u0000", ""));
            } else if (value instanceof Struct) {
                cleanedStruct.put(field.name(), cleanStruct((Struct) value));
            } else {
                cleanedStruct.put(field.name(), value);
            }
        }
        return cleanedStruct;
    }

    @Override
    public ConfigDef config() {
        return new ConfigDef();
    }

    @Override
    public void close() {}
}

3. 打包部署

  • 将代码打包成Jar包,注意使用provided scope避免重复引入Kafka Connector核心依赖。
  • 将Jar包放入Kafka Connector的插件目录(对应配置中的plugin.path路径)。

4. 连接器配置

在Connector配置文件中添加SMT配置:

transforms=removeUnicodeNull
transforms.removeUnicodeNull.type=你的完整包名.RemoveUnicodeNull

内容的提问来源于stack exchange,提问作者Italo Vieira

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 07:53:16