You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.15 14:56:14