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

GCP PubSub订阅消费报错:'generator'对象无add_done_callback属性

GCP PubSub容器部署消费报错:'generator' object has no attribute 'add_done_callback'

问题详情

本地环境可正常消费GCP PubSub订阅消息,但部署到云端容器(预发布环境)时触发如下错误:

2024-05-14 15:50:15,379  ERROR [google.api_core.bidi:_thread_main:678] [pid=1] [tname=Thread-ConsumeBidirectionalStream] [cluster=mt2] [1.65.1] Thread-ConsumeBidirectionalStream caught unexpected exception 'generator' object has no attribute 'add_done_callback' and will exit.
Traceback (most recent call last):
  File "/opt/project/lib/python3.9/site-packages/google/api_core/bidi.py", line 644, in _thread_main
    self._bidi_rpc.open()
  File "/opt/project/lib/python3.9/site-packages/google/api_core/bidi.py", line 294, in open
    call._wrapped.add_done_callback(self._on_call_done)
AttributeError: 'generator' object has no attribute 'add_done_callback'

代码结构

部署时通过Docker容器bash入口启动初始化脚本,核心代码如下:

init.py

def main():
    """
    Entrypoint for all consumers
    """

    args = get_args()
    config_settings = get_config_section(args.config_filename)
    consumer = BaseGooglePubSubConsumer(config_settings)
    consumer.run()


if __name__ == "__main__":
    main()

consumer.py

from functools import cached_property

from google.cloud.pubsub_v1 import SubscriberClient

class BaseGooglePubSubConsumer(BaseConsumer):

    @cached_property
    def consumer(self) -> SubscriberClient:
        json_account_info = self._get_service_account_info()
        credentials = service_account.Credentials.from_service_account_info(json_account_info)
        return SubscriberClient(credentials=credentials)

    def run(self):
        """
        Consumes messages from a Pub/Sub Subscriber.
        """
        streaming_pull_future = self.consumer.subscribe(self._get_subscription_path(), callback=self.process_message)
        with self.consumer:
            try:
                streaming_pull_future.result()
            except Exception as exc:
                le.errors = "Failed to process message"
                logger.error(le, exc_info=exc)
                streaming_pull_future.cancel()
                streaming_pull_future.result()
            finally:
                le.end()
    def process_message(self, message: Message) -> None:
        # message.ack()
        logger.info(message)

问题根源

该报错与OpenTelemetry的PubSub instrumentation相关:云端环境可能自动注入了OpenTelemetry监控,其instrumentation代码错误地将SubscriberClient.subscribe方法的返回值(预期为StreamingPullFuture)包装成了生成器,导致后续调用add_done_callback时触发属性不存在的异常。本地环境未启用该instrumentation,因此无此问题。

修复方案

试试以下几种解决方式:

  1. 禁用PubSub的OpenTelemetry instrumentation
    在代码启动时添加禁用逻辑:

    import opentelemetry.instrumentation.google_cloud_pubsub_v1
    opentelemetry.instrumentation.google_cloud_pubsub_v1.uninstrument()
    

    或者通过环境变量全局禁用:

    OTEL_PYTHON_DISABLED_INSTRUMENTATIONS=google_cloud_pubsub_v1
    
  2. 调整客户端初始化方式
    将consumer的@cached_property初始化改为在__init__中直接初始化,避免instrumentation代理延迟初始化对象时出现异常:

    class BaseGooglePubSubConsumer(BaseConsumer):
        def __init__(self, config_settings):
            super().__init__(config_settings)
            json_account_info = self._get_service_account_info()
            credentials = service_account.Credentials.from_service_account_info(json_account_info)
            self.consumer = SubscriberClient(credentials=credentials)
    
  3. 对齐依赖版本
    检查并升级/降级相关依赖到兼容版本:

    • google-cloud-pubsub 升级至 >=2.18.0
    • opentelemetry-instrumentation-google-cloud-pubsub-v1 升级至 >=0.41b0
  4. 手动修正返回值类型
    如果确定返回的是生成器,手动取出实际的StreamingPullFuture对象:

    streaming_pull_future = self.consumer.subscribe(self._get_subscription_path(), callback=self.process_message)
    # 处理生成器包装情况
    if hasattr(streaming_pull_future, "__next__"):
        streaming_pull_future = next(streaming_pull_future)
    

内容的提问来源于stack exchange,提问作者juan jose Palacio

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 22:42:04