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

Google Pub/Sub Python客户端:多订阅者优雅关闭实现疑问

多Google Pub/Sub订阅优雅关闭解决方案

原代码问题分析

  1. 循环阻塞:future.result()是阻塞方法,循环调用时会卡在第一个订阅的result()上,后续订阅的等待逻辑根本无法执行。
  2. 重复触发关闭逻辑:同时注册了SIGINT信号处理器和try-except KeyboardInterrupt,导致按下Ctrl+C后,信号处理器执行cancel操作,但主线程会回到循环继续执行future.result(),进而重复打印"Result"和"Cancelling"。
  3. 可能的遗漏:创建的订阅future(f1、f2)未添加到subscription_futures列表,导致信号处理器无法找到要取消的对象。

修正后的实现代码

import signal
from concurrent.futures import wait

# 全局标志,防止重复触发关闭流程
shutdown_initiated = False

def shutdown_handler(signum, frame):
    global shutdown_initiated
    if shutdown_initiated:
        return
    shutdown_initiated = True
    
    print("Shutting down gracefully...")
    # 仅取消未完成的future
    for future in subscription_futures:
        if not future.done():
            print("Cancelling subscription future")
            future.cancel()

# 存储所有订阅的future对象
subscription_futures = []

# 注册信号处理器,同时处理Ctrl+C(SIGINT)和容器停止信号(SIGTERM)
signal.signal(signal.SIGINT, shutdown_handler)
signal.signal(signal.SIGTERM, shutdown_handler)

# 初始化Pub/Sub订阅客户端(假设已完成配置)
# subscriber = pubsub_v1.SubscriberClient()

# 创建订阅并将future加入列表
f1 = subscriber.subscribe(subscription_1, callback=call_back1)
subscription_futures.append(f1)

f2 = subscriber.subscribe(subscription_2, callback=call_back2)
subscription_futures.append(f2)

# 等待所有订阅future完成或被取消
try:
    print("Subscriptions started, press Ctrl+C to exit")
    wait(subscription_futures)
except Exception as e:
    # 仅处理非主动关闭的异常
    if not shutdown_initiated:
        print(f"Unexpected error occurred: {e}")

print("All subscriptions have stopped")

关键改动说明

  1. 用wait替代循环result():concurrent.futures.wait()会阻塞主线程,同时监听所有订阅future的状态,不会卡在单个订阅上。
  2. 防止重复关闭:通过shutdown_initiated全局标志,避免多次触发关闭逻辑(比如用户连续按Ctrl+C)。
  3. 扩展信号支持:新增SIGTERM信号处理,适配Docker等容器环境的优雅停止需求。
  4. 精准取消future:仅对未完成的future执行cancel操作,避免无效调用。

额外注意事项

  • 回调函数中尽量避免长时间阻塞操作,若必须执行耗时任务,建议设置超时或在回调内检查shutdown_initiated标志,确保能快速响应关闭指令。
  • 取消future后,正在处理的消息会根据Pub/Sub的ack_deadline配置重新入队,若需要保证消息不丢失,可在回调内完成当前消息处理后再退出。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 15:25:55