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

Airflow中使用Python向Pub/Sub发布消息报错及等待方案咨询

报错原因

Google Cloud Pub/Sub 客户端的publish()方法返回的是自定义实现的pubsub_v1.publisher.futures.Future类,并非Python标准库concurrent.futures模块提供的标准Future类型。标准库的futures.wait()方法仅支持处理自身模块的Future实例,执行逻辑依赖标准Future的私有属性_condition实现锁控制,而Pub/Sub自定义Future没有该属性,因此触发AttributeError。

替代等待发布完成的方案
  • 方案1:遍历发布Future逐个等待结果
    Pub/Sub自定义Future本身实现了result()方法,直接遍历调用即可等待所有发布任务完成,适配代码如下:
# 替换原代码中的 futures.wait 部分
for publish_future in publish_futures:
    try:
        # 可根据需求调整超时时间,无需超时可省略timeout参数
        publish_future.result(timeout=60)
    except Exception as e:
        print(f"消息发布失败:{str(e)}")

该方案不会影响你已经添加的add_done_callback回调逻辑,回调会在发布完成后正常触发。

  • 方案2:调用PublisherClient自带的shutdown方法
    如果当前PublisherClient实例提交完所有发布请求后不再使用,可以直接调用shutdown()方法,方法会自动阻塞直到所有已提交的发布请求全部处理完成,不需要手动维护Future列表:
# 所有publish()调用执行完成后执行即可
publisher.shutdown(timeout=60)

内容的提问来源于stack exchange,提问作者Javier Muñoz

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 18:57:01