news 2026/9/17 5:06:08

Strimzi Kafka Operator 中 Topic Operator 扩展性性能测试实战:TopicOperatorScalabilityPerformance 测试套件全解

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Strimzi Kafka Operator 中 Topic Operator 扩展性性能测试实战:TopicOperatorScalabilityPerformance 测试套件全解

Strimzi Kafka Operator 中 Topic Operator 扩展性性能测试实战:TopicOperatorScalabilityPerformance 测试套件全解

【免费下载链接】strimzi-kafka-operatorApache Kafka® running on Kubernetes项目地址: https://gitcode.com/GitHub_Trending/st/strimzi-kafka-operator

本篇围绕 Strimzi Kafka Operator 系统测试框架中的TopicOperatorScalabilityPerformance性能测试套件展开,覆盖其“部署受控资源限额的 Topic Operator → 对 10/100/500/1000 个 Topic 做并发全生命周期压测 → 持久化并汇总性能数据”的完整流程。读完本文,你将掌握该套件如何配置 Topic Operator 的批处理参数、如何以“每个 Topic 一个线程”的模式测量吞吐(throughput),以及性能报告如何落盘与解析。

测试套件定位:测量的是吞吐而非延迟

TopicOperatorScalabilityPerformance是 Strimzi 系统测试(systemtest)模块中的性能测试套件,源码位于 TopicOperatorScalabilityPerformance.java,继承自AbstractST,并被标注为:

  • @Tag(PERFORMANCE)@Tag(SCALABILITY):用于 JUnit 5 标签过滤,只运行性能/扩展性相关测试;
  • @IsolatedTest:该测试需要独占式环境(独占命名空间与集群资源),避免与其他测试并行干扰测量结果;
  • 套件文档标签(@SuiteDoc)声明其描述为 “Test suite for measuring Topic Operator scalability under concurrent topic operations”,并归属 topic-operator 标签页——该标签页汇总了所有由 Topic Operator 管理的 KafkaTopic 相关测试。

套件中最关键的测试方法是testScalability()。它在源码的@TestDoc注解中明确声明:“This test measures throughput (time to process N topics in parallel), NOT latency (response time for a single topic)”——即测量的是“并行处理 N 个 Topic 的总耗时”(吞吐量),而不是“单个 Topic 请求的响应时间”(延迟)。这是理解整套测试数据解读的前提:报告中随 Topic 数量增长而增长的总时间,反映的是 Operator 侧事件队列与调谐(reconciliation)流水线在并发压力下的处理能力,而不是单次 API 调用的快慢。

测试前置条件:受控资源限额的集群与 Topic Operator 配置

文档中的 “Before test execution steps” 要求:部署一个 Topic Operator 带有特定资源限额与批处理配置的 Kafka 集群。这一前置条件由测试类的@BeforeAll setUp()方法(见 TopicOperatorScalabilityPerformance.java)完整实现,分为三步:

1. 安装 Cluster Operator 并创建命名空间

SetupClusterOperator.getInstance() .withDefaultConfiguration() .install(); suiteTestStorage = new TestStorage(KubeResourceManager.get().getTestContext(), TestConstants.CO_NAMESPACE);

2. 部署 Kafka 节点池(NodePool 模式)

测试使用较新的 KafkaNodePool 模板分别创建 3 个 broker 节点和 3 个 controller 节点(KRaft 架构下的 broker/controller 分离布局),两个节点池都设置相同的资源 request/limit:内存 768Mi、CPU 750m:

KafkaNodePoolTemplates.brokerPoolPersistentStorage(namespace, brokerPoolName, clusterName, 3) .editSpec() .withResources(new ResourceRequirementsBuilder() .addToLimits("memory", new Quantity("768Mi")) .addToLimits("cpu", new Quantity("750m")) .addToRequests("memory", new Quantity("768Mi")) .addToRequests("cpu", new Quantity("750m")) .build()) .endSpec() // 同样的资源限制也应用到 KafkaNodePoolTemplates.controllerPoolPersistentStorage(...)

固定 broker/controller 的资源规格是性能测试的关键手段——它排除了“资源充足所以跑得快”这一变量,使测量结果可以在相同硬件预算下横向对比。

3. 创建带 Entity Operator 的 Kafka 集群,并注入 Topic Operator 调优参数

KafkaTemplates.kafka(namespace, clusterName, 3)基础上,测试对spec.entityOperator.topicOperator做了两类配置:

a) 常规 Spec 配置:调谐间隔 10 秒、资源限制 768Mi/750m:

.editEntityOperator() .editTopicOperator() .withReconciliationIntervalMs(10_000L) // 10s 调谐间隔 .withResources(/* 768Mi / 750m */) .endTopicOperator()

b) 容器级环境变量(批处理三参数):通过entityOperator.template.topicOperatorContainer.env注入:

| 环境变量 | 测试中取值 | 对应源码常量 | 含义 | | - | - | - | - | |STRIMZI_MAX_BATCH_SIZE|100|maxBatchSize = 100| Topic Operator 事件批处理的最大批次大小 | |MAX_BATCH_LINGER_MS|100|maxBatchLingerMs = 100| 批次等待(linger)时间,单位毫秒 | |STRIMZI_MAX_QUEUE_SIZE|Integer.MAX_VALUE|maxQueueSize = Integer.MAX_VALUE| 事件队列最大长度,测试中不限制队列容量 |

这三个参数是 Topic Operator 处理外部事件(KafkaTopic 增删改)的背压/吞吐控制旋钮:队列无上限意味着压力全部转移到调谐循环的批处理效率上。测试结束时,这三个“输入参数”会作为属性写回性能报告,保证结果可追溯、可复现。

testScalability:多档位 Topic 数量的全生命周期并发压测

测试矩阵:10 / 100 / 500 / 1000 个事件

testScalability()的核心逻辑是对一组“事件批次”循环执行(源码见 TopicOperatorScalabilityPerformance.java):

private final List<Integer> eventBatches = List.of(10, 100, 500, 1000); eventBatches.forEach(numEvents -> { final int eventPerTask = 4; final int numberOfTasks = numEvents / eventPerTask; final int numSpareEvents = numEvents % eventPerTask; this.reconciliationTimeMs = TopicOperatorPerformanceUtils .processAllTopicsConcurrently(suiteTestStorage, numberOfTasks, numSpareEvents, 0); // finally 中执行清理 + 性能数据落盘 });

这里的“事件”与“Topic 生命周期”之间有一个换算关系:每个 Task(一个 Topic 的完整生命周期)按 4 个事件计(文档步骤中的 CREATE、MODIFY、DELETE,加上内部等待/确认开销;源码中以eventPerTask = 4换算)。因此实际并发创建的 Topic 数为:

| 事件档位numEvents| Topic 数numberOfTasks| 剩余事件numSpareEvents| | - | - | - | | 10 | 2 | 2 | | 100 | 25 | 0 | | 500 | 125 | 0 | | 1000 | 250 | 0 |

numSpareEvents用于补齐不足一个完整生命周期的零头事件,在并发执行层中以“已完成的空 future”占位,保证事件总数与档位一致。

六个执行步骤(继承自文档步骤表)

文档中列出的 6 个步骤在源码中一一对应:

  1. 为每个 KafkaTopic 生成一个并发线程,各自执行完整生命周期processAllTopicsConcurrently最终委托给共享的 PerformanceTestExecutorService.processResourcesConcurrently,在System.nanoTime()计时开始处为每个资源索引提交一个CompletableFuture.runAsync任务;
  2. CREATE:每个线程创建指定分区数与副本数的 KafkaTopic,并等待其状态就绪;
  3. MODIFY:更新 Topic 配置并等待调谐完成;
  4. DELETE:删除 KafkaTopic 并等待删除完成;
  5. 等待所有线程完成并统计总耗时CompletableFuture.allOf(futures).join()后返回Duration.ofNanos(...).toMillis(),即 N 个 Topic 全部走完生命周期的总时长(吞吐口径);
  6. 清理残余 Topic 并收集指标:在finally块中执行“安全网”清理(见下文),并把性能数据持久化到 topic-operator 报告目录。

并发模型:固定线程池与“每 Topic 一线程”

PerformanceTestExecutorService 是 Topic Operator 与 User Operator 性能测试共享的并发基础设施:

  • 线程池大小CONCURRENCY_HINT = Runtime.getRuntime().availableProcessors() * 10,即 CPU 核数的 10 倍,确保 250 个 Topic 的并发任务不会被线程池成为瓶颈;
  • 执行器是静态共享的单例,若检测到已shutdown会自动重建,并在stopExecutor()中执行优雅停机(先shutdown等待 5 秒,超时则shutdownNow);
  • 若配置了 warm-up 任务(本测试传0),结束后会额外执行 5 秒冷却(COOLDOWN_PERIOD_MS = 5_000)以降低批次间干扰。

单 Topic 生命周期:创建、修改配置、删除

每个 Topic 任务的实际动作由 TopicOperatorPerformanceUtils.performFullLifecycle 串联三步,且每一步都是“操作 + 等待收敛”:

创建performCreationWithWait):通过KafkaTopicScalabilityUtils.createTopicsViaK8s(namespace, clusterName, topicPrefix, start, end, 12, 3, 2)以 K8s API 创建 Topic(12 分区、3 副本、min.insync.replicas=2),随后waitForTopicStatus等待所有 Topic 的 CustomResource 状态达到Ready且 Condition 为True——这一步会真实触发 Topic Operator 的调谐。

修改performModificationWithWait):将一组 10 项 Topic 配置写入KafkaTopicSpec.config并等待调谐生效,配置集合在源码中定义为KAFKA_TOPIC_CONFIG_TO_MODIFY

Map.of( "compression.type", "gzip", "cleanup.policy", "delete", "min.insync.replicas", 2, "max.compaction.lag.ms", 54321L, "min.compaction.lag.ms", 54L, "retention.ms", 3690L, "segment.ms", 123456L, "retention.bytes", 9876543L, "segment.bytes", 321654L, "flush.messages", 456123L )

覆盖压缩、清理策略、副本一致性、段切分、保留策略、compact 滞后等典型 Topic 配置维度,waitForTopicsContainConfig会验证这些配置已实际反映到 Topic 上。

删除performDeletionWithWait):KafkaTopicUtils.deleteKafkaTopicsInRange+waitForTopicWithPrefixDeletion,等待按前缀的 Topic 全部消失。

一个实现细节值得注意:由于生命周期任务运行在线程池的新线程中,而 JUnit 的ExtensionContext是线程绑定的,源码在每一步操作前都调用KubeResourceManager.get().setTestContext(currentContext)重新注入上下文,否则会触发 NPE(源码注释对此有明确说明)。

清理安全网与性能数据落盘

文档第 6 步“Clean up any remaining topics and collect performance metrics”对应testScalability()中的finally块,它保证无论压测成功与否都会执行:

} finally { List<KafkaTopic> kafkaTopics = CrdClients.kafkaTopicClient() .inNamespace(suiteTestStorage.getNamespaceName()).list().getItems(); KubeResourceManager.get().deleteResourceAsyncWait(kafkaTopics.toArray(new KafkaTopic[0])); KafkaTopicUtils.waitForTopicWithPrefixDeletion(namespace, topicName); // 组装 performanceAttributes 并落盘 }

性能属性表(LinkedHashMap)记录本次运行的完整输入与输出:

| 属性常量(PerformanceConstants) | 本测试的取值 | | - | - | |TOPIC_OPERATOR_IN_MAX_QUEUE_SIZE|maxQueueSize(Integer.MAX_VALUE) | |TOPIC_OPERATOR_IN_MAX_BATCH_SIZE|100| |TOPIC_OPERATOR_IN_NUMBER_OF_TOPICS|numberOfTasks(2/25/125/250) | |TOPIC_OPERATOR_IN_NUMBER_OF_EVENTS|numberOfTasks * 3 + numSpareEvents| |TOPIC_OPERATOR_IN_MAX_BATCH_LINGER_MS|100| |TOPIC_OPERATOR_IN_PROCESS_TYPE|"TOPIC-CONCURRENT"(每 Topic 并发,区别于按批次并发) | |OPERATOR_OUT_RECONCILIATION_INTERVAL|reconciliationTimeMs(本档位的总调谐耗时,毫秒) |

随后TopicOperatorPerformanceReporter.logPerformanceData将数据写入REPORT_DIRECTORY("topic-operator") + "/" + GENERAL_SCALABILITY_USE_CASE("scalabilityUseCase")目录,根目录取自环境变量PERFORMANCE_DIR,默认值为${user.home}/../systemtest/target/performance/(见 Environment.java)。

报告目录的命名规则

TopicOperatorPerformanceReporter.resolveComponentUseCasePathDir 会根据输入属性拼出带语义的子目录,例如本测试numEvents=1000的一档,目录名形如:

topic-operator/scalabilityUseCase/max-batch-size-100-max-linger-time-100-with-clients-false-number-of-topics-250

从源码结构看,目录名依次编码:用例名(scalabilityUseCase)→ 最大批次大小(100)→ 最大 linger 时间(100ms)→ 是否启用客户端实例(本测试为 false)→ Topic 数量;只有当用例是TOPIC_OPERATOR_FIXED_SIZE_OF_EVENTS_USE_CASE时才会追加-process-type-topic-concurrent后缀。这种命名方式让同一硬件上不同批次参数、不同规模的结果天然形成可对比的文件树。

收尾:自动生成绩标汇总表

@AfterAll tearDown()在全部档位跑完后调用:

BasePerformanceMetricsParser.main(new String[]{PerformanceConstants.TOPIC_OPERATOR_PARSER});

即运行 BasePerformanceMetricsParser 的topic-operator解析器,把落盘的原始数据解析成指标表格输出到测试日志——对应文档中“collect performance metrics, including total reconciliation time”的结果项。

运行方式与解读要点

运行方式:该套件属于 systemtest 模块,可通过 systemtest/Makefile 与 systemtest/scripts/run_tests.sh 提供的系统测试入口执行,用 JUnit 5 标签performancescalability过滤到本套件。适用前提:需要一个可用的 Kubernetes 集群、可访问的 Strimzi Operator 与 Kafka 镜像,且因为@IsolatedTest的存在,该测试需要独占运行环境(它会自行安装 Cluster Operator、创建/销毁 Kafka 集群与数百个 Topic),不适合与功能测试混跑。

结果解读要点

  1. 同一硬件预算下(broker/controller/Topic Operator 统一 768Mi/750m),比较 2/25/125/250 个 Topic 各档的OPERATOR_OUT_RECONCILIATION_INTERVAL,可以观察吞吐随并发规模的增长曲线——线性增长说明瓶颈在单 Topic 处理耗时,明显超线性则说明队列、批处理或 K8s API 调用成为瓶颈;
  2. Topic Operator 的reconciliationIntervalMs=10s是吞吐的下限约束之一:事件被排队等待调谐轮次时,总耗时中会包含等待成分;
  3. 报告中的输入参数(batch size=100、linger=100ms、queue=MAX)与目录名一一对应,修改这三个环境变量后重跑即可得到同坐标系下的对照实验。

小结

TopicOperatorScalabilityPerformance套件用一套可复现的受控环境(统一资源限额、固定的批处理参数与调谐间隔)+ “每 Topic 一线程”的全生命周期并发模型 + 带语义的报告目录与自动汇总解析,回答了“Topic Operator 在 10 到 1000 个事件规模下能并行消化多快的 Topic 增删改流量”这一问题。核心源码链路为:测试入口 → 并发生命周期工具 → 共享线程池 → 报告写入 → 指标解析。若需要同类视角的 User Operator 压测,可对照 UserOperatorScalabilityPerformance 文档。

【免费下载链接】strimzi-kafka-operatorApache Kafka® running on Kubernetes项目地址: https://gitcode.com/GitHub_Trending/st/strimzi-kafka-operator

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/17 5:05:33

STM32CubeIDE Attach调试:不复位不烧录,直接接管运行中目标

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/17 5:03:56

Python零基础入门:条件循环与数据结构实战

1. 项目概述&#xff1a;零基础Python入门第三课"0基础Python-003"这个标题背后&#xff0c;是一个面向编程新手的Python入门系列课程。作为该系列的第三课&#xff0c;它通常承担着承前启后的关键作用——在学员掌握了基础语法和简单逻辑后&#xff0c;开始接触更贴…

作者头像 李华
网站建设 2026/9/17 5:03:54

WorkBuddy技术拆解:AI工作台的产品化壁垒与工程实践

WorkBuddy这段时间讨论度确实高。我从它刚火的时候开始折腾&#xff0c;装客户端、配本地模型、挂SkillHub里的各种技能&#xff0c;也把网上那些“从入门到精通”的实操手册翻了个遍。一个很直接的感受是&#xff1a;大家把它想得太神秘了。拆到技术层&#xff0c;它就是LLM调…

作者头像 李华