ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

Strimzi Kafka Operator 中 MirrorMaker 2 系统测试深度解析:跨集群复制、安全认证与滚动更新验证

Strimzi Kafka Operator 中 MirrorMaker 2 系统测试深度解析:跨集群复制、安全认证与滚动更新验证 Strimzi Kafka Operator 中 MirrorMaker 2 系统测试深度解析跨集群复制、安全认证与滚动更新验证【免费下载链接】strimzi-kafka-operatorApache Kafka® running on Kubernetes项目地址: https://gitcode.com/GitHub_Trending/st/strimzi-kafka-operator本文基于 Strimzi 仓库中的系统测试套件 MirrorMaker2ST 及其测试文档 development-docs/systemtests/io.strimzi.systemtest.mirrormaker.MirrorMaker2ST.md 展开逐条剖析该套件验证的 9 个核心场景基础消息镜像、IdentityReplicationPolicy、TLS/mTLS 与 SCRAM-SHA-512 认证、active-active 模式下的偏移量检查点恢复、Scale-to-Zero 伸缩、连接器状态机与错误处理以及 Secret/证书变更驱动的滚动更新。读完本文你可以理解 Strimzi 如何为KafkaMirrorMaker2资源构建端到端的验收标准并能将其作为生产环境故障排查连接器 FAILED 状态、认证失败、证书轮换后 Pod 未滚动等的对照手册。测试套件总览MirrorMaker2ST是一个继承自AbstractST的 JUnit 5 测试类位于 systemtest/src/test/java/io/strimzi/systemtest/mirrormaker/MirrorMaker2ST.java。类级注解表明它同时带有三个标签REGRESSION、MIRROR_MAKER2与CONNECT_COMPONENTS即它既属于回归测试集也属于 MirrorMaker 2 专项与 Connect 组件类别。BeforeAll的setup()方法第 1413-1422 行会先安装带自定义操作超时的 Cluster Operator这是所有用例的共同前置条件。该套件官方文档labels/mirror-maker-2.md对其覆盖范围的描述是验证 Kafka MirrorMaker 2 跨集群复制。覆盖安全连接TLSmTLS、TLSSCRAM、消息/Header 镜像、伸缩含 scale-to-zero、topic 分区传播、identity replication policy、active-active 场景下的 offset checkpoint/restore、连接器状态迁移failed - running、pause/resume、手动与 secret/cert 驱动的滚动更新以及通过 status/conditions 与 offset 内省实现的错误处理。套件中每个用例都使用KubeResourceManager.get().createResourceWithWait(...)创建 Kafka 集群由KafkaNodePoolTemplates的 broker/controller 节点池 KafkaTemplates的Kafka资源构成采用 KRaft 模式下的节点池架构然后通过KafkaMirrorMaker2Templates构建KafkaMirrorMaker2自定义资源。MM2 的默认模板见 KafkaMirrorMaker2Templates.java其生成的资源与官方示例 kafka-mirror-maker-2.yaml 结构一致apiVersion: kafka.strimzi.io/v1 kind: KafkaMirrorMaker2 metadata: name: my-mirror-maker-2 spec: version: 4.3.1 replicas: 1 target: alias: cluster-b bootstrapServers: cluster-b-kafka-bootstrap:9092 groupId: my-mirror-maker-2-group # internal topics should be prefixed with __ to be excluded from mirroring by default configStorageTopic: __my-mirror-maker-2-config offsetStorageTopic: __my-mirror-maker-2-offset statusStorageTopic: __my-mirror-maker-2-status config: # -1 means it will use the default replication factor configured in the broker config.storage.replication.factor: -1 offset.storage.replication.factor: -1 status.storage.replication.factor: -1 mirrors: - source: bootstrapServers: cluster-a-kafka-bootstrap:9092 alias: cluster-a sourceConnector: tasksMax: 1 config: replication.factor: -1 offset-syncs.topic.replication.factor: -1 sync.topic.acls.enabled: false checkpointConnector: tasksMax: 1 config: checkpoints.topic.replication.factor: -1 sync.group.offsets.enabled: false refresh.groups.interval.seconds: 600 topicsPattern: .* groupsPattern: .*从源码结构看每个镜像对mirror由三类连接器组成sourceConnector镜像消息、checkpointConnector镜像消费组检查点、以及隐含的Heartbeat机制target中的configStorageTopic、offsetStorageTopic、statusStorageTopic三个内部 topic 则分别存储 Connect 配置、镜像偏移量与心跳状态。示例中的注释也提示将内部 topic 以__前缀命名可以让 MM2 默认将其排除在镜像范围之外。以下按文档中的测试条目逐一展开并在每节中补充源码级实现细节。testMirrorMaker2基础镜像、手动滚动更新与分区数传播这是套件的主干用例验证 MM2 在两个集群间镜像消息并处理滚动更新事件。源码实现见 第 126-245 行。步骤与结果继承自测试文档步骤操作预期结果1以默认配置部署源/目标 Kafka 集群Kafka 集群就绪2部署 MirrorMaker 2 与一个源 topicMM2 部署完成且初始 topic 存在3在源集群生产并消费消息客户端成功完成收发4检查 MM2 的 ConfigMap、Pod 标签等元数据配置与标签符合预期5验证消息已镜像到目标集群目标集群消费者能消费到镜像消息6通过注解手动触发 MM2 滚动更新MM2 Pod 滚动且状态保持7更新 topic 分区数验证分区数传播镜像 topic 的分区数同步更新源码级细节ConfigMap 内容断言。测试构造了一份期望的kafka-connect.properties内容第 131-141 行逐项校验 MM2 生成的 ConfigMapbootstrap.servers指向目标集群的 plain bootstrap 地址group.idmirrormaker2-clusterkey/value/header 三类 converter 均为ByteArrayConverter这是 MM2 作为“透明复制管道”的关键——不做任何格式转换三个内部 topicmirrormaker2-cluster-configs/-status/-offsets三个 replication factor 均为-1使用 broker 默认值。手动滚动更新通过注解strimzi.io/manual-rolling-update: true打在 MM2 对应的 StrimziPodSet 上触发第 223-226 行。该注解常量定义在 ResourceAnnotations.javapublic static final String ANNO_STRIMZI_IO_MANUAL_ROLLING_UPDATE STRIMZI_DOMAIN manual-rolling-update;滚动完成后测试先用 Admin Client 确认镜像 topic 分区数为 3再执行KafkaTopicUtils.replace(... setPartitions(8))将源 topic 扩到 8 分区并断言目标集群中镜像 topic 的分区数最终变为 8——这验证了 MM2 的refresh.topics.interval.seconds1配置下 topic 变更的传播能力。测试还通过VerificationUtils系列方法校验了 Pod/Service/ConfigMap/ServiceAccount 的标签一致性第 195-200 行这些正是 CustomResourceStatusST 中同类断言的复用。testIdentityReplicationPolicy保留源 topic 名称的镜像策略描述验证 MM2 使用IdentityReplicationPolicy使镜像 topic 在两个集群间保持完全相同的名称源码注释明确写道“This is what should be used by all new users”第 659-661 行。步骤与结果步骤操作预期结果1部署源/目标 Kafka 集群与一个 scraper Pod集群与 scraper Pod 就绪2部署使用 IdentityReplicationPolicy 的 MM2MM2 使用 identity 策略3在源集群生产/消费消息消息成功收发4在目标集群消费镜像 topic镜像 topic 名称与源 topic 完全一致消息消费成功源码级细节关键配置是打在sourceConnector上的两项第 696-705 行.editSourceConnector() .addToConfig(replication.policy.class, org.apache.kafka.connect.mirror.IdentityReplicationPolicy) .addToConfig(refresh.topics.interval.seconds, 1) .endSourceConnector()replication.policy.class指定使用 Kafka 原生的IdentityReplicationPolicy即目标集群 topic 名 源 topic 名而非默认的source-alias_topic前缀式重命名refresh.topics.interval.seconds1把 topic 刷新周期压到 1 秒保证测试中新建的源 topic 能被 MM2 快速发现生产示例 kafka-mirror-maker-2-custom-replication-policy.yaml 中该值为 600 秒说明这是测试专用的加速参数与默认策略不同identity 策略下target侧无需为每个源 alias 维护前缀映射镜像 topic 名即源 topic 名这正是测试第 4 步能直接以testStorage.getTopicName()而非getMirroredSourceTopicName()在目标集群消费的原因。testMirrorMaker2TlsAndScramSha512AuthTLS SCRAM-SHA-512 安全镜像描述验证消息经由 TLS 通道、以 SCRAM-SHA-512 认证完成镜像。步骤与结果步骤操作预期结果1部署带 TLS 监听器 SCRAM-SHA-512 认证的源/目标集群集群与 SCRAM 用户就绪2经 TLSSCRAM 在源集群收发客户端使用 SCRAM-SHA-512 over TLS 成功操作3部署带 SCRAM-SHA-512 凭证与受信任证书的 MM2MM2 以 SCRAM-SHA-512 连接运行就绪4在目标集群消费镜像 topic使用 SCRAM-SHA-512 认证成功消费5检查目标集群镜像 topic 分区数与源 topic 分区数一致源码级细节第 433-578 行两个 Kafka 集群均配置internal类型、tls: true、认证为KafkaListenerAuthenticationScramSha512的 9093 端口监听器并各创建一个KafkaUserTemplates.scramShaUser(...)用户MM2 的 source/target 侧分别声明kafkaClientAuthenticationScramSha512引用用户名 PasswordSecretSource密码 key 固定为password和tls.trustedCertificates引用各集群的 cluster CA 证书 Secret证书 key 为ca.crt# 等价 YAML由测试代码构建器生成 source: bootstrapServers: source-kafka-bootstrap:9093 authentication: type: scram-sha-512 username: source-user passwordSecret: secretName: source-user password: password tls: trustedCertificates: - certificate: ca.crt secretName: source-cluster-ca-cert最终通过 TLS Admin Client 连接目标集群断言镜像 topic 存在且partitionCount()等于源 topic 的 3 个分区。testMirrorMaker2TlsAndTlsClientAuthTLS mTLS 双向认证镜像描述验证消息经由 TLS 通道、以 mTLS双向证书认证完成镜像。该用例被额外打上ACCEPTANCE标签第 252 行。步骤与结果步骤操作预期结果1部署带 TLS 监听器 mTLS 认证的源/目标集群集群与用户按 TLS/mTLS 部署2部署 topic 与 TLS 用户topic 与 TLS 用户存在3在源集群经 TLS 收发TLS 客户端操作成功4部署带 TLSmTLS 配置与受信任证书的 MM2MM2 以 mTLS 连接运行就绪5在目标集群消费镜像 topic使用 TLS 成功消费镜像消息6检查目标集群镜像 topic 分区数与源 topic 分区数一致源码级细节第 267-413 行与 SCRAM 用例的差异在于客户端身份由证书提供。MM2 的 source/target 侧使用KafkaClientAuthenticationTls并指向 KafkaUser Secret 中的user.crt/user.key.withNewKafkaClientAuthenticationTls() .withNewCertificateAndKey() .withSecretName(testStorage.getSourceUsername()) // KafkaUser 名即 Secret 名 .withCertificate(user.crt) .withKey(user.key) .endCertificateAndKey() .endKafkaClientAuthenticationTls() .withNewTls() .withTrustedCertificates(certSecretSource) // cluster CA .endTls()从源码结构看mTLS 场景中 MM2 的“客户端证书”来自KafkaUser生成的 Secret而“受信任 CA”来自各 Kafka 集群的cluster-ca-certSecret两者缺一Kafka 监听器都会拒绝连接。测试最后同样以 Admin ClientTLS 认证断言镜像 topic 分区数为 3。testRestoreOffsetsInConsumerGroupactive-active 下的消费组偏移量同步描述验证 MM2 active-active双向模式下消费组偏移量的 checkpoint/restore 机制能够防止消息重复消费。步骤与结果步骤操作预期结果1部署 Kafka 集群与 active-active 配置的 MM2Active-active MM2 就绪2在两个 Kafka 集群间生产/消费消息消息正常收发3生产新消息后一部分从源集群消费、一部分从目标集群消费两集群偏移量出现分歧同步机制被验证4验证偏移量检查点防止重复消费消费任务在空偏移量上如期超时源码级细节第 745-878 行双向镜像部署两个KafkaMirrorMaker2实例——A→B 与 B→A每个实例均设置topicsPattern: .*与groupsPattern: .*且 source/checkpoint 连接器配置刻意调低间隔refresh.topics.interval.seconds1、sync.group.offsets.interval.seconds1、emit.checkpoints.interval.seconds1等以加速收敛MapString, Object checkpointConnectorConfig new HashMap(); checkpointConnectorConfig.put(refresh.groups.interval.seconds, 1); checkpointConnectorConfig.put(sync.group.offsets.enabled, true); checkpointConnectorConfig.put(sync.group.offsets.interval.seconds, 1); checkpointConnectorConfig.put(emit.checkpoints.enabled, true); checkpointConnectorConfig.put(emit.checkpoints.interval.seconds, 1); checkpointConnectorConfig.put(checkpoints.topic.replication.factor, 1);这与官方示例 kafka-mirror-maker-2-sync-groups.yaml 中的sync.group.offsets.enabled: true场景对应。偏移量分歧构造同一consumerGroup的消费者先消费源集群全部消息“Producer A”再消费目标集群镜像 topic“Producer B”随后源集群新增 50 条消息消费者故意从源集群消费 10 条并附加max.poll.records10再从目标集群消费 40 条——人为制造两集群间组偏移量的分歧。无重复消费断言再次从目标集群以及源集群尝试至少消费 1 条消息时任务会因“无消息可读”而超时测试用ClientUtils.waitForClientTimeout断言这种失败如期发生——即 checkpoint 机制保证了 50 条消息恰好被消费一次既不丢也不重。testScaleMirrorMaker2UpAndDownToZero扩缩容与缩容到零描述验证 MM2 的扩容、缩容尤其是 scale-to-zero。步骤与结果步骤操作预期结果1部署源/目标 Kafka 集群集群就绪2以初始副本数部署 MM2MM2 正常启动3向上扩容 MM2校验 observedGeneration 与 Pod 命名Pod 数量增加且新 Pod 命名正确4将 MM2 缩容至 0 副本等待 status URL 变为 null所有 Pod 被移除replicas0状态反映关闭URL 为 null源码级细节第 594-657 行扩容使用 Kubernetes 的 scale 子资源scaleByName(KafkaMirrorMaker2.RESOURCE_KIND .kafka.strimzi.io, name, 2)即kubectl scale --resource-version ... KafkaMirrorMaker2的等效 API 调用随后断言spec.replicas、status.replicas与实际 Pod 数三者一致且metadata.generationobserved generation大于扩容前记录值——说明 Operator 已完成一次成功协调缩容到零通过KafkaMirrorMaker2Utils.replace(... setReplicas(0))实现然后依次断言MM2 选择器下 Pod 数为 0、status conditions 中Ready条件存在、generation 变化并轮询等待status.url null——Connect REST API 的 URL 字段在组件关闭后被清空这是“彻底下线”的可观测信号。testKafkaMirrorMaker2ConnectorsStateAndOffsetManagement连接器状态机与偏移量管理这是套件中最复杂的用例验证连接器从 FAILED 到 RUNNING 的故障恢复、pause/resume 语义以及通过注解导出偏移量的“offset 内省”能力。步骤与结果步骤操作预期结果1以错误配置错误 bootstrap 地址部署 Kafka 集群与 MM2强制连接器失败MM2 显示 NotReady 并带错误信息2修正 MM2 配置修复 bootstrap 地址消除连接器失败连接器从 FAILED 迁移到 RUNNING3暂停/恢复连接器并验证状态迁移连接器状态与消息镜像行为符合预期4用 scraper 与 ConfigMap 校验偏移量外部存储中的偏移量值正确源码级细节第 893-1018 行强制失败MM2 的source.bootstrapServers被故意写成source-kafka-bootstrap.:9092多了一个点MM2 因此进入NotReadystatus message 为One or more connectors are in FAILED state测试用KafkaMirrorMaker2Utils.waitForKafkaMirrorMaker2StatusMessage等待该状态出现。修复与状态迁移replace修正bootstrapServers后测试等待 MM2 恢复 Ready并断言 conditions 中不再有该错误 message——对应 Connect 层面MirrorSourceConnector内部名为source-target.MirrorSourceConnector从 FAILED 回到 RUNNING。pause/resume通过spec.mirrors[0].sourceConnector.state在PAUSED与RUNNING之间切换ConnectorState枚举定义在 api 模块。暂停期间源集群生产/消费正常但目标集群消费任务按预期超时镜像中断恢复后目标消费者成功收到积压消息。偏移量内省MM2 的 source connector 上预先声明了listOffsets.toConfigMap写入cluster-offsets-listConfigMap恢复 RUNNING 后先经 scraper Pod 的 Connect REST API/offsets/0/offset/offset路径确认偏移量为 99再通过两个注解把偏移量导出到 ConfigMap 并解析 JSON 断言同值mm2.getMetadata().getAnnotations().putAll(Map.of( Annotations.ANNO_STRIMZI_IO_CONNECTOR_OFFSETS, list, Annotations.ANNO_STRIMZI_IO_MIRRORMAKER_CONNECTOR, sourceConnectorName));这两个注解分别对应 ResourceAnnotations.java 中的strimzi.io/connector-offsets与strimzi.io/mirrormaker-connector。注意源码中对连接器名做了两次转义处理URL 中-编码为%2D%3EConfigMap key 中替换为--这是操作者手工查询偏移量时容易踩到的细节。testKMM2RollAfterSecretsCertsUpdateScramShaSCRAM 密码变更后 MM2 自动滚动描述验证修改 SCRAM-SHA 用户密码 Secret 后MM2 Pod 自动完成滚动更新并继续镜像。步骤与结果步骤操作预期结果1部署带 SCRAM-SHA 的 Kafka 集群与用户用户与集群就绪2部署带 SCRAM-SHA 与 CA 凭证的 MM2MM2 正常镜像消息3更新源/目标用户密码验证 MM2 Pod 滚动Secret 更新后 MM2 完成滚动4滚动更新后生产/消费Secret 变更后镜像持续可用源码级细节第 1038-1197 行MM2 的 target 侧trustedCertificates使用了pattern 形式setPattern(*.crt)引用集群 CA Secret验证了通配符匹配证书文件这一配置方式在真实滚动场景中可用密码更新调用KafkaUserUtils.modifyKafkaUserPasswordWithNewSecret(...)两次源用户、目标用户各一次每次更新后通过RollingUpdateUtils.waitTillComponentHasRolledAndPodsReady断言 MM2 Pod 完成一次滚动Pod 快照比对——因为 MM2 的sasl.jaas.config是从用户 Secret 渲染出来的Secret 内容变化必然触发 MM2 的滚动更新滚动完成后测试注释明确说明“passwords have changed, we need to change the authentication as well (to get new sasl.jaas.config from the users secrets)”第 1187 行重新以新凭证生产到源、消费自目标闭环验证镜像链路在凭证轮换后依旧完好。testKMM2RollAfterSecretsCertsUpdateTLS证书轮换后的级联滚动描述验证更换 TLS 用户 Secret 与集群证书后MM2 与 Kafka Pod 的滚动更新行为。步骤与结果步骤操作预期结果1部署带 TLS 的 Kafka 集群与用户TLS 用户与集群就绪2部署带 TLS 凭证与受信任证书的 MM2MM2 正常镜像消息3更新客户端 CA 与集群 CA Secret验证 MM2 与 Kafka Pod 滚动CA 更新后 Pod 完成滚动4滚动更新后生产/消费证书/Secret 变更后镜像持续可用源码级细节第 1213-1411 行该用例是四阶段“证书轮换 → 级联滚动”的完整演练每个阶段都通过strimzi.io/force-renew注解常量见 ResourceAnnotations.java驱动源集群客户端 CA 轮换对source-clients-ca-certSecret 打strimzi.io/force-renew: true然后等待源 broker Pod 与 MM2 Pod 同时完成滚动——说明 MM2 因挂载/信任客户端 CA 而必须跟进目标集群客户端 CA 轮换同样断言目标 broker 与 MM2 滚动源集群集群 CA 轮换等待源 controller/broker Pod 与 Entity Operator Deployment 滚动MM2 再次滚动目标集群集群 CA 轮换目标 controller/broker 与 Entity Operator 滚动MM2 最后一次滚动。每轮滚动后都重新跑一轮“源生产 → 目标消费”验证镜像链路且最后两轮使用TestConstants.GLOBAL_TIMEOUT_LONG的加长超时注释解释原因是“Extend the timeout for clients to be sure that all messages are synced by MM2”——滚动后 MM2 需要时间重建连接与追平偏移量。从源码结构看这个用例实际上把 MM2 的证书依赖链clients CA、cluster CA、用户证书 Secret完整映射到了 Pod 滚动行为上可作为生产环境执行 CA 轮换前的行为基线。如何在仓库中定位与运行这些测试测试源码systemtest/src/test/java/io/strimzi/systemtest/mirrormaker/MirrorMaker2ST.javaCRD 模板与默认参数KafkaMirrorMaker2Templates.java其中默认sourceConnector配置为replication.factor-1、sync.topic.acls.enabledfalse、refresh.topics.interval.seconds600checkpointConnector配置为sync.group.offsets.enabledfalse生产可用示例examples/mirror-maker/ 目录下有基础、TLS、自定义复制策略、同步消费组四类 YAML均可直接参考修改相关测试文档同属mirror-maker-2标签的还有状态/错误用例CustomResourceStatusST、日志LogSettingST、指标MetricsST与 Leader 选举LeaderElectionST中的 MM2 相关测试运行方式整个 systemtest 模块通过 systemtest/Makefile 驱动运行约定与前置条件见 development-docs/TESTING.md。这些用例均要求真实 Kubernetes/OpenShift 集群与已部署的 Strimzi Operator由SetupClusterOperator在安装阶段完成属于集群级集成测试而非单元测试。小结MirrorMaker2ST套件以 9 个场景系统化覆盖了 StrimziKafkaMirrorMaker2的验收维度维度覆盖用例关键机制基础镜像与元数据testMirrorMaker2Connect 配置断言、标签校验、手动滚动注解命名策略testIdentityReplicationPolicyIdentityReplicationPolicytopic 名跨集群不变安全认证testMirrorMaker2TlsAndScramSha512Auth / testMirrorMaker2TlsAndTlsClientAuthSCRAM-SHA-512 密码 Secret、mTLS 用户证书 集群 CA组偏移量同步testRestoreOffsetsInConsumerGroupactive-active 双向镜像、checkpoint 防重复消费伸缩testScaleMirrorMaker2UpAndDownToZeroscale 子资源、replicas0 时 status.url 为 null故障与状态机testKafkaMirrorMaker2ConnectorsStateAndOffsetManagementFAILED→RUNNING 恢复、pause/resume、偏移量导出注解凭证轮换testKMM2RollAfterSecretsCertsUpdateScramSha / testKMM2RollAfterSecretsCertsUpdateTLSSecret/CA 变更触发 MM2 与 Kafka Pod 级联滚动对运维人员而言这套测试同时是一份“行为契约文档”例如遇到 MM2 status 出现One or more connectors are in FAILED state对应的恢复路径就是修正bootstrapServers并等待条件消失需要临时停止镜像时可将sourceConnector.state置为PAUSED需要审计镜像位置时可借助strimzi.io/connector-offsets注解将偏移量导出到 ConfigMap。【免费下载链接】strimzi-kafka-operatorApache Kafka® running on Kubernetes项目地址: https://gitcode.com/GitHub_Trending/st/strimzi-kafka-operator创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表