ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Strimzi Kafka Operator:Topic Operator 容量压测(testCapacity)全解析——从 CR 配置到指标采集与报告

Strimzi Kafka Operator:Topic Operator 容量压测(testCapacity)全解析——从 CR 配置到指标采集与报告 Strimzi Kafka OperatorTopic Operator 容量压测testCapacity全解析——从 CR 配置到指标采集与报告【免费下载链接】strimzi-kafka-operatorApache Kafka® running on Kubernetes项目地址: https://gitcode.com/GitHub_Trending/st/strimzi-kafka-operator本文以 Strimzi 系统测试框架中的TopicOperatorPerformance测试套件文档为主线逐步骤还原testCapacity容量测试的完整执行流程并结合topic-operator模块源码与systemtest模块实现深入讲解批处理参数STRIMZI_MAX_BATCH_SIZE、MAX_BATCH_LINGER_MS、STRIMZI_MAX_QUEUE_SIZE的实际作用、指标采集机制以及性能数据的持久化方式。读完本文你可以理解该仓库如何测量 Topic Operator 可管理的 KafkaTopic 数量上限以及如何在自己的环境中参照这套方法复现容量压测。一、测试套件定位文档骨架关联文档 TopicOperatorPerformance 是 Strimzi 系统测试文档development-docs/systemtests/下由测试注解自动生成中的一篇其核心内容如下套件描述Test suite for measuring Topic Operator capacity and performance limits.用于测量 Topic Operator 容量与性能上限的测试套件标签topic-operator —— 该标签页 topic-operator 汇总了所有覆盖 KafkaTopic 资源管理的测试包括 topic 的创建、更新、删除行为、配额quota强制执行与指标校验testCapacity即其中之一。测试方法testCapacity—— This test measures the maximum capacity of KafkaTopics that can be managed by the Topic Operator by incrementally creating topics until failure.通过递增式创建 topic 直至失败测量 Topic Operator 可管理的 KafkaTopic 最大容量。文档给出的 6 个执行步骤构成全文骨架| Step | Action | Result | | - | - | - | | 1. | 部署一个 Kafka 集群Topic Operator 按指定 batch size 与 linger time 配置。 | Kafka 集群与 Topic Operator 部署就绪。 | | 2. | 开始采集 Topic Operator 指标。 | 指标采集中。 | | 3. | 每批创建 100 个 KafkaTopic每个 topic 12 分区、3 副本。 | Topics 创建完成并进入 Ready 状态。 | | 4. | 持续批量创建直到 Topic Operator 调谐reconcile失败。 | 达到最大容量失败被检测到。 | | 5. | 使用 TestLogCollector 并指定自定义资源列表收集范围化日志pods、deployments、configmaps、Kafka CR避免收集成千上万个 KafkaTopic CR 和 Secret。 | 日志被收集用于定位瓶颈。 | | 6. | 清理所有 KafkaTopic并持久化性能指标。 | 命名空间被清理性能数据保存到 topic-operator 报告目录。 |下文将逐步骤对照源码 TopicOperatorPerformance.java 展开。二、集群拓扑与 Kafka CR 配置步骤 1测试使用双角色dual-roleNodePool 部署 KafkaBroker 池KafkaNodePoolTemplates.brokerPoolPersistentStorage(...)3 个副本brokerReplicas 3持久化存储 50Gi PVCwithSize(50Gi)、withDeleteClaim(true)Controller 池KafkaNodePoolTemplates.controllerPoolPersistentStorage(...)3 个副本controllerReplicas 35Gi PVC。随后创建的 Kafka CRKafkaTemplates.kafkaWithMetrics在spec.entityOperator.topicOperator层面做了三处关键配置源码见 TopicOperatorPerformance.java.editEntityOperator() .editTopicOperator() .withReconciliationIntervalMs(10_000L) // 周期调谐间隔 10 秒 .endTopicOperator() .editOrNewTemplate() .editOrNewTopicOperatorContainer() // Finalizers ensure orderly and controlled deletion of KafkaTopic resources. // In this case we would delete them automatically via ResourceManager .addNewEnv().withName(STRIMZI_USE_FINALIZERS).withValue(false).endEnv() .addNewEnv().withName(STRIMZI_ENABLE_ADDITIONAL_METRICS).withValue(true).endEnv() .addNewEnv().withName(STRIMZI_MAX_QUEUE_SIZE).withValue(maxQueueSize).endEnv() // Integer.MAX_VALUE .addNewEnv().withName(STRIMZI_MAX_BATCH_SIZE).withValue(maxBatchSize).endEnv() // 100 .addNewEnv().withName(MAX_BATCH_LINGER_MS).withValue(maxBatchLingerMs).endEnv() // 100 .endTopicOperatorContainer() .endTemplate() .endEntityOperator()每个环境变量的含义、默认值可以在 Topic Operator 的配置类 TopicOperatorConfig.java 中找到定义| 环境变量 | 作用 | 默认值源码 | 本测试取值 | | - | - | - | - | |STRIMZI_USE_FINALIZERS| 是否为 KafkaTopic 资源添加 finalizer |true|false| |STRIMZI_ENABLE_ADDITIONAL_METRICS| 是否启用针对外部服务Kafka、Kubernetes、Cruise Control请求的附加指标 |false|true| |STRIMZI_MAX_QUEUE_SIZE| topic 事件队列最大长度 |1024|Integer.MAX_VALUE| |STRIMZI_MAX_BATCH_SIZE| 单个批中允许的最大 topic 事件数 |100| 参数化默认100| |MAX_BATCH_LINGER_MS| 批在开始处理前最多等待的毫秒数 |100源码中配置键名为STRIMZI_MAX_BATCH_LINGER_MS | 参数化默认100|这里有两个值得注意的实现细节为何关闭 finalizer源码注释说明 finalizer 用于保证 KafkaTopic 的有序、受控删除但容量测试中的清理是由 ResourceManager 统一删除的关闭 finalizer 可以加快大规模删除速度与第六节的批量删除策略呼应。环境变量命名差异从源码结构看TopicOperatorConfig.java 中注册的配置键是STRIMZI_MAX_BATCH_LINGER_MS默认 100ms而测试代码与官方文档 con-tuning-topic-request-batches.adoc 中使用的都是MAX_BATCH_LINGER_MS缺少STRIMZI_前缀。复现该测试时建议以实际部署版本中 Topic Operator 容器实际读取的键名为准两个性能测试类 TopicOperatorPerformance 与 TopicOperatorScalabilityPerformance 均沿用MAX_BATCH_LINGER_MS这一写法。批处理参数在官方文档 con-tuning-topic-request-batches.adoc 中的描述与测试参数直接对应Topic Operator 利用 Kafka Admin API 的请求批处理能力来操作 topic 资源STRIMZI_MAX_QUEUE_SIZE默认 1024设置事件队列上限STRIMZI_MAX_BATCH_SIZE默认 100限制单批事件数MAX_BATCH_LINGER_MS默认 100ms控制批的等待时间。若请求批队列超过上限Topic Operator 会直接关闭并重启文档因此建议根据典型负载调整STRIMZI_MAX_QUEUE_SIZE。而本容量测试恰恰把队列上限拉到Integer.MAX_VALUE目的就是排除队列溢出触发重启这一变量让测试真正跑向容量极限。另外测试还会部署一个scraper podScraperTemplates.scraperPod它从外部抓取 Topic Operator 的 Prometheus 指标端点是第二步指标采集的基础设施。三、参数化批处理配置的扫描入口testCapacity是一个ParameterizedTest配置来源是provideConfigurationsForCapacity()源码见 TopicOperatorPerformance.javaprivate static StreamArguments provideConfigurationsForCapacity() { return Stream.of( Arguments.of(100, 100) // Default configuration ); }每对参数依次为maxBatchSize与maxBatchLingerMs注释明确说明该方法的意图是测试不同 max batch sizes 和 linger times 对系统吞吐量和响应性的影响当前只跑默认的平衡配置。这也意味着该测试天然支持扩展往流中追加更多参数组即可扫描不同批处理配置下的容量边界。测试方法上标注了Tag(TOPIC_CAPACITY)标签便于按标签筛选执行。四、核心压测循环批量创建直到失败步骤 3、4指标采集启动后测试进入无限循环按固定参数批量创建 topicfinal int batchSize 100; final int topicPartitions 12; final int topicReplicas 3; final int minInSyncReplicas 2; while (true) { // Endless loop int start successfulCreations; int end successfulCreations batchSize; try { // 通过 Kubernetes API 批量创建 KafkaTopic CR KafkaTopicScalabilityUtils.createTopicsViaK8s(namespace, clusterName, topicName, start, end, topicPartitions, topicReplicas, minInSyncReplicas); // 等待本批全部进入 Ready 状态 KafkaTopicScalabilityUtils.waitForTopicStatus(namespace, topicName, start, end, CustomResourceStatus.Ready, ConditionStatus.True); successfulCreations batchSize; LOGGER.info(Successfully created and verified batch from {} to {}, start, end); } catch (WaitException e) { LOGGER.error(Failed to create Kafka topics from index {} to {}: {}, start, end, e.getMessage()); // 收集范围化日志后退出循环 ... break; } }这段循环体现了容量测试的方法论以 Ready 状态为验收标准每批不仅创建 CR还通过waitForTopicStatus等待所有 topic 的 condition 达到ReadyTrue即 topic 真正在 Kafka 集群中创建完成失败即边界WaitException来自 kubetest4j 的等待框架意味着某一批 topic 在超时窗口内没有就绪此时successfulCreations的值就是本次运行测得的容量边界——该值会作为输出指标OUT: Successful KafkaTopics Created写入报告创建路径选择createTopicsViaK8s通过 Kubernetes API 创建 KafkaTopic CR由 Topic Operator 调谐落地这与测量 Topic Operator 管理容量的目标一致。五、失败时的范围化日志采集步骤 5失败分支中使用了带自定义资源列表的TestLogCollector源码见 TopicOperatorPerformance.javathis.logCollector TestLogCollector.of( // Pod logs and descriptions are always collected by LogCollector automatically. // Here we scope the additional namespaced resources to only deployments, configmaps, and Kafka CRs, // avoiding the default full set which would include thousands of KafkaTopic CRs. TestLogCollector.defaultLogCollectorBuilder() .withNamespacedResources( TestConstants.DEPLOYMENT.toLowerCase(Locale.ROOT), TestConstants.CONFIG_MAP.toLowerCase(Locale.ROOT), Kafka.RESOURCE_SINGULAR ).build()); this.logCollector.collectLogs();这对应文档步骤 5 的关键设计Pod 日志与描述由 LogCollector 自动收集额外采集的命名空间资源被刻意收窄为deployments、configmaps 和 Kafka CR三类。因为容量测试结束时命名空间里可能有数千个 KafkaTopic CR 和对应 Secret若使用默认的全量资源列表日志包会大到无法分析。这一范围化设计正是定位瓶颈例如 Topic Operator 的调谐日志、Kafka CR 状态而不被海量 topic CR 淹没的关键。六、清理策略与结果持久化步骤 6finally块保证无论成功与否都会执行清理与报告ListKafkaTopic kafkaTopics CrdClients.kafkaTopicClient().inNamespace(namespace).list().getItems(); KubeResourceManager.get().deleteResourceAsyncWait(kafkaTopics.toArray(new KafkaTopic[0])); KafkaTopicUtils.waitForTopicWithPrefixDeletion(namespace, topicName);源码注释解释了为何一次性批量删除而非逐个删除作者观察到逐个删除 KafkaTopic 时每个资源可能产生约 10 秒的删除延迟在数千 topic 的规模下清理时间不可接受deleteResourceAsyncWait并发发起删除并统一等待配合第五节提到的关闭 finalizer共同加速了大规模回收。随后构建报告属性并写入报告目录REPORT_DIRECTORY topic-operatoruse case 名为PerformanceConstants.GENERAL_CAPACITY_USE_CASE即capacityUseCaseperformanceAttributes.put(TOPIC_OPERATOR_IN_MAX_QUEUE_SIZE, maxQueueSize); performanceAttributes.put(TOPIC_OPERATOR_IN_MAX_BATCH_SIZE, maxBatchSize); performanceAttributes.put(TOPIC_OPERATOR_IN_MAX_BATCH_LINGER_MS, maxBatchLingerMs); performanceAttributes.put(TOPIC_OPERATOR_OUT_SUCCESSFUL_KAFKA_TOPICS_CREATED, successfulCreations); performanceAttributes.put(METRICS_HISTORY, topicOperatorMetricsGatherer.getMetricsStore()); topicOperatorPerformanceReporter.logPerformanceData(testStorage, performanceAttributes, REPORT_DIRECTORY / PerformanceConstants.GENERAL_CAPACITY_USE_CASE, TimeHolder.getActualTime(), Environment.PERFORMANCE_DIR);报告目录结构由 TopicOperatorPerformanceReporter.java 解析以capacityUseCase/max-batch-size-N-max-linger-time-M-with-clients-bool形式命名使不同批处理配置的结果可以并列对比。所有 IN/OUT 常量如IN: MAX BATCH SIZE (ms)、OUT: Successful KafkaTopics Created、Metrics History集中定义在 PerformanceConstants.java。七、指标采集体系5 秒轮询一次的操作与 JVM 全景第二步Start collecting Topic Operator metrics的实现在测试中对应this.topicOperatorCollector new TopicOperatorMetricsCollector.Builder() .withScraperPodName(this.testStorage.getScraperPodName()) .withNamespaceName(this.testStorage.getNamespaceName()) .withComponent(TopicOperatorMetricsComponent.create(namespace, clusterName)) .build(); this.topicOperatorMetricsGatherer TopicOperatorMetricsCollectionScheduler .getInstance(this.topicOperatorCollector, strimzi.io/cluster clusterName); this.topicOperatorMetricsGatherer.startCollecting();采集器 TopicOperatorMetricsCollector.java 依赖STRIMZI_ENABLE_ADDITIONAL_METRICStrue打开的附加指标抓取的核心指标组名称见 PerformanceConstants.java操作耗时sum 与 max 两种聚合strimzi_create_topics_duration_seconds、strimzi_delete_topics_duration_seconds、strimzi_describe_topics_duration_seconds、strimzi_alter_configs_duration_seconds、strimzi_describe_configs_duration_seconds、strimzi_update_status_duration_seconds、strimzi_create_partitions_duration_seconds调谐层strimzi_reconciliations_duration_seconds、strimzi_reconciliations_total/successful_total/failed_total/locked_total、strimzi_reconciliations_max_queue_size、strimzi_reconciliations_max_batch_sizefinalizer 操作strimzi_add_finalizer_duration_seconds、strimzi_remove_finalizer_duration_seconds资源数strimzi_resourceskindKafkaTopic直接反映受管 topic 数量曲线JVM 与系统jvm_gc_memory_allocated_bytes_total、jvm_memory_used_bytes、jvm_threads_live_threads、system_cpu_usage、process_cpu_usage等。调度器 TopicOperatorMetricsCollectionScheduler.java 以PerformanceConstants.DEFAULT_METRICS_POLLING_INTERVAL_SEC 5秒为周期轮询并按时间戳存入metricsStoreMapLong, MapString, ListDouble。它还计算了两个派生指标strimzi_total_time_spend_on_uto_event_queue_duration_seconds用调谐总耗时减去各内部操作create/delete/describe/alter/update status/finalizer 等耗时之和估算事件在队列中等待的时间——这是判断瓶颈在排队还是在执行的关键量system_load_average_per_core1 分钟负载均值除以 CPU 核数。采集选择器strimzi.io/cluster clusterName用于从 Prometheus 文本中按集群标签过滤出目标 Topic Operator 实例的指标。八、测试生命周期与结果表格套件级前置BeforeAll中通过SetupClusterOperator.getInstance().withDefaultConfiguration().install()安装 Cluster OperatorBeforeEach为每个测试实例创建TestStorage命名空间为TestConstants.CO_NAMESPACE上下文。套件级后置AfterAll中调用BasePerformanceMetricsParser.main(new String[]{PerformanceConstants.TOPIC_OPERATOR_PARSER})把test-performance-metrics文件中的记录解析渲染成结果表格供人工与 CI 对比。九、姊妹测试可扩展性场景延伸阅读同一标签下的 TopicOperatorScalabilityPerformance对应文档 TopicOperatorScalabilityPerformance与本文的容量测试互补testScalability不是测能管多少 topic而是对 10/100/500/1000 个 topic 各起一个线程并发执行 create-modify-delete 完整生命周期测量的是吞吐N 个 topic 全部完成所需总时间而非单 topic 延迟。两者共用同一套 Reporter 与解析器区别在于 use case 名分别为capacityUseCase与scalabilityUseCase且可扩展性测试显式设置了 768Mi 内存 / 750m CPU 的资源上下限以固定变量。十、小结这套容量测试可复用的方法论综合文档与源码testCapacity给出了一套可迁移的 Operator 容量压测方法固定集群拓扑与 Operator 参数双角色 NodePool、10s 调谐间隔、队列上限拉满、批参数参数化把容量极限与队列溢出重启两类失败解耦——后者在 con-tuning-topic-request-batches.adoc 中是明确记录的运维行为以 Ready 状态而非 API 接受为准每批验证后再继续失败批次的首个索引即为容量边界全程 5 秒粒度采集操作耗时、调谐队列指标与 JVM 全景指标并用派生指标区分排队时间与执行时间范围化日志 批量异步清理让数千资源规模的失败现场仍可诊断、可回收。关键参考路径汇总测试文档 TopicOperatorPerformance、测试实现 TopicOperatorPerformance.java、Operator 配置 TopicOperatorConfig.java、指标采集 TopicOperatorMetricsCollector.java、调度器 TopicOperatorMetricsCollectionScheduler.java、常量 PerformanceConstants.java、报告 TopicOperatorPerformanceReporter.java、调优文档 con-tuning-topic-request-batches.adoc。【免费下载链接】strimzi-kafka-operatorApache Kafka® running on Kubernetes项目地址: https://gitcode.com/GitHub_Trending/st/strimzi-kafka-operator创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表