Strimzi Cluster Operator 的 StrimziPodSet 独立调和验证:解析 STRIMZI_POD_SET_RECONCILIATION_ONLY 系统测试
【免费下载链接】strimzi-kafka-operatorApache Kafka® running on Kubernetes项目地址: https://gitcode.com/GitHub_Trending/st/strimzi-kafka-operator
本文基于 Strimzi Kafka Operator 系统测试套件中的PodSetST测试文档展开,说明该套件如何验证StrimziPodSet资源的调和(reconciliation)行为,并结合仓库源码解析STRIMZI_POD_SET_RECONCILIATION_ONLY环境变量的底层实现:开启该模式后 Cluster Operator 只调和StrimziPodSet资源,Kafka等自定义资源的配置变更不再触发 Pod 滚动更新,而 Pod 缺失仍会被自动重建。读完后你能理解这一"仅调和 Pod 集合"模式的适用场景、生效机制,以及如何用系统测试中的步骤复现和验证该行为。
PodSetST 测试套件概述
PodSetST是 Strimzi 系统测试(systemtest)模块中的一个测试套件,其定位在 测试文档 中描述为:
Test suite for
StrimziPodSetrelated functionality and features, which verifies pod set reconciliation behavior.
对应的测试类位于 PodSetST.java,它继承自AbstractST,使用 JUnit 5 的@Tag(REGRESSION)标注为回归测试。套件当前只包含一个核心测试方法testPodSetOnlyReconciliation,并且通过@IsolatedTest注解声明该测试会修改 Cluster Operator 的环境变量,因此需要与其他测试隔离运行:
@Tag(REGRESSION) @SuiteDoc( description = @Desc("Test suite for `StrimziPodSet` related functionality and features, which verifies pod set reconciliation behavior.") ) public class PodSetST extends AbstractST { ... @IsolatedTest("We are changing CO env variables in this test") @TestDoc( description = @Desc("This test verifies that when the `STRIMZI_POD_SET_RECONCILIATION_ONLY` environment variable is enabled, only `StrimziPodSet` resources are reconciled, and Kafka configuration changes do not trigger rolling updates of pods."), ... ) void testPodSetOnlyReconciliation() { ... } }值得注意的是,类上的@SuiteDoc、@TestDoc、@Step、@Label等注解来自io.skodjob库,用于在测试代码中声明式地描述测试意图,development-docs/systemtests/目录下的 Markdown 文档即由这些注解生成,保证文档与代码保持同步。
STRIMZI_POD_SET_RECONCILIATION_ONLY 模式:文档与实现
PodSetST围绕的核心功能是 Cluster Operator 的一个运行模式。官方文档 con-configuring-cluster-operator.adoc 对该环境变量的说明是:
STRIMZI_POD_SET_RECONCILIATION_ONLY:: Optional, defaultfalse. When set totrue, the Cluster Operator reconciles only theStrimziPodSetresources and any changes to the other custom resources (Kafka,KafkaConnect, and so on) are ignored. This mode is useful for ensuring that your pods are recreated if needed, but no other changes happen to the clusters.
也就是说,该模式的价值在于:让 Operator 继续守护 Pod 的生命周期(Pod 被删除后会被重建),但冻结所有其他变更(不滚动更新、不修改集群配置),是一种用于运维管控、灾难恢复或灰度冻结的"最小干预"模式。
在源码中,该配置项定义于 ClusterOperatorConfig.java:
public static final ConfigParameter<Boolean> POD_SET_RECONCILIATION_ONLY = new ConfigParameter<>("STRIMZI_POD_SET_RECONCILIATION_ONLY", BOOLEAN, "false", CONFIG_VALUES); public boolean isPodSetReconciliationOnly() { return get(POD_SET_RECONCILIATION_ONLY); }isPodSetReconciliationOnly()在 ClusterOperator.java 中三处被检查,这三处正是"仅调和 PodSet"模式的落点:
启动阶段(
start方法,约 L94):StrimziPodSetController始终先启动(startFutures.add(startStrimziPodSetController())),只有当该模式为false时,才会为KafkaAssemblyOperator、KafkaConnectAssemblyOperator、KafkaBridgeAssemblyOperator、KafkaMirrorMaker2AssemblyOperator、KafkaRebalanceAssemblyOperator创建 watch,并额外开启 KafkaNodePool 和 KafkaConnector 的 watch。开启模式后,这些 watch 全部不建立,Operator 不再感知其他自定义资源的事件。周期性调和(约 L124-L129):
vertx.setPeriodic设置的定时任务内,只有模式为false时才会调用reconcileAll("timer")。这意味着即使 Operator 丢失了某些事件,定时器也不会再去兜底调和Kafka等资源。reconcileAll方法(约 L176-L185):作为最后一道保险,即便被触发,方法内部仍会判断isPodSetReconciliationOnly(),为true时直接跳过对所有 assembly operator 的调和调用。
// ClusterOperator.java(节选) startFutures.add(startStrimziPodSetController()); if (!config.isPodSetReconciliationOnly()) { List<AbstractOperator<?, ?, ?, ?>> operators = new ArrayList<>(asList( kafkaAssemblyOperator, kafkaConnectAssemblyOperator, kafkaBridgeAssemblyOperator, kafkaMirrorMaker2AssemblyOperator, kafkaRebalanceAssemblyOperator)); ... }而StrimziPodSetController的创建路径与上述分支无关,因此该模式下它仍然独立运行。从 StrimziPodSetController.java 的结构看,它是一个基于 Informer(podInformer、strimziPodSetInformer)与BlockingQueue<SimplifiedReconciliation>工作队列的独立控制器,拥有自己的调和循环,不依赖 assembly operator 的 watch,这正是"Pod 缺失仍可重建"这一行为的来源。
testPodSetOnlyReconciliation 测试场景解析
测试文档中的步骤表完整描述了测试流程,下面结合 PodSetST.java 的实现逐一说明。
| Step | Action | Result | | - | - | - | | 1. | 部署一个 broker 与 controller 节点池各 3 副本的 Kafka 集群。 | 带节点池和 topic 的 Kafka 集群部署成功。 | | 2. | 启动持续生产/消费客户端。 | 客户端持续生产并消费消息。 | | 3. | 在 Cluster Operator 上启用STRIMZI_POD_SET_RECONCILIATION_ONLY环境变量。 | Cluster Operator 携带新环境变量重启。 | | 4. | 修改Kafka资源中的 readiness 探针超时。 | 尽管配置变更,不发生 Pod 滚动更新。 | | 5. | 删除一个 Kafka Pod。 | Pod 由StrimziPodSet控制器重建。 | | 6. | 从 Cluster Operator 移除STRIMZI_POD_SET_RECONCILIATION_ONLY环境变量。 | Cluster Operator 再次重启。 | | 7. | 验证发生 Pod 滚动更新。 | 由于存在挂起的配置变更,Kafka Pod 被滚动重启。 | | 8. | 验证StrimziPodSet状态。 | 所有StrimziPodSet就绪,Pod 数量匹配。 | | 9. | 验证消息连续性。 | 客户端持续成功生产和消费全部消息。 |
阶段一:高可用集群与持续流量
测试首先部署一个具备数据保护能力的集群:replicas = 3、default.replication.factor = 3、min.insync.replicas = 2,并同时创建 3 副本的 broker 节点池和 controller 节点池(KRaft 分离角色模式):
final int replicas = 3; final int probeTimeoutSeconds = 6; EnvVar reconciliationEnv = new EnvVar(Environment.STRIMZI_POD_SET_RECONCILIATION_ONLY_ENV, "true", null); KubeResourceManager.get().createResourceWithWait( KafkaNodePoolTemplates.brokerPoolPersistentStorage(...).build(), KafkaNodePoolTemplates.controllerPoolPersistentStorage(...).build() ); KubeResourceManager.get().createResourceWithWait(KafkaTemplates.kafka(testStorage.getNamespaceName(), testStorage.getClusterName(), replicas) .editSpec() .editOrNewKafka() .addToConfig("default.replication.factor", 3) .addToConfig("min.insync.replicas", 2) .endKafka() .endSpec() .build());其中STRIMZI_POD_SET_RECONCILIATION_ONLY的名称统一由 Environment.java 中的常量管理:
public static final String STRIMZI_POD_SET_RECONCILIATION_ONLY_ENV = "STRIMZI_POD_SET_RECONCILIATION_ONLY";随后测试创建一个 3 副本的持续 topic,并启动KafkaProducerConsumer(每 1000ms 发送一条消息的 Producer/Consumer Job),为后续验证"变更不影响数据面"建立基线:
final KafkaProducerConsumer continuousKafkaProducerConsumer = new KafkaProducerConsumerBuilder() .withProducerName(testStorage.getContinuousProducerName()) .withConsumerName(testStorage.getContinuousConsumerName()) ... .withBootstrapAddress(KafkaResources.plainBootstrapAddress(testStorage.getClusterName())) .withMessageCount(testStorage.getContinuousMessageCount()) .withDelayMs(1000) .build();阶段二:开启模式并验证"配置变更不触发滚动"
测试直接读取 Cluster Operator Deployment 容器内的现有 env 列表,追加新变量后整体替换 Deployment,并等待 Operator 滚动完成:
DeploymentUtils.replace(ns, coName, coDep -> coDep.getSpec().getTemplate().getSpec().getContainers().get(0).setEnv(envVars)); DeploymentUtils.waitTillDepHasRolled(ns, coName, 1, coPod);滚动完成后先对 broker Pod 做一次快照,然后修改Kafka资源的 readiness 探针超时(一个必然会引起 Pod 模板变化、正常情况下会触发滚动更新的配置项),并断言没有发生滚动:
Map<String, String> brokerPods = PodUtils.podSnapshot(testStorage.getNamespaceName(), testStorage.getBrokerSelector()); KafkaUtils.replace(testStorage.getNamespaceName(), testStorage.getClusterName(), kafka -> { kafka.getSpec().getKafka().setReadinessProbe( new StrimziProbeBuilder().withTimeoutSeconds(probeTimeoutSeconds).build()); }); RollingUpdateUtils.waitForNoKafkaRollingUpdate(testStorage.getNamespaceName(), testStorage.getClusterName(), brokerPods);这一步验证的正是ClusterOperator.start()中的行为:模式开启后KafkaAssemblyOperator的 watch 未建立,Kafka资源上的探针变更不会产生任何调和动作。
阶段三:删除 Pod 并验证 StrimziPodSet 控制器仍在工作
紧接着测试删除一个 broker Pod,并等待 Pod 数量恢复到 3:
String kafkaPodName = PodUtils.listPodNames(testStorage.getNamespaceName(), testStorage.getBrokerSelector()).get(0); KubeResourceManager.get().deleteResourceWithoutWait(...pods()...withName(kafkaPodName).get()); PodUtils.waitForPodsReady(testStorage.getNamespaceName(), testStorage.getBrokerSelector(), replicas, true);这一步是整个测试的关键对照:assembly operator 的调和被禁用的同时,StrimziPodSetController(见 StrimziPodSetController.java)依旧通过 Informer 感知 Pod 消失事件并通过工作队列完成重建。这也印证了官方文档对该模式的描述——"ensuring that your pods are recreated if needed, but no other changes happen to the clusters"。
阶段四:关闭模式并验证挂起变更生效
移除环境变量、等待 Operator 再次滚动后,测试断言之前"挂起"的探针配置变更此时会触发滚动更新:
envVars.remove(reconciliationEnv); DeploymentUtils.replace(ns, coName, ...); DeploymentUtils.waitTillDepHasRolled(ns, coName, 1, coPod); // 配置在模式期间被修改过,恢复调和后 Pod 应当滚动 RollingUpdateUtils.waitTillComponentHasRolledAndPodsReady(ns, brokerSelector, replicas, brokerPods); // 等待所有 StrimziPodSet 状态与就绪 Pod 数量一致 StrimziPodSetUtils.waitForAllStrimziPodSetAndPodsReady(ns, clusterName, brokerComponentName, 3); ClientUtils.waitForClientsSuccess(ns, consumerName, producerName, messageCount);StrimziPodSetUtils.waitForAllStrimziPodSetAndPodsReady(定义于 StrimziPodSetUtils.java)会核对每个StrimziPodSet的status中的 Pod 数量与实际就绪 Pod 数量是否匹配,覆盖文档步骤 8。最后通过ClientUtils.waitForClientsSuccess校验 Producer 与 Consumer 统计的消息总数一致,证明从开启模式、Pod 删除重建到滚动更新的全过程中没有消息丢失(步骤 9)。
测试辅助设施与可复现性
该测试的可复现性建立在几个系统测试辅助类之上,读者若要自行验证相同行为,可以参照它们封装的等待逻辑:
- SetupClusterOperator:
@BeforeAll中通过withDefaultConfiguration().install()安装默认配置的 Cluster Operator,是所有 Operator 类测试的公共前置; DeploymentUtils:封装 Deployment 快照、替换与滚动等待,用于安全地变更 CO 环境变量并确认重启完成;RollingUpdateUtils:提供waitForNoKafkaRollingUpdate与waitTillComponentHasRolledAndPodsReady这对"断言不滚动/断言滚动"的工具,是本测试正反向验证的核心;StrimziPodSetUtils与PodUtils:分别负责StrimziPodSet状态核对和 Pod 快照、就绪等待。
小结
PodSetST虽然只有一个测试方法,但它完整覆盖了STRIMZI_POD_SET_RECONCILIATION_ONLY模式的行为边界:
- 模式开启时:
Kafka资源的配置变更被完全忽略(无滚动更新),但StrimziPodSet控制器独立工作,被删除的 Pod 会被自动重建——这一点由ClusterOperator.start()中 assembly operator watch 的条件注册和StrimziPodSetController的独立工作队列共同保证; - 模式关闭后:之前挂起的配置变更立即被调和,Pod 按预期滚动,且全程消息收发无丢失;
- 测试方法本身:通过修改真实 Cluster Operator Deployment 的环境变量并在滚动后做正/反向断言,验证了该模式在生产形态下的生效路径,而非仅验证单元测试中的 mock 行为。
对于需要在不改变集群状态的前提下修复 Pod 异常(例如清理僵尸 Pod、重建异常节点)的运维场景,这正是官方推荐的 Cluster Operator 配置开关,其完整说明可参考 con-configuring-cluster-operator.adoc。
【免费下载链接】strimzi-kafka-operatorApache Kafka® running on Kubernetes项目地址: https://gitcode.com/GitHub_Trending/st/strimzi-kafka-operator
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考