尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Strimzi Kafka Operator 系统测试解析:Kafka Connect 全生命周期与安全特性的端到端验证

Strimzi Kafka Operator 系统测试解析:Kafka Connect 全生命周期与安全特性的端到端验证 Strimzi Kafka Operator 系统测试解析Kafka Connect 全生命周期与安全特性的端到端验证【免费下载链接】strimzi-kafka-operatorApache Kafka® running on Kubernetes项目地址: https://gitcode.com/GitHub_Trending/st/strimzi-kafka-operator导读本文深入解析 Strimzi Kafka Operator 仓库中 ConnectST 系统测试套件对应实现见 ConnectST.java该套件针对 Kafka Connect 组件进行端到端验证覆盖部署与手动滚动更新、连接器KafkaConnector生命周期管理pause/stop/run、连接器 Offset 的 list/alter/reset 管理、任务自动重启、TLS 与 SCRAM-SHA 认证、Secret/ConfigMap 挂载、JVM 参数与资源配额、以及从扩容缩容到缩为零副本的完整缩放场景。读完本文你将掌握 ConnectST 每个测试用例的验证目标、操作步骤与预期结果并能对照源码理解 Kafka Connect 相关 API 与注解在真实集群中的工作方式。ConnectST 测试套件概览ConnectST 是 Strimzi 系统测试体系中 Kafka Connect 功能验证的核心套件。根据测试文档的描述其总体目标是验证 Kafka Connect 组件的部署、手动滚动更新与卸载undeployment并在执行任何测试之前先部署一个用于访问其他所有 Pod 的 scraper Pod测试探针容器。从套件注解ConnectST.java L116-L128可以看到类级标签REGRESSION回归、CONNECT、CONNECT_COMPONENTS套件级前置步骤部署 scraper Pod 用于访问所有其他 Pod预期结果为 scraper Pod 成功部署套件文档标签归类为 connect。该标签文档说明这些测试验证 Kafka Connect 组件确保 Kafka 与外部系统之间通过连接器可靠集成覆盖插件管理、构建过程、网络配置及各种安全协议等场景其正确性对于流生态系统中数据一致性与可用性的维护至关重要。ConnectST 共包含 19 个测试方法全部使用ParallelNamespaceTest注解以支持并行命名空间执行。测试的通用资源编排模式是通过KafkaNodePoolTemplates.brokerPool(...)/controllerPool(...)创建 KRaft 模式下的 broker 与 controller 节点池再创建 Kafka 集群、KafkaTopic、KafkaUser、KafkaConnect 与 KafkaConnector并以KubeResourceManager.get().createResourceWithWait(...)等待资源就绪。部署、手动滚动更新与配置校验testDeployRollUndeploy部署—滚动—验证全流程该用例ConnectST.java L148-L190验证 Kafka Connect 的部署、手动滚动更新及资源标签一致性主要步骤如下初始化 TestStorage每个并行命名空间测试的上下文容器设定 Kafka Connect 副本数为 2通过StUtils.loadProperties定义期望的 Connect 配置包括bootstrap.servers指向 TLS bootstrap 地址、group.id、key/value converter 均为JsonConverter、三个存储 topic 的副本因子为-1以及 config/status/offset 存储 topic 名称创建 broker/controller KafkaNodePools各 3 副本与 Kafka 集群、2 副本的 KafkaConnect手动滚动更新先对 Pod 做快照PodUtils.podSnapshot然后使用StrimziPodSetUtils.annotateStrimziPodSet为 StrimziPodSet 打上strimzi.io/manual-rolling-updatetrue注解再通过RollingUpdateUtils.waitTillComponentHasRolled等待组件完成滚动Pod 快照发生变化校验生成的 ConfigMap 中kafka-connect.properties包含全部期望配置项校验 Kafka Connect 容器镜像、Pod/Service/ConfigMap/ServiceAccount 的标签是否符合预期。这里使用的strimzi.io/manual-rolling-update注解定义于 ResourceAnnotations.java L36。它是 Strimzi 触发组件手动滚动更新的标准手段——用户无需删除 Pod只需在 StrimziPodSet 上添加该注解Cluster Operator 即会在下一次协调时逐个重建 Pod。testCustomAndUpdatedValues自定义环境变量与探针的可更新性该用例ConnectST.java L869-L959验证 Kafka Connect 的自定义环境变量与 readiness/liveness 探针在更新前后的正确性初始阶段通过spec.template.connectContainer.env注入TEST_ENV_1test.env.one、TEST_ENV_2test.env.two设置 readiness/liveness 探针initialDelaySeconds30、timeoutSeconds10、periodSeconds10、successThreshold1、failureThreshold3验证初始环境变量与探针参数确实反映在 Pod 上更新阶段将环境变量改为TEST_ENV_2updated.test.env.two与新增TEST_ENV_3同时更新spec.config三个存储 replication factor 设为-1并调整探针的 initialDelay31s、timeout11s、period5s、failureThreshold1等待组件完成滚动更新后再次验证 Pod 上的环境变量、探针参数以及 ConfigMap 中生成的kafka-connect.properties已包含更新后的配置。该测试的工程意义在于Strimzi 的KafkaConnectTemplate允许用户以声明式方式定制容器环境变量与健康探针而任何 spec 变更都会触发协调与滚动这一机制是生产环境热调整 Connect 配置的基础。testJvmAndResourcesJVM 参数与资源配额落地验证该用例ConnectST.java L445-L478确保 JVM 选项与资源请求/限制被正确应用到 Kafka Connect 组件通过spec.resources设置 limits 为memory400M、cpu2requests 为memory300M、cpu1通过spec.jvmOptions设置-Xmx200m、-Xms200m并追加-XX:UseG1GCJVM 选项以 Map 形式提供验证手段VerificationUtils.assertPodResourceRequests检查 Pod 上实际生效的资源请求/限制VerificationUtils.assertJvmOptions检查容器启动参数中确实包含-Xmx200m、-Xms200m、-XX:UseG1GC。KafkaConnector 生命周期与 Offset 管理testKafkaConnectAndConnectorStateWithFileSinkPluginpause/stop/run 状态控制该用例ConnectST.java L207-L252通过 FileSink 连接器验证 KafkaConnector 的暂停、停止与恢复能力是 SANITY/SMOKE 级别的核心用例。其实现借助verifySinkConnectorByBlockAndUnblock辅助方法ConnectST.java L1956-L1989部署 KafkaConnect 时开启strimzi.io/use-connector-resources: true注解定义于 ResourceAnnotations.java L61使 KafkaConnector 自定义资源生效同时将 key/value converter 配置为StringConverter并关闭 schema分别验证三条路径通过spec.state设置为STOPPED/RUNNING停止与恢复连接器通过spec.state设置为PAUSED/RUNNING暂停与恢复连接器文档明确说明spec.state的优先级高于spec.pause验证逻辑阻塞连接器后清空 FileSink 目标文件向 topic 生产新消息并确认消息不会出现在目标文件中assertThrows等待超时解除阻塞后再次确认消息最终写入目标文件证明连接器真正恢复消费。spec.state的定义位于 AbstractConnectorSpec.java L129-L146其字段注释明确说明连接器应处于的状态默认为 running。ConnectorState枚举提供RUNNING、PAUSED、STOPPED三种取值。testConnectorOffsetManagement连接器 Offset 的 list/alter/reset 三步曲这是 ConnectST 中最具代表性的高级用例ConnectST.java L1809-L1939完整验证连接器 Offset 管理功能的三类操作其核心机制是strimzi.io/connector-offsets注解定义于 ResourceAnnotations.java L91。前置配置KafkaConnector 的 spec 中需要同时指定两处 ConfigMap 引用——.spec.listOffsets.toConfigMaplist 阶段写入的 ConfigMap与.spec.alterOffsets.fromConfigMapalter 阶段读取的 ConfigMap。在 AbstractConnectorSpec.java L148-L184 中可以看到这两个字段的类型分别为ListOffsets与AlterOffsets。测试中两者指向同一个名为cluster-offsets的 ConfigMap。执行流程部署 KafkaNodePools、Kafka、启用连接器资源的 KafkaConnectFile Sink 插件、KafkaConnector以及用于访问 Connect API 的 scraper Pod 与 NetworkPolicy生产并消费 100 条消息等待连接器 Offset 更新list 阶段为 KafkaConnector 打上strimzi.io/connector-offsetslist注解等待 Cluster Operator 创建包含 offsets 的 ConfigMap随后解析 ConfigMap 数据中offsets.json的/offsets/0/offset/kafka_offset断言其等于消息数 100alter 阶段先把 ConfigMap 中的kafka_offset改为 20再先将连接器状态置为STOPPEDConnectST.java L1896-L1897更新 ConfigMap 后打上strimzi.io/connector-offsetsalter注解等待注解被 Operator 移除表示 alter 完成然后调用 Connect API 的 offsets 端点确认真实 Offset 已变为 20reset 阶段连接器仍处于 stopped 状态打上strimzi.io/connector-offsetsreset注解等待注解移除后再次查询 Connect API断言 offsets 返回[]空。该测试揭示了 Offset 管理的两个关键约束连接器必须处于 stopped 状态才能执行 alter 与 reset。文档明确指出若连接器未停止Connector 的 status 部分会抛出警告且 Offset 不会被更新每次操作后strimzi.io/connector-offsets注解会被 Operator 移除因此注解的消失是操作完成的判定信号。testConnectorTaskAutoRestart任务失败自动重启该用例ConnectST.java L764-L853验证 Kafka Connect 任务失败后的自动重启功能注意该测试标注了MicroShiftNotSupported因为它依赖 Connect Build 功能通过spec.build.plugins使用 Connect Build 特性从 Jar 工件含 sha512 校验和构建 EchoSink 插件创建 EchoSink KafkaConnectorspec 中启用autoRestartnew AutoRestartBuilder().withEnabled().build()对应 AutoRestart.javatasksMax1并通过fail.task.after.records5让连接器在处理第 5 条消息时主动抛错分批发送消息的设计细节先发 4 条、再发 1 条。测试代码注释解释了原因——若一次性发送全部消息整批会在put()方法中处理并抛异常可能导致 Offset 无法提交进而引发 EchoSink 任务无限重启验证点任务首次失败后立即被自动重启autoRestartCount变为 1waitForConnectorAutoRestartCount(..., 1)任务恢复到 RUNNING 状态并稳定约 2 分钟后自动重启计数被重置为 0。该机制的价值在于spec.autoRestart让 Operator 能够在连接器任务因瞬时错误失败时自动恢复而无需人工介入同时避免无限重启风暴。认证与安全场景testConnectScramShaAuthWithWeirdUserName 与 testConnectTlsAuthWithWeirdUserName特殊用户名两个用例分别验证 SCRAM-SHA-512 与 TLS 认证下使用怪异用户名的端到端可用性SCRAM 用例ConnectST.java L1155-L1230用户名包含点号且长度约 92 字符超过常规 64 字符约束。Kafka 集群使用 TLS listener SCRAM-SHA-512 认证KafkaConnect 通过kafkaClientAuthenticationScramSha512引用该用户并配置信任证书TLS 用例ConnectST.java L1059-L1134用户名包含点号且恰为 64 字符KafkaConnect 使用 TLS 客户端认证certificateAndKey指向用户 Secret 中的user.crt/user.key。两者均通过 FileStreamSink 连接器验证消息链路使用对应认证方式的生产者向 topic 写入消息然后检查 Kafka Connect FileSink 文件中出现了等量消息。这些用例验证了 Strimzi 对 Kafka 用户名的边界约束处理点号与长度不会破坏认证与连接器消费链路。testKafkaConnectWithPlainAndScramShaAuthenticationPLAIN/SCRAM 混合认证该用例ConnectST.java L273-L353验证 Kafka Connect 在 Plain 与 SCRAM-SHA 认证下的功能Kafka 仅开放 internal plain listener端口 9092认证方式为 SCRAM-SHA-512KafkaConnect 通过kafkaClientAuthenticationScramSha512携带用户凭据用户名 引用用户 Secret 中 password 字段的PasswordSecretSource并显式清空 TLS 配置通过 NetworkPolicy 放开 scraper Pod 对 KafkaConnect 的访问经 Connect REST API 创建 FileStreamSink 连接器Kafka 客户端以 SCRAM-SHA-PLAIN 方式生产/消费消息最后验证 FileSink 文件中消息数量符合预期。testSecretsWithKafkaConnectWithTlsAndScramShaAuthentication / TLS 客户端认证两个用例分别验证 Kafka Connect 组合使用 TLS 传输加密与 SCRAM-SHA-512ConnectST.java L661-L741、以及 TLS 传输加密与 TLS 客户端认证ConnectST.java L551-L638的完整链路Kafka 集群配置 TLS listener端口 9093认证分别为 SCRAM-SHA-512 与 TLSKafkaConnect 配置tls.trustedCertificates引用cluster-cluster-ca-certSecret 中的ca.crt与对应的客户端认证TLS 客户端认证用例还额外验证Operator 会为 KafkaConnect 创建内部 truststore SecretKafkaConnectResources.internalTlsTrustedCertsSecretName其中包含合并后的ca.crt且该数据非空通过 scraper Pod 创建 FileStreamSink 连接器TLS 客户端完成生产与消费最终验证消息落入 FileSink。testKafkaConnectWithScramShaAuthenticationRolledAfterPasswordChanged改密后滚动更新该用例ConnectST.java L1659-L1761验证 SCRAM-SHA 密码变更后 Kafka Connect 的滚动更新与功能连续性通过KafkaUser.spec.authentication.password.valueFrom.secretKeyRef将用户密码绑定到外部 Secretcustom-pwd-secret的pwd键密码明文长度满足 FIPS 要求KafkaConnect 使用该 SCRAM-SHA-512 用户凭据部署等待 REST API 可用创建包含新密码的第二个 Secretnew-custom-pwd-secret更新 KafkaUser 引用新 Secret由于密码引用的 Secret 内容变化Operator 检测到变化后触发 KafkaConnect 滚动更新waitTillComponentHasRolledAndPodsReady滚动完成后再次等待 REST API 可用证明组件在改密 滚动后仍正常工作。这一场景直接对应生产运维中的轮换密码操作密码存放在 Secret 中更新 Secret 即可驱动组件滚动无需删除资源。testMountingSecretAndConfigMapAsVolumesAndEnvVars配置挂载该用例ConnectST.java L1445-L1641验证 Secret 与 ConfigMap 可同时以卷和容器环境变量两种方式挂载到 Kafka Connect创建常规与含点号名称connect.config.map、connect.secret的 ConfigMap/Secret 各一通过spec.template.pod.volumes定义四个卷Secret 卷、ConfigMap 卷、含点号 Secret/ConfigMap 卷注意含点号资源的卷名称必须去掉点号Kubernetes 命名约束测试中使用doted-configmap-volume-name等替代名通过spec.template.connectContainer.env与valueFrom.secretKeyRef/configMapKeyRef将 Secret 键与 ConfigMap 键注入环境变量验证方式在 Connect Pod 内执行printenv检查环境变量值、cat挂载路径下的文件内容均与预期一致含点号名称的 Secret/ConfigMap 同样生效。缩放场景testKafkaConnectScaleUpScaleDown常规扩缩容该用例ConnectST.java L499-L529验证 Kafka Connect 副本的常规缩放以 1 副本部署 KafkaConnect确认初始副本数为 1通过KafkaConnectUtils.replace将spec.replicas更新为 4等待 4 个 Pod 全部就绪并断言 Pod 数量再将 replicas 改回 1等待就绪并断言 Pod 数量恢复。整个缩放过程由 Cluster Operator 协调完成无需人工管理 Pod。testScaleConnectWithConnectorToZero 与 testScaleConnectWithoutConnectorToZero缩到零副本两个用例验证缩容到 0 副本这一特殊场景的差异无连接器场景ConnectST.java L1250-L12752 副本 KafkaConnect 直接缩到 0等待 Connect 状态变为 ReadyKafkaConnectUtils.waitForConnectReady断言 Pod 数量为 0 且 status condition 为Ready有连接器场景ConnectST.java L1298-L13412 副本 KafkaConnect 上运行着 KafkaConnector缩到 0 后KafkaConnect 本身仍为 Ready0 副本但 KafkaConnector 进入NotReady且 condition message 包含 has 0 replicas——因为连接器依赖 Connect 集群提供运行环境。这组用例澄清了一个重要的行为边界缩容到零时 KafkaConnect 资源本身可以保持 Ready但挂载其上的 KafkaConnector 会因无工作节点而转为 NotReady。testScaleConnectAndConnectorSubresourcesubresource 缩放该用例ConnectST.java L1361-L1427验证通过 Kubernetes 标准scalesubresource 进行声明式缩放使用kubectl scale KafkaConnect.kafka.strimzi.io name --replicas4等价操作测试通过scaleByName实现缩放 KafkaConnect验证 Pod 数、spec.replicas、status.replicas均为 4且observedGeneration大于缩放前的值证明协调确实执行使用同样的方式对 KafkaConnector 执行 scale将tasks.max缩放为 4验证spec.tasksMax与status.tasksMax均为 4最后通过 Connect REST API逐个 Connect Pod检查连接器实际配置中的tasks.max也变为 4。此用例说明 Strimzi 为 KafkaConnect 与 KafkaConnector 都实现了标准的/scalesubresource从而可以无缝接入 HPA、kubectl scale等 Kubernetes 原生缩放工作流。多节点集群与端到端消息链路testMultiNodeKafkaConnectWithConnectorCreation多节点 Connect 集群该用例ConnectST.java L980-L1041验证多节点 Kafka Connect 集群上的连接器创建与消息处理且属于CONNECTOR_OPERATOR与ACCEPTANCE级别创建 3 副本的 KafkaConnect并显式指定自定义的groupId与 config/offset/status 存储 topic 名称避免多集群冲突创建 FileStreamSink KafkaConnector 后通过GET /connectors/name/status从 Connect API 获取worker_id解析出实际承载该连接器任务的 Connect Pod 名部署生产者/消费者客户端完成消息生产消费后在解析出的具体 Pod 上验证 FileSink 文件包含全部消息。testKafkaConnectAndConnectorFileSinkPluginFile Sink 插件与 REST 配置核对该用例ConnectST.java L374-L428在并行命名空间下验证 File Sink 插件的功能与连接器配置的透明度部署带 File 插件的 KafkaConnect开启use-connector-resources、scraper Pod、NetworkPolicy创建名为license-source的 FileStreamSource KafkaConnectortasksMax2topic 指向测试 topic通过 scraper Pod 执行curl service/connectors/license-source断言响应 JSON 中包含连接器名、connector.classorg.apache.kafka.connect.file.FileStreamSourceConnector、tasks.max2与正确的 topic同时创建消费者客户端消费消息确认数据链路完整。测试基础设施与运行机制ConnectST 体现了 Strimzi 系统测试systemtest 模块的核心基础设施设计理解这些机制有助于复现与扩展测试TestStorage每个ParallelNamespaceTest用例都会初始化一个 TestStorage 实例集中管理命名空间、集群名、broker/controller 节点池名、用户名、topic 名、消息数等上下文保证并行命名空间之间相互隔离KafkaNodePools所有用例均基于 KRaft 模式创建独立的 broker 与 controller 节点池替代传统的 ZooKeeper 集群scraper Pod 与 NetworkPolicyscraper Pod 作为测试探针通过curl访问 KafkaConnect 的 REST API默认端口 8083同时部署 NetworkPolicy 以放开 scraper 到 KafkaConnect 的访问路径模拟可控的网络边界KubeResourceManager统一管理资源的创建、更新与等待就绪并在测试结束后完成清理strimzi.io/use-connector-resources注解在 KafkaConnect 上启用该注解后KafkaConnector 自定义资源才会被 Connector Operator 管理见 ResourceAnnotations.java L61ConnectST 中大多数连接器用例都依赖它验证工具类VerificationUtils镜像、标签、探针、资源配额、JVM 参数、KafkaConnectUtilsREST API 可用性、FileSink 文件消息数、KafkaConnectorUtilsOffset 查询、注解等待、任务状态、autoRestart 计数等封装了各测试共用的断言逻辑。小结ConnectST 覆盖了 Kafka Connect 从部署、滚动更新、连接器状态控制、Offset 管理、自动重启、安全认证到缩放的完整功能矩阵。对照源码可以看到每个测试用例都对应着可验证的用户行为strimzi.io/manual-rolling-update触发滚动、strimzi.io/connector-offsets注解驱动 Offset 的 list/alter/reset、spec.state控制连接器启停、spec.autoRestart保障任务自愈、spec.template定制容器环境与挂载、scalesubresource 接入原生 Kubernetes 缩放。对于希望在真实集群中验证 Strimzi Kafka Connect 行为的开发者而言这套测试既是行为规范文档也是可直接参考的端到端操作清单对应实现位于 ConnectST.javaCRD 模型定义见 KafkaConnectorSpec.java 与 AbstractConnectorSpec.java注解常量统一维护在 ResourceAnnotations.java。【免费下载链接】strimzi-kafka-operatorApache Kafka® running on Kubernetes项目地址: https://gitcode.com/GitHub_Trending/st/strimzi-kafka-operator创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表