如何在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中指定版本号,需通过插件目录隔离+类加载器隔离实现多版本共存:
- 拆分插件目录
将两个版本的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文件
- 配置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会看到两个版本的插件,此时它们的类加载器完全隔离,不会互相干扰。
- 为不同版本插件设置别名(可选)
在每个插件目录下创建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.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
相关产品推荐
相关产品推荐

