Strimzi Kafka 多版本系统测试深度解析:KafkaVersionsST 如何验证每个受支持的 Kafka 版本
【免费下载链接】strimzi-kafka-operatorApache Kafka® running on Kubernetes项目地址: https://gitcode.com/GitHub_Trending/st/strimzi-kafka-operator
导读
KafkaVersionsST 是 Strimzi Kafka Operator 系统测试套件(systemtest)中负责"按受支持 Kafka 版本逐一遍历、验证集群基础功能"的冒烟级(smoke)测试套件。它以 kafka-versions.yaml 中标记为supported: true的每个 Kafka 版本为参数,完整验证"集群部署 → Topic Operator / User Operator 就绪 → PLAIN+SCRAM-SHA 与 TLS 双通道消息收发"这条生产核心链路。阅读本文后,你将理解 Strimzi 多版本兼容性测试的整体设计、testKafkaWithVersion的每一步断言逻辑、Kafka 版本清单文件的结构与解析机制,并能独立在本地运行该套件。
一、套件定位:面向"每个受支持版本"的基础功能冒烟验证
KafkaVersionsST.md 将本套件的职责定义为:验证每个受支持的 Kafka 版本的基础功能(basic functionality)。它归属于 kafka 测试标签,该标签下的测试共同覆盖动态配置、监听器、节点池、版本升级、配额、分层存储等核心 Kafka 能力。
对应到源码,套件声明于 KafkaVersionsST.java:
@Tag(KAFKA_SMOKE) @SuiteDoc( description = @Desc("Verifies the basic functionality for each supported Kafka version."), beforeTestSteps = { @Step(value = "Deploy Cluster Operator with default installation.", expected = "Cluster Operator is deployed.") }, labels = { @Label(value = TestDocsLabels.KAFKA) } ) public class KafkaVersionsST extends AbstractST {其中KAFKA_SMOKE定义于 TestTags.java,值为"kafkasmoke",表明这是一组需要在功能变更合入前快速跑通的冒烟测试——因为"新版本能否正常部署"是所有其他能力验证的前提。
测试前置条件
根据文档中"Before test execution steps"表格,套件唯一的全局前置步骤是:
| 步骤 | 动作 | 预期结果 | | - | - | - | | 1 | 使用默认配置部署 Cluster Operator | Cluster Operator 已成功部署 |
该步骤由@BeforeAll void setup()实现(KafkaVersionsST.java),使用SetupClusterOperator.getInstance().withDefaultConfiguration().install()完成,即沿用系统测试框架的默认安装配置,不做额外定制。
二、testKafkaWithVersion:逐版本验证核心测试方法的全流程
testKafkaWithVersion是本套件的唯一测试方法(KafkaVersionsST.java),其定义中的关键注解是:
@ParameterizedTest(name = "Kafka version: {0}.version()") @MethodSource("io.strimzi.systemtest.utils.TestKafkaVersion#getSupportedKafkaVersions")这是一个 JUnit 5 参数化测试:参数来源TestKafkaVersion#getSupportedKafkaVersions会返回kafka-versions.yaml中所有supported: true的版本,因此每新增/移除一个受支持版本,该测试的用例数会自动随之增减,无需修改测试代码。测试名会带上具体版本,例如Kafka version: 4.3.1.version(),便于在 CI 报告中快速定位失败版本。
文档给出了该方法的 5 步验证流程:
| 步骤 | 动作 | 预期结果 | | - | - | - | | 1 | 以指定版本部署 Kafka 集群 | Kafka 集群无任何问题地完成部署 | | 2 | 验证 Topic Operator 创建 | Topic Operator 工作正常 | | 3 | 验证 User Operator 创建 | User Operator 与 SCRAM-SHA 及 ACL 配合工作正常 | | 4 | 通过 PLAIN + SCRAM-SHA 发送并接收消息 | 消息收发成功 | | 5 | 通过 TLS 发送并接收消息 | 消息收发成功 |
下面结合源码逐步拆解这 5 步背后的具体实现。
第 1 步:以指定版本部署 Kafka 集群(含 Node Pool)
测试首先创建 3 副本的 Broker 节点池与 Controller 节点池(KafkaVersionsST.java):
KubeResourceManager.get().createResourceWithWait( KafkaNodePoolTemplates.brokerPool(testStorage.getNamespaceName(), testStorage.getBrokerPoolName(), testStorage.getClusterName(), 3).build(), KafkaNodePoolTemplates.controllerPool(testStorage.getNamespaceName(), testStorage.getControllerPoolName(), testStorage.getClusterName(), 3).build() );随后创建Kafka自定义资源,关键配置点如下(KafkaVersionsST.java):
KafkaTemplates.kafka(testStorage.getNamespaceName(), testStorage.getClusterName(), 3) .editOrNewSpec() .editOrNewKafka() .withVersion(testKafkaVersion.version()) // 参数化注入目标版本 .addToConfig("auto.create.topics.enable", "true") // 允许客户端自动建 topic .withNewKafkaAuthorizationSimple() // 启用 Simple 授权 .endKafkaAuthorizationSimple() .withListeners( new GenericKafkaListenerBuilder() .withName(TestConstants.PLAIN_LISTENER_DEFAULT_NAME) // plain 监听器 .withPort(9092) .withType(KafkaListenerType.INTERNAL) .withTls(false) .withNewKafkaListenerAuthenticationScramSha512Auth() .endKafkaListenerAuthenticationScramSha512Auth() .build(), new GenericKafkaListenerBuilder() .withName(TestConstants.TLS_LISTENER_DEFAULT_NAME) // tls 监听器 .withPort(9093) .withType(KafkaListenerType.INTERNAL) .withTls(true) .withNewKafkaListenerAuthenticationTlsAuth() .endKafkaListenerAuthenticationTlsAuth() .build() ) .endKafka() .endSpec() .build();需要特别说明的是withVersion(testKafkaVersion.version()):这里写入Kafka资源的spec.kafka.version字段,是 Strimzi 选择容器镜像与运行参数的核心依据。createResourceWithWait会阻塞等待集群真正达到 Ready 状态,因此"部署成功"本身就隐含了版本兼容性检查——例如该版本的 Kafka 镜像能否正常启动、Controller/Broker 节点能否形成 KRaft 仲裁组等。
从运行时实现看,Cluster Operator 在 KafkaVersion.java 的Lookup类中加载 classpath 上的kafka-versions.yaml,并据此校验Kafka资源里声明的版本是否受支持、映射到哪个容器镜像。同时parseKafkaVersions会强制约束:版本清单中不允许出现重复版本,也不允许出现多个default: true(KafkaVersion.java)。
第 2 步:验证 Topic Operator
测试通过KafkaTopicTemplates.topic(testStorage).build()创建 Topic 资源并等待就绪(KafkaVersionsST.java)。Topic Operator 需要把KafkaTopic自定义资源翻译成真实的 Kafka topic,这一步验证的是"Entity Operator 中的 Topic Operator 在该 Kafka 版本上能正常同步 topic 元数据"。
第 3 步:验证 User Operator(SCRAM-SHA 凭据 + ACL)
测试创建了三个具有不同权限模型的KafkaUser(KafkaVersionsST.java),用于后续两路消息收发的鉴权验证:
- writeUser:SCRAM-SHA 用户,拥有对目标 topic 的
WRITE、DESCRIBE、CREATE操作(CREATE是为配合auto.create.topics.enable: true让生产者在 topic 不存在时自动创建); - readUser:SCRAM-SHA 用户,拥有对目标 topic 的
READ、DESCRIBE操作,以及对消费组的READ操作; - tlsReadWriteUser:TLS 证书认证用户,同时具备读写权限。
这一步验证了 User Operator 在该 Kafka 版本上能够正确生成 SCRAM-SHA 凭据 Secret、签发 TLS 客户端证书,并将 ACL 规则同步到 Kafka 的 Simple Authorizer 中。
第 4 步:PLAIN + SCRAM-SHA 消息收发
测试通过系统测试专用的 Kafka 客户端构建器创建生产/消费 Job(KafkaVersionsST.java):
final KafkaProducerClient kafkaProducerPlainScramShaWrite = new KafkaProducerClientBuilder() .withName(testStorage.getProducerName()) .withNamespaceName(testStorage.getNamespaceName()) .withTopicName(testStorage.getTopicName()) .withBootstrapAddress(KafkaResources.plainBootstrapAddress(testStorage.getClusterName())) .withMessageCount(testStorage.getMessageCount()) .withAuthentication(ClientsAuthentication.configurePlainScramSha(testStorage.getNamespaceName(), kafkaUserWrite)) .build();生产者使用writeUser的 SCRAM-SHA-512 凭据连接 PLAIN(9092)监听器写入messageCount条消息,随后消费者使用readUser在随机生成的消费组上读取并断言全部消息成功消费:
ClientUtils.waitForClientSuccess(testStorage.getNamespaceName(), testStorage.getProducerName(), testStorage.getMessageCount()); ClientUtils.waitForClientSuccess(testStorage.getNamespaceName(), testStorage.getConsumerName(), testStorage.getMessageCount());这一步同时覆盖了 PLAIN 监听器、SCRAM-SHA 认证、Simple 授权 ACL 在目标 Kafka 版本上的端到端可用性。
第 5 步:TLS 消息收发
最后使用tlsReadWriteUser,通过KafkaProducerConsumerBuilder同时创建生产与消费 Job,连接 TLS(9093)监听器完成读写(KafkaVersionsST.java),并以waitForClientsSuccess断言两侧同时成功。这一步验证了 TLS 监听器、TLS 客户端证书认证在该版本上的可用性。
三、版本数据源:kafka-versions.yaml 与版本解析机制
所有受支持版本集中维护在仓库根目录的 kafka-versions.yaml 中。该文件注释明确说明它是"变更受支持 Kafka 版本时唯一需要更新的单一事实来源",同时影响编译期与运行期:Docker 镜像构建、Helm chart 中的KAFKA_IMAGE_MAP、配置模型生成、Cluster Operator 运行时加载的KafkaVersion、以及make docu_versions生成的文档片段。
字段说明
| 字段 | 含义 | | - | - | |version| 在Kafka/KafkaConnect/KafkaMirrorMaker2自定义资源中使用的版本号 | |maven-version| 构建时使用的 Apache Kafka Maven 制品版本;未设置时默认与version相同,用于自定义 Kafka 构建场景 | |metadata| 该版本默认使用的 Kafka 协议 metadata 版本 | |url/checksum| Kafka 二进制发行包的下载地址及 SHA512 校验和 | |third-party-libs| 该版本使用的第三方依赖库目录(如4.3.x) | |supported|true表示当前受支持;不支持的版本保留在文件中用于历史追溯 | |default|true表示用户未显式指定版本时的默认版本,全文件只允许一个 |
以当前仓库为例,标记为supported: true的是 Kafka 4.2.0、4.2.1、4.3.0,以及default: true的 4.3.1(kafka-versions.yaml);更早的 2.1.0 到 4.1.2 等版本均以supported: false保留。
测试侧的版本解析
系统测试通过 TestKafkaVersion.java 读取该 YAML 文件:静态块在类加载时解析kafka-versions.yaml,过滤出所有supported: true的版本并排序(TestKafkaVersion.java),getSupportedKafkaVersions()即返回该列表作为参数化测试的数据源。若文件中没有任何受支持版本,测试框架会直接抛异常拒绝启动。此外该类还实现了点分版本比较(compareDottedVersions)与isUpgrade判断,为升级类测试提供基础。
运行期校验
在 Cluster Operator 侧,KafkaVersion.java 运行时解析同一份清单,为每个Kafka资源匹配版本对应的镜像并完成默认版本的唯一性校验。因此KafkaVersionsST的每个参数化用例,本质上也在验证"运行期版本映射与镜像拉取"这一完整链路。
四、支持的版本策略与版本增删
关于"哪些版本应该被支持",仓库根目录的 KAFKA_VERSION_SUPPORT.md 给出了两条维护规则:
- 至少支持最近两个 major/minor 版本的 Apache Kafka(例如:加入 2.4.x 支持时可以移除 2.2.x);
- 相邻两个 Strimzi 大/小版本之间至少保留一个共同的 Kafka 版本,以保证 Strimzi 与 Kafka 版本都能平滑升级。
当需要新增或移除版本时,development-docs/KAFKA_VERSIONS.md 提供了完整操作清单:新增版本需更新kafka-versions.yaml、为新 minor 版本准备docker-images下的第三方依赖目录、同步文档attributes.adoc与示例资源、更新pom.xml中的 Kafka 版本等;移除版本则只需将supported置为false(不要删除条目)。该文件还介绍了自 Strimzi 1.1.0 起支持的自定义 Kafka 版本能力——通过在kafka-versions.yaml中声明maven-version、自定义url/checksum即可让集群 Operator 使用带补丁的 Kafka 构建。可以看到,无论是新增、移除还是自定义版本,KafkaVersionsST都是版本变更后必须通过的回归验证关卡。
五、如何运行 KafkaVersionsST
系统测试基于 Maven 与 JUnit 5,套件本身继承自AbstractST。运行该套件前需确保已具备可用的 Kubernetes/OpenShift 集群,可参考仓库内的 TESTING.md 了解系统测试的环境准备与执行方式。典型的定向执行命令如下(在systemtest模块目录下):
mvn verify -pl systemtest -Dtest=KafkaVersionsST或按标签过滤冒烟类测试:
mvn verify -pl systemtest -Dgroups=kafkasmoke参数化测试会在 JUnit 报告中为每个受支持版本生成独立用例(如Kafka version: 4.3.1.version()),便于区分失败项。若新增版本后套件失败,通常需要回溯检查:镜像是否已构建并推送到仓库、kafka-versions.yaml中supported/default标记是否正确、对应版本的第三方依赖目录是否就绪。
六、小结
KafkaVersionsST 是 Strimzi 多版本兼容性保障体系中的第一道防线:它以参数化测试的形式,把 kafka-versions.yaml 中每个受支持版本都跑一遍"部署 → 双 Operator 就绪 → PLAIN/SCRAM-SHA 与 TLS 双通道收发"的完整链路。结合 TestKafkaVersion.java 的版本解析与 KafkaVersion.java 的运行时校验,开发者可以在合入新版本支持前快速获得"该版本在 Kubernetes 上能否正常运转"的明确信号——这正是 Strimzi 承诺"Apache Kafka running on Kubernetes"时最底层的质量保障。
延伸阅读
- kafka 测试标签文档:了解 kafka 标签下全部测试的完整清单
- KafkaVersionsST 源码:本文所有断言的实现出处
- kafka-versions.yaml:受支持版本的单一事实来源
- KAFKA_VERSION_SUPPORT.md:版本支持策略
- development-docs/KAFKA_VERSIONS.md:新增/移除/自定义 Kafka 版本的操作手册
【免费下载链接】strimzi-kafka-operatorApache Kafka® running on Kubernetes项目地址: https://gitcode.com/GitHub_Trending/st/strimzi-kafka-operator
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考