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

使用Python异步Pub/Sub Lite客户端订阅时遇间歇性认证错误求助

Python异步Pub/Sub Lite订阅间歇性凭据错误排查

我用Python异步版Google Cloud Pub/Sub Lite客户端订阅消息时,间歇性遇到以下错误:

Traceback (most recent call last): 
File "/usr/local/lib/python3.10/site-packages/grpc/_plugin_wrapping.py", line 89, in __call__
 self._metadata_plugin( 
File "/usr/local/lib/python3.10/site-packages/google/auth/transport/grpc.py", line 101, in __call__
 callback(self._get_authorization_headers(context), None) 
File "/usr/local/lib/python3.10/site-packages/google/auth/transport/grpc.py", line 87, in _get_authorization_headers
 self._credentials.before_request( 
File "/usr/local/lib/python3.10/site-packages/google/auth/credentials.py", line 134, in before_request
 self.apply(headers) 
File "/usr/local/lib/python3.10/site-packages/google/auth/credentials.py", line 110, in apply
 _helpers.from_bytes(token or self.token) 
File "/usr/local/lib/python3.10/site-packages/google/auth/_helpers.py", line 130, in from_bytes 
raise ValueError("{0!r} could not be converted to unicode".format(value))
ValueError: None could not be converted to unicode" 

我未通过GOOGLE_APPLICATION_CREDENTIALS环境变量指定凭据,而是直接用ServiceAccount信息构建(避免在AWS主机上写入凭据文件),代码如下:

import asyncio

from google.cloud.pubsublite.cloudpubsub import AsyncSubscriberClient
from google.cloud.pubsublite.types import (
    CloudRegion,
    CloudZone,
    FlowControlSettings,
    SubscriptionPath,
)
from google.oauth2 import service_account


class AsyncTimedIterable:
    def __init__(self, iterable, poll_timeout=90):
        class AsyncTimedIterator:
            def __init__(self):
                self._iterator = iterable.__aiter__()

            async def __anext__(self):
                try:
                    result = await asyncio.wait_for(
                        self._iterator.__anext__(), int(poll_timeout)
                    )
                    if not result:
                        raise StopAsyncIteration
                    return result

                except asyncio.TimeoutError as e:
                    raise e

        self._factory = AsyncTimedIterator

    def __aiter__(self):
        return self._factory()


# TODO add project info below
location = CloudZone(CloudRegion("region"), "zone")

subscription_path = SubscriptionPath("project_number", location, "subscription_id")

# TODO add service account details
gcp_creds = {}


async def async_receive_from_subscription(per_partition_count=100):
    # Configure when to pause the message stream for more incoming messages based on the
    # maximum size or number of messages that a single-partition subscriber has received,
    # whichever condition is met first.
    per_partition_flow_control_settings = FlowControlSettings(
        # 1,000 outstanding messages. Must be >0.
        messages_outstanding=per_partition_count,
        # 10 MiB. Must be greater than the allowed size of the largest message (1 MiB).
        bytes_outstanding=10 * 1024 * 1024,
    )

    async with AsyncSubscriberClient(
        credentials=service_account.Credentials.from_service_account_info(gcp_creds)
    ) as async_subscriber_client:
        message_iterator = await async_subscriber_client.subscribe(
            subscription_path,
            per_partition_flow_control_settings=per_partition_flow_control_settings,
        )

        timed_iter = AsyncTimedIterable(message_iterator, 90)
        async for message in timed_iter:
            yield message


async def main():
    async for message in async_receive_from_subscription(per_partition_count=100_000):
        print(message.data)


if __name__ == "__main__":
    asyncio.run(main())

查看堆栈跟踪里的文件,看到这段注释:

The plugin may be invoked on a thread created by Core, which will not
have the context propagated. This context is stored and installed in
the thread invoking the plugin.

请问是不是因为我设置的凭据没传递到Core创建的新线程里,才导致了这个问题?


问题分析与解决方案

是的,你的猜测完全正确。这个间歇性错误的核心原因就是凭据的线程上下文没有被正确传播到gRPC Core创建的后台线程中。

当gRPC需要刷新过期的凭据时,会在自身创建的后台线程中调用元数据插件,但直接构建的service_account.Credentials实例没有绑定可跨线程复用的请求上下文,导致后台线程中尝试获取token时返回None,最终触发ValueError。

修复步骤

修改凭据初始化逻辑,为凭据绑定一个可复用的请求对象,确保token刷新能跨线程正常执行:

# 新增导入
from google.auth.transport.requests import Request

# 替换原凭据构建代码
gcp_creds = {
    # 填入完整的ServiceAccount字段:type、project_id、private_key、client_email等
}

# 创建可跨线程复用的请求对象
request = Request()
# 构建凭据并提前刷新初始token
credentials = service_account.Credentials.from_service_account_info(gcp_creds)
credentials.refresh(request)

然后在创建AsyncSubscriberClient时传入这个预处理后的凭据:

async with AsyncSubscriberClient(credentials=credentials) as async_subscriber_client:
    # 后续订阅逻辑保持不变

额外优化建议

  • 将凭据初始化逻辑放在全局作用域,复用同一个凭据实例,避免重复构建和刷新token
  • 确保gcp_creds字典包含完整的ServiceAccount必填字段,缺失字段也可能导致token获取失败

内容的提问来源于stack exchange,提问作者Vasu Kandimalla

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 17:45:40