- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
导读
本指南围绕 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.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。
整体架构:客户端、控制器与监控器如何协作
仿真的工作流可以概括为一条"控制链路":
- 在若干台压测机器上分别启动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_COMMAND | 0 | 修改既有主题的速率/消息大小 |
STOP_COMMAND | 1 | 停止(关闭)一个主题 |
TRADE_COMMAND | 2 | 创建"生产者 + 消费者"对 |
CHANGE_GROUP_COMMAND | 3 | 批量修改一个组内所有主题 |
STOP_GROUP_COMMAND | 4 | 批量停止一个组内所有主题 |
FIND_COMMAND | 5 | 查询某个主题当前由哪个客户端持有 |
由于大负载往往需要多台压测机器协同,用户不与客户端直接交互,而是把请求统一委托给控制器,由控制器把指令分发到各客户端——这正是这套架构的核心设计动机(文档原文亦有说明)。
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 的组合)为单位管理每个主题,内部持有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_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_topic:
- tenant 名称是
my_tenant - namespace 名称是
my_namespace - topic 名称是
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 上的历史数据仿真负载(应为正在仿真的同一个 ZooKeeper)
simulate <tenant> <zookeeper> [--rate-multiplier value]流式读取给定活跃 ZooKeeper 的最新数据,仿真该 ZooKeeper 的实时负载
stream <tenant> <zookeeper> [--rate-multiplier value]所有 ZooKeeper 参数的格式均为zookeeper_host:port。
命令选项与默认值
下表汇总了 Shell 命令支持的全部选项(对应ShellArguments的 JCommander 定义):
| 选项 | 默认值 | 说明 |
|---|---|---|
--rate | 1(条/秒) | 每秒消息数 |
--rand-rate | 空 | 从两个逗号分隔值构成的区间内均匀随机选取消息速率,会覆盖--rate;源码中先取两者min/max,再random.nextDouble() * (max - min) + min |
--size | 1024(字节) | 消息大小 |
--separation | 0(毫秒) | trade_group创建各主题之间的间隔,0 表示不间隔(源码中每次trade后Thread.sleep(separation)) |
--topics-per-namespace | 1 | trade_group中每个 namespace 创建的主题数(主题总数 = num_namespaces × topics_per_namespace) |
--rate-multiplier | 1 | copy/simulate/stream的负载比例系数 |
分组命令的语义
命令中的 group 参数允许用户一次操作多个主题。分组在调用trade_group时创建:源码中 namespace 命名为<group>-<namespace索引>,主题名为索引字符串,即最终主题为persistent://<tenant>/<cluster>/<group>-<i>/<j>(i ∈ [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是要复制的源 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()的实现,simulate与stream实际只接收一个参数(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 字符以便阅读。
一个完整的仿真流程示例
把上述组件串起来,一次典型的负载仿真实验大致如下:
- 在压测机 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-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
相关推荐
Apache Pulsar 负载仿真测试工具实战指南:Simulation Client / Controller 与 Broker Monitor 全解析
Apache Pulsar 负载仿真测试工具实战指南:Simulation Client / Controller 与 Broker Monitor 全解析 导
消息队列后端流处理Apache Pulsar 负载模拟工具实战:Simulation Client / Controller 与 Broker Monitor
Apache Pulsar 负载模拟工具实战:Simulation Client / Controller 与 Broker Monitor 本篇技术指南围绕
消息队列后端流处理axum 提取器(Extractor)完全指南:从内置类型到自定义实现
axum 提取器(Extractor)完全指南:从内置类型到自定义实现 在 axum 中, 提取器(Extractor) 是处理 HTTP 请求的核心抽象:一个
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考