如何在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
相关产品推荐
相关产品推荐

