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

如何让Google Pub/Sub与Locust兼容运行?

解决Locust + Google Pub/Sub的gRPC与gevent兼容性问题

问题背景

尝试为Locust实现GCloudSubUser自定义用户类,用于向Google Pub/Sub主题发送消息,但遇到gRPC与gevent的兼容性冲突,已尝试常规解决步骤仍无法正常运行。

可行解决思路

  • 严格控制初始化顺序:gevent猴子补丁必须在所有gRPC、Pub/Sub相关模块导入前执行,确保gRPC正确适配gevent协程环境
  • 复用Pub/Sub客户端实例:避免在每次任务中重复创建PublisherClient,重复创建会导致连接池混乱,应在用户初始化时创建一次并复用
  • 异步处理发布结果:避免直接调用future.result()阻塞协程,改用回调方式处理结果,减少gevent上下文阻塞
  • 验证gRPC gevent补丁生效:确保gRPC的gevent扩展正确初始化,避免遗漏关键步骤

修改后的代码示例

# 第一步:先执行gevent猴子补丁,必须在所有其他导入之前
from gevent import monkey
# 补丁时排除thread/threading,避免与gRPC的gevent补丁冲突
monkey.patch_all(thread=False, threading=False)

# 第二步:初始化gRPC的gevent支持
import grpc.experimental.gevent as grpc_gevent
grpc_gevent.init_gevent()

# 第三步:导入其他依赖模块
import os
from google.cloud import pubsub_v1
from locust import User, task
from google.api_core.exceptions import GoogleAPICallError


class GCloudSubUser(User):
    abstract = True

    def __init__(self, environment):
        super().__init__(environment)
        self.client = GCloudSubClient(environment)


class GCloudSubClient:
    def __init__(self, environment):
        self.environment = environment
        # 初始化一次PublisherClient,复用连接池
        self.publisher = pubsub_v1.PublisherClient()
        self.project_id = os.getenv("PUBSUB_PROJECT_ID")
        self.topic_id = os.getenv("PUBSUB_TOPIC_ID")
        self.topic_path = self.publisher.topic_path(self.project_id, self.topic_id)

    def _handle_publish_result(self, future, message_num):
        """异步处理发布结果的回调函数"""
        try:
            message_id = future.result()
            self.environment.events.request.fire(
                request_type="pubsub_publish",
                name=self.topic_id,
                response_time=0,  # Pub/Sub无直接响应时间,可按需统计
                response_length=0,
                exception=None,
            )
            print(f"Message {message_num} published with ID: {message_id}")
        except GoogleAPICallError as e:
            self.environment.events.request.fire(
                request_type="pubsub_publish",
                name=self.topic_id,
                response_time=0,
                response_length=0,
                exception=e,
            )
            print(f"Failed to publish message {message_num}: {str(e)}")

    def publish_messages(self):
        for n in range(1, 10):
            data_str = f"Message number {n}"
            data = data_str.encode("utf-8")
            # 发布消息,添加回调而非直接阻塞等待结果
            future = self.publisher.publish(self.topic_path, data)
            future.add_done_callback(lambda fut, num=n: self._handle_publish_result(fut, num))
        print(f"Started publishing messages to {self.topic_path}")

额外验证步骤

  1. 确认依赖版本兼容:使用grpcio>=1.40.0、google-cloud-pubsub>=2.10.0、gevent>=21.12.0
  2. 检查环境变量:确保PUBSUB_PROJECT_ID和PUBSUB_TOPIC_ID已正确配置
  3. 先测试低并发:先运行少量用户验证无阻塞或崩溃问题,再逐步扩大并发量

内容的提问来源于stack exchange,提问作者user24477818

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 20:43:15