首页/新闻资讯/正文详情

Apache Pulsar 负载仿真工具详解:Simulation Client / Controller 与 Broker Monitor 使用指南

发布时间:2026/9/24 15:33:22 来源:云帆数科 栏目:资讯中心
Apache Pulsar 负载仿真工具详解:Simulation Client / Controller 与 Broker Monitor 使用指南
消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载导读本指南围绕 Apache Pulsar 官方文档中的Simulation tools负载仿真工具展开系统讲解如何通过pulsar-perf脚本提供的三个子命令——simulation-client、simulation-controller与monitor-brokers——搭建一个可编程控制的人工负载测试环境用于观察 Load Manager负载管理器在处理大规模流量时的表现。读完本文你将掌握仿真客户端/控制器的完整启动方式、控制器交互 Shell 的全部命令与参数语义以及 Broker 监控器的输出解读方法并能结合仓库源码理解其底层实现机制。为什么需要负载仿真工具在生产或预发环境正式部署前往往需要人为制造负载来观察系统行为。Apache Pulsar 的负载管理器Load Manager负责在多个 Broker 之间均衡分配 bundle 负载而验证其行为最直接的方式就是人为制造可控的消息流量观察负载管理器如何调度与均衡。为此Pulsar 在pulsar-testclient模块中提供了三件配套工具对应官方文档 site2/docs/developing-tools.mdSimulation Client仿真客户端真正制造流量的一端按可配置的消息速率与消息大小创建并订阅主题Simulation Controller仿真控制器面向用户的交互 Shell负责向一个或多个仿真客户端下发指令控制负载的创建、变更与停止Broker MonitorBroker 监控器持续从 ZooKeeper 读取各 Broker 的负载数据以表格形式打印到控制台用于观察负载管理器在仿真过程中的行为。三个组件对应的核心实现类都位于pulsar-testclient/src/main/java/org/apache/pulsar/testclient/目录下分别是 LoadSimulationClient.java、LoadSimulationController.java 与 BrokerMonitor.java。整体架构客户端、控制器与监控器如何协作仿真的工作流可以概括为一条控制链路在若干台压测机器上分别启动Simulation Client每个客户端持有一个ServerSocket监听指定端口等待指令见LoadSimulationClient.run()中new ServerSocket(port)的循环accept。启动Simulation Controller它会根据--clients传入的主机名列表通过Socket与所有客户端建立连接见LoadSimulationController构造器中new Socket(clients[i], clientPort)并为每个客户端维护一对DataInputStream/DataOutputStream。用户在控制器 Shell 中敲入命令控制器将命令编码后写入对应客户端的输出流客户端解码后执行创建生产者/消费者、调整速率、停止主题等动作。在任意机器上启动Broker Monitor它通过 ZooKeeper Watcher 监听各 Broker 的负载数据节点一旦数据更新就刷新控制台表格。从源码看客户端与控制器之间的命令协议由LoadSimulationClient中的一组字节码常量定义常量值含义CHANGE_COMMAND0修改既有主题的速率/消息大小STOP_COMMAND1停止关闭一个主题TRADE_COMMAND2创建生产者 消费者对CHANGE_GROUP_COMMAND3批量修改一个组内所有主题STOP_GROUP_COMMAND4批量停止一个组内所有主题FIND_COMMAND5查询某个主题当前由哪个客户端持有由于大负载往往需要多台压测机器协同用户不与客户端直接交互而是把请求统一委托给控制器由控制器把指令分发到各客户端——这正是这套架构的核心设计动机文档原文亦有说明。Simulation Client仿真客户端职责仿真客户端是一台按配置制造流量的机器它创建并订阅主题以可配置的消息速率和消息大小持续发送/接收消息。出于性能考虑客户端会预先缓存各消息大小的byte[]载荷源码中的payloadCache字段注释明确说明为每条消息新建 byte[] 对压测机压力很大。启动方式使用pulsar-perf脚本pulsar-testclient模块随发行包提供的 CLI 入口启动pulsar-perf simulation-client --port listen port --service-url pulsar service url启动后客户端即进入监听状态等待控制器连接与指令。命令行参数对应源码中LoadSimulationClient.MainArguments均必填参数必填说明--port是客户端用于接收控制器指令的监听端口--service-url是Pulsar Service URL如pulsar://broker-host:6650客户端内部实现要点从源码可以确认以下实现细节对应LoadSimulationClient.java连接配置构造器使用PulsarClient.builder()创建客户端设置了memoryLimit(0, SizeUnit.BYTES)不设内存上限、connectionsPerBroker(4)、ioThreads取 CPU 核数、关闭 stats 输出同时用PulsarAdmin负责后续自动创建 namespace。TradeUnit客户端以TradeUnit一个 Consumer 与 Producer 的组合为单位管理每个主题内部持有RateLimiterGuava控制发送速率并用AtomicReferencebyte[]包装消息载荷以便运行时无痛更换消息大小。创建主题收到TRADE_COMMAND后客户端会用admin.namespaces().createNamespace()自动创建对应 namespace已存在时捕获ConflictException忽略消费者以Subscriber- topic作为订阅名通过subscribeAsync()异步订阅并使用Consumer::acknowledgeAsync作为轻量级确认监听器。消息发送发送循环中调用producer.sendAsync(payload.get())后由rateLimiter.acquire()限速生产者在sendTimeout(0, TimeUnit.SECONDS)不超时下创建。若发送过程出现异常例如 Broker 重启exceptionHandler会置位健康标志客户端随后调用getNewProducer()无限重试——该方法在创建失败时休眠 10 秒后重试保证 Broker 重启后流量自动恢复。默认值TradeConfiguration的默认消息速率为 100 条/秒、默认消息大小为 1024 字节未指定时按此执行。命令处理handle(byte command, ...)方法按字节码分发到创建/修改/停止等分支其中组命令CHANGE_GROUP_COMMAND、STOP_GROUP_COMMAND通过正则.*://tenant/.*/group-.*/.*匹配组内全部主题组名作为 namespace 前缀FIND_COMMAND则向控制器回写布尔值指示主题是否由本客户端持有。Simulation Controller仿真控制器职责仿真控制器负责向仿真客户端下发指令创建新主题、停止旧主题、调整主题负载以及其他若干任务。它以交互 Shell 形式呈现给用户实现类为LoadSimulationController。启动方式pulsar-perf simulation-controller --cluster cluster to simulate on --client-port listen port for clients --clients comma-separated list of client host names必须先启动所有仿真客户端再启动控制器控制器启动时会立即连接各客户端见LoadSimulationController构造器中逐台打印Connected to client的逻辑。启动后将出现一个简单的提示符在此输入命令即可向客户端下发指令。命令行参数对应源码中LoadSimulationController.MainArguments参数必填说明--cluster是要仿真的集群名称用于拼装persistent://tenant/cluster/namespace/topic形式的主题全名见源码makeTopic()--client-port是仿真客户端正在监听的端口--clients是客户端主机名的逗号分隔列表命名约定始终使用 BASE 名称Shell 命令中的参数经常涉及 tenant、namespace、topic。在所有情况下命令都只使用 tenant、namespace、topic 的 BASE 名称。例如对主题persistent://my_tenant/my_cluster/my_namespace/my_topictenant 名称是my_tenantnamespace 名称是my_namespacetopic 名称是my_topic控制器内部会用makeTopic()把三者拼装为完整主题名再下发给客户端。控制器命令大全控制器支持以下动作命令之后的参数均来自源码ShellArguments与read()分发逻辑创建单个主题含生产者和消费者trade tenant namespace topic [--rate message rate per second] [--rand-rate lower bound,upper bound] [--size message size in bytes]批量创建一组主题每个主题含生产者和消费者trade_group tenant group num_namespaces [--rate message rate per second] [--rand-rate lower bound,upper bound] [--separation separation between creating topics in ms] [--size message size in bytes] [--topics-per-namespace number of topics to create per namespace]修改既有主题的配置change tenant namespace topic [--rate message rate per second] [--rand-rate lower bound,upper bound] [--size message size in bytes]批量修改一组主题的配置change_group tenant group [--rate message rate per second] [--rand-rate lower bound,upper bound] [--size message size in bytes] [--topics-per-namespace number of topics to create per namespace]停止关闭一个已创建的主题stop tenant namespace topic批量停止一组已创建的主题stop_group tenant group将历史数据从一个 ZooKeeper 复制到另一个并按历史中的消息速率与大小进行仿真copy tenant source zookeeper target zookeeper [--rate-multiplier value]基于当前 ZooKeeper 上的历史数据仿真负载应为正在仿真的同一个 ZooKeepersimulate tenant zookeeper [--rate-multiplier value]流式读取给定活跃 ZooKeeper 的最新数据仿真该 ZooKeeper 的实时负载stream tenant zookeeper [--rate-multiplier value]所有 ZooKeeper 参数的格式均为zookeeper_host:port。命令选项与默认值下表汇总了 Shell 命令支持的全部选项对应ShellArguments的 JCommander 定义选项默认值说明--rate1条/秒每秒消息数--rand-rate空从两个逗号分隔值构成的区间内均匀随机选取消息速率会覆盖--rate源码中先取两者min/max再random.nextDouble() * (max - min) min--size1024字节消息大小--separation0毫秒trade_group创建各主题之间的间隔0 表示不间隔源码中每次trade后Thread.sleep(separation)--topics-per-namespace1trade_group中每个 namespace 创建的主题数主题总数 num_namespaces × topics_per_namespace--rate-multiplier1copy/simulate/stream的负载比例系数分组命令的语义命令中的 group 参数允许用户一次操作多个主题。分组在调用trade_group时创建源码中 namespace 命名为group-namespace索引主题名为索引字符串即最终主题为persistent://tenant/cluster/group-i/ji ∈ [0, num_namespaces)j ∈ [0, topics_per_namespace)。之后change_group与stop_group通过正则.*://tenant/.*/group-.*/.*匹配并批量修改/停止该组下所有主题。脚本命令与退出除上述命令外控制器 Shell 还支持两个实用命令源码read()中实现script script_name从脚本文件逐行读取并执行命令按空白拆分后递归调用read()便于批量下发预置的仿真指令quit/exit退出控制器。copy、simulate、stream 的区别copy、simulate、stream三条命令非常相似但语义有显著差异copy用于在目标 ZooKeeper上仿真一个静态的外部 ZooKeeper的负载。因此source zookeeper是要复制的源 ZooKeepertarget zookeeper是正在仿真的目标 ZooKeeper。命令会递归读取源端/loadbalance/resource-quota/namespace下的全部资源配额getResourceQuotas()把源端的完整历史数据以两种格式写入目标端使负载管理器无论新旧实现都能获得完整的历史数据收益。源码中为了区分不同集群/租户下的同名 namespace会将 namespace 重组为sourceCluster-sourceTenant-keyRange-namespace的形式同时写入旧 API 路径/loadbalance/resource-quota/namespace/...与新 API 路径/loadbalance/bundle-data/...。simulate只接收一个 ZooKeeper 参数即正在仿真的那个适用于该 ZooKeeper 上已有SimpleLoadManagerImpl的历史数据旧 API 的资源配额的场景它据此为ModularLoadManagerImpl新 API生成等价的历史数据BundleData再由客户端按历史数据仿真负载。写入时若节点已存在则用setData覆盖。stream接收一个与正在仿真目标不同的、活跃的 ZooKeeper通过BrokerWatcher/LoadReportWatcher两个 Watcher 监听/loadbalance/broker节点及其下的负载报告每当 Broker 上线或负载报告更新就按msgRateIn/msgRateOut平均值估算消息速率、按吞吐量除以速率估算消息大小再调用changeOrCreate()实时创建或调整客户端上的主题从而仿真目标集群的实时负载。需要说明以上三条命令中copy需要三个参数tenant、源 ZK、目标 ZK而根据当前仓库源码handleSimulate()与handleStream()的实现simulate与stream实际只接收一个参数ZooKeeper 连接串文档中示例的tenant参数在当前实现中并未使用使用时请以实际命令行为准。三者的共同点是都支持可选的--rate-multiplier参数例如--rate-multiplier 0.05会让消息以被仿真负载 5% 的速率发送方便按比例缩放负载强度。Broker MonitorBroker 监控器职责要在仿真中观察负载管理器的行为可以使用 Broker Monitor实现类BrokerMonitor。它通过 ZooKeeper Watcher 监听各 Broker 的负载数据节点一旦数据更新就以表格形式把负载数据打印到控制台。启动方式使用pulsar-perf脚本中的monitor-brokers子命令pulsar-perf monitor-brokers --connect-string zookeeper host:port其中--connect-string为 ZooKeeper 连接串必填源码中 ZooKeeper 连接超时 30 秒。启动后控制台将持续打印负载数据直到进程被中断CtrlC。输出内容解读监控器支持两种负载管理器实现输出格式随之变化源码通过 JSON 中是否包含allocated字符串来区分SimpleLoadManagerImpl旧实现输出LoadReport格式包含 COUNT 行TOPIC / BUNDLE / PRODUCER / CONSUMER / BUNDLE / BUNDLE-、RAW SYSTEM 与 ALLOC SYSTEM 行CPU % / MEMORY % / DIRECT % / BW IN % / BW OUT % / MAX %、RAW MSG 与 ALLOC MSG 行MSG/S IN / MSG/S OUT / TOTAL / KB/S IN / KB/S OUT / TOTAL。ModularLoadManagerImpl新实现除了读取 Broker 的LocalBrokerData还会读取/loadbalance/broker-time-average下的TimeAverageBrokerData输出 SYSTEM 行、COUNT 行、以及 LATEST / SHORT / LONG 三组消息速率与吞吐行便于同时观察实时、短期平均与长期平均三个时间尺度的负载。除了单 Broker 的明细表格监控器每 60 秒源码GLOBAL_STATS_PRINT_PERIOD_MILLIS 60000还会打印一次全局统计表printGlobalData()汇总每个 Broker 的 BUNDLE 数、总消息速率、KB/S 吞吐、长期速率与最大资源使用率MAX %并给出 TOTAL 汇总行有 Broker 上线/下线时也会打印Gained broker/Lost broker日志并实时更新统计。表格渲染使用FixedColumnLengthTableMaker元素宽度 14、小数格式%.2f表格宽度约 120 字符全局表则用*边框、Broker 列加宽到 60 字符以便阅读。一个完整的仿真流程示例把上述组件串起来一次典型的负载仿真实验大致如下在压测机 A例如10.0.0.1启动仿真客户端pulsar-perf simulation-client --port 7777 --service-url pulsar://broker-1:6650在压测机 B例如10.0.0.2再启动一个客户端以扩大负载容量可选。在控制台机器启动仿真控制器--clients按逗号分隔列出所有客户端主机名pulsar-perf simulation-controller --cluster my_cluster --client-port 7777 --clients 10.0.0.1,10.0.0.2在控制器提示符下创建主题负载 trade my_tenant my_namespace my_topic --rate 1000 --size 512 trade_group my_tenant group_a 10 --rate 500 --topics-per-namespace 2 --separation 50观察负载管理器表现随时调整或停止负载 change my_tenant my_namespace my_topic --rate 5000 --size 1024 stop_group my_tenant group_a在另一终端启动 Broker Monitor 观察各 Broker 负载变化pulsar-perf monitor-brokers --connect-string zk-1:2181实验结束后在控制器输入quit退出。小结Pulsar 的负载仿真工具链simulation-clientsimulation-controllermonitor-brokers为验证负载管理器行为提供了一套控制端 执行端 观测端的完整方案控制器负责编排客户端负责制造可控流量监控器则借助 ZooKeeper Watcher 实时呈现负载数据。其核心实现全部位于 pulsar-testclient 模块官方文档 site2/docs/developing-tools.md 是理解三者用法的第一手资料配合本仓库源码LoadSimulationClient.java、LoadSimulationController.java、BrokerMonitor.java可以进一步掌握命令协议、速率控制、分组匹配与 ZooKeeper 数据读写等底层细节为搭建自己的大规模负载仿真实验提供可复用的实践依据。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar 负载仿真测试工具实战指南Simulation Client / Controller 与 Broker Monitor 全解析Apache Pulsar 负载仿真测试工具实战指南Simulation Client / Controller 与 Broker Monitor 全解析 导消息队列后端流处理Apache Pulsar 负载模拟工具实战Simulation Client / Controller 与 Broker MonitorApache Pulsar 负载模拟工具实战Simulation Client / Controller 与 Broker Monitor 本篇技术指南围绕消息队列后端流处理axum 提取器Extractor完全指南从内置类型到自定义实现axum 提取器Extractor完全指南从内置类型到自定义实现 在 axum 中 提取器Extractor 是处理 HTTP 请求的核心抽象一个消息队列后端流处理创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

相关推荐

stats4cj分位数完全指南:percentile、四分位数、IQR与三均值逐层详解
stats4cj分位数完全指南:percentile、四分位数、IQR与三均值逐层详解

stats4cj分位数完全指南:percentile、四分位数、IQR与三均值逐层详解 【免费下载链接】stats4cj stats4cj是一个仓颉实现的数学统计库,包括总体/样本均值、总体/样本方差、分位数、统计分布等多种数理统计函数。 项目地址: https://gitcode.com/Cangji… · 2026/9/24 15:33:22

鸿蒙远控横评:ToDesk/向日葵/RayLink,功能完整度才是真正的分水岭
鸿蒙远控横评:ToDesk/向日葵/RayLink,功能完整度才是真正的分水岭

鸿蒙远控,功能完整度才是真正的分水岭鸿蒙系统从发布至今,远程控制软件的适配一直是个痛点。不少用户反馈,安卓端能用的功能到了鸿蒙端要么缺失,要么体验打折扣。市面上主流的远控软件——ToDesk、向日葵、RayLink——在鸿蒙端的表… · 2026/9/24 15:33:16

Sinon `assert.calledOnceWithMatch` 断言完全指南:精确匹配一次调用
Sinon `assert.calledOnceWithMatch` 断言完全指南:精确匹配一次调用

测试开发工具 【免费下载链接】sinon Test spies, stubs and mocks for JavaScript. 项目地址: https://gitcode.com/gh_mirrors/si/sinon 点击查看 免费下载 本文基于 Sinon.JS 官方文档(docs/concepts/assertions/api/called-once-with-match.md&… · 2026/9/24 15:33:16

mcp-use TypeScript 全栈 MCP 框架实战:构建 Agent、MCP Server 与跨客户端 MCP Apps
mcp-use TypeScript 全栈 MCP 框架实战:构建 Agent、MCP Server 与跨客户端 MCP Apps

后端MCP 服务MCP ClientsAI Agent人工智能 【免费下载链接】mcp-use The fullstack MCP framework to develop MCP Apps for ChatGPT / Claude & MCP Servers for AI Agents. 项目地址: https://gitcode.com/gh_mirrors/mc/mcp-use 点击查看 免费下载 mcp-use … · 2026/9/24 16:05:59

62.qt quick-QML虚拟软键盘V2版本(手机键盘弹出机制)-支持换肤、动态加载移除语言、发布linux软键盘程序、支持qt6
62.qt quick-QML虚拟软键盘V2版本(手机键盘弹出机制)-支持换肤、动态加载移除语言、发布linux软键盘程序、支持qt6

最新版本已更新,请用最新版本: 108.qt quick-QML虚拟软键盘V3版本-新增点击空白区域收回键盘、支持ListView、Flickable自动布局-CSDN博客 在上章我们学习了45.qt quick-qml虚拟软键盘详解(一)_诺谦的博客-CSDN博客46.qt quick-自定义非常好看的qml虚拟软键盘-支持换肤、动态… · 2026/9/24 16:05:52

AI黄瓜病虫害防治机器人 QT 信创完整项目
AI黄瓜病虫害防治机器人 QT 信创完整项目

# AI黄瓜病虫害防治机器人 QT 信创完整项目 ## 项目说明 1. 平台:Qt5.15 / Qt6 兼容(适配银河麒麟、统信UOS信创操作系统) 2. 功能:AI图像识别黄瓜病虫害、机器人运动控制、病害数据库、喷洒作业调度、日志记录、本地模型推理 3. 架构:主窗口+AI推理模块+串口机器人控制+… · 2026/9/24 16:05:46

力扣刷题总结(内容简单,个人记录,有问题请各位大佬评论区指出)
力扣刷题总结(内容简单,个人记录,有问题请各位大佬评论区指出)

1. 二分法简单题给定一个 n 个元素有序的(升序)整型数组 nums 和一个目标值 target ,写一个函数搜索 nums 中的 target,如果目标值存在返回下标,否则返回 -1。示例 1:输入: nums [-1,0,3,5,9,12], target 9 输出: 4… · 2026/9/24 16:05:46

Prisma CLI 集群管理实战:`prisma cluster list` 命令详解与集群注册表机制剖析
Prisma CLI 集群管理实战:`prisma cluster list` 命令详解与集群注册表机制剖析

Prisma CLI 集群管理实战:prisma cluster list 命令详解与集群注册表机制剖析 【免费下载链接】prisma1 💾 Database Tools incl. ORM, Migrations and Admin UI (Postgres, MySQL & MongoDB) [deprecated] 项目地址: https://gitcode.com/gh_mirr… · 2026/9/24 16:05:40

Basic Computer Games 之 Weekday 的 MiniScript 移植:安装运行指南与格里高利历算法源码解析
Basic Computer Games 之 Weekday 的 MiniScript 移植:安装运行指南与格里高利历算法源码解析

示例工程 【免费下载链接】basic-computer-games An updated version of the classic "Basic Computer Games" book, with well-written examples in a variety of common MEMORY SAFE, SCRIPTING programming languages. See https://coding-horror.github.io/basic… · 2026/9/24 16:05:40

基于YOLOv8的渔船作业监控系统:从环境搭建到边缘部署全流程
基于YOLOv8的渔船作业监控系统:从环境搭建到边缘部署全流程

简介:这是一套面向计算机、人工智能、自动化等专业学生与教师的毕业设计级项目资源,围绕YOLOv8实现渔船作业监控系统,可用于毕设、课程设计、大作业或项目立项演示。压缩包共97个文件,约24.21MB,以70个Python源码文件为… · 2026/9/24 0:00:13

1D-CNN时间序列建模实战:从Conv1d原理到工业落地
1D-CNN时间序列建模实战:从Conv1d原理到工业落地

简介:面向时间序列数据建模的一维卷积神经网络完整实现,适合深度学习入门者及需要快速验证时序模型的研究者,能够从音频、文本、传感器或股价等序列中挖掘局部特征与时间依赖。压缩包体积很小,只有3KB,内含3个Python脚… · 2026/9/24 0:00:26

柔软的L:汉语语流中被忽视的舌肌张力控制
柔软的L:汉语语流中被忽视的舌肌张力控制

1. 这个“L”不是字母表里的L,而是舌尖上的L最近在几个方言群和语音教学社群里,反复看到有人发一句:“也说字母L:柔软的长舌”。初看以为是英语发音课笔记,点开才发现全是方言爱好者、播音系学生、语言康复师甚至戏曲演… · 2026/9/24 0:00:44

了解更多?预约专属演示

我们的顾问将为您一对一讲解产品与方案

企业微信二维码