Strimzi Kafka Connect 构建与插件管理深度解读:ConnectBuilderST 系统测试全解析
2026/9/17 5:37:23 网站建设 项目流程

Strimzi Kafka Connect 构建与插件管理深度解读:ConnectBuilderST 系统测试全解析

【免费下载链接】strimzi-kafka-operatorApache Kafka® running on Kubernetes项目地址: https://gitcode.com/GitHub_Trending/st/strimzi-kafka-operator

ConnectBuilderST 是 Strimzi Kafka Operator 仓库中专门针对Kafka Connect 镜像构建(Build)与连接器插件管理的系统测试套件,覆盖了从 Maven 坐标、Jar/Tgz/Zip 制品、校验和校验、镜像卷挂载到 OpenShift ImageStream 推送等完整构建场景。本文以该测试套件的 7 个用例为骨架,结合 ConnectBuilderST.java 源码、build 相关 API 模型 与真实配置示例,帮助读者掌握 Kafka Connect 自定义镜像构建的配置方式、失败恢复机制与端到端验证方法。

一、测试套件概览:验证什么、如何组织

根据 ConnectBuilderST.md 的描述,该套件的核心主题是"Testing Kafka Connect build and plugin management"(测试 Kafka Connect 构建与插件管理)。它属于系统测试(System Test,简称 ST)体系,运行在真实或近真实的 Kubernetes/OpenShift 集群之上,通过KubeResourceManager创建真实的 CRD 资源并等待其收敛。

从源码(ConnectBuilderST.java)可以看到套件级别的注解定义:

@Tag(REGRESSION) @Tag(CONNECT_COMPONENTS) @Tag(CONNECT) @MicroShiftNotSupported @SuiteDoc( description = @Desc("Testing Kafka Connect build and plugin management."), labels = { @Label(value = TestDocsLabels.CONNECT) } ) class ConnectBuilderST extends AbstractST { ... }
  • 测试标签:标记为REGRESSION(回归)、CONNECT_COMPONENTS(Connect 组件)、CONNECT(Connect 领域),其中testBuildPluginUsingMavenCoordinatesArtifacts额外带有SANITYACCEPTANCE标签,属于冒烟/验收级别的关键用例。
  • 标签索引connect标签的说明文档见 labels/connect.md,它概括了 Connect 测试的目标——"确保 Kafka Connect 与外部系统之间通过连接器实现可靠集成",本套件聚焦其中的"插件管理、构建流程"子集。
  • 前置环境@BeforeAll中会安装带CO_OPERATION_TIMEOUT_SHORT(短操作超时)配置的 Cluster Operator,并创建一个 3 broker + 3 controller 的 KRaft Kafka 集群(ConnectBuilderST.java)。
  • 并行执行:所有用例标注@ParallelTest,可并行运行;testMountPluginUsingImageVolume要求 Kubernetes API 版本 ≥ 1.35(Kubernetes Image Volume 特性),testPushIntoImageStream标注@OpenShiftOnly仅在 OpenShift 上运行。

二、构建模型基础:spec.build的结构与五种制品类型

要理解这些测试,先要理解 Strimzi 的 Kafka Connect 构建模型。KafkaConnectspec.build字段由 Build.java 建模,包含三个核心字段:

| 字段 | 是否必填 | 说明 | | - | - | - | |output| 是 | 新构建镜像的存储位置(Docker 仓库或 OpenShift ImageStream) | |plugins| 是 | 要打进镜像的连接器插件列表 | |resources| 否 | 构建 Pod 预留的 CPU 与内存资源 |

plugins中的每个插件由 Plugin.java 建模,包含nameartifacts两个必填字段。插件名遵循正则^[a-z0-9][-_a-z0-9]*[a-z0-9]$,且在同一个 KafkaConnect 资源内必须唯一;它会用于生成插件在容器内的存储路径(测试中可见plugins/plugin-with-other-type/*这样的路径)。

artifacts是多态类型,由 Artifact.java 定义,通过type字段区分,目前支持5 种制品类型

| type | 制品类 | 说明 | | - | - | - | |jar| JarArtifact | 单个 JAR 文件 | |tgz| TgzArtifact | tar.gz 压缩包 | |zip| ZipArtifact | zip 压缩包 | |maven| MavenArtifact | 从 Maven 仓库拉取坐标制品 | |other| OtherArtifact | 其他任意类型,可指定落盘文件名 |

其中jar/tgz/zip/other继承自 DownloadableArtifact.java,公共字段包括:

  • url(必填):制品下载地址,支持httphttpsftp协议;
  • sha512sum(可选):制品 SHA-512 校验和,指定后构建时会校验,不指定则不校验。模型注释特别提醒:Strimzi 不对下载制品做安全扫描,生产环境应先在本地人工验证制品并配置校验和;
  • insecure(可选):置为true时跳过所有 TLS 校验,允许从不可信证书的服务器下载。

测试中使用的 EchoSink 制品常量定义在 TestConstants.java:例如 EchoSink 连接器类为cz.scholz.kafka.connect.echosink.EchoSinkConnector,JAR 下载地址为https://github.com/scholzj/echo-sink/releases/download/1.6.0/echo-sink-1.6.0.jar,并带有固定的 SHA-512 校验和。

一个可直接参考的完整 YAML 示例位于 examples/connect/kafka-connect-build.yaml:

spec: build: output: type: docker image: ttl.sh/strimzi-connect-example-4.3.0:24h plugins: - name: kafka-connect-file artifacts: - type: maven group: org.apache.kafka artifact: connect-file version: 4.3.1

此外,系统测试使用的基础模板 connect-build-template.yaml 展示了构建输出与构建容器的可调点,包括build.output.additionalBuildOptions(如--log-format=json)以及template.buildContainer.securityContext(如runAsUser: 1000并授予SETUID/SETGID/DAC_OVERRIDE/SYS_ADMIN能力)。注意 Strimzi 在推送镜像时使用 tag,但在拉取时使用 digest,确保拉取到的是本次构建的确切镜像(见示例文件注释)。

三、用例详解一:校验和错误导致构建失败与自动恢复

用例:testBuildFailsWithWrongChecksumOfArtifact

该用例验证"Kafka Connect 构建在制品校验和错误时失败,并在更正校验和后恢复"。这正是sha512sum字段在真实构建链路中生效的证明。源码见 ConnectBuilderST.java,流程如下:

| 步骤 | 动作 | 预期结果 | | - | - | - | | 1 | 初始化 TestStorage 并获取测试镜像名 | TestStorage 实例创建,镜像名就绪 | | 2 | 用错误校验和创建 Plugin,并以此构建 KafkaConnect 资源 | 资源创建成功,但构建因校验和错误而失败 | | 3 | 部署 Scraper pod(带特定配置) | Scraper pod 成功部署 | | 4 | 等待 Kafka Connect 状态指示构建失败 | 状态包含构建失败消息 | | 5 | 为 Kafka Connect 部署网络策略 | 网络策略部署成功 | | 6 | 将校验和替换为正确值并更新资源 | 资源以正确校验和更新 | | 7 | 等待 Kafka Connect 就绪 | Kafka Connect 变为 Ready | | 8 | 通过 Kafka Connect API 验证 EchoSink 连接器可用 | API 返回 EchoSink 连接器 | | 9 | 验证 EchoSink 连接器出现在资源 status 中 | status 中列出 EchoSink 连接器 |

实现要点(源码级)

  • 测试先构造带ECHO_SINK_JAR_WRONG_CHECKSUMJarArtifact构建资源,并用createResourceWithoutWait提交(不等待就绪),随后调用KafkaConnectUtils.waitForConnectNotReady(...)waitUntilKafkaConnectStatusConditionContainsMessage(..., "The Kafka Connect build failed(.*)?")断言失败状态;
  • 通过kafkaConnect.getStatus().getConditions()断言状态条件消息匹配The Kafka Connect build failed(.*)?,且条件类型为NotReady
  • 恢复手段是调用KafkaConnectUtils.replace(...)移除列表中的错误插件并追加校验和正确的插件,之后waitForConnectReady(...)等待恢复;
  • 最终验证分两层:一是通过 Scraper pod 内curl http://<connect-service>:8083/connector-plugins查询 Kafka Connect REST API,断言响应包含 EchoSink 类名;二是读取KafkaConnect.status.connectorPlugins,断言其中的connectorClass包含 EchoSink 类名。这两层验证分别对应"运行时 API 可见"与"资源状态可见"。

四、用例详解二:Jar + Tgz + Zip 混合制品构建与消息收发验证

用例:testBuildWithJarTgzAndZip

该用例验证"混合 jar、tar.gz、zip 三种制品的 Kafka Connect 镜像构建,并验证消息发送-接收功能"(ConnectBuilderST.java)。测试同时覆盖了Docker 输出(push into Docker output)路径。

| 步骤 | 动作 | 预期结果 | | - | - | - | | 1 | 创建 TestStorage 对象 | 实例创建并携带上下文 | | 2 | 获取测试用例镜像名 | 镜像名获取成功 | | 3 | 创建 Kafka Topic 资源 | 资源创建并等待就绪 | | 4 | 创建 Kafka Connect 资源 | 资源创建并等待就绪 | | 5 | 配置 Kafka Connector | 连接器配置并创建完成 | | 6 | 验证 Connector 类名 | 与ECHO_SINK_CLASS_NAME一致 | | 7 | 创建 Kafka 客户端并发送消息 | 消息发送并验证成功 | | 8 | 检查日志中的消息 | 日志包含预期的接收消息 |

实现要点

  • 该测试一次性声明两个插件(ConnectBuilderST.java):connector-with-tar-and-jar由一个JarArtifact(EchoSink JAR)和一个TgzArtifact(EchoSink 源码 tar.gz)组成;connector-from-zip由一个ZipArtifact(Camel HTTP connector 的-package.zip,版本 0.7.0)组成,均带sha512sum
  • 构建资源上同时配置了StringConverter与关闭 schema 的 connect 配置(key.converter/value.converter均为org.apache.kafka.connect.storage.StringConverterschemas.enable=false),以及 inline logging 将 rootLogger 设为 INFO;
  • 验证手段:通过KafkaProducerClient的 Job 发送testStorage.getMessageCount()条消息,ClientUtils.waitForClientSuccess等待发送成功,最后PodUtils.waitUntilMessageIsInPodLogs(...)在 Connect Pod 日志中查找Received message with key 'null' and value 'Hello world - 99',证明 EchoSink 连接器真正消费到了消息。

五、用例详解三:Maven 坐标制品构建(SANITY/ACCEPTANCE 级)

用例:testBuildPluginUsingMavenCoordinatesArtifacts

这是本套件中唯一同时带SANITYACCEPTANCE标签的用例,验证"使用 Maven 坐标制品构建插件"(ConnectBuilderST.java)。

| 步骤 | 动作 | 预期结果 | | - | - | - | | 1 | 创建测试存储对象 | 对象创建成功 | | 2 | 生成测试用例镜像名 | 镜像名生成成功 | | 3 | 创建 Kafka Topic 与带 mvn 坐标插件配置的 Kafka Connect 资源 | 资源创建并可用 | | 4 | 配置并部署 Kafka Connector | 连接器以正确配置部署 | | 5 | 创建 Kafka consumer 并开始消费 | 消费者开始消费 | | 6 | 验证消费者收到消息 | 消费者收到预期消息 |

实现要点

  • 插件构建使用MavenArtifactBuilder(ConnectBuilderST.java):group=org.apache.camel.kafkaconnectorartifact=camel-timer-kafka-connectorversion=0.9.0;若环境变量ST_MAVEN_MIRROR_URL存在,还会追加MavenMirror指向镜像仓库;
  • 用例开头有一条针对Kind 集群 + 未启用 Buildah环境的assumeFalse跳过假设(源码注释说明该假设可在 Buildah 进入 GA 后移除),说明 Maven 制品构建依赖特定构建器配置;
  • 连接器使用 Camel Timer Source 连接器(org.apache.camel.kafkaconnector.timer.CamelTimerSourceConnector),配置camel.source.path.timerName
  • 验证手段与上一用例相反——用KafkaConsumerClient的 Job 消费消息,ClientUtils.waitForClientSuccess确认消费者收到testStorage.getMessageCount()条消息。

关于 Maven 制品的字段,MavenArtifact.java 给出了完整定义:groupartifactversion为坐标三元组;repository默认https://repo1.maven.org/maven2/mirrors会将构建期间(含插件仓库与 Maven Central)的所有仓库请求重定向到配置的镜像;insecure可关闭 TLS 校验;includeScope可取compile/provided/runtime/test/system,未配置时包含所有依赖。

六、用例详解四:other类型制品的文件名与哈希命名行为

用例:testBuildOtherPluginTypeWithAndWithoutFileName

该用例验证"不同插件类型在有/无文件名时的行为"(ConnectBuilderST.java)。

| 步骤 | 动作 | 预期结果 | | - | - | - | | 1 | 初始化测试存储与 topic | 命名空间与 topic 创建成功 | | 2 | 以指定插件与构建配置创建 Kafka Connect | 部署并配置正确 | | 3 | 快照当前 Connect Pod 并验证插件文件名 | 文件名与预期一致 | | 4 | 改用无文件名的插件并触发滚动更新 | 无文件名插件更新成功 | | 5 | 基于插件哈希验证更新后的文件名 | 文件名与之前不同且匹配哈希 |

实现要点

  • OtherArtifact独有的fileName字段(见 OtherArtifact.java)指定制品落盘名称。测试第一阶段构造fileName=echo-sink-test.jar(即TestConstants.ECHO_SINK_FILE_NAME),并用辅助方法getPluginFileNameFromConnectPod在 Connect Pod 内执行ls plugins/plugin-with-other-type/*断言文件名一致;
  • 第二阶段通过KafkaConnectUtils.replace(...)移除fileName后触发滚动更新(RollingUpdateUtils.waitTillComponentHasRolledAndPodsReady对比 Pod 快照),此时断言文件名变为Util.hashStub(ECHO_SINK_JAR_URL)——即未指定文件名时,Strimzi 使用制品 URL 的哈希作为落盘文件名,保证不同来源制品不会互相覆盖。

七、用例详解五:通过 Kubernetes Image Volume 挂载 OCI 制品插件

用例:testMountPluginUsingImageVolume

该用例验证"通过 Kubernetes Image Volume 从 OCI 制品挂载 Kafka Connect 插件,并验证消息收发功能"(ConnectBuilderST.java)。它要求 Kubernetes API 版本 ≥ 1.35,属于较新的镜像卷特性。

| 步骤 | 动作 | 预期结果 | | - | - | - | | 1 | 创建 TestStorage 对象 | 实例创建并携带上下文 | | 2 | 创建 Kafka Topic 资源 | Topic 创建成功 | | 3 | 创建带 Image Volume 挂载 EchoSink 插件的 Kafka Connect | 资源创建并使用挂载的插件 | | 4 | 创建 Kafka Connector | 连接器创建成功 | | 5 | 验证 Connector 类名 | 与ECHO_SINK_CLASS_NAME一致 | | 6 | 创建 Kafka 客户端并发送消息 | 消息发送并验证 | | 7 | 检查日志中的接收消息 | 日志包含预期消息 |

实现要点

  • 该用例与构建用例不同,不走spec.build流程,而是走MountedPlugin(见 MountedPlugin.java):name=connector-from-image-volume,制品为ImageArtifact(见 ImageArtifact.java),reference=ghcr.io/scholzj/echo-sink:latest
  • ImageArtifact支持pullPolicyAlways(始终拉取,失败则容器创建失败)、Never(仅用本地镜像)、IfNotPresent(本地无则拉取);默认对:latest标签为Always,否则IfNotPresent
  • 数据面验证与混合制品用例一致:Producer Job 发送消息 →waitForClientSuccess→ Connect Pod 日志中出现Received message ... 'Hello world - 99',证明以卷挂载方式提供的插件同样可被 Kafka Connect 加载并执行。

八、用例详解六:构建产物推送到 OpenShift ImageStream

用例:testPushIntoImageStream(仅 OpenShift)

该用例验证"KafkaConnect 构建成功推送到 OpenShift ImageStream"(ConnectBuilderST.java)。

| 步骤 | 动作 | 预期结果 | | - | - | - | | 1 | 初始化测试存储 | 以测试上下文初始化 | | 2 | 创建 ImageStream | 在指定命名空间创建成功 | | 3 | 部署带 ImageStream 输出的 KafkaConnect | 以预期构建配置部署 | | 4 | 验证构建制品与状态 | 两个插件、使用 ImageStream 输出、状态为 Ready |

实现要点

  • 测试先用ImageStreamBuilder创建名为custom-image-stream的 ImageStream,再通过 OpenShiftClient 提交;
  • spec.build.output使用ImageStreamOutput(见 ImageStreamOutput.java):type: imagestreamimage: custom-image-stream:latest
  • 插件复用PLUGIN_WITH_TAR_AND_JAR(2 个制品),断言包括:spec.build.plugins[0].artifacts.size()==2output.type=="imagestream"output.image=="custom-image-stream:latest"、status 条件类型为Readystatus.connectorPlugins非空且包含 EchoSink 类名。

作为对照,Docker 输出类型由 DockerOutput.java 建模,字段包括image(完整镜像名,如quay.io/my-organization/my-custom-connect:latest)、pushSecret(推送凭证 Secret)、additionalBuildOptionsadditionalPushOptions(分别透传给 Kaniko/Buildah 的buildpush命令,仅 Kubernetes 平台生效、OpenShift 忽略,且变更这些字段不会触发镜像重建)。允许的选项在白名单中明确列出,例如 Kaniko 的--insecure--log-format--reproducible等,Buildah 的--authfile--creds--tls-verify等;相关校验逻辑位于 KafkaConnectBuild.java。

九、用例详解七:向已有 Connect 追加插件与滚动更新

用例:testUpdateConnectWithAnotherPlugin

该用例验证"用另一个插件更新 Kafka Connect 并验证"(ConnectBuilderST.java),覆盖了"先构建一个插件、运行、再动态追加第二个插件"的运行时扩展场景。

| 步骤 | 动作 | 预期结果 | | - | - | - | | 1 | 创建 TestStorage 实例 | 实例创建成功 | | 2 | 生成随机 topic 名并创建 Kafka topic | Topic 创建成功 | | 3 | 为 KafkaConnect 部署网络策略 | 网络策略部署成功 | | 4 | 创建 EchoSink KafkaConnector | 创建并验证成功 | | 5 | 向 Kafka Connect 添加第二个插件并滚动更新 | 第二插件添加并完成滚动更新 | | 6 | 创建 Camel-HTTP-Sink KafkaConnector | 创建并验证成功 | | 7 | 验证两个连接器与插件同时存在 | 均验证成功 |

实现要点

  • 初始插件为PLUGIN_WITH_TAR_AND_JAR;第二阶段通过KafkaConnectUtils.replace(...)spec.build.plugins追加Camel HTTP Sink 插件(TgzArtifactcamel-http-kafka-connector-0.7.0-package.tar.gz,校验和d0bb8c...51b5);
  • 追加前先用PodUtils.podSnapshot记录 Pod 状态,追加后RollingUpdateUtils.waitTillComponentHasRolledAndPodsReady等待滚动更新完成——构建配置变更会触发 Connect 滚动更新以加载新镜像
  • 验证包含三层:追加前通过curl .../connector-plugins断言存在 EchoSink 类名且不存在Camel HTTP Sink 类名;追加后创建 Camel-HTTP-Sink 连接器;最后断言spec.build.plugins.size()==2status.connectorPlugins同时包含两个连接器类名。

十、测试工程实践:镜像命名、基础设施与运行前提

从套件源码可以总结出几条可复用的系统测试工程实践:

  1. 隔离的镜像命名getImageNameForTestCase()(ConnectBuilderST.java)用随机数生成strimzi-sts-connect-build镜像的 registry 输出地址,避免并行测试间镜像相互覆盖;
  2. TestStorage 上下文:每个用例创建独立的TestStorage(命名空间、cluster 名、topic 名、scraper 名等均从测试上下文派生),配合KubeResourceManager统一管理资源生命周期与清理;
  3. 构建失败的可观测性:套件反复使用KafkaConnectUtils.waitForConnectNotReadywaitUntilKafkaConnectStatusConditionContainsMessagewaitForConnectStatusContainsPlugins等工具,验证构建状态条件(NotReady/Ready)与connectorPlugins状态字段,这也是排查线上构建问题的关键入口;
  4. 网络策略联动:多个用例调用NetworkPolicyUtils.deployNetworkPolicyForResource为 KafkaConnect 部署 NetworkPolicy,确保启用网络策略的集群中 Scraper 仍能访问 Connect REST API;
  5. 运行前提:需要在已就绪的 Kubernetes(≥ 1.35 才能跑 Image Volume 用例)或 OpenShift 集群上先安装 Cluster Operator 并准备 KRaft Kafka 集群;Maven 制品用例在 Kind 且未启用 Buildah 时会自动跳过。

综上,ConnectBuilderST 完整覆盖了 Kafka Connect 自定义镜像构建从"制品获取(URL/校验和/Maven 坐标/镜像卷)→ 构建输出(Docker/ImageStream)→ 失败恢复 → 运行时插件扩展"的全链路,是理解 Strimzi Kafka Connect Build 特性行为与验证配置正确性的最佳参考起点。配合 connect-build-template.yaml 与 kafka-connect-build.yaml 两个模板,开发者可以快速复刻同样的构建配置到自己的生产集群。

【免费下载链接】strimzi-kafka-operatorApache Kafka® running on Kubernetes项目地址: https://gitcode.com/GitHub_Trending/st/strimzi-kafka-operator

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询