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

Python Flask应用中如何保持Pub/Sub消费者持续监听消息?

问题原因

你当前的代码逻辑存在两个核心问题,导致消费者无法在Flask运行周期内持续监听:

  • 执行顺序错误:consumer()是同步阻塞调用,只有当consumer函数完全执行完毕(即触发超时、关闭订阅客户端之后),代码才会走到app.run()逻辑启动Flask服务。也就是说Flask启动的时候,Pub/Sub消费者已经被主动关停,自然收不到消息。
  • 消费者逻辑中给streaming_pull_future.result()设置了固定timeout参数,超时就会主动触发订阅流取消、关闭客户端,本身就不支持长期驻留监听。
解决方案

轻量场景下可以通过后台守护线程运行消费者,让消费逻辑和Flask Web服务在同一个进程内并行执行,同时移除消费者的超时退出逻辑,保证其持续驻留。
修正后的可运行代码如下:

import threading
from flask import Flask
from google.cloud import pubsub_v1

app = Flask(__name__)
# 替换为实际项目配置
project_id = "apnatime-fbc72"
subscription_id = "new-job-application-sub-event-local"

def callback(message):
    # 替换为实际的消息处理逻辑
    print(f"Received message: {message.data.decode('utf-8')}")
    # 消息处理完成后必须调用ack,否则消息会重新投递
    message.ack()

def consumer():
    subscriber = pubsub_v1.SubscriberClient()
    subscription_path = subscriber.subscription_path(project_id, subscription_id)
    streaming_pull_future = subscriber.subscribe(subscription_path, callback=callback)
    print(f"Listening for messages on {subscription_path}..\n")
    with subscriber:
        try:
            # 不传timeout参数,会持续阻塞监听消息,直到主动取消或抛出异常
            streaming_pull_future.result()
        except Exception as e:
            print(f"Pub/Sub subscriber encountered error: {str(e)}")
            streaming_pull_future.cancel()

if __name__ == "__main__":
    # 启动守护线程运行消费者,线程随主进程退出自动销毁
    consumer_thread = threading.Thread(target=consumer, daemon=True)
    consumer_thread.start()
    # 主线程直接启动Flask服务
    app.run(debug=False, port=5000)
生产环境注意事项
  • 不要开启Flask debug模式:debug=True会启动两个重载进程,导致消费者实例重复运行,引发重复消费问题。
  • 如果用Gunicorn等生产级WSGI服务器部署,需要把worker数量设置为1,否则每个worker进程都会启动一个消费者实例,同样会导致重复消费。
  • 更稳妥的生产实践是把Pub/Sub消费逻辑拆成独立进程单独部署,和Flask Web服务解耦,避免Web服务的进程调度、扩缩容影响消费稳定性。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 12:03:20