
Strimzi Kafka Connect 构建与插件管理深度解读ConnectBuilderST 系统测试全解析【免费下载链接】strimzi-kafka-operatorApache Kafka® running on Kubernetes项目地址: https://gitcode.com/GitHub_Trending/st/strimzi-kafka-operatorConnectBuilderST 是 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_COMPONENTSConnect 组件、CONNECTConnect 领域其中testBuildPluginUsingMavenCoordinatesArtifacts额外带有SANITY与ACCEPTANCE标签属于冒烟/验收级别的关键用例。标签索引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.35Kubernetes Image Volume 特性testPushIntoImageStream标注OpenShiftOnly仅在 OpenShift 上运行。二、构建模型基础spec.build的结构与五种制品类型要理解这些测试先要理解 Strimzi 的 Kafka Connect 构建模型。KafkaConnect的spec.build字段由 Build.java 建模包含三个核心字段| 字段 | 是否必填 | 说明 | | - | - | - | |output| 是 | 新构建镜像的存储位置Docker 仓库或 OpenShift ImageStream | |plugins| 是 | 要打进镜像的连接器插件列表 | |resources| 否 | 构建 Pod 预留的 CPU 与内存资源 |plugins中的每个插件由 Plugin.java 建模包含name与artifacts两个必填字段。插件名遵循正则^[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必填制品下载地址支持http、https、ftp协议sha512sum可选制品 SHA-512 校验和指定后构建时会校验不指定则不校验。模型注释特别提醒Strimzi 不对下载制品做安全扫描生产环境应先在本地人工验证制品并配置校验和insecure可选置为true时跳过所有 TLS 校验允许从不可信证书的服务器下载。测试中使用的 EchoSink 制品常量定义在 TestConstants.java例如 EchoSink 连接器类为cz.scholz.kafka.connect.echosink.EchoSinkConnectorJAR 下载地址为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.yamlspec: 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-formatjson以及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_CHECKSUM的JarArtifact构建资源并用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.javaconnector-with-tar-and-jar由一个JarArtifactEchoSink JAR和一个TgzArtifactEchoSink 源码 tar.gz组成connector-from-zip由一个ZipArtifactCamel HTTP connector 的-package.zip版本 0.7.0组成均带sha512sum构建资源上同时配置了StringConverter与关闭 schema 的 connect 配置key.converter/value.converter均为org.apache.kafka.connect.storage.StringConverterschemas.enablefalse以及 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这是本套件中唯一同时带SANITY与ACCEPTANCE标签的用例验证使用 Maven 坐标制品构建插件ConnectBuilderST.java。| 步骤 | 动作 | 预期结果 | | - | - | - | | 1 | 创建测试存储对象 | 对象创建成功 | | 2 | 生成测试用例镜像名 | 镜像名生成成功 | | 3 | 创建 Kafka Topic 与带 mvn 坐标插件配置的 Kafka Connect 资源 | 资源创建并可用 | | 4 | 配置并部署 Kafka Connector | 连接器以正确配置部署 | | 5 | 创建 Kafka consumer 并开始消费 | 消费者开始消费 | | 6 | 验证消费者收到消息 | 消费者收到预期消息 |实现要点插件构建使用MavenArtifactBuilderConnectBuilderST.javagrouporg.apache.camel.kafkaconnector、artifactcamel-timer-kafka-connector、version0.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 给出了完整定义group、artifact、version为坐标三元组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指定制品落盘名称。测试第一阶段构造fileNameecho-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.javanameconnector-from-image-volume制品为ImageArtifact见 ImageArtifact.javareferenceghcr.io/scholzj/echo-sink:latestImageArtifact支持pullPolicyAlways始终拉取失败则容器创建失败、Never仅用本地镜像、IfNotPresent本地无则拉取默认对:latest标签为Always否则IfNotPresent数据面验证与混合制品用例一致Producer Job 发送消息 →waitForClientSuccess→ Connect Pod 日志中出现Received message ... Hello world - 99证明以卷挂载方式提供的插件同样可被 Kafka Connect 加载并执行。八、用例详解六构建产物推送到 OpenShift ImageStream用例testPushIntoImageStream仅 OpenShift该用例验证KafkaConnect 构建成功推送到 OpenShift ImageStreamConnectBuilderST.java。| 步骤 | 动作 | 预期结果 | | - | - | - | | 1 | 初始化测试存储 | 以测试上下文初始化 | | 2 | 创建 ImageStream | 在指定命名空间创建成功 | | 3 | 部署带 ImageStream 输出的 KafkaConnect | 以预期构建配置部署 | | 4 | 验证构建制品与状态 | 两个插件、使用 ImageStream 输出、状态为 Ready |实现要点测试先用ImageStreamBuilder创建名为custom-image-stream的 ImageStream再通过 OpenShiftClient 提交spec.build.output使用ImageStreamOutput见 ImageStreamOutput.javatype: imagestream、image: custom-image-stream:latest插件复用PLUGIN_WITH_TAR_AND_JAR2 个制品断言包括spec.build.plugins[0].artifacts.size()2、output.typeimagestream、output.imagecustom-image-stream:latest、status 条件类型为Ready、status.connectorPlugins非空且包含 EchoSink 类名。作为对照Docker 输出类型由 DockerOutput.java 建模字段包括image完整镜像名如quay.io/my-organization/my-custom-connect:latest、pushSecret推送凭证 Secret、additionalBuildOptions与additionalPushOptions分别透传给 Kaniko/Buildah 的build与push命令仅 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()2且status.connectorPlugins同时包含两个连接器类名。十、测试工程实践镜像命名、基础设施与运行前提从套件源码可以总结出几条可复用的系统测试工程实践隔离的镜像命名getImageNameForTestCase()ConnectBuilderST.java用随机数生成strimzi-sts-connect-build镜像的 registry 输出地址避免并行测试间镜像相互覆盖TestStorage 上下文每个用例创建独立的TestStorage命名空间、cluster 名、topic 名、scraper 名等均从测试上下文派生配合KubeResourceManager统一管理资源生命周期与清理构建失败的可观测性套件反复使用KafkaConnectUtils.waitForConnectNotReady、waitUntilKafkaConnectStatusConditionContainsMessage、waitForConnectStatusContainsPlugins等工具验证构建状态条件NotReady/Ready与connectorPlugins状态字段这也是排查线上构建问题的关键入口网络策略联动多个用例调用NetworkPolicyUtils.deployNetworkPolicyForResource为 KafkaConnect 部署 NetworkPolicy确保启用网络策略的集群中 Scraper 仍能访问 Connect REST API运行前提需要在已就绪的 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),仅供参考