消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载在 Pulsar 集群中命名空间 bundle 如何分配到哪个 broker直接决定了消息流量的均衡程度与集群稳定性。Pulsar 的负载管理器load manager正是负责这一决策的组件。本文围绕 Pulsar 文档《Modular load manager》展开先讲清模块化负载管理器ModularLoadManagerImpl的定位与启用方式再给出三种可落地的验证手段最后深入到 ModularLoadManagerImpl.java 等源码解析其数据模型LocalBrokerData / TimeAverageBrokerData / BundleData、双时间窗口的采样机制以及LeastLongTermMessageRate放置策略的评分与过载处理逻辑。读完本文你不仅能切换并验证负载管理器还能从源码层面解释 Pulsar 该把 bundle 分给谁的每一步决策。一、模块化负载管理器是什么模块化负载管理器由 ModularLoadManagerImpl.java 实现它是早期 SimpleLoadManagerImpl.java 的灵活替代方案。两者的定位差异在于SimpleLoadManagerImpl实现相对简单直接基于 broker 上报的系统资源使用率CPU、内存、带宽等配合一组可配置的资源权重对候选 broker 打分排序选出剩余容量最大的 brokerModularLoadManagerImpl在简化负载管理方式的同时引入抽象层把候选 broker 过滤BrokerFilter、放置决策ModularLoadManagerStrategy、卸载策略LoadSheddingStrategy拆成可独立替换的组件便于后续实现更复杂的负载管理策略。从源码结构看这种模块化体现在构造阶段initialize()方法中依次装配了BundleSplitStrategybundle 拆分策略、ModularLoadManagerStrategy.create(conf)放置策略、filterPipeline含BrokerVersionFilter的 broker 过滤管线以及loadSheddingPipeline默认OverloadShedder可通过配置替换见 ModularLoadManagerImpl.java#L246-L284。另外要注意一个重要的架构特征模块化负载管理器是集中式的centralized。所有 bundle 分配请求——无论该 bundle 是首次出现还是之前已被分配过——都只会由lead broker领导者 broker可随时间变化处理。要查看当前 lead broker可检查 ZooKeeper 中的/loadbalance/leader节点。二、启用方式有两种方式启用模块化负载管理器方式一修改 broker.conf 静态配置在 conf/broker.conf 中将loadManagerClassName参数的值从org.apache.pulsar.broker.loadbalance.impl.SimpleLoadManagerImpl改为org.apache.pulsar.broker.loadbalance.impl.ModularLoadManagerImplloadManagerClassNameorg.apache.pulsar.broker.loadbalance.impl.ModularLoadManagerImpl说明在当前仓库版本的 ServiceConfiguration.java 中loadManagerClassName的默认值已经是org.apache.pulsar.broker.loadbalance.impl.ModularLoadManagerImpl即较新版本默认启用模块化负载管理器而在 2.3.0 版本语境下它是需要显式切换的替代方案。无论哪个版本配置错误时的行为一致一旦指定了不存在的类Pulsar 会回退到SimpleLoadManagerImpl。方式二使用 pulsar-admin 动态配置不重启 broker通过pulsar-admin动态更新$ pulsar-admin update-dynamic-config \ --config loadManagerClassName \ --value org.apache.pulsar.broker.loadbalance.impl.ModularLoadManagerImpl用同样的方法也可以改回原值。对应到 broker 端动态配置的接收入口在 BrokersBase.java 的updateDynamicConfiguration()方法而loadManagerClassName被注册为可动态变更的配置项并带有合法性校验器校验类名是否存在且实现LoadManager接口见 BrokerService.java#L2271-L2293。三、验证当前使用的负载管理器文档给出了三种验证手段均可直接照做。1. 查看动态配置项$ bin/pulsar-admin brokers get-all-dynamic-config { loadManagerClassName : org.apache.pulsar.broker.loadbalance.impl.ModularLoadManagerImpl }如果输出中没有loadManagerClassName元素则说明当前使用的是默认负载管理器即未做动态覆盖取ServiceConfiguration中该字段的默认值。2. 对比 ZooKeeper 中的负载报告两种负载管理器写入/loadbalance/brokers/...节点的负载报告Load Report结构不同这是最直观的区分方式模块化负载管理器的负载报告——系统资源使用项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 } }而 Simple 负载管理器的负载报告是嵌套在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 } } }从源码看这种结构差异源于两个实现写入的不同数据模型模块化负载管理器写的是 LocalBrokerData.java字段平铺另附 bundle 列表与每个 bundle 的最新统计而 Simple 负载管理器写的是LoadReport/SystemResourceUsage结构资源项嵌套在systemResourceUsage下。3. 观察 broker monitor 命令行输出pulsar-admin clusters ...之外的另一条线索是bin/pulsar-admin brokers相关的 monitor 命令其输出格式随负载管理器不同而变化。模块化负载管理器的示例输出包含 SYSTEM / COUNT / LATEST / SHORT / LONG 等区块 ||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 || Simple 负载管理器的示例输出包含 COUNT / RAW SYSTEM / ALLOC SYSTEM / RAW MSG / ALLOC MSG 等区块 ||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 || 两种输出中 SHORT / LONG或 ALLOC区块的语义正好对应下一节要讲的短期/长期双时间窗口数据。四、数据模型负载管理器监控了什么模块化负载管理器监控的数据整体包含在 LoadData.java 中按 broker 数据与 bundle 数据两大类组织。4.1 Broker 数据Broker 数据由 BrokerData.java 承载进一步细分为两部分本地数据Local Broker Data每个 broker 各自独立写入 ZooKeeper与历史 broker 数据Historical Broker Data由 lead broker 写入 ZooKeeper。本地 Broker 数据本地 broker 数据由 LocalBrokerData.java 承载提供以下资源信息CPU 使用率JVM 堆内存使用率直接内存Direct Memory使用率入/出带宽使用率所有 bundle 最近一次的总消息速率in/outtopic、bundle、producer、consumer 的总数量分配到本 broker 的所有 bundle 名称本 broker 最近发生的 bundle 分配变更本地数据的更新周期由服务配置loadBalancerReportUpdateMaxIntervalMinutes控制broker.conf 默认值为 15 分钟ServiceConfiguration.java 中同样默认 15。实际上写入并非纯定时源码 ModularLoadManagerImpl.java 的 needBrokerDataUpdate() 显示除了超过最大间隔强制写入外还会对比上一次写入的数据当最大资源使用率、消息速率、消息吞吐量或 bundle 数量中任意一项的变化幅度超过loadBalancerReportUpdateThresholdPercentagebroker.conf 默认 10%时提前写入。任一 broker 更新本地数据后lead broker 会通过 ZooKeeper watch 立即感知——本地数据读取自 ZooKeeper 节点/loadbalance/brokers/broker host/port即LoadManager.LOADBALANCE_BROKERS_ROOT下的每个 broker 锁节点。历史 Broker 数据历史 broker 数据由 TimeAverageBrokerData.java 承载。为兼顾稳态下的良好决策与危急场景下的快速反应历史数据被拆分为两部分短期数据用于反应式决策与长期数据用于稳态决策。两个时间窗都维护整个 broker 的消息速率in/out整个 broker 的消息吞吐量in/out与 bundle 数据不同broker 数据不单独维护全局消息速率与吞吐量的采样——因为这些值会随 bundle 的增删而天然漂移broker 级别的短期/长期值实际上是对其所承载 bundle 数据的聚合因此理解 broker 数据的前提是先理解 bundle 数据见下节。历史 broker 数据的更新链路是任一 broker 把本地数据写入 ZooKeeper 后lead broker 在内存中更新各 broker 的历史数据随后 lead broker 再按配置loadBalancerResourceQuotaUpdateIntervalMinutesbroker.conf 与 ServiceConfiguration.java 中默认均为 15 分钟周期性地把历史数据写回 ZooKeeper。源码中对应的实现是 ModularLoadManagerImpl.java 的 writeBundleDataOnZooKeeper()批量把所有 bundle 的BundleData与每个 broker 的TimeAverageBrokerData写入元数据存储broker 级路径为/loadbalance/broker-time-average/broker见 ModularLoadManagerImpl.java#L116。Bundle 数据Bundle 数据由 BundleData.java 承载与历史 broker 数据类似同样分为短期与长期两个时间窗。每个时间窗维护该 bundle 的消息速率in/out该 bundle 的消息吞吐量in/out该 bundle 当前的样本数时间窗的实现方式是在有限数量的样本集合上对速率/吞吐值求滑动平均样本取自本地数据中的消息速率与吞吐量。举例若本地数据更新间隔为 2 分钟、短期样本数为 10、长期样本数为 1000则短期数据覆盖10 samples * 2 minutes/sample 20 minutes长期数据同理覆盖 2000 分钟。当样本不足以填满某个时间窗时仅对已有样本求平均当完全没有样本时使用默认值直到第一个真实样本到来后被覆盖。源码 ModularLoadManagerImpl.java#L99-L110 确认了这些默认值与样本数常量// Default message rate to assume for unseen bundles. public static final double DEFAULT_MESSAGE_RATE 50; // Default message throughput to assume for unseen bundles. public static final double DEFAULT_MESSAGE_THROUGHPUT 50000; // 50KB/s public static final int NUM_LONG_SAMPLES 1000; public static final int NUM_SHORT_SAMPLES 10;即当前默认值为消息速率in/out50 msg/s消息吞吐量in/out50 KB/s50000 B/s隐含的默认消息大小 50000/50 1KB注意一个源码细节吞吐量默认值常量命名为DEFAULT_MESSAGE_THROUGHPUT 50000字节/秒即文档所述 50KB/s。bundle 数据同样由 lead broker 在任一 broker 写入本地数据时更新内存副本并按loadBalancerResourceQuotaUpdateIntervalMinutes周期性、与历史 broker 数据同时写入 ZooKeeper/loadbalance/bundle-data/bundle见 ModularLoadManagerImpl.java#L97。五、流量分配策略Least Long Term Message Rate模块化负载管理器通过 ModularLoadManagerStrategy.java 抽象做出 bundle 分配决策。策略的入口方法签名为OptionalString selectBroker(SetString candidates, BundleData bundleToAssign, LoadData loadData, ServiceConfiguration conf);即策略的决策输入包括服务配置、全量负载数据loadData以及待分配 bundle 自身的BundleData。当前唯一受支持的策略是 LeastLongTermMessageRate.javaModularLoadManagerStrategy.create(conf)工厂方法目前固定返回它接口注释表明未来允许用户注入自己的策略。Least Long Term Message Rate 策略原理顾名思义该策略试图把 bundle 分布到各 broker 上使每个 broker 在长期时间窗口内的消息速率大致相等。但仅按消息速率均衡无法处理每条消息在各 broker 上产生的非对称资源负担问题——同样的消息速率在资源吃紧的机器上代价更高。因此分配过程还纳入了系统资源使用率CPU、内存、直接内存、入带宽、出带宽文档给出的概念公式是按1 / (overload_threshold - max_usage)对最终消息速率进行加权其中overload_threshold对应配置loadBalancerBrokerOverloadedThresholdPercentagebroker.conf 默认 85即 85%max_usage是候选 broker 各项系统资源中的最大使用比例。这个乘子的效果是承受同等消息速率时资源负担更重的机器会少分负载从而尽量保证如果有一台机器过载那么所有机器都差不多同时过载。当某 broker 的 max usage 超过过载阈值时该 broker 不参与分配候选若所有 broker 都过载则随机分配。对照源码 LeastLongTermMessageRate.java#L53-L136实际打分逻辑与上述概念一致、实现上更可执行getScore()对每个候选 broker 打分若maxUsage overloadThreshold直接记为POSITIVE_INFINITY并打 warn 日志输出 CPU/MEMORY/DIRECT/BW IN/BW OUT 明细否则累加预分配 bundle 的长期消息速率preallocatedBundleData中各 bundle 的getLongTermData()消息速率 inout加上 broker 时间平均数据的长期速率getLongTermMsgRateIn() getLongTermMsgRateOut()得分越低越优先selectBroker()维护所有并列最优 broker 的列表最后随机选取其一避免所有请求扎堆到同一台分数最低的机器若所有候选得分都是无穷大全部过载回退为在全体候选中随机分配分配完成后ModularLoadManagerImpl会调用preallocateBundle()把 bundle 记入预分配表使其立即计入该 broker 的后续评分直到真实数据把它覆盖。此外ModularLoadManagerImpl.selectBroker(ServiceUnitId)在委托放置策略前还有一层候选收敛管线见 ModularLoadManagerImpl.java#L852-L921先按命名空间策略applyNamespacePolicies圈定候选再依次过滤 topic 数超过loadBalancerBrokerMaxTopics默认 50000的 broker、按 anti-affinity group 与故障域failure domain过滤、按命名空间维度剔除已承载最多的 broker最后再走BrokerFilter管线默认含BrokerVersionFilter若任何过滤把候选清空则回退到完整的候选集合。六、相关配置参数速查结合 conf/broker.conf 与 ServiceConfiguration.java与本文主题直接相关的参数如下默认值均来自当前仓库参数默认值作用loadManagerClassName...impl.ModularLoadManagerImpl指定负载管理器实现类配置错误时回退到SimpleLoadManagerImplloadBalancerReportUpdateMaxIntervalMinutes15本地 broker 数据写入的最大间隔分钟loadBalancerReportUpdateThresholdPercentage10本地数据相对上次写入的最大允许变化百分比超过则提前写入loadBalancerReportUpdateMinIntervalMillis5000本地数据两次写入之间的最小间隔毫秒loadBalancerResourceQuotaUpdateIntervalMinutes15lead broker 把 bundle 数据与历史 broker 数据周期性地写回 ZooKeeper 的间隔分钟loadBalancerBrokerOverloadedThresholdPercentage85过载判定阈值%max usage 超过该值的 broker 不参分配候选全部过载时随机分配loadBalancerHostUsageCheckIntervalMinutes1主机资源使用率的采集间隔分钟loadBalancerBrokerMaxTopics50000候选过滤阈值topic 数超过该值的 broker 被移出候选此外模块化负载管理器还自带一组周期性任务与可插拔策略均受对应开关控制load shedding默认OverloadShedder由loadBalancerSheddingEnabled、loadBalancerSheddingIntervalMinutes、loadBalancerSheddingGracePeriodMinutes等控制、自动 bundle 拆分loadBalancerAutoBundleSplitEnabled等、以及 Prometheus 侧的负载均衡指标如brk_lb_cpu_usage、brk_lb_memory_usage等见 ModularLoadManagerImpl.java#L1043-L1058。这些机制共享同一套 LoadData 数据底座是模块化设计数据与策略解耦的直接收益。七、小结与延伸阅读回顾全文的关键脉络启用静态改loadManagerClassName或用pulsar-admin update-dynamic-config动态切换配置错误一律回退SimpleLoadManagerImpl验证get-all-dynamic-config看配置、对比/loadbalance/brokers/...负载报告结构顶层平铺 vssystemResourceUsage嵌套、观察 monitor 输出格式差异数据本地数据各 broker 自写、watch 驱动 历史数据lead broker 聚合、双时间窗口滑动平均未观测 bundle 按 50 msg/s、50KB/s 默认 bundle 数据每 bundle 的短期/长期速率与吞吐决策集中式 lead broker 依据LeastLongTermMessageRate策略按长期消息速率 资源使用惩罚选 broker过载者出局、全过载则随机。想进一步深入建议按以下路径阅读源码与测试核心实现ModularLoadManagerImpl.java、BrokerData.java、BundleData.java、TimeAverageBrokerData.java策略与过滤LeastLongTermMessageRate.java、BrokerFilter.java、ModularLoadManagerStrategy.java行为验证ModularLoadManagerImplTest.java 覆盖了分配、预分配、负载报告与负载卸载等核心路径适合对照本文各结论逐一验证配置全集conf/broker.conf 中loadBalancer*系列参数理解这套数据上报 → 双窗口聚合 → 策略打分 → 预分配占位的闭环之后你就能针对具体集群的流量形态突发热点、大小 broker 混布等判断默认策略是否够用以及自定义ModularLoadManagerStrategy时应当使用哪些LoadData中的既有数据。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar 模块化负载管理器Modular Load Manager深入指南启用、验证与实现原理Apache Pulsar 模块化负载管理器Modular Load Manager深入指南启用、验证与实现原理 Apache Pulsar 的负载均衡机消息队列后端流处理Apache Pulsar 模块化负载管理器Modular Load Manager深入解析启用、验证、数据模型与 Bundle 分配策略Apache Pulsar 模块化负载管理器Modular Load Manager深入解析启用、验证、数据模型与 Bundle 分配策略 本指南以 Pu消息队列后端流处理Apache Pulsar 模块化负载均衡器Modular Load Manager开发与运维指南架构、数据模型与流量分配策略Apache Pulsar 模块化负载均衡器Modular Load Manager开发与运维指南架构、数据模型与流量分配策略 导读 本文以 Apache消息队列后端流处理上一篇WebGL 2 示例项目推荐下一篇InsightFace_Pytorch常见问题解决从环境配置到性能调优创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
企业数字化 ERP 产品动态
相关推荐
安卓逆向助手:抓包脱壳反编译全流程脚本化实战 简介:安卓逆向助手是一款面向Android应用开发者与安全研究人员的图形化逆向工具,旨在降低APK反编译与分析门槛,让初学者也能快速理解应用内部结构。它集成dex2jar、JD-GUI、apktool、baksmali等常用组件,支持一键将Dalvik字节码转… · 2026/9/25 7:50:58
通信驱动的CRM工作台:DeskcommCRM设计思路与落地实践 最近大半年我在推进一个项目,内部代号 DeskcommCRM,聊的人不多,但用起来确实和传统 CRM 是两个思路。它不是那种把客户信息塞进数据库就完事的系统,而是把“客户关系”这件事重新拉回到桌面上——电话、邮件、会话、跟进记录&… · 2026/9/25 7:50:45
SPL 迁移到 Axiom APL 实战指南:基于 spl-to-apl 技能的完整查询翻译手册 后端前端AI 技能AI 插件搜索引擎 【免费下载链接】clawhub Skill Plugin Registry for OpenClaw 项目地址: https://gitcode.com/gh_mirrors/mo/clawhub 点击查看 免费下载 本指南围绕本仓库 .agents/skills/spl-to-apl/ 目录下的 SPL→APL 翻译技能展开ÿ… · 2026/9/25 7:50:39
IT技术岗转网络安全值得吗?成本、路线与就业全景解析 我经常在后台收到类似的提问:干了几年IT技术岗,到底要不要转网络安全?说实话,每次看到这种问题,我都能大概猜到提问者的处境——现有工作不算差,但天花板感越来越明显;网络安全听起来热门、有技… · 2026/9/25 8:18:44
Moto 中 Amazon Managed Prometheus(amp)服务的模拟实现与实战指南 Mock测试 【免费下载链接】moto A library that allows you to easily mock out tests based on AWS infrastructure. 项目地址: https://gitcode.com/gh_mirrors/mo/moto 点击查看 免费下载 Amazon Managed Prometheus(AMP,AWS 的托管 Prom… · 2026/9/25 8:18:31
平头哥倚天720/730/750三代CPU规划解读:微架构迭代与ARM服务器落地实践 /* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 8:18:31
SQL Explorer 实战指南:在 RocketRide 管道中浏览数据库、编写 SQL 与解读查询计划 【免费下载链接】rocketride-server High-performance AI pipeline engine with a C core and 50 Python-extensible nodes. Build, debug, and scale LLM workflows with 13 model providers, 8 vector databases, and agent orchestration, all from your IDE. Includes VS C… · 2026/9/25 8:18:31
XAgent utils 模块深度解析:Token 计数、文本裁剪、状态码枚举、任务数据结构与单例元类 AI Agent大模型后端任务调度 【免费下载链接】XAgent An Autonomous LLM Agent for Complex Task Solving 项目地址: https://gitcode.com/gh_mirrors/xa/XAgent 点击查看 免费下载 XAgent/utils.py 是 XAgent 框架中一个"小而关键"的基础模块࿱… · 2026/9/25 8:18:31
如何用Auto-Empirical-Research-Skills在10分钟跑通第一篇DID论文?新手快速上手教程 如何用Auto-Empirical-Research-Skills在10分钟跑通第一篇DID论文?新手快速上手教程 【免费下载链接】Auto-Empirical-Research-Skills 🔬 A curated collection of 23,000 agent skills for empirical research across 8 social science disciplines. |… · 2026/9/25 8:18:31
创维E900V22D刷机全攻略:S905L3SB芯片兼容性解析与救砖实战 /* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 1:00:31
MQTT协议原理与Broker服务器搭建实战:从Mosquitto到EMQX /* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/25 1:00:37