Locust中Google Pub/Sub客户端无法订阅问题求助
问题分析与解决:Locust环境下Google Pub/Sub订阅阻塞问题
问题描述
尝试为Python性能测试库Locust实现自定义GCloudSubUser,需求是完成两个操作:
- 向Google Pub/Sub的指定主题发送消息
- 等待API响应(响应会发送至输出主题)
当前代码实现如下:
import os from google.cloud import pubsub_v1 from locust import User class GCloudSubUser(User): abstract = True def __init__(self, environment) -> None: super().__init__(environment) self.client = GCloudSubClient(environment) class GCloudSubClient: def __init__(self, environment): self.environment = environment def publish(self): # ...some code to publish to a topic... project_id = os.getenv("PUBSUB_PROJECT_ID") subscription_id = os.getenv("PUBSUB_TOPIC_ID") subscriber = pubsub_v1.SubscriberClient() subscription_path = subscriber.subscription_path(project_id, subscription_id) future = subscriber.subscribe(subscription_path, callback=callback) print("I reached this point") def callback(message: pubsub_v1.subscriber.message.Message) -> None: print(f"Received {message}.") message.ack()
遇到的问题:Locust环境下执行future = subscriber.subscribe(subscription_path, callback=callback)时会永久阻塞,无法打印后续日志;但脱离Locust环境时,可正常订阅输出主题并接收消息。
原因分析
Locust基于gevent协程运行,会对Python标准库做猴子补丁(monkey-patching)实现异步IO。而Google Pub/Sub同步客户端的subscribe()方法默认用原生线程池处理订阅逻辑,gevent的协程模型和原生线程池存在兼容性冲突:线程被阻塞后无法正常进行协程切换,导致整个Locust用户进程卡住。
另外,subscribe()本身不会阻塞,但它启动的后台线程在gevent打补丁的环境中,无法正确处理事件循环,导致后续代码(比如print("I reached this point"))无法执行。
解决方案
1. 使用Pub/Sub异步客户端(推荐)
Google Cloud提供基于asyncio的异步Pub/Sub客户端google.cloud.pubsub_v1.aio,能完美适配Locust的协程模型。修改后的代码示例:
import os import asyncio from google.cloud import pubsub_v1 from google.cloud.pubsub_v1.aio import SubscriberClient from locust import User, task class GCloudSubUser(User): abstract = True def __init__(self, environment) -> None: super().__init__(environment) self.client = GCloudSubClient(environment) class GCloudSubClient: def __init__(self, environment): self.environment = environment self.project_id = os.getenv("PUBSUB_PROJECT_ID") self.subscription_id = os.getenv("PUBSUB_TOPIC_ID") self.input_topic_id = os.getenv("PUBSUB_INPUT_TOPIC_ID") async def publish_and_wait(self): # 发布消息到输入主题 publisher = pubsub_v1.PublisherClient() topic_path = publisher.topic_path(self.project_id, self.input_topic_id) publish_future = publisher.publish(topic_path, b"test message") await publish_future.result() # 异步订阅输出主题 async with SubscriberClient() as subscriber: subscription_path = subscriber.subscription_path(self.project_id, self.subscription_id) # 定义异步回调,收到消息后取消订阅 async def callback(message): print(f"Received {message}.") message.ack() future.cancel() future = subscriber.subscribe(subscription_path, callback=callback) print("I reached this point") try: # 等待消息或超时 await future except asyncio.CancelledError: print("Subscription cancelled after receiving message") class TestUser(GCloudSubUser): @task async def test_pubsub_flow(self): await self.client.publish_and_wait()
2. 禁用gevent线程猴子补丁(不推荐)
如果必须使用同步客户端,可在Locust启动脚本开头禁用gevent对线程的补丁,但这会降低Locust的性能测试效率:
from gevent import monkey monkey.patch_all(thread=False)
关键注意事项
- Locust中涉及异步操作的任务方法,需用
async def定义,并搭配@task装饰器 - 使用异步客户端时,所有IO操作都要采用异步方式,避免阻塞协程
- 订阅完成后要及时取消订阅,防止资源泄漏
内容的提问来源于stack exchange,提问作者JPFrancoia
相关产品推荐
相关产品推荐

