ARTICLE DETAIL

资讯详情

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

深入解析 Strimzi Kafka 配额机制:基于 QuotasST 系统测试的磁盘保护与带宽限流实战

深入解析 Strimzi Kafka 配额机制:基于 QuotasST 系统测试的磁盘保护与带宽限流实战 深入解析 Strimzi Kafka 配额机制基于 QuotasST 系统测试的磁盘保护与带宽限流实战【免费下载链接】strimzi-kafka-operatorApache Kafka® running on Kubernetes项目地址: https://gitcode.com/GitHub_Trending/st/strimzi-kafka-operator导读本文以 Strimzi Kafka Operator 仓库中的QuotasST系统测试套件为主线完整剖析 Kafka 集群配额Quotas的两大核心能力通过strimzi配额插件保护磁盘空间、通过kafka内置插件做带宽与 CPU 限流。你将掌握Kafka与KafkaUser两种自定义资源中配额字段的配置方法、excludedPrincipals白名单机制的实现原理以及如何利用系统测试验证配额是否真正生效——所有结论均以仓库源码、测试用例与官方文档为依据。配额机制在 Strimzi 中的定位配额Quotas是 Kafka 多租户场景下最重要的保护手段之一它控制客户端吞吐量、限制 CPU 占用、约束分区变更频率并能在磁盘空间不足时阻止生产者继续写入从而避免一个失控客户端拖垮整个集群。在 Strimzi 中配额能力通过配额插件Quota Plugin暴露在Kafka自定义资源的spec.kafka.quotas字段中。从 QuotasPlugin.java 的源码可见该抽象类的 Javadoc 明确了两类插件kafka使用 Kafka 自带的配额插件QuotasPluginKafka提供静态的、按用户/按 broker 的限流能力strimzi使用 Strimzi 配额插件QuotasPluginStrimzi提供集群整体的聚合限流与磁盘保护能力。两类插件互斥同一时间只能启用一个kafka插件为默认启用启用strimzi插件后内置插件即被禁用见 con-choosing-a-quota-plugin.adoc。配额能力分为两个层面集群级配额在Kafka资源的quotas字段配置作用于所有客户端与用户级配额在KafkaUser资源的spec.quotas字段配置作用于单个用户可覆盖 broker 级默认值。本文的 QuotasST 测试关注前者但两者共享 Kafka 底层相同的 quota 机制。测试套件总览QuotasST 在验证什么QuotasST.java 是 Strimzi 系统测试System TestsST中专门验证配额插件的测试类归属kafka测试标签见 labels/kafka.md。该类包含两个测试方法| 测试方法 | 验证重点 | | - | - | |testKafkaQuotasPluginIntegration|strimzi插件基于磁盘可用字节数minAvailableBytesPerVolume的配额强制与excludedPrincipals白名单行为 | |testKafkaQuotasPluginWithBandwidthLimitation|strimzi插件基于带宽producerByteRate/consumerByteRate的限流效果对比普通用户与白名单用户的生产耗时 |前置条件为什么 Minikube 上跑不通测试套件在类级 Javadoc 中明确给出重要提示QuotasST 无法在 Minikube以及可能使用本地存储的其他集群上正确运行。原因在于磁盘空间的计算基于本地存储而 Minikube 的本地存储可能在多个 Docker 容器之间共享导致当前已用存储的计算失真。因此运行该套件必须使用具备独立、真实存储的 Kubernetes 集群。这与磁盘配额测试的语义直接相关testKafkaQuotasPluginIntegration需要通过真实的持久化卷PVC来观察剩余可用字节数低于阈值后生产者被停止的现象本地共享存储会破坏这一前提。运行前置步骤无论执行哪个测试都需要先完成统一准备| 步骤 | 动作 | 预期结果 | | - | - | - | | 1 | 以默认配置部署 Cluster OperatorSetupClusterOperator.getInstance().withDefaultConfiguration().install() | Cluster Operator 正常部署 |测试类在BeforeAll中执行这一步骤每个测试方法使用ParallelNamespaceTest在独立命名空间并行运行。测试一磁盘空间配额强制testKafkaQuotasPluginIntegration该测试验证的是strimzi配额插件最独特的能力——磁盘空间保护。核心思路配置一个极低的最小可用字节数阈值让普通生产者在达到阈值后被迫停止而白名单用户不受影响。测试执行步骤| 步骤 | 动作 | 结果 | | - | - | - | | 1 | 跳过 Minikube / MicroShift 环境assumeFalse(cluster.isMinikube() || cluster.isMicroShift()) | 集群适合测试 | | 2 | 创建 KafkaNodePoolsbroker controller各 1 副本与持久化存储 Kafka 集群配置配额插件、字节速率限制与白名单 principal | Kafka 与 KafkaNodePools 创建成功配额与白名单正确应用 | | 3 | 不使用任何用户身份发送消息 | 生产者达到最小可用字节数后停止 | | 4 | 检查 Kafka broker 日志中的配额强制消息 | 日志包含预期的配额强制消息 | | 5 | 使用白名单用户发送消息 | 消息正常发送未触发配额 | | 6 | 清理资源 | 资源删除成功 |关键资源配置解读测试通过 Java Builder API 构造 Kafka CR等价于以下 YAML 配置字段均来自 QuotasPluginStrimzi.java 的定义apiVersion: kafka.strimzi.io/v1beta2 kind: Kafka metadata: name: my-cluster spec: kafka: # ... config: client.quota.callback.static.storage.check-interval: 5 # 每 5 秒检查一次磁盘 quotas: type: strimzi producerByteRate: 10000000 # 10 MB/s 聚合生产限流 consumerByteRate: 10000000 # 10 MB/s 聚合消费限流 minAvailableBytesPerVolume: 800000000 # 每个卷剩余可用字节低于该值时停止生产 excludedPrincipals: - User:my-user # 白名单 principal 绕过配额 listeners: - name: plain port: 9092 type: internal tls: false - name: scramsha port: 9095 type: internal tls: false authentication: type: scram-sha-512字段语义与取值边界从 QuotasPluginStrimzi.java 的源码注释可以逐项确认字段含义producerByteRate/consumerByteRatebroker 级的字节速率配额与客户端数量无关。当多个非白名单生产者以最大速度生产时配额在所有生产客户端间均分否则按各客户端实际生产速率比例动态分配源码描述为 If clients produce at maximum speed, the quota is shared equally... Otherwise, the quota is divided based on each clients production rate。字段类型Long最小值为 0。minAvailableBytesPerVolume当任意存储卷的可用字节数 ≤ 该值时停止消息生产。字段类型Long最小值为 0。minAvailableRatioPerVolume以百分比小数形式如0.1表示 10%设置可用空间比例阈值语义同上。类型Double取值范围[0, 1]源码标注Minimum(0)与Maximum(1)。excludedPrincipals被豁免配额的用户列表。源码特别强调 principal必须以User:前缀开头例如User:my-user;User:CNmy-other-user。测试中构造的excludedPrincipal User: username正是这一规则的直接体现。磁盘检查频率由 Kafka 原生配置控制注意测试在spec.kafka.config中额外设置了client.quota.callback.static.storage.check-interval: 5。这是 Strimzi 配额插件暴露给 Kafka broker 的静态存储检查配置单位为秒用于控制磁盘可用空间的轮询频率。该参数直接决定磁盘配额强制的响应速度——间隔越短磁盘告警后生产者停止得越快但 broker 的检查开销也越高。如何验证配额真的生效测试用三个步骤形成完整的证据链生产者任务失败信号构造一个发送 100,000 条消息、每条 10,000 个#字符约 10 KB的生产者 Job并配置delivery.timeout.ms10005、request.timeout.ms10000、linger.ms5使失败快速暴露。随后通过JobUtils.waitForJobContainingLogMessage(..., Failed to send messages)等待 Job 日志中出现发送失败记录——这正是配额强制后的表现。Broker 日志取证抓取 broker Pod 日志并断言其中包含below the limit of 800000000字样源码第 141 行构造的日志匹配串。这句日志直接来自配额插件在磁盘剩余字节低于阈值时输出的警告是配额被触发的硬证据。白名单对照切换 bootstrap 地址到scramsha监听器端口 9095使用 SCRAM-SHA-512 认证的白名单用户身份重新发送相同规模的消息通过ClientUtils.waitForClientSuccess断言全部成功——证明同一磁盘条件下白名单用户完全不受配额约束。测试二带宽限流效果验证testKafkaQuotasPluginWithBandwidthLimitation第二个测试验证strimzi插件的带宽限流能力核心手段是对比普通用户与白名单用户发送同样数量消息的耗时。测试执行步骤| 步骤 | 动作 | 结果 | | - | - | - | | 1 | 设置白名单 principal | principal 已设置 | | 2 | 创建 KafkaNodePools 与持久化 Kafka启用配额并配置字节速率限制与白名单 | Kafka 资源创建成功 | | 3 | 创建 Kafka topic 与 SCRAM-SHA 认证用户 | 创建成功 | | 4 | 普通用户发送消息并计时 | 记录耗时 | | 5 | 白名单用户发送消息并计时 | 记录耗时 | | 6 | 断言普通用户耗时 白名单用户耗时 | 断言通过 |限流强度的关键设计与第一个测试不同这里把配额压到了极端值.withProducerByteRate(1000L) // 1 kB/s .withConsumerByteRate(1000L) // 1 kB/s而消息总量为 100 条 × 2,000 个#约 2 KB/条合计约 200 KB。在 1 kB/s 的配额下普通用户理论上需要约 200 秒才能发完而白名单用户未认证或使用 SCRAM 认证不受限流可在数秒内完成。测试通过System.currentTimeMillis()前后差值分别计算durationNormal与durationExcluded最终断言assertThat(Time taken for normal user should be greater than time taken for excluded user, durationNormal, Matchers.greaterThan(durationExcluded));这个测试设计揭示了strimzi配额插件的另一个特性带宽配额是聚合性的、动态分配的见 con-choosing-a-quota-plugin.adoc 中的说明。例如设置了 40 MBps 的生产者总配额若某生产者只用了 10 MBps其余 30 MBps 会自动释放给其他生产者——而非静态均分。这也是它和kafka插件按用户静态限流的本质区别。监听器与认证设计两个测试都配置了双监听器plain9092无认证用于匿名发送验证配额强制与scramsha9095SCRAM-SHA-512 认证用于白名单用户发送。这种设计确保了无用户 → 受配额约束与白名单用户 → 不受约束两组实验可以在同一个集群上并行对比且能通过认证身份明确区分 principal。白名单机制excludedPrincipals 的语义与坑excludedPrincipals是 Strimzi 配额插件区别于 Kafka 内置配额的核心能力之一适用于需要运维账号不受限、业务账号受限的场景例如监控、备份、管理类客户端。使用时有三个关键点均有源码依据必须带User:前缀源码 QuotasPluginStrimzi.java 明确要求 principal 形如User:my-user或User:CNmy-other-user。多值之间用;分隔。与认证类型强绑定白名单按 principal 匹配因此匿名连接无认证无法命中白名单测试中白名单用户通过 SCRAM-SHA-512 认证建立身份后才能被豁免。豁免范围白名单同时豁免带宽限流与磁盘限流。测试一中白名单用户在磁盘低于阈值的情况下依然能成功发送全部消息测试二中白名单用户不受 1 kB/s 带宽约束。用户级配额KafkaUser 中的 spec.quotas虽然 QuotasST 聚焦集群级配额但理解用户级配额有助于建立完整认知。KafkaUserQuotas.java 定义了KafkaUser资源中可用的四个配额字段与kafka插件字段一一对应| 字段 | 类型 | 含义 | 最小值 | | - | - | - | - | |producerByteRate| Integer | 客户端组每秒可发布到 broker 的最大字节数按 broker 计 | 0 | |consumerByteRate| Integer | 客户端组每秒可从 broker 拉取的最大字节数按 broker 计 | 0 | |requestPercentage| Integer | 客户端组占用网络和 I/O 线程的 CPU 百分比上限 | 0 | |controllerMutationRate| Double | 每秒可接受的创建 topic、创建分区、删除 topic 变更速率按分区数累计 | 0 |用户级配额的配置示例根据 con-configuring-client-quotas.adoc配置方式如下apiVersion: kafka.strimzi.io/v1beta2 kind: KafkaUser metadata: name: my-user labels: strimzi.io/cluster: my-cluster spec: # ... quotas: producerByteRate: 1048576 # 每秒最多推送 1 MiB 数据 consumerByteRate: 2097152 # 每秒最多拉取 2 MiB 数据 requestPercentage: 55 # CPU 占用上限 55% controllerMutationRate: 10 # 每秒最多 10 次分区变更操作注意Strimzi 支持用户级user-level配额但不支持客户端级client-level配额——文档中对此有明确说明。底层实现User Operator 的 QuotasOperator用户级配额由 User Operator 中的 QuotasOperator.java 负责调协。其核心逻辑值得关注通过QuotasCache缓存各用户的当前配额基于 Kafka Admin API 查询按cacheRefresh周期刷新避免每次调协都访问 Kafka通过QuotasBatchReconciler进行**微批量micro-batching**提交ClientQuotaAlteration提升大量用户配额更新时的吞吐批量大小与阻塞时间由getBatchQueueSize、getBatchMaxBlockSize、getBatchMaxBlockTime配置reconcile()方法按期望为 null 当前为 null → NoOp / 期望为 null 当前非空 → 删除 / 期望非空 当前为 null → 新增 / 两者不同 → 更新 / 相同 → NoOp五路分支决策并同步维护缓存支持STRIMZI_IGNORE_USERS白名单匹配的用户即使 Kafka 中已有配额也会被忽略。集群级配额kafka 插件与 DefaultKafkaQuotasManager当使用默认的kafka配额插件时集群级默认配额通过 Cluster Operator 在 Kafka 中设置。配置示例来自 proc-setting-broker-limits.adocapiVersion: kafka.strimzi.io/v1beta2 kind: Kafka metadata: name: my-cluster spec: kafka: # ... quotas: type: kafka producerByteRate: 1000000 consumerByteRate: 1000000 requestPercentage: 55 controllerMutationRate: 50字段语义源码确认| 字段 | 类型 | 含义来自 QuotasPluginKafka.java | 最小值 | | - | - | - | - | |producerByteRate| Long | 每个客户端每秒可向每个 broker 发布的最大字节数超过即被节流按 broker 计 | 0 | |consumerByteRate| Long | 每个客户端每秒可从每个 broker 拉取的最大字节数超过即被节流按 broker 计 | 0 | |requestPercentage| Integer | 每个客户端可占用的 broker 网络与 I/O 线程 CPU 百分比上限按 broker 计 | 0 | |controllerMutationRate| Double | 每个 broker 对 create topic、create partition、delete topic 请求的每秒变更接受速率按创建/删除的分区数计量 | 0 |底层实现Admin API 调协默认用户配额Cluster Operator 通过 DefaultKafkaQuotasManager.java 将上述配置落地到 KafkaprepareQuotaConfigurationRequest()将 CR 中的字段转换为 Kafka Admin API 的ClientQuotaAlteration.Op列表key 分别为producer_byte_rate、consumer_byte_rate、request_percentage、controller_mutation_rate通过describeClientQuotas查询 Kafka 中当前的默认用户配额DEFAULT_USER_ENTITY即 USER 类型为 null 的实体shouldAlterDefaultQuotasConfig()做差异比较当前配额不存在且未配置任何字段时跳过若期望配置包含非空值或与当前不同则执行alterClientQuotas更新。值得注意的是当 CR 中使用的是strimzi插件而非kafka时reconcileDefaultUserQuotas会用空配额覆盖默认用户配额emptyQuotasPluginKafka()确保两插件互斥、不留残留配置。测试 DefaultKafkaQuotasManagerTest.java 对这一逻辑有完整覆盖。两类插件的选型决策综合 con-choosing-a-quota-plugin.adoc 与本文的测试证据选型建议如下选择strimzi插件对应 QuotasST 的验证场景适用场景防止 broker 磁盘耗尽minAvailableBytesPerVolume/minAvailableRatioPerVolume按每个磁盘卷独立生效跨所有客户端的聚合吞吐限制配额是总量而非人均吞吐配额在活跃客户端间动态共享豁免特定用户excludedPrincipals白名单。注意使用strimzi插件时只能看到聚合的配额指标看不到按客户端细分的指标。选择kafka插件适用场景按用户、按 broker 的静态限流broker 级限额可被用户级配额覆盖控制客户端请求的 CPU 占用requestPercentage限制分区变更速率controllerMutationRate。需要特别警惕默认kafka配额同样作用于内部组件包括 Topic Operator 和 Cruise Control。若配额过紧可能导致 Topic Operator 被节流、无法及时创建 topic或 Cruise Control 再平衡任务因吞吐不足而失败。文档给出的实操建议是监控 controller 指标显式设置足以容纳 topic 操作的controllerMutationRate小集群3 broker建议 producer/consumer 配额至少设为 1 KB/s更大或更活跃的集群应相应提高。磁盘保护阈值的两种配置strimzi插件提供两种互斥的磁盘阈值配置文档与源码均强调只能配置其中一个quotas: type: strimzi minAvailableBytesPerVolume: 500000000000 # 绝对字节数方式剩余 ≤ 500 GB 停止生产 # minAvailableRatioPerVolume: 0.1 # 比例方式可用空间 ≤ 10% 停止生产与上者互斥从测试到实践如何在真实集群复现与验证QuotasST 的测试方法完全可以迁移为生产环境的验证手段建议按以下步骤操作准备集群使用具备独立持久化存储的集群Kubernetes 或 OpenShift不要使用 Minikube/MicroShift共享本地存储会破坏磁盘计算。部署 Cluster Operator参照 install/cluster-operator 目录下的 YAML 安装。创建带配额插件的 Kafka 集群按本文测试一/测试二中的配置定义spec.kafka.quotas建议先设一个宽松阈值观察行为再逐步收紧。构造发送任务用kafka-console-producer或自建 producer如测试中的KafkaProducerClient大批量发送消息观察普通用户是否被节流/停止broker 日志是否出现below the limit of bytes之类的配额警告白名单用户是否畅通无阻。反向验证用户级配额为单个KafkaUser配置spec.quotas对比该用户与其他用户的实际吞吐差异。结语配额是 Kafka 集群稳定性的最后一道防线。Strimzi 通过strimzi与kafka双插件设计把防磁盘耗尽与防带宽/CPU 滥用两类诉求解耦分别以聚合动态限流和按用户静态限流的方式落地。QuotasST 这套系统测试用最直观的停不下来 / 停下来了 / 白名单没事三段式实验把配额机制的每个细节都变成了可验证的事实。无论是排查生产事故还是为新集群做容量规划理解本文的配置语义与测试证据链都能帮你少走弯路。【免费下载链接】strimzi-kafka-operatorApache Kafka® running on Kubernetes项目地址: https://gitcode.com/GitHub_Trending/st/strimzi-kafka-operator创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表