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

