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

自定义Debezium SMT加载失败,请求排查Kafka Connect配置问题

自定义Debezium SMT加载失败排查与解决

问题现象

为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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 09:54:55