Kafka Connect自定义Transform类无法加载报错求助
针对你遇到的Kafka Connect自定义Transform类加载失败的问题,我从你的错误日志、代码和配置里找到了核心问题点,下面一步步帮你排查解决:
问题背景回顾
你已经将自定义Transform的jar放入Kafka Connect插件目录,但启动时持续报错找不到transforms.RemoveWeirdCharacters类,核心错误、代码及配置如下:
核心错误日志
[2025-02-10 14:54:53,530] INFO AbstractConfig values: (org.apache.kafka.common.config.AbstractConfig:370) [2025-02-10 14:54:53,639] ERROR Failed to create connector for .\kafka-connect\oracle-sink.properties (org.apache.kafka. connect.cli.ConnectStandalone:74) [2025-02-10 14:54:53,642] ERROR Stopping after connector error (org.apache.kafka.connect.cli.ConnectStandalone:84) java.util.concurrent.ExecutionException: org.apache.kafka.connect.runtime.rest.errors.BadRequestException: Connector con figuration is invalid and contains the following 2 error(s): Invalid value transforms.RemoveWeirdCharacters for configuration transforms.removeWeirdChars.type: Class transforms.Remo veWeirdCharacters could not be found. Invalid value null for configuration transforms.removeWeirdChars.type: Not a Transformation You can also find the above list of errors at the endpoint `/connector-plugins/{connectorType}/config/validate` at org.apache.kafka.connect.util.ConvertingFutureCallback.result(ConvertingFutureCallback.java:123) at org.apache.kafka.connect.util.ConvertingFutureCallback.get(ConvertingFutureCallback.java:107) at org.apache.kafka.connect.cli.ConnectStandalone.processExtraArgs(ConnectStandalone.java:81) at org.apache.kafka.connect.cli.AbstractConnectCli.startConnect(AbstractConnectCli.java:150) at org.apache.kafka.connect.cli.AbstractConnectCli.run(AbstractConnectCli.java:94) at org.apache.kafka.connect.cli.ConnectStandalone.main(ConnectStandalone.java:112) Caused by: org.apache.kafka.connect.runtime.rest.errors.BadRequestException: Connector configuration is invalid and cont ains the following 2 error(s): Invalid value transforms.RemoveWeirdCharacters for configuration transforms.removeWeirdChars.type: Class transforms.Remo veWeirdCharacters could not be found. Invalid value null for configuration transforms.removeWeirdChars.type: Not a Transformation You can also find the above list of errors at the endpoint `/connector-plugins/{connectorType}/config/validate` at org.apache.kafka.connect.runtime.AbstractHerder.maybeAddConfigErrors(AbstractHerder.java:754) at org.apache.kafka.connect.runtime.standalone.StandaloneHerder.putConnectorConfig(StandaloneHerder.java:204) at org.apache.kafka.connect.runtime.standalone.StandaloneHerder.lambda$null$0(StandaloneHerder.java:190) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:539) at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) at java.base/java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecuto r.java:304) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635) at java.base/java.lang.Thread.run(Thread.java:833)
自定义Transform代码
package transforms; import org.apache.kafka.common.config.ConfigDef; import org.apache.kafka.connect.connector.ConnectRecord; import org.apache.kafka.connect.transforms.Transformation; import java.util.Map; public class RemoveWeirdCharacters<R extends ConnectRecord<R>> implements Transformation<R> { @Override public R apply(R record) { if (record.value() instanceof String value) { value = value.replaceAll("^[^a-zA-Z0-9]+", ""); return record.newRecord( record.topic(), record.kafkaPartition(), record.keySchema(), record.key(), record.valueSchema(), value, record.timestamp() ); } return record; } @Override public ConfigDef config() { return new ConfigDef(); } @Override public void configure(Map<String, ?> configs) { } @Override public void close() { } }
pom.xml配置
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd"> <modelVersion>4.0.0</modelVersion> <groupId>com.example</groupId> <artifactId>custom-transforms</artifactId> <name>custom-transforms</name> <version>1.0.0</version> <properties> <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding> <junit.version>5.9.2</junit.version> </properties> <dependencies> <!-- Kafka Connect dependencies --> <dependency> <groupId>org.apache.kafka</groupId> <artifactId>connect-api</artifactId> <version>3.7.1</version> </dependency> <dependency> <groupId>org.apache.kafka</groupId> <artifactId>connect-transforms</artifactId> <version>3.7.1</version> </dependency> <dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>3.7.1</version> </dependency> <!-- JUnit for testing --> <dependency> <groupId>org.junit.jupiter</groupId> <artifactId>junit-jupiter-api</artifactId> <version>${junit.version}</version> <scope>test</scope> </dependency> <dependency> <groupId>org.junit.jupiter</groupId> <artifactId>junit-jupiter-engine</artifactId> <version>${junit.version}</version> <scope>test</scope> </dependency> </dependencies> <build> <plugins> <!-- Compiler plugin --> <plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-compiler-plugin</artifactId> <version>3.11.0</version> <configuration> <source>21</source> <target>21</target> </configuration> </plugin> <!-- Surefire plugin for running tests --> <plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-surefire-plugin</artifactId> <version>3.0.0-M7</version> </plugin> </plugins> </build> </project>
Kafka Connect配置
name=oracle-kafka-connector connector.class=io.debezium.connector.oracle.OracleConnector tasks.max=1 topic.prefix=prefix database.hostname=host database.port=port database.user=user database.password=password database.dbname=dbname schema.include.list=schema table.include.list=schema.table schema.history.internal=io.debezium.storage.file.history.FileSchemaHistory schema.history.internal.file.filename=path transforms=reroute,unwrap,removeWeirdChars transforms.reroute.type=io.debezium.transforms.ByLogicalTableRouter transforms.reroute.topic.regex=(.*).schema.table(.*) transforms.reroute.topic.replacement=$1 snapshot.mode=configuration_based snapshot.include.collection.list =dbname.schema.table schema.history.internal.store.only.captured.tables.ddl=true snapshot.mode.configuration.based.snapshot.schema=true snapshot.mode.configuration.based.start.stream=true tombstones.on.delete=false transforms.unwrap.type=io.debezium.transforms.ExtractNewRecordState transforms.unwrap.drop.tombstones=false transforms.unwrap.delete.handling.mode=rewrite transforms.unwrap.add.fields=table,lsn transforms.removeWeirdChars.type=transforms.RemoveWeirdCharacters
解决方案
核心问题出在Kafka Connect的类加载机制和Maven打包配置上,按以下步骤操作即可解决:
1. 修复Kafka Connect插件目录结构
Kafka Connect不会直接加载plugins根目录下的jar包,必须把每个插件放在独立的子目录中:
- 在你的Kafka Connect插件目录(比如
./kafka-connect/plugins)下新建custom-transforms子目录 - 将你打包好的
custom-transforms-1.0.0.jar移动到这个子目录 - 重启Kafka Connect服务
2. 调整Maven依赖避免冲突
你的pom.xml中Kafka Connect相关依赖用了compile scope,会导致这些依赖被打包进jar,和Connect runtime的依赖冲突,修改为provided scope(Connect runtime已自带这些依赖):
<dependency> <groupId>org.apache.kafka</groupId> <artifactId>connect-api</artifactId> <version>3.7.1</version> <scope>provided</scope> </dependency> <dependency> <groupId>org.apache.kafka</groupId> <artifactId>connect-transforms</artifactId> <version>3.7.1</version> <scope>provided</scope> </dependency> <dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>3.7.1</version> <scope>provided</scope> </dependency>
然后重新打包:
mvn clean package
把新生成的jar替换到custom-transforms子目录中,重启Connect。
3. 验证类是否正确打包
用以下命令检查jar中是否包含你的自定义类:
jar tf custom-transforms-1.0.0.jar
确认输出中存在transforms/RemoveWeirdCharacters.class,如果没有,检查Maven编译配置或类的包名是否正确。
4. 快速验证配置
用Connect的配置验证端点提前排查问题(替换为你的Connect地址和端口):
curl -X POST http://<connect-host>:<connect-port>/connector-plugins/OracleConnector/config/validate \ -H "Content-Type: application/json" \ -d '{ "name": "oracle-kafka-connector", "config": { "connector.class": "io.debezium.connector.oracle.OracleConnector", "tasks.max": "1", "topic.prefix": "prefix", "database.hostname": "host", "database.port": "port", "database.user": "user", "database.password": "password", "database.dbname": "dbname", "schema.include.list": "schema", "table.include.list": "schema.table", "schema.history.internal": "io.debezium.storage.file.history.FileSchemaHistory", "schema.history.internal.file.filename": "path", "transforms": "reroute,unwrap,removeWeirdChars", "transforms.reroute.type": "io.debezium.transforms.ByLogicalTableRouter", "transforms.reroute.topic.regex": "(.*).schema.table(.*)", "transforms.reroute.topic.replacement": "$1", "snapshot.mode": "configuration_based", "snapshot.include.collection.list": "dbname.schema.table", "schema.history.internal.store.only.captured.tables.ddl": "true", "snapshot.mode.configuration.based.snapshot.schema": "true", "snapshot.mode.configuration.based.start.stream": "true", "tombstones.on.delete": "false", "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState", "transforms.unwrap.drop.tombstones": "false", "transforms.unwrap.delete.handling.mode": "rewrite", "transforms.unwrap.add.fields": "table,lsn", "transforms.removeWeirdChars.type": "transforms.RemoveWeirdCharacters" } }'
如果返回"error_count":0,说明配置已经没问题了。
补充:你的Transform逻辑是正确的
你写的正则^[^a-zA-Z0-9]+能准确去掉消息开头的非字母数字字符(包括空格),完全符合你解决Elasticsearch展示异常的需求,这部分不需要修改。
备注:内容来源于stack exchange,提问作者Killer Queen

