
Strimzi Kafka Access Operator 端到端系统测试详解从 KafkaAccess 到 Secret 的消息收发实战【免费下载链接】strimzi-kafka-operatorApache Kafka® running on Kubernetes项目地址: https://gitcode.com/GitHub_Trending/st/strimzi-kafka-operator本文基于 Strimzi Kafka Operator 仓库中的系统测试文档AccessOperatorST讲解 Kafka Access Operator下称 KAO与真实 Kafka 集群、KafkaUser 资源的集成测试testAccessOperator的完整设计与实现。读完后你将理解 KAO 的安装方式、KafkaAccess自定义资源CR的字段语义与 Secret 产出结构并掌握如何用测试产出的凭据完成 Kafka 生产/消费的端到端验证。测试目标为什么需要 testAccessOperator仓库中的测试文档 AccessOperatorST.md 对testAccessOperator的定义如下该测试在真实环境中验证 Kafka Access Operator 与 Kafka、KafkaUser CR 协同工作的功能。同时验证基于 Kafka 集群的凭据与信息Kafka 客户端能够成功连接并完成消息收发。KAO 的定位可以从仓库的安装文件与示例中看出它是一个独立的轻量级 Operator读取KafkaAccess资源后查找指定的 Kafka 实例并创建一个包含连接该 Listener 所需详情bootstrap 地址、信任链以及可选的用户证书的 Kubernetes Secret。示例文件 kafka-access-with-user.yaml 的注释进一步说明当指定了user字段时KAO 会查找对应的 KafkaUser并检查其认证方式与 Listener 是否匹配若匹配则将用户凭据一并写入所创建的 Secret。也就是说KAO 的价值在于把「连接哪个集群 用哪个用户」这两个声明合并成一份可直接消费的标准 Secret。测试被标注了 kafka-access 标签该标签下的测试专门验证「KAO 提供的数据足以连接真实 Kafka 集群并完成消息收发」对应源码实现位于 AccessOperatorST.javaTag(REGRESSION)属于回归测试集。五步流程总览原文档给出的完整步骤如下后文逐步结合源码展开| Step | Action | Result | | - | - | - | | 1. | Deploy Kafka Access Operator using the installation files from packaging folder. | Kafka Access Operator is successfully deployed. | | 2. | Deploy and create NodePools, Kafka with configured plain and TLS listeners, TLS KafkaUser, and KafkaTopic where we will produce the messages. | All of the resources are successfully deployed/created. | | 3. | Create KafkaAccess resource pointing to our Kafka cluster (and TLS listener) and TLS KafkaUser; Wait for KafkaAccess Secret creation. | Both KafkaAccess resource and KafkaAccess Secret is created. | | 4. | Collect KafkaAccess Secret and its data. | Data successfully collected. | | 5. | Use the data from KafkaAccess Secret in producer and consumer clients, do message transmission. | With the data provided by KAO, both clients can connect to Kafka and do the message transmission. |步骤 1从安装文件部署 Kafka Access Operator测试第一步调用SetupAccessOperator.install(namespace)见 AccessOperatorST.java#L72-L73。安装逻辑实现在 SetupAccessOperator.java安装文件目录为packaging/install/access-operator/即仓库中的 install/access-operator/ 目录按文件名顺序遍历 YAML 文件跳过 Namespace 文件测试使用独立命名空间按资源类型分别注入命名空间后应用ServiceAccount、ClusterRole、ClusterRoleBinding、Deployment、CRD。该目录下共有 6 个安装文件| 文件 | 内容 | | - | - | | 000-Namespace.yaml |strimzi-access-operator命名空间 | | 010-ServiceAccount.yaml | KAO 运行身份 | | 020-ClusterRole.yaml | RBAC 权限定义 | | 030-ClusterRoleBinding.yaml | 权限绑定 | | 040-Crd-kafkaaccess.yaml |kafkaaccesses.access.strimzi.ioCRD | | 050-Deployment.yaml |strimzi-access-operatorDeployment |020-ClusterRole.yaml 的权限范围恰好反映了 KAO 的工作面对access.strimzi.io组的kafkaaccesses及其status子资源拥有完整读写权限作为它管理的 CR对kafka.strimzi.io组的kafkas、kafkausers仅有get/list/watch只读权限KAO 只读取目标集群和用户不修改它们对核心组secrets拥有完整读写权限创建/更新/删除结果 Secret 是其核心产物。050-Deployment.yaml 显示 KAO 以单副本 Recreate 策略运行镜像为quay.io/strimzi/access-operator:0.3.0入口脚本为/opt/strimzi/bin/access_operator_run.sh并通过 8080 端口的/healthy、/ready探针做存活与就绪检查。步骤 2搭建带 plain/TLS 双监听器的 Kafka 集群测试的BeforeAll先通过SetupClusterOperator以默认配置安装 Cluster Operator然后在命名空间内创建见 AccessOperatorST.java#L76-L106一个控制器节点池KafkaNodePoolTemplates.controllerPool3 副本和一个 broker 节点池KafkaNodePoolTemplates.brokerPool3 副本一个 3 副本的 Kafka 集群显式配置两个internal监听器plain9092 端口tls: false即TestConstants.PLAIN_LISTENER_DEFAULT_NAMEtls9093 端口tls: true并启用KafkaListenerAuthenticationTlsAuthTLS 客户端认证即TestConstants.TLS_LISTENER_DEFAULT_NAME一个 KafkaTopic 与一个 TLS KafkaUserKafkaTopicTemplates.topic和KafkaUserTemplates.tlsUser。之所以要求 TLS Listener 开启客户端认证是因为后续 KafkaAccess 引用的 KafkaUser 是 TLS 用户——KAO 需要验证用户认证方式与 Listener 匹配后才能把用户证书写进 Secret这正对应 kafka-access-with-user.yaml 注释中描述的匹配检查逻辑。步骤 3创建 KafkaAccess 资源并等待 Secret测试随后用 Fabric8 构建器创建 KafkaAccess见 AccessOperatorST.java#L108-L130new KafkaAccessBuilder() .editOrNewMetadata() .withName(kafkaAccessName) // 测试存储的集群名 -access .withNamespace(testStorage.getNamespaceName()) .endMetadata() .editOrNewSpec() .withNewKafka() .withName(testStorage.getClusterName()) .withNamespace(testStorage.getNamespaceName()) .withListener(TestConstants.TLS_LISTENER_DEFAULT_NAME) // 指向 TLS 监听器 .endKafka() .withNewUser() .withName(testStorage.getUsername()) .withNamespace(testStorage.getNamespaceName()) .withApiGroup(KafkaUser.RESOURCE_GROUP) // kafka.strimzi.io .withKind(KafkaUser.RESOURCE_KIND) // KafkaUser .endUser() .endSpec() .build()随后通过SecretUtils.waitForSecretReady等待名为kafkaAccessName的 Secret 就绪——即验证「KafkaAccess 资源与其 Secret 都被创建」。KafkaAccess的 Schema 由 040-Crd-kafkaaccess.yaml 定义关键字段如下v1alpha1Namespaced 作用域短名kaCRD 上还带有servicebinding.io/provisioned-service: true标签表明它面向 Service Binding 场景| 字段 | 说明 | | - | - | |spec.kafka.name| 目标 Kafka 集群名称必填 | |spec.kafka.namespace| 目标集群所在命名空间 | |spec.kafka.listener| 指定使用的 Listener 名称CR 注释说明若未指定Operator 会自行选择一个优先 internal 监听器 | |spec.user.kind/spec.user.apiGroup/spec.user.name| 关联用户的 Kind如KafkaUser、API Group如kafka.strimzi.io与名称三者均必填 | |spec.user.namespace| 用户所在命名空间 | |spec.secretName| 自定义结果 Secret 名称测试中未设置默认与 KafkaAccess 同名 | |spec.template.secret.metadata| 为生成的 Secret 附加 labels / annotations | |status.conditions| 状态条件CRD 的 printer column 展示Ready条件状态 |kubectl get kafkaaccesses时CRD 的 additionalPrinterColumns 会直接展示Listener、Cluster、User与Ready四列便于运维快速核对。步骤 4收集 KafkaAccess Secret 的数据测试获取 Secret 后见 AccessOperatorST.java#L130-L163从中提取了三个关键要素bootstrap 地址Util.decodeFromBase64(accessSecret.getData().get(bootstrapServers))——Secret 的bootstrapServers键以 Base64 存储需解码后使用TLS 证书三件套通过环境变量间接引用同一 Secret 的三个键ssl.truststore.crt→ 注入为环境变量CA_CRT集群 CAssl.keystore.crt→ 注入为USER_CRT用户证书ssl.keystore.key→ 注入为USER_KEY用户私钥。这说明 KAO 产出的 Secret 是「自包含」的客户端无需再分别去查询 Kafka 集群的 CA Secret 和 KafkaUser 的 Secret仅凭这一个 Secret 即可完成 TLS 双向认证连接。步骤 5用 KAO 数据完成消息收发验证最后测试用KafkaProducerConsumerBuilder组装生产者与消费者 Job见 AccessOperatorST.java#L165-L183withBootstrapAddress(bootstrapServer)使用步骤 4 解码出的地址withAdditionalEnvVars(tlsEnvVarsForKafkaAccess)将 CA/用户证书/私钥三个环境变量各自引用 KafkaAccess Secret 中的对应键传入客户端容器withMessageCount(TestConstants.MESSAGE_COUNT)指定消息数量并随机生成消费组。客户端 Job 创建后ClientUtils.waitForClientsSuccess(...)等待生产者 Job 成功、消费者读到全部消息完成闭环验证。这一步证明KAO 汇总的凭据在真实 TLS 客户端认证场景下可用客户端连接与消息收发成功。手工复现从 KafkaAccess 到 Secret 的最小实践若不在系统测试环境而在真实集群中复现该流程可以按以下顺序操作依次应用 install/access-operator/ 下的 6 个 YAML注意050-Deployment.yaml中命名空间为strimzi-access-operator可按需调整部署好 Strimzi Kafka 集群含 TLS Listener 与 KafkaUser然后提交如下最小 KafkaAccess对照 kafka-access.yaml 与 kafka-access-with-user.yamlapiVersion: access.strimzi.io/v1alpha1 kind: KafkaAccess metadata: name: my-kafka-access spec: kafka: name: my-cluster namespace: kafka listener: tls user: kind: KafkaUser apiGroup: kafka.strimzi.io name: my-user namespace: kafka等待同名 Secret 创建可用kubectl get kafkaaccesses观察Ready列其中bootstrapServers需 Base64 解码若用户为 TLS 用户Secret 中还包含ssl.truststore.crt、ssl.keystore.crt、ssl.keystore.key等键将 Secret 键挂载或注入到客户端容器如上文测试中的CA_CRT/USER_CRT/USER_KEY环境变量即可建立连接。如何运行与定位该测试测试类systemtest/src/test/java/io/strimzi/systemtest/specific/AccessOperatorST.java方法testAccessOperator标注IsolatedTest与Tag(REGRESSION)属于独立运行的回归测试运行环境依赖需要可用的 Kubernetes 集群、已构建的 KAO 镜像quay.io/strimzi/access-operator:0.3.0以及systemtest模块的测试基础设施AbstractST提供的存储与资源管理安装路径逻辑SetupAccessOperator.java#L29 中PATH_TO_KAO_CONFIG指向packaging/install/access-operator/与仓库install/access-operator/目录内容对应相关文档标签kafka-access 标签说明 声明该标签下所有测试均验证「使用 KAO 提供的数据连接真实 Kafka 集群并完成消息收发」。小结testAccessOperator以五步链路完整覆盖了 Kafka Access Operator 的核心价值主张以KafkaAccess这一个 CR 为入口把 Kafka 集群的 bootstrap 信息、Listener 信任链与 KafkaUser 的 TLS 凭据聚合为一个标准 Kubernetes Secret再由真实生产者/消费者 Job 验证该凭据可用。从 RBAC对kafkaaccesses读写、对kafkas/kafkausers只读、对secrets读写到 CRD 字段spec.kafka必填、spec.user三要素必填、可选secretName与模板再到 Secret 键名bootstrapServers、ssl.truststore.crt、ssl.keystore.crt、ssl.keystore.key测试与安装文件、示例文件相互印证构成了对 KAO 端到端行为的可靠回归验证。【免费下载链接】strimzi-kafka-operatorApache Kafka® running on Kubernetes项目地址: https://gitcode.com/GitHub_Trending/st/strimzi-kafka-operator创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考