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包,注意使用
providedscope避免重复引入Kafka Connector核心依赖。 - 将Jar包放入Kafka Connector的插件目录(对应配置中的
plugin.path路径)。
4. 连接器配置
在Connector配置文件中添加SMT配置:
transforms=removeUnicodeNull transforms.removeUnicodeNull.type=你的完整包名.RemoveUnicodeNull
内容的提问来源于stack exchange,提问作者Italo Vieira
相关产品推荐
相关产品推荐

