Python操作Pulsar:单生产者单分区下send()方法QPS约500是否正常?
Pulsar单生产者单分区下send()方法QPS约500是否正常?
一、QPS500的合理性判断
单生产者单分区使用同步send()方法时,QPS约500属于可能出现的情况,但并非最优状态,核心原因在于:
send()是同步阻塞调用,每发送一条消息都要等待Broker的确认响应,QPS直接受网络往返延迟(RTT)、Broker处理能力、客户端配置影响。如果测试环境是跨机房网络,RTT较高(比如几十ms),500的QPS符合同步调用的性能逻辑;但如果是本地Broker环境,这个数值还有明显优化空间。- 你的测试代码完全是单条消息逐一发送,没有利用Pulsar的批量发送机制,导致网络请求开销占比极高,进一步拉低了QPS。
二、兼顾有序性、失败即停与性能的优化方案
你需要的有序性+失败即停+高性能,完全可以通过调整生产者配置实现,不需要放弃同步调用的可靠性:
1. 开启生产者批量发送(核心优化)
Pulsar Python客户端默认未开启批量发送,你可以在创建生产者时启用批量配置,让同步send()自动将多条消息打包发送,大幅减少网络请求次数:
from pulsar import Client import time service_url='' topic = '' data = '0xf9ba0955b0509ac6138908ccc50fasdfa5bd296e48d7d,0x74b1d2f771adsfasdfasebdbea80fe4013bdace4fc4b653322c2895c' client = Client(service_url) # 开启批量发送,调整批量参数以适配你的场景 producer = client.create_producer( topic, batching_enabled=True, batching_max_messages=1000, # 批量最多包含1000条消息 batching_max_delay_ms=10, # 最多等待10ms触发批量 send_timeout=30000 # 设置发送超时 ) try: for i in range(18000,18100): start_time = time.time() for j in range(0,10000): producer.send(data.encode('utf-8')) end_time = time.time() elapsed_time = end_time - start_time print(f"Block {i} takes {elapsed_time:.2f} seconds to write, QPS: {10000/elapsed_time:.2f}") except Exception as e: print(f"发送失败,终止后续发送: {str(e)}") producer.close() client.close()
开启批量后,同步send()的QPS能提升数倍甚至一个数量级,同时依然保证消息的有序性——Pulsar的批量发送会严格保留消息的发送顺序。
2. 失败即停的实现
同步send()方法在发送失败时会直接抛出异常,你只需要用try-except捕获异常,一旦触发就关闭生产者、终止循环,完全符合你“失败后立即停止后续发送”的需求,且天然保证消息有序。
3. 其他辅助优化
- 确认Broker的ack配置:Pulsar默认的
send()会等待所有副本确认(类似Kafka的ack=-1),如果Broker副本数过多或同步副本配置严格,可根据业务需求调整,但不建议为了性能降低可靠性。 - 调整客户端线程池:如果后续扩展多生产者,可设置
io_threads参数提升并发,但单生产者场景下影响不大。
三、关于send_async()的补充说明
send_async()性能高是因为异步非阻塞,但并非无法保证可靠性和有序性:你可以在回调函数中记录发送状态,一旦检测到失败就设置全局标志停止后续发送,但需要处理线程安全问题。不过对于你的需求而言,开启批量的同步send()是更简单直接的方案,既能满足可靠性和有序性,又能大幅提升性能。
内容的提问来源于stack exchange,提问作者ke du
相关产品推荐
相关产品推荐

