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

如何在FastAPI应用中持续监听Google Cloud Pub/Sub消息?

问题描述

我正在用Google Scheduler向Pub/Sub主题发送消息,希望在FastAPI应用中持续监听这些消息,但当前代码仅执行一次,无法实现持续监听。

main.py

from fastapi import FastAPI, Depends
from typing import List

from core.config import get_db
from sqlalchemy.orm import Session

app = FastAPI()

from concurrent.futures import TimeoutError
from google.cloud import pubsub_v1

subscriber = pubsub_v1.SubscriberClient()
subscription_path = subscriber.subscription_path("project_id", "subscription_id")

def callback(message: pubsub_v1.subscriber.message.Message) -> None:
    print(f"Received {message}.")
    message.ack()

streaming_pull_future = subscriber.subscribe(subscription_path, callback=callback)
print(f"Listening for messages on {subscription_path}..\n")

with subscriber:
    try:
        streaming_pull_future.result(timeout=5)
    except TimeoutError:
        streaming_pull_future.cancel()  # Trigger the shutdown.
        streaming_pull_future.result()  # Block until the shutdown is complete.


@app.get("/")
def home(db: Session = Depends(get_db)):
    return {
        "message": "Welcome!"
    }

请问是否有方法在FastAPI应用中持续监听Pub/Sub消息?


解决方案

原代码无法持续监听的核心原因是设置了timeout=5,5秒后监听就被主动取消了。要在FastAPI中持续监听Pub/Sub消息,需要将监听逻辑放到后台线程中执行,避免阻塞FastAPI的主线程(主线程需要处理HTTP请求)。

具体修改后的代码示例:

from fastapi import FastAPI, Depends
from typing import List
import threading
from core.config import get_db
from sqlalchemy.orm import Session
from concurrent.futures import TimeoutError
from google.cloud import pubsub_v1

app = FastAPI()

subscriber = pubsub_v1.SubscriberClient()
subscription_path = subscriber.subscription_path("project_id", "subscription_id")

def callback(message: pubsub_v1.subscriber.message.Message) -> None:
    print(f"Received {message}.")
    message.ack()

def run_pubsub_listener():
    """在后台线程运行Pub/Sub监听"""
    streaming_pull_future = subscriber.subscribe(subscription_path, callback=callback)
    print(f"Listening for messages on {subscription_path}..\n")
    
    with subscriber:
        try:
            # 不设置超时,让监听持续运行
            streaming_pull_future.result()
        except Exception as e:
            print(f"监听异常终止: {e}")
            streaming_pull_future.cancel()
            streaming_pull_future.result()

# 在FastAPI启动时启动后台监听线程
@app.on_event("startup")
def startup_event():
    listener_thread = threading.Thread(target=run_pubsub_listener, daemon=True)
    listener_thread.start()

@app.get("/")
def home(db: Session = Depends(get_db)):
    return {
        "message": "Welcome!"
    }

关键说明

  • 使用@app.on_event("startup")装饰器,在FastAPI启动时自动触发监听线程的启动
  • 设置线程为daemon=True,保证FastAPI应用关闭时,后台监听线程会自动退出
  • 移除timeout=5参数,让streaming_pull_future.result()持续阻塞,实现永久监听
  • 增加异常捕获逻辑,确保监听异常终止时能正确清理资源

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 14:20:40