Apache Pulsar 负载仿真工具详解:Simulation Client / Controller 与 Broker Monitor 使用指南
2026/9/24 15:33:19 网站建设 项目流程
  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

项目地址:https://gitcode.com/gh_mirrors/pulsar28/pulsar
点击查看免费下载

导读

本指南围绕 Apache Pulsar 官方文档中的Simulation tools(负载仿真工具)展开,系统讲解如何通过pulsar-perf脚本提供的三个子命令——simulation-clientsimulation-controllermonitor-brokers——搭建一个可编程控制的人工负载测试环境,用于观察 Load Manager(负载管理器)在处理大规模流量时的表现。读完本文,你将掌握仿真客户端/控制器的完整启动方式、控制器交互 Shell 的全部命令与参数语义,以及 Broker 监控器的输出解读方法,并能结合仓库源码理解其底层实现机制。

为什么需要负载仿真工具

在生产或预发环境正式部署前,往往需要人为制造负载来观察系统行为。Apache Pulsar 的负载管理器(Load Manager)负责在多个 Broker 之间均衡分配 bundle 负载,而验证其行为最直接的方式,就是人为制造可控的消息流量,观察负载管理器如何调度与均衡。

为此,Pulsar 在pulsar-testclient模块中提供了三件配套工具(对应官方文档 site2/docs/developing-tools.md):

  • Simulation Client(仿真客户端):真正制造流量的一端,按可配置的消息速率与消息大小创建并订阅主题;
  • Simulation Controller(仿真控制器):面向用户的交互 Shell,负责向一个或多个仿真客户端下发指令,控制负载的创建、变更与停止;
  • Broker Monitor(Broker 监控器):持续从 ZooKeeper 读取各 Broker 的负载数据,以表格形式打印到控制台,用于观察负载管理器在仿真过程中的行为。

三个组件对应的核心实现类都位于pulsar-testclient/src/main/java/org/apache/pulsar/testclient/目录下,分别是 LoadSimulationClient.java、LoadSimulationController.java 与 BrokerMonitor.java。

整体架构:客户端、控制器与监控器如何协作

仿真的工作流可以概括为一条"控制链路":

  1. 在若干台压测机器上分别启动Simulation Client,每个客户端持有一个ServerSocket,监听指定端口等待指令(见LoadSimulationClient.run()new ServerSocket(port)的循环accept)。
  2. 启动Simulation Controller,它会根据--clients传入的主机名列表,通过Socket与所有客户端建立连接(见LoadSimulationController构造器中new Socket(clients[i], clientPort)),并为每个客户端维护一对DataInputStream/DataOutputStream
  3. 用户在控制器 Shell 中敲入命令,控制器将命令编码后写入对应客户端的输出流;客户端解码后执行"创建生产者/消费者、调整速率、停止主题"等动作。
  4. 在任意机器上启动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-urlPulsar 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 的组合)为单位管理每个主题,内部持有RateLimiter(Guava)控制发送速率,并用AtomicReference<byte[]>包装消息载荷以便运行时无痛更换消息大小。
  • 创建主题:收到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_COMMANDSTOP_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_topic

  • tenant 名称是my_tenant
  • namespace 名称是my_namespace
  • topic 名称是my_topic

控制器内部会用makeTopic()把三者拼装为完整主题名再下发给客户端。

控制器命令大全

控制器支持以下动作(命令之后的参数均来自源码ShellArgumentsread()分发逻辑):

创建单个主题(含生产者和消费者)

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 上的历史数据仿真负载(应为正在仿真的同一个 ZooKeeper)

simulate <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 表示不间隔(源码中每次tradeThread.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>/<j>(i ∈ [0, num_namespaces),j ∈ [0, topics_per_namespace))。之后change_groupstop_group通过正则.*://<tenant>/.*/<group>-.*/.*匹配并批量修改/停止该组下所有主题。

脚本命令与退出

除上述命令外,控制器 Shell 还支持两个实用命令(源码read()中实现):

  • script <script_name>:从脚本文件逐行读取并执行命令(按空白拆分后递归调用read()),便于批量下发预置的仿真指令;
  • quit/exit:退出控制器。

copy、simulate、stream 的区别

copysimulatestream三条命令非常相似,但语义有显著差异:

  • copy:用于在目标 ZooKeeper上仿真一个静态的外部 ZooKeeper的负载。因此source zookeeper是要复制的源 ZooKeeper,target 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()的实现,simulatestream实际只接收一个参数(ZooKeeper 连接串),文档中示例的<tenant>参数在当前实现中并未使用,使用时请以实际命令行为准。

三者的共同点是都支持可选的--rate-multiplier参数:例如--rate-multiplier 0.05会让消息以被仿真负载 5% 的速率发送,方便按比例缩放负载强度。

Broker Monitor:Broker 监控器

职责

要在仿真中观察负载管理器的行为,可以使用 Broker Monitor(实现类BrokerMonitor)。它通过 ZooKeeper Watcher 监听各 Broker 的负载数据节点,一旦数据更新,就以表格形式把负载数据打印到控制台

启动方式

使用pulsar-perf脚本中的monitor-brokers子命令:

pulsar-perf monitor-brokers --connect-string <zookeeper host:port>

其中--connect-string为 ZooKeeper 连接串(必填,源码中 ZooKeeper 连接超时 30 秒)。启动后控制台将持续打印负载数据,直到进程被中断(Ctrl+C)。

输出内容解读

监控器支持两种负载管理器实现,输出格式随之变化(源码通过 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 字符以便阅读。

一个完整的仿真流程示例

把上述组件串起来,一次典型的负载仿真实验大致如下:

  1. 在压测机 A(例如10.0.0.1)启动仿真客户端:
pulsar-perf simulation-client --port 7777 --service-url pulsar://broker-1:6650
  1. 在压测机 B(例如10.0.0.2)再启动一个客户端,以扩大负载容量(可选)。

  2. 在控制台机器启动仿真控制器(--clients按逗号分隔列出所有客户端主机名):

pulsar-perf simulation-controller --cluster my_cluster --client-port 7777 --clients 10.0.0.1,10.0.0.2
  1. 在控制器>提示符下创建主题负载:
> 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
  1. 观察负载管理器表现,随时调整或停止负载:
> change my_tenant my_namespace my_topic --rate 5000 --size 1024 > stop_group my_tenant group_a
  1. 在另一终端启动 Broker Monitor 观察各 Broker 负载变化:
pulsar-perf monitor-brokers --connect-string zk-1:2181
  1. 实验结束后在控制器输入quit退出。

小结

Pulsar 的负载仿真工具链(simulation-client+simulation-controller+monitor-brokers)为验证负载管理器行为提供了一套"控制端 + 执行端 + 观测端"的完整方案:控制器负责编排,客户端负责制造可控流量,监控器则借助 ZooKeeper Watcher 实时呈现负载数据。其核心实现全部位于 pulsar-testclient 模块,官方文档 site2/docs/developing-tools.md 是理解三者用法的第一手资料,配合本仓库源码(LoadSimulationClient.java、LoadSimulationController.java、BrokerMonitor.java)可以进一步掌握命令协议、速率控制、分组匹配与 ZooKeeper 数据读写等底层细节,为搭建自己的大规模负载仿真实验提供可复用的实践依据。

  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

项目地址:https://gitcode.com/gh_mirrors/pulsar28/pulsar
点击查看免费下载

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询