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

如何在Kafka Connect中指定插件版本并兼容新旧版本运行?

问题描述

我的Kafka Connect Worker的plugin.path中存在两个版本的Debezium PostgreSQL连接器,调用GET /connector-plugins返回如下结果:

[
    {
        "class": "io.debezium.connector.postgresql.PostgresConnector",
        "type": "source",
        "version": "1.9.6.Final"
    },
    {
        "class": "io.debezium.connector.postgresql.PostgresConnector",
        "type": "source",
        "version": "2.5.1.Final"
    }
]

需要在连接器配置的connector.class中指定插件版本,让依赖1.9.6版本的旧连接器和使用2.5.1版本的新连接器在同一Worker中同时正常运行。若不指定版本,旧连接器会出现报错:

{
    "old-connector": {
        "status": {
            "name": "old-connector",
            "connector": {
                "state": "RUNNING",
                "worker_id": "1.2.3.4:8083"
            },
            "tasks": [
                {
                    "id": 0,
                    "state": "FAILED",
                    "worker_id": "1.2.3.4:8083",
                    "trace": "org.apache.kafka.connect.errors.ConnectException: Error configuring an instance of PostgresConnectorTask; check the logs for details\n\tat io.debezium.connector.common.BaseSourceTask.start(BaseSourceTask.java:131)\n\tat org.apache.kafka.connect.runtime.WorkerSourceTask.initializeAndStart(WorkerSourceTask.java:226)\n\tat org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:186)\n\tat org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:243)\n\tat java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515)\n\tat java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)\n\tat java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)\n\tat java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)\n\tat java.base/java.lang.Thread.run(Thread.java:829)\n"
                }
            ],
            "type": "source"
        }
    }
}
解决方法

Kafka Connect本身不支持直接在connector.class中指定版本号,需通过插件目录隔离+类加载器隔离实现多版本共存:

  1. 拆分插件目录
    将两个版本的Debezium PostgreSQL连接器分别放在独立目录中,确保每个目录仅包含对应版本的所有依赖jar包:
# 示例目录结构
/kafka/plugins/
  ├── debezium-postgres-1.9.6/
  │   ├── debezium-connector-postgres-1.9.6.Final.jar
  │   └── 所有1.9.6版本的依赖jar文件
  └── debezium-postgres-2.5.1/
      ├── debezium-connector-postgres-2.5.1.Final.jar
      └── 所有2.5.1版本的依赖jar文件
  1. 配置Worker的插件隔离
    修改Kafka Connect Worker的配置文件,更新plugin.path包含两个独立目录,并开启插件级类加载隔离:
# 配置插件路径,包含两个版本的目录
plugin.path=/kafka/plugins/debezium-postgres-1.9.6,/kafka/plugins/debezium-postgres-2.5.1
# 开启插件级类加载隔离,每个插件目录使用独立的类加载器
plugin.isolation=plugin

重启Worker后,再次调用GET /connector-plugins会看到两个版本的插件,此时它们的类加载器完全隔离,不会互相干扰。

  1. 为不同版本插件设置别名(可选)
    在每个插件目录下创建plugin.properties文件,为不同版本的插件设置唯一别名,方便后续配置区分:
  • 在debezium-postgres-1.9.6目录下创建plugin.properties:
    plugin.name=debezium-postgres-v1
    
  • 在debezium-postgres-2.5.1目录下创建plugin.properties:
    plugin.name=debezium-postgres-v2
    
  1. 指定连接器使用的版本
  • 旧连接器(依赖1.9.6版本):保持原有connector.class配置即可,重启Worker后会自动关联到对应版本的插件目录,无需额外修改。
  • 新连接器(使用2.5.1版本):在配置中指定connector.class为类名即可,Worker会默认使用最新版本的插件;若设置了别名,也可通过plugin.name明确指定:
    {
        "name": "new-connector",
        "config": {
            "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
            // 可选:若设置了别名,添加此参数指定版本
            "plugin.name": "debezium-postgres-v2",
            // 其他连接器配置...
        }
    }
    
关键说明
  • plugin.isolation=plugin是实现多版本共存的核心配置,它让每个插件目录拥有独立的类加载器,避免不同版本的类互相冲突。
  • 拆分目录后,旧连接器无需修改配置即可继续使用原有版本,新连接器默认优先使用最新版本插件。
  • 确保每个插件目录的jar包完整,不要混合不同版本的依赖,否则会引发类加载异常。

内容的提问来源于stack exchange,提问作者user22659171

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 20:00:18