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

