Apache Pulsar 模块化负载管理器Modular Load Manager深入指南启用、验证与实现原理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址: https://gitcode.com/gh_mirrors/pulsar28/pulsarApache Pulsar 的负载均衡机制负责在多个 broker 之间合理分配 namespace bundle是集群稳定性的基石。本文基于 Pulsar 源码仓库中的 developing-load-manager.md 文档系统讲解模块化负载管理器Modular Load Manager的启用方式、验证手段、数据模型与流量分配策略并结合 ModularLoadManagerImpl.java 等核心源码与 conf/broker.conf 真实配置从配置、命令行到源码级原理逐层深入。读完本文你将掌握如何在集群中切换负载管理器、如何确认当前生效的实现、如何解读负载报告以及LeastLongTermMessageRate策略背后完整的决策逻辑。什么是模块化负载管理器模块化负载管理器Modular Load Manager是 Pulsar 对早期 SimpleLoadManagerImpl 的灵活替代方案。它试图简化负载管理方式同时通过抽象接口支持实现更复杂的负载管理策略。两者在架构上的本质区别在于SimpleLoadManagerImpl每个 broker 独立基于自己的资源使用情况做决策负载报告以systemResourceUsage嵌套结构写入 ZooKeeper。ModularLoadManagerImpl采用集中式centralized设计所有 bundle 分配请求无论是新出现的 bundle 还是曾经见过的 bundle统一由lead broker处理lead broker 随时间推移可能发生变化。判断当前 lead broker 的方法查看 ZooKeeper 中的/loadbalance/leader节点。核心实现在ModularLoadManagerImpl中它实现了 ModularLoadManager 接口。该接口定义了几类关键职责selectBrokerForAssignment(ServiceUnitId)作为 leader broker为给定的 bundle 寻找合适的 brokerdoLoadShedding()作为 leader broker选择需要卸载以重新分配的 bundlecheckNamespaceBundleSplit()作为 leader broker自动检测并拆分热点 bundledisableBroker()/start()/stop()/initialize(PulsarService)生命周期管理。从源码看ModularLoadManagerImpl 内部维护了本地数据LocalBrokerData、全集群负载数据LoadData、预分配 bundle 映射preallocatedBundleToBroker、broker 过滤器管道filterPipeline和负载卸载管道loadSheddingPipeline等核心状态。启用模块化负载管理器启用 Modular Load Manager 有两种方式任选其一即可。方式一修改 broker 配置文件编辑conf/broker.conf将loadManagerClassName参数的值从org.apache.pulsar.broker.loadbalance.impl.SimpleLoadManagerImpl改为org.apache.pulsar.broker.loadbalance.impl.ModularLoadManagerImpl在当前的仓库配置中conf/broker.conf 的默认值已经就是模块化实现# Name of load manager to use loadManagerClassNameorg.apache.pulsar.broker.loadbalance.impl.ModularLoadManagerImpl也就是说从当前仓库对应的版本起Modular Load Manager 已经成为默认负载管理器无需额外修改即可生效。方式二使用 pulsar-admin 动态修改配置使用pulsar-admin的update-dynamic-config命令无需重启 broker 即可切换$ pulsar-admin update-dynamic-config \ --config loadManagerClassName \ --value org.apache.pulsar.broker.loadbalance.impl.ModularLoadManagerImpl用同样的方法可以切回原来的值$ pulsar-admin update-dynamic-config \ --config loadManagerClassName \ --value org.apache.pulsar.broker.loadbalance.impl.SimpleLoadManagerImpl重要无论通过哪种方式只要loadManagerClassName的值填写错误Pulsar 都会回退到默认的SimpleLoadManagerImpl而不会直接报错因此在切换后务必按下一节的方法验证实际生效的实现。验证当前生效的负载管理器有以下几种方式可以确认集群当前正在使用哪种负载管理器。方法一通过 pulsar-admin 检查动态配置$ bin/pulsar-admin brokers get-all-dynamic-config { loadManagerClassName : org.apache.pulsar.broker.loadbalance.impl.ModularLoadManagerImpl }如果返回结果中没有loadManagerClassName元素说明当前使用的是默认负载管理器即SimpleLoadManagerImpl。方法二对比 ZooKeeper 中的负载报告结构两种负载管理器写入 ZooKeeper 的负载报告结构差异明显。模块化负载管理器的负载报告位于/loadbalance/brokers/...其中的systemResourceUsage子元素bandwidthIn、bandwidthOut等被提升到了顶层例如{ bandwidthIn: { limit: 10240000.0, usage: 4.256510416666667 }, bandwidthOut: { limit: 10240000.0, usage: 5.287239583333333 }, bundles: [], cpu: { limit: 2400.0, usage: 5.7353247655435915 }, directMemory: { limit: 16384.0, usage: 1.0 } }而简单负载管理器写入/loadbalance/brokers/...的报告则将这些指标嵌套在systemResourceUsage之下{ systemResourceUsage: { bandwidthIn: { limit: 10240000.0, usage: 0.0 }, bandwidthOut: { limit: 10240000.0, usage: 0.0 }, cpu: { limit: 2400.0, usage: 0.0 }, directMemory: { limit: 16384.0, usage: 1.0 }, memory: { limit: 8192.0, usage: 3903.0 } } }两者对比可以归纳为对比维度Modular Load ManagerSimple Load Manager资源指标位置顶层平铺cpu、memory、bandwidthIn等直接可见嵌套在systemResourceUsage之下字段差异负载报告为LocalBrokerData序列化结果负载报告为LoadReport额外包含overLoaded、underLoaded等标记报告内容含 bundle 列表、资源、消息速率等多维数据以系统资源使用为核心这一差异在源码层面有明确对应模块化实现直接使用LocalBrokerData作为 ZooKeeper 节点的数据模型ZooKeeper 路径为/loadbalance/brokers/broker host/port而简单实现在 SimpleLoadManagerImpl 中通过getSystemResourceUsage()组装SystemResourceUsage并调用loadReport.setSystemResourceUsage(...)写入。方法三对比 pulsar-perf monitor-brokers 输出命令行工具 broker monitor即pulsar-perf monitor-brokers会持续接收 broker 数据与负载报告其输出格式会因负载管理器实现的不同而有明显差别。模块化负载管理器的输出示例 ||SYSTEM |CPU % |MEMORY % |DIRECT % |BW IN % |BW OUT % |MAX % || || |0.00 |48.33 |0.01 |0.00 |0.00 |48.33 || ||COUNT |TOPIC |BUNDLE |PRODUCER |CONSUMER |BUNDLE |BUNDLE - || || |4 |4 |0 |2 |4 |0 || ||LATEST |MSG/S IN |MSG/S OUT |TOTAL |KB/S IN |KB/S OUT |TOTAL || || |0.00 |0.00 |0.00 |0.00 |0.00 |0.00 || ||SHORT |MSG/S IN |MSG/S OUT |TOTAL |KB/S IN |KB/S OUT |TOTAL || || |0.00 |0.00 |0.00 |0.00 |0.00 |0.00 || ||LONG |MSG/S IN |MSG/S OUT |TOTAL |KB/S IN |KB/S OUT |TOTAL || || |0.00 |0.00 |0.00 |0.00 |0.00 |0.00 || 简单负载管理器的输出示例 ||COUNT |TOPIC |BUNDLE |PRODUCER |CONSUMER |BUNDLE |BUNDLE - || || |4 |4 |0 |2 |0 |0 || ||RAW SYSTEM |CPU % |MEMORY % |DIRECT % |BW IN % |BW OUT % |MAX % || || |0.25 |47.94 |0.01 |0.00 |0.00 |47.94 || ||ALLOC SYSTEM |CPU % |MEMORY % |DIRECT % |BW IN % |BW OUT % |MAX % || || |0.20 |1.89 | |1.27 |3.21 |3.21 || ||RAW MSG |MSG/S IN |MSG/S OUT |TOTAL |KB/S IN |KB/S OUT |TOTAL || || |0.00 |0.00 |0.00 |0.01 |0.01 |0.01 || ||ALLOC MSG |MSG/S IN |MSG/S OUT |TOTAL |KB/S IN |KB/S OUT |TOTAL || || |54.84 |134.48 |189.31 |126.54 |320.96 |447.50 || 差异要点模块化实现按SYSTEM / COUNT / LATEST / SHORT / LONG五个区块组织输出其中 SHORT 与 LONG 对应短周期与长周期的历史聚合数据简单实现则按COUNT / RAW SYSTEM / ALLOC SYSTEM / RAW MSG / ALLOC MSG组织强调原始值与分配后值的对比。该命令的用法详见 reference-cli-tools.md$ pulsar-perf monitor-brokers --connect-string zk-connect-string主要参数--connect-string用于指定一个或多个 ZooKeeper 服务器的连接串。数据模型Modular Load Manager 如何组织负载数据模块化负载管理器监控的全部数据封装在LoadData类中。它把可用数据划分为两部分bundle 数据与broker 数据。Broker 数据broker 数据封装在BrokerData类中进一步细分为两部分本地数据Local Broker Data每个 broker 各自写入 ZooKeeper历史数据Historical Broker Data由 leader broker 写入 ZooKeeper。本地 Broker 数据本地 broker 数据封装在LocalBrokerData中提供以下资源信息CPU 使用率JVM 堆内存使用率Direct memory堆外直接内存使用率带宽入/出使用率所有 bundle 的最近总消息速率入/出主题topic、bundle、生产者producer、消费者consumer的总数量分配给该 broker 的所有 bundle 名称该 broker 最近一次 bundle 分配的变化。本地数据按照服务配置loadBalancerReportUpdateMaxIntervalMinutes周期性更新。每当某个 broker 更新其本地数据后leader broker 会通过ZooKeeper watch立即收到更新通知并从 ZooKeeper 节点/loadbalance/brokers/broker host/port读取本地数据。对应的配置默认值见 conf/broker.conf# Percentage of change to trigger load report update loadBalancerReportUpdateThresholdPercentage10 # minimum interval to update load report loadBalancerReportUpdateMinIntervalMillis5000 # maximum interval to update load report loadBalancerReportUpdateMaxIntervalMinutes15 # Frequency of report to collect loadBalancerHostUsageCheckIntervalMinutes1即在默认配置下负载报告在变化超过 10% 时触发更新、最短间隔 5 秒、最长间隔 15 分钟强制更新主机资源使用情况每分钟采集一次。系统资源的采集在 ModularLoadManagerImpl 初始化时根据操作系统选择实现Linux 上使用LinuxBrokerHostUsageImpl其他平台回退到GenericBrokerHostUsageImpl。历史 Broker 数据历史 broker 数据封装在TimeAverageBrokerData类中。为了在稳态场景做出良好决策、同时在临界场景做出快速响应历史数据被拆分为两部分短期数据short-term用于响应式决策长期数据long-term用于稳态决策。两个时间框架都维护以下信息整个 broker 的消息速率入/出整个 broker 的消息吞吐入/出。与 bundle 数据不同broker 数据不维护全局消息速率与吞吐量的样本——因为随着 bundle 被添加或移除这些全局值本身不会保持稳定。相反这些数据是通过对各个 bundle 的短期与长期数据聚合得到的聚合方式见下文 bundle 数据部分。历史 broker 数据由 leader broker 在任何 broker 向 ZooKeeper 写入本地数据时于内存中更新然后按照配置loadBalancerResourceQuotaUpdateIntervalMinutes周期性地写入 ZooKeeper默认值为 15 分钟见 conf/broker.conf。Bundle 数据bundle 数据封装在BundleData类中。与历史 broker 数据类似bundle 数据也分为短期与长期两个时间框架每个时间框架维护该 bundle 的消息速率入/出该 bundle 的消息吞吐入/出该 bundle 当前的样本数量。时间框架通过在一组有限数量的样本上维护均值来实现样本来自本地数据中的消息速率与吞吐值。在 ModularLoadManagerImpl 源码中样本数量是硬编码的常量// The number of effective samples to keep for observing long term data. public static final int NUM_LONG_SAMPLES 1000; // The number of effective samples to keep for observing short term data. public static final int NUM_SHORT_SAMPLES 10;因此如果本地数据更新间隔为 2 分钟、短样本数为 10、长样本数为 1000那么短期数据覆盖的时段为10 样本 × 2 分钟/样本 20 分钟长期数据覆盖的时段为1000 样本 × 2 分钟/样本 2000 分钟。当某个时间框架的样本不足时均值只基于现有样本计算当完全没有样本时使用默认值直到第一个样本覆盖它。当前默认值同样定义在 ModularLoadManagerImpl 中// Default message rate to assume for unseen bundles. public static final double DEFAULT_MESSAGE_RATE 50; // Default message throughput to assume for unseen bundles. // Note that the default message size is implicitly defined as // DEFAULT_MESSAGE_THROUGHPUT / DEFAULT_MESSAGE_RATE. public static final double DEFAULT_MESSAGE_THROUGHPUT 50000;即消息速率默认入/出各 50 条/秒消息吞吐默认入/出各 50KB/秒隐含的默认消息大小为 50KB / 50 1KB/条。这些默认值在initialize()阶段被写入NamespaceBundleStats defaultStats字段用于初始化未见过的 bundle 的历史数据。bundle 数据同样由 leader broker 在任何 broker 写入本地数据时于内存中更新并按照loadBalancerResourceQuotaUpdateIntervalMinutes配置默认 15 分钟与历史 broker 数据在同一时刻写入 ZooKeeper。模块化负载管理器在 ZooKeeper 中使用的关键路径见 ModularLoadManagerImpl 常量包括ZooKeeper 路径存储内容/loadbalance/brokers各 broker 的LocalBrokerData本地数据/loadbalance/broker-time-average各 broker 的TimeAverageBrokerData历史数据/loadbalance/bundle-data各 bundle 的BundleData数据/loadbalance/leader当前 lead broker 信息/loadbalance/resource-quota/namespace旧版 ResourceQuota 数据其中/loadbalance/brokers的根路径常量定义在 LoadManager.java 中LOADBALANCE_BROKERS_ROOT /loadbalance/brokers并被模块化实现复用。流量分配ModularLoadManagerStrategy 与 LeastLongTermMessageRate模块化负载管理器通过ModularLoadManagerStrategy提供的抽象做出 bundle 分配决策。策略在做决策时会综合考虑服务配置ServiceConfiguration、完整的负载数据LoadData以及待分配 bundle 的 bundle 数据。目前唯一支持的策略是LeastLongTermMessageRate其类注释明确说明这是基于哪个 broker 拥有最小长期消息速率来选择 broker 的放置策略。最小长期消息速率策略Least Long Term Message Rate正如其名称所示该策略试图将 bundle 分配到各 broker使每个 broker长期时间窗口内的消息速率大致相等。它通过ModularLoadManagerStrategy接口的selectBroker(candidates, bundleToAssign, loadData, conf)方法实现。评分公式对每个候选 broker策略通过getScore()计算得分final double overloadThreshold conf.getLoadBalancerBrokerOverloadedThresholdPercentage() / 100.0; final double maxUsage brokerData.getLocalData().getMaxResourceUsage(); if (maxUsage overloadThreshold) { // 返回 POSITIVE_INFINITY表示该 broker 过载不参与分配 }具体评分逻辑为对 broker 的预分配 bundle 数据preallocatedBundleData求所有长期消息速率入 出之和加上该 broker 自身的长期消息速率来自TimeAverageBrokerData本身是所有已分配 bundle 长期消息速率之和两者相加得到该 broker 的最终估计分数。资源权重与过载处理仅仅按消息速率均衡负载无法处理每条消息在不同 broker 上资源消耗不对称的问题。因此系统资源使用率——CPU、内存、direct memory、带宽入、带宽出——也被纳入分配过程。策略按以下方式对最终消息速率加权1 / (overload_threshold - max_usage)其中overload_threshold对应配置loadBalancerBrokerOverloadedThresholdPercentage默认 85即 85%max_usage是候选 broker 各系统资源中被利用的最高比例由LocalBrokerData.getMaxResourceUsage()计算。这个乘数确保在相同消息速率下承受更大资源压力的机器会被分配更少的负载。特别地它力求达到一种状态——如果一台机器过载那么所有机器都近似过载。处理规则与 LeastLongTermMessageRate 源码一致若 broker 的maxUsage超过过载阈值该 broker 得分为POSITIVE_INFINITY不会被考虑用于 bundle 分配同时记录一条包含 CPU/MEMORY/DIRECT/BANDWIDTH 各项使用率的告警日志得分最低的 broker 胜出多个 broker 得分相同平局时全部保留如果所有 broker 都过载即没有有限得分则从所有候选中随机分配一个 broker如果最终候选列表为空例如当前确实没有可用 broker返回Optional.empty()。策略在选择最佳 broker 时使用ThreadLocalRandom从并列最优的 broker 中随机挑选这有助于避免多个 bundle 连续压到同一台机器上。候选 broker 的过滤在策略打分之前候选集合还会经过 ModularLoadManagerImpl 中的 broker 过滤器管道filterPipeline例如BrokerVersionFilter以及持久化/非持久化主题可用性检查BrokerTopicLoadingPredicate确保只有符合版本与能力要求的 broker 进入评分环节。与负载均衡相关的配套配置除了loadManagerClassName之外conf/broker.conf 中还有一批与模块化负载管理器协同工作的关键配置理解它们有助于更好地调优集群# Enable load balancer loadBalancerEnabledtrue # Usage threshold to determine a broker as over-loaded loadBalancerBrokerOverloadedThresholdPercentage85 # Interval to flush dynamic resource quota to ZooKeeper loadBalancerResourceQuotaUpdateIntervalMinutes15 # Enable/disable automatic bundle unloading for load-shedding loadBalancerSheddingEnabledtrue # Load shedding interval loadBalancerSheddingIntervalMinutes1 # Prevent the same topics to be shed within this timeframe loadBalancerSheddingGracePeriodMinutes30 # load shedding strategy, support OverloadShedder and ThresholdShedder loadBalancerLoadSheddingStrategyorg.apache.pulsar.broker.loadbalance.impl.ThresholdShedder # Usage threshold to allocate max number of topics to broker loadBalancerBrokerMaxTopics50000关键说明loadBalancerBrokerOverloadedThresholdPercentage85即overload_threshold正是LeastLongTermMessageRate评分公式中用来判定 broker 过载的阈值loadBalancerLoadSheddingStrategy当前仓库版本2.10.0 之后默认使用ThresholdShedder它与LeastLongTermMessageRate配合分别负责卸载过载 broker 上的 bundle与新 bundle 的放置选择loadBalancerEnabledtrue负载均衡总开关若置为false则负载管理器不再执行分配与卸载。小结模块化负载管理器通过集中式的 leader 决策、清晰的LocalBrokerData / TimeAverageBrokerData / BundleData数据分层以及可插拔的ModularLoadManagerStrategy策略抽象将 Pulsar 的负载均衡从每台 broker 各自为政升级为全局视角统一调度。目前默认的LeastLongTermMessageRate策略在长期消息速率均衡的基础上叠加了系统资源使用率加权与过载保护兼顾了稳态公平与临界响应。实践中你可以通过pulsar-admin动态切换实现、通过动态配置或 ZooKeeper 负载报告结构快速确认生效的实现并借助pulsar-perf monitor-brokers持续观察集群负载分布。若需进一步深入源码推荐按以下顺序阅读ModularLoadManagerImpl.java整体实现包含 ZooKeeper 路径、样本常量、默认统计值与核心调度逻辑LeastLongTermMessageRate.java分配策略的评分与选择逻辑LocalBrokerData.java本地负载数据模型conf/broker.conf全部负载均衡相关配置及其默认值。【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址: https://gitcode.com/gh_mirrors/pulsar28/pulsar创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考