关于Strimzi Kafka Connect连接器插件版本的疑问
场景背景
场景一:未指定镜像的Kafka Connect实例
创建未指定自定义镜像的Strimzi Kafka Connect实例,配置如下:
apiVersion: kafka.strimzi.io/v1beta2 kind: KafkaConnect metadata: name: my-connect-cluster-connect-build namespace: first-kafka-namespace labels: my-connect: my-connect-cluster-connect-build annotations: strimzi.io/use-connector-resources: "true" spec: version: 3.8.0 replicas: 1 bootstrapServers: "first-kafka-cluster-kafka-bootstrap:9092" config: group.id: connect-cluster offset.storage.topic: connect-cluster-offsets config.storage.topic: connect-cluster-configs status.storage.topic: connect-cluster-status config.storage.replication.factor: 1 offset.storage.replication.factor: 1 status.storage.replication.factor: 1
执行查询命令:
curl -s http://localhost:8083/connector-plugins | jq .
返回结果仅包含3个Mirror系列连接器,未显示默认的FileStream连接器:
[ { "class": "org.apache.kafka.connect.mirror.MirrorCheckpointConnector", "type": "source", "version": "3.8.0" }, { "class": "org.apache.kafka.connect.mirror.MirrorHeartbeatConnector", "type": "source", "version": "3.8.0" }, { "class": "org.apache.kafka.connect.mirror.MirrorSourceConnector", "type": "source", "version": "3.8.0" } ]
场景二:自定义镜像构建后的实例
先创建用于构建镜像的Kafka Connect实例,配置如下:
apiVersion: kafka.strimzi.io/v1beta2 kind: KafkaConnect metadata: name: kafka-connect-cluster namespace: kafka annotations: strimzi.io/use-connector-resources: "true" spec: version: 3.9.0 replicas: 1 bootstrapServers: "kafka-cluster-kafka-bootstrap:9092" config: group.id: connect-cluster offset.storage.topic: connect-cluster-offsets config.storage.topic: connect-cluster-configs status.storage.topic: connect-cluster-status config.storage.replication.factor: 1 offset.storage.replication.factor: 1 status.storage.replication.factor: 1 build: output: type: docker image: <GHCR Registry> pushSecret: ghcr-push-secret plugins: - name: file-source-connector artifacts: - type: jar url: https://repo1.maven.org/maven2/org/apache/kafka/connect-file/3.6.0/connect-file-3.6.0.jar
镜像成功推送后,使用该镜像创建测试实例:
apiVersion: kafka.strimzi.io/v1beta2 kind: KafkaConnect metadata: name: my-connect-cluster-connect-build namespace: first-kafka-namespace labels: my-connect: my-connect-cluster-connect-build annotations: strimzi.io/use-connector-resources: "true" spec: image: <image from GHCR> replicas: 1 bootstrapServers: "first-kafka-cluster-kafka-bootstrap:9092" config: group.id: connect-cluster offset.storage.topic: connect-cluster-offsets config.storage.topic: connect-cluster-configs status.storage.topic: connect-cluster-status config.storage.replication.factor: 1 offset.storage.replication.factor: 1 status.storage.replication.factor: 1 template: pod: imagePullSecrets: - name: ghcr-pull-secret
执行查询后返回结果显示FileStream连接器版本为3.9.0(而非预期的3.6.0),且该连接器正常显示:
[ { "class": "org.apache.kafka.connect.file.FileStreamSinkConnector", "type": "sink", "version": "3.9.0" }, { "class": "org.apache.kafka.connect.file.FileStreamSourceConnector", "type": "source", "version": "3.9.0" }, { "class": "org.apache.kafka.connect.mirror.MirrorCheckpointConnector", "type": "source", "version": "3.9.0" }, { "class": "org.apache.kafka.connect.mirror.MirrorHeartbeatConnector", "type": "source", "version": "3.9.0" }, { "class": "org.apache.kafka.connect.mirror.MirrorSourceConnector", "type": "source", "version": "3.9.0" } ]
疑问解答
1. 为何指定的3.6.0版本FileStream连接器被覆盖为3.9.0?
Strimzi的镜像构建基于你指定的spec.version对应的官方基础镜像,该镜像已内置与Kafka版本匹配的核心连接器组件(包括FileStream),这些组件存放在Kafka的核心类路径目录中,类加载优先级高于你手动添加的自定义插件jar。
当你添加3.6.0版本的connect-file jar时,Kafka Connect会优先加载基础镜像中自带的3.9.0版本组件,导致你手动添加的低版本jar无法生效。另外,FileStream属于Kafka Connect核心内置连接器,无需手动添加,基础镜像已默认包含。
2. 为何首次创建的实例中未列出默认的FileStream连接器?
Strimzi官方默认的Kafka Connect镜像,其plugin.path配置仅指向/opt/kafka/plugins目录,而FileStream等核心连接器的jar存放在/opt/kafka/libs目录,不在默认的插件扫描路径内。
Kafka Connect仅会扫描plugin.path配置的目录来识别可用连接器,因此即使镜像中存在该组件,也不会被检测到。而使用build功能构建镜像时,Strimzi会自动调整插件路径配置,将核心类路径目录纳入扫描范围,因此FileStream连接器会被正常识别并显示。
内容的提问来源于stack exchange,提问作者Longclaw

