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

Kafka Connect自定义Transform类无法加载报错求助

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 09:54:35