自定义Debezium SMT加载失败,请求排查Kafka Connect配置问题
问题现象
为Debezium MySQL Source Connector开发了自定义Single Message Transforms(SMT),将其与MySQL连接器Jar包放入Kafka Connect Docker镜像并配置插件路径后,注册连接器时出现以下错误:
{ "error_code": 400, "message": "Connector configuration is invalid and contains the following 2 error(s):\nInvalid value com.github.gunbos.kafka.connect.smt.CustomTransformer for configuration transforms.router.type: Class com.github.gunbos.kafka.connect.smt.CustomTransformer could not be found.\nInvalid value null for configuration transforms.router.type: Not a Transformation\nYou can also find the above list of errors at the endpoint `/connector-plugins/{connectorType}/config/validate`" }
移除SMT相关配置后,MySQL Source Connector可正常运行。
相关配置
Dockerfile
FROM confluentinc/cp-kafka-connect-base:7.2.6 COPY debezium-connector-mysql/ /opt/custom-connector/debezium-connector-mysql COPY transformer/ /opt/custom-transformer/custom-event-router
docker-compose.yml
version: '2' services: zookeeper: container_name: zookeeper image: confluentinc/cp-zookeeper:7.2.6 environment: ZOOKEEPER_CLIENT_PORT: 2181 ZOOKEEPER_TICK_TIME: 2000 ports: - "2181:2181" kafka: container_name: kafka image: confluentinc/cp-kafka:7.2.6 depends_on: - zookeeper ports: - "9092:9092" environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:29092,PLAINTEXT_HOST://localhost:9092 KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 kafka_ui: image: provectuslabs/kafka-ui:latest depends_on: - kafka ports: - "8080:8080" environment: KAFKA_CLUSTERS_0_ZOOKEEPER: zookeeper:2181 KAFKA_CLUSTERS_0_NAME: local KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS: kafka:29092 mysql: image: "mysql:8.0.31" ports: - "3306:3306" env_file: - .env restart: always kafka_connect: container_name: custom-connect image: "my-custom-connect" build: dockerfile: ./Dockerfile depends_on: - kafka ports: - "8083:8083" environment: CONNECT_BOOTSTRAP_SERVERS: "kafka:29092" CONNECT_REST_PORT: "8083" CONNECT_GROUP_ID: "outbox" CONNECT_CONFIG_STORAGE_TOPIC: "outbox-config" CONNECT_CONFIG_STORAGE_REPLICATION_FACTOR: 1 CONNECT_OFFSET_STORAGE_TOPIC: "outbox-offset" CONNECT_OFFSET_STORAGE_REPLICATION_FACTOR: 1 CONNECT_STATUS_STORAGE_TOPIC: "outbox-status" CONNECT_STATUS_STORAGE_REPLICATION_FACTOR: 1 CONNECT_KEY_CONVERTER: "org.apache.kafka.connect.json.JsonConverter" CONNECT_KEY_CONVERTER_SCHEMAS_ENABLE: "false" CONNECT_VALUE_CONVERTER: "org.apache.kafka.connect.json.JsonConverter" CONNECT_VALUE_CONVERTER_SCHEMAS_ENABLE: "false" CONNECT_INTERNAL_KEY_CONVERTER: "org.apache.kafka.connect.json.JsonConverter" CONNECT_INTERNAL_VALUE_CONVERTER: "org.apache.kafka.connect.json.JsonConverter" CONNECT_REST_ADVERTISED_HOST_NAME: "localhost" CONNECT_PLUGIN_PATH: "/usr/share/java, /opt/custom-connector/, /opt/custom-transformer/"
连接器注册配置
{ "name": "outbox-connector", "config": { "connector.class": "io.debezium.connector.mysql.MySqlConnector", "tasks.max": "1", "database.hostname": "mysql", "database.port": "3306", "database.user": "cain", "database.password": "cain", "topic.prefix": "test", "database.server.id": "123456", "database.include.list": "kafka_connect", "table.include.list": "kafka_connect.outbox", "topic.creation.default.replication.factor" : 1, "topic.creation.default.partitions" : 1, "schema.history.internal.kafka.topic": "schemaHistory.outbox", "schema.history.internal.kafka.bootstrap.servers": "kafka:29092", "transforms": "router", "transforms.router.type": "com.github.gunbos.kafka.connect.smt.CustomTransformer" } }
排查与解决步骤
1. 修正插件目录结构与Jar包复制方式
Kafka Connect要求每个插件(包括SMT)必须放在独立目录下,且目录直接包含编译好的Jar包,不能嵌套源码目录或多层子目录。当前Dockerfile复制的是源码目录transformer/到子目录,导致Connect无法识别SMT类。
修改Dockerfile:先在本地执行mvn package编译自定义SMT项目生成Jar包,再将Jar包直接复制到插件路径目录:
FROM confluentinc/cp-kafka-connect-base:7.2.6 # 复制Debezium连接器(保持原有结构) COPY debezium-connector-mysql/ /opt/custom-connector/debezium-connector-mysql # 复制编译好的自定义SMT Jar包到插件目录 COPY custom-transformer/target/custom-transformer-1.0-SNAPSHOT.jar /opt/custom-transformer/
2. 验证插件加载状态
启动容器后,调用Kafka Connect的插件列表接口,确认自定义SMT是否被加载:
curl http://localhost:8083/connector-plugins
若返回结果包含com.github.gunbos.kafka.connect.smt.CustomTransformer,说明插件加载成功;若未出现,需检查Jar包是否正确复制、插件路径配置是否正确。
3. 确认依赖配置正确性
自定义SMT需依赖Kafka Connect API,确保pom.xml中依赖配置正确,使用provided scope避免与Connect镜像中的依赖冲突:
<dependency> <groupId>org.apache.kafka</groupId> <artifactId>connect-api</artifactId> <version>7.2.6</version> <scope>provided</scope> </dependency>
4. 核对类名与包名
确认代码中的包名com.github.gunbos.kafka.connect.smt和类名CustomTransformer与连接器配置中的transforms.router.type完全一致,避免拼写错误。
内容的提问来源于stack exchange,提问作者gun bos

