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

FastAPI(Cloud Run)向GCP PubSub发消息时无响应报504问题排查

问题诊断与解决方案

核心问题分析

  1. PubSub客户端重复初始化:当前通过Depends(get_pub_sub)在每次请求时新建PubSub实例,并发场景下会导致多个客户端初始化操作竞争资源,在Cloud Run容器环境中易引发阻塞。
  2. 同步调用阻塞事件循环:publish_message中调用future.result()是同步阻塞操作,在FastAPI的异步路由中会阻塞整个事件循环,导致后续请求排队超时返回504。
  3. 客户端日志缺失:默认未开启PubSub客户端的底层日志,无法定位初始化卡住的具体原因。

一、修复PubSub客户端并发与阻塞问题

1. 全局单例复用PubSub客户端

Google Cloud客户端(含PubSub)天生线程安全,全局复用可避免重复初始化的资源竞争:

from google.cloud import pubsub
import json
from fastapi import Depends

# 全局单例存储客户端实例
_pub_sub_client = None

def get_pub_sub():
    global _pub_sub_client
    if _pub_sub_client is None:
        _pub_sub_client = PubSub("primary-care-378721")
    return _pub_sub_client

class PubSub:
    def __init__(self, project_name, enable_message_ordering=True) -> None:
        print("Instantiating PubSub client...")
        publisher_options = pubsub.types.PublisherOptions(
            enable_message_ordering=enable_message_ordering
        )
        self.project_name = project_name
        self.publisher = pubsub.PublisherClient(publisher_options=publisher_options)
        print("Done")

    def publish_message(
        self, topic_name: str, message: dict, ordering_key: str
    ):
        topic_path = self.publisher.topic_path(self.project_name, topic_name)
        data = json.dumps(message).encode("utf-8")
        future = self.publisher.publish(
            topic=topic_path, data=data, ordering_key=ordering_key
        )
        return future

2. 异步化处理发布请求

方案A:用后台任务隔离同步操作

将PubSub发布逻辑放入FastAPI后台任务,避免阻塞Webhook响应:

from fastapi import Request, status, Depends, BackgroundTasks
from fastapi.routing import APIRouter

router = APIRouter()

@router.post("/incoming-message", status_code=status.HTTP_204_NO_CONTENT)
async def receive_incoming_message(
    request: Request, 
    background_tasks: BackgroundTasks,
    pub_sub_client=Depends(get_pub_sub)
):
    print("Processing request...")
    form_data = await request.form()
    form_dict = form_data._dict
    ordering_key = form_dict.get("From")
    topic_name = "twilio_incoming_message-de6f415"

    # 后台执行发布逻辑
    def publish_task():
        try:
            future = pub_sub_client.publish_message(
                topic_name=topic_name, message=form_dict, ordering_key=ordering_key
            )
            rs = future.result()
            print(f"Message {rs} published successfully to topic {topic_name} with ordering key {ordering_key}")
        except Exception as e:
            print(f"Publish failed: {str(e)}")

    background_tasks.add_task(publish_task)

方案B:使用PubSub异步客户端(推荐)

适配FastAPI异步环境,直接使用官方异步客户端:

from google.cloud.pubsub_v1.publisher.async_client import AsyncPublisherClient
from google.cloud.pubsub_v1.types import PublisherOptions

class PubSub:
    def __init__(self, project_name, enable_message_ordering=True) -> None:
        print("Instantiating Async PubSub client...")
        publisher_options = PublisherOptions(
            enable_message_ordering=enable_message_ordering
        )
        self.project_name = project_name
        self.publisher = AsyncPublisherClient(publisher_options=publisher_options)
        print("Done")

    async def publish_message(
        self, topic_name: str, message: dict, ordering_key: str
    ):
        topic_path = self.publisher.topic_path(self.project_name, topic_name)
        data = json.dumps(message).encode("utf-8")
        future = await self.publisher.publish(
            topic=topic_path, data=data, ordering_key=ordering_key
        )
        return await future

路由中直接异步调用:

@router.post("/incoming-message", status_code=status.HTTP_204_NO_CONTENT)
async def receive_incoming_message(
    request: Request, 
    pub_sub_client=Depends(get_pub_sub)
):
    print("Processing request...")
    form_data = await request.form()
    form_dict = form_data._dict
    ordering_key = form_dict.get("From")
    topic_name = "twilio_incoming_message-de6f415"

    try:
        rs = await pub_sub_client.publish_message(
            topic_name=topic_name, message=form_dict, ordering_key=ordering_key
        )
        print(f"Message {rs} published successfully to topic {topic_name} with ordering key {ordering_key}")
    except Exception as e:
        print(f"Publish failed: {str(e)}")

二、开启PubSub客户端日志

通过Python logging模块开启底层调试日志,在应用启动时添加:

import logging

# 开启Google Cloud客户端调试日志
logging.basicConfig(level=logging.DEBUG)
logging.getLogger("google.cloud.pubsub").setLevel(logging.DEBUG)
logging.getLogger("google.api_core").setLevel(logging.DEBUG)

日志会输出到Cloud Run的日志系统中,可查看客户端初始化、网络请求等细节,定位阻塞原因。


三、额外优化建议

  • Cloud Run配置调整:提升容器CPU/内存配额,避免资源不足;设置最小实例数,减少冷启动时的并发压力。
  • 异常捕获增强:在客户端初始化和发布逻辑中添加全局异常捕获,避免未处理异常导致进程崩溃。
  • Twilio超时适配:Twilio Webhook默认超时15秒,后台任务方案可确保快速返回响应,避免触发超时。

内容的提问来源于stack exchange,提问作者p.magalhaes

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 16:39:51