Strimzi Kafka Operator 的 User Operator 缓存机制:用本地缓存化解 Kafka Admin API 逐用户查询的性能困境

Strimzi Kafka Operator 的 User Operator 缓存机制:用本地缓存化解 Kafka Admin API 逐用户查询的性能困境 Strimzi Kafka Operator 的 User Operator 缓存机制用本地缓存化解 Kafka Admin API 逐用户查询的性能困境【免费下载链接】strimzi-kafka-operatorApache Kafka® running on Kubernetes项目地址: https://gitcode.com/GitHub_Trending/st/strimzi-kafka-operatorUser Operator 需要持续比对 Kafka 中每个用户的 ACL、Quotas 与 SCRAM-SHA 凭据的当前状态但 Kafka Admin API 并不支持批量 describe 查询逐用户查询又慢又重。Strimzi 的 User Operator 为此设计了一套基于ConcurrentHashMap的周期性全量缓存机制以用户名为键用单次 Admin API 调用拉取全量数据、原子替换本地 Map再由各 Operator 在对账写入时同步回写缓存。本文基于仓库中 cache 包的设计文档 与完整实现源码拆解这套缓存的动机、架构、刷新机制、对账联动与已知局限。问题背景为什么 User Operator 需要缓存Kafka 的 Admin API 在查询当前状态这件事上并不友好。设计文档开篇就点明了核心矛盾The Kafka Admin API does not make it easy to query it to find the current state of things. It does not allow to batch the describe requests for some of the things the UO cares about, so we cannot use the micro-batching as we do for update requests.具体来说写请求可以微批处理读请求却不行。User Operator 对 ACL 的增删请求已经通过 micro-batching reconciler见 batching 包合并下发但 describe 类查询无法这样批量合并。逐用户查询代价高。Kafka Admin API 支持查询全量数据所有用户的所有 Quotas、所有用户的所有 ACLs、所有带 SCRAM-SHA 凭据的用户但如果按用户逐个查询速度慢、对 Broker 压力大还会导致超时。因此 User Operator 的策略是利用 Admin API 支持全量查询的能力一次请求拿到所有数据本地缓存按用户按键检索。这正是 DESIGN.md 描述的整个方案。总体架构一个抽象基类 三个具体缓存缓存代码位于io.strimzi.operator.user.operator.cache包共 4 个核心文件文件职责AbstractCache.java泛型抽象基类提供定时刷新、原子替换、通用读写方法AclCache.java缓存每个用户的 ACL 规则集合QuotasCache.java缓存每个用户的客户端配额ScramShaCredentialsCache.java标记每个用户是否存在 SCRAM-SHA 凭据三个具体缓存都是AbstractCacheT的子类只是缓存值的类型不同AclCache extends AbstractCacheSetSimpleAclRule一个用户可能有多条 ACL 规则缓存值是该用户全部规则的集合。设计文档特别指出AclCache会把一个用户的多条 ACL 规则聚合成缓存中的单个条目collate the ACL rules for a single user as a single item。QuotasCache extends AbstractCacheKafkaUserQuotas值为 Strimzi 的KafkaUserQuotas模型对象。ScramShaCredentialsCache extends AbstractCacheBoolean由于 Kafka 无法列出具体的 SCRAM 凭据值只能知道哪些用户名下有凭据所以值只是一个Boolean表示该用户凭据是否存在。这套结构由 KafkaUserOperator 统一编排其start()方法依次启动三个 OperatorquotasOperator、aclOperator、scramCredentialsOperator每个 Operator 在自己的构造器中创建对应缓存stop()时对称关闭。AbstractCache 的实现细节设计文档描述AbstractCache利用 Java 定时任务机制按可配置间隔周期刷新缓存。对照当前源码AbstractCache.java实际实现要点如下1. 单线程调度器驱动周期刷新构造器接收缓存名与刷新间隔毫秒创建一个以缓存名-cache命名的单线程ScheduledExecutorServicepublic AbstractCache(String name, long refreshIntervalMs) { this.refreshIntervalMs refreshIntervalMs; this.scheduledExecutor Executors.newSingleThreadScheduledExecutor(r - new Thread(r, name -cache)); }start()先同步执行一次initialize()即立即updateCache()做首次装载然后按固定速率调度后续刷新public void start() { initialize(); scheduledExecutor.scheduleAtFixedRate(this::updateCache, refreshIntervalMs, refreshIntervalMs, TimeUnit.MILLISECONDS); }也就是说每个缓存实例有自己独立的工作线程如ACL-cache、Quotas-cache、ScramShaCredentials-cache刷新互不阻塞。2. 原子替换而非逐条更新缓存本体是一个volatile的ConcurrentHashMap引用。刷新时不是原地修改旧 Map而是由子类loadCache()构建一个全新的 Map后整体换引用private void updateCache() { try { cache loadCache(); } catch (Exception e) { LOGGER.error({} failed to update, this.getClass().getSimpleName(), e); cache null; // Reset the cache } }这种构建新 Map volatile 原子换引用的方式保证了读线程对账线程要么看到完整的旧快照要么看到完整的新快照不会读到半成品。同时读操作内部仍依赖ConcurrentHashMap自身的并发安全。注意异常分支任何一次刷新失败都会把cache置回null此后所有读写方法都会抛出RuntimeException(XXX is not ready!)。这是刻意的 fail-fast 设计——与其拿着可能过期的数据继续对账产生错误操作不如让错误显式暴露。3. 通用访问方法基类提供get、getOrDefault、put、remove、keys五个方法全部带 not ready 检查stop()则关闭调度器并把cache置空public void stop() { scheduledExecutor.shutdownNow(); cache null; }唯一的抽象方法loadCache()留给子类实现约定返回一份包含最新数据的 Map。三个缓存各自的 loadCache一次 Admin API 调用拉全量AclCache全量 ACL 按用户聚合AclCache.loadCache() 用AclBindingFilter.ANY一次性描述全部 ACL超时 1 分钟KafkaFutureCollectionAclBinding futureAcls adminClient.describeAcls(AclBindingFilter.ANY).values(); CollectionAclBinding aclsBindings futureAcls.get(1, TimeUnit.MINUTES);然后遍历结果解析每条AclBinding的 principal只保留 principal 类型为User的规则按用户名computeIfAbsent聚合进HashSetSimpleAclRule。源码注释解释了初始容量的取值// Each user can have multiple ACL rules. So the size of the map will not directly correspond to the number // of rules. But we size it for 3-5 rules per user to give us at least some start ... ConcurrentHashMapString, SetSimpleAclRule map new ConcurrentHashMap(aclsBindings.size() / 3);即按每用户平均 3 条规则估算 Map 大小避免 Java 默认初始容量带来的反复扩容。QuotasCache全量 Quotas 与 NPE 陷阱QuotasCache.loadCache() 调用describeClientQuotas(ClientQuotaFilter.all())拿全部配额实体。这里有一个容易踩坑的细节被源码注释明确说明ClientQuotaFilter.all()返回中可能包含默认用户配额default user quota其USER键存在但值为null直接取值插入 Map 会抛 NPEif (entry.getKey().entries().get(ClientQuotaEntity.USER) ! null) { map.put(entry.getKey().entries().get(ClientQuotaEntity.USER), QuotaUtils.fromClientQuota(entry.getValue())); }ScramShaCredentialsCache只存存在性ScramShaCredentialsCache.loadCache() 调用describeUserScramCredentials().users()拿到用户名列表然后逐个写入Boolean值true。这印证了设计文档的说法凭据本身无法被列出缓存只能记录该用户是否有凭据这一事实。缓存与对账流程的联动写入方回写减少陈旧缓存的副作用设计文档强调了一点缓存虽由定时任务刷新但各 Operator 在对账成功写入 Kafka之后也会立即更新本地缓存目的是避免陈旧缓存引发无谓操作。文档给出的典型故障场景是缓存还没刷新认为某条 ACL 不存在于是 Operator 反复创建一条 Kafka 中实际已存在的 ACL形成无效操作循环。源码中这条联动链路清晰可见SimpleAclOperator 的reconcile()第一步就是cache.getOrDefault(username, Set.of())取当前状态与期望状态做差集决定走 NoOp / 创建 / 删除 / 增量更新哪条路径随后在internalCreate、internalUpdate成功后cache.put(username, desired)在internalDelete成功后cache.remove(username)。QuotasOperator 同样是cache.get(username)起步写 Kafka 成功后cache.put/cache.remove回写。ScramCredentialsOperator 用Boolean.TRUE.equals(cache.get(username))判断凭据是否已存在创建/删除后同步put(username, true)或remove(username)。SimpleAclOperator.getAllUsers() 则通过cache.keys()遍历缓存键找出所有持有 ACL 的用户并过滤掉ignoredUsersPattern匹配的用户、对CN前缀做解码。这样形成读走缓存、写后回写、定时兜底的闭环Operator 自己的写操作让缓存即时保持一致周期刷新则负责兜住缓存中不来自 Operator的变化例如用户直接在 Kafka 侧改了 ACL。刷新间隔的配置三个缓存共用同一个刷新间隔配置项定义在 UserOperatorConfig.java/** * Refresh interval for the cache storing the resources from the Kafka Admin API */ public static final ConfigParameterLong CACHE_REFRESH_INTERVAL_MS new ConfigParameter(STRIMZI_CACHE_REFRESH_INTERVAL_MS, LONG, 15000, CONFIG_VALUES);环境变量STRIMZI_CACHE_REFRESH_INTERVAL_MS类型毫秒LONG默认值15000即 15 秒读取入口UserOperatorConfig#getCacheRefresh()各 Operator 在构造缓存时通过config.getCacheRefresh()传入例如new AclCache(adminClient, config.getCacheRefresh())。调小该值可提高状态一致性、缩短陈旧窗口但会提高 Admin API 全量查询的频率调大则反之需要在 Broker 负载与收敛速度之间权衡。已知局限Admin API 没有分页DESIGN.md 的 Limitations 一节指出了一个规模化的硬约束Since we are currently using the Kafka Admin API to get all data in a single query, we might run into problems in big clusters where the response would not fit into a single response.由于所有数据压在一次请求/响应里ACL 风险最高——单个用户可能拥有大量 ACL 规则全量响应可能超过max.request.size/max.partition.fetch.bytes等消息大小限制Quotas 和 SCRAM-SHA 凭据的用户级数据量非常有限通常不是问题Kafka Admin API 本身不支持分页文档给出的当前唯一解法是调大 Broker 的消息大小限制如max.message.bytes相关配置这也是部署超大用户规模集群时值得提前评估的点。此外从源码还能补充两个隐含限制单次loadCache对 Admin 调用设置了 1 分钟超时超时即抛异常并触发前文所述的cache null降级而刷新失败期间的所有对账读操作都会因 not ready 而快速失败依赖上层对账重试机制兜底。未来方向从缓存构建 watch/informer摆脱周期对账设计文档的 Future possibilities 部分讨论了更进一步的可能把缓存做成类似 Kubernetes 的 watch/informer 机制事件驱动地通知 Operator从而彻底免除周期性轮询。但作者论证了当前这条路走不通SCRAM-SHA 凭据没有变更信号——Admin API 只能告诉你凭据是否存在无法告诉你它何时、为何变化只有存在性这一种状态这一缺失能力由上游 JIRA 问题 KAFKA-14356 跟踪即使未来对 Quotas 和 ACLs 实现了基于事件的监听由于 SCRAM-SHA 凭据仍需周期对账兜底整体架构也无法去掉定时刷新。结论是在 Kafka 补齐凭据变更通知之前周期刷新 写后回写的缓存方案就是 User Operator 的最优解。小结Strimzi User Operator 的缓存方案是一个典型的以空间换查询、以快照换一致设计值得借鉴的工程决策包括利用 API 的全量查询能力而非逐条查询单次describeAcls(ANY)/describeClientQuotas(all)/describeUserScramCredentials()拉全量摊薄了 N 个用户查询的成本volatile 整体换 Map 的原子快照读线程永远看到完整一致的旧或新数据无需加锁写后回写缓存Operator 自身的 Kafka 写操作成功后立即同步本地缓存把缓存滞后导致的无效重试窗口压到最小失败即置空的 fail-fast刷新失败宁可让后续操作快速报错也不用脏数据做对账决策明确的局限意识文档坦承无分页的 Admin API 在超大集群下的响应体风险并给出调大消息大小的现实解法同时论证了 watch/informer 方案因 SCRAM 凭据无变更信号而暂不可行。这套机制的实现集中在 user-operator/src/main/java/io/strimzi/operator/user/operator/cache 目录约 200 行基类 3 个各百行以内的实现配合 SimpleAclOperator、QuotasOperator、ScramCredentialsOperator 中的读写联动构成了 User Operator 与 Kafka 之间状态同步的核心路径。【免费下载链接】strimzi-kafka-operatorApache Kafka® running on Kubernetes项目地址: https://gitcode.com/GitHub_Trending/st/strimzi-kafka-operator创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考