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

GCP Pub/Sub Lite异步订阅报错:同一SubscriberId重复连接分区

问题描述

使用Python异步版Pub/Sub Lite客户端订阅消息,对应主题包含10个分区。同时创建5个客户端,每个负责消费2个分区(逻辑正常运行),但偶尔部分订阅者会出现**"A second subscriber connected with the same SubscriberId and topic partition"**错误。

错误日志

JSON格式日志

{
    "time": 1662717414,
    "level": "ERROR",
    "message": "Future exception was never retrieved\nfuture: <Future finished exception=Aborted('A second subscriber connected with the same SubscriberId and topic partition.')>",
    "type": "applog",
    "name": "asyncio",
    "exc_info": "Traceback (most recent call last):\n  File \"/usr/local/lib/python3.10/site-packages/google/api_core/grpc_helpers_async.py\", line 102, in _wrapped_aiter\n    async for response in self._call:  # pragma: no branch\n  File \"/usr/local/lib/python3.10/site-packages/grpc/aio/_call.py\", line 326, in _fetch_stream_responses\n    await self._raise_for_status()\n  File \"/usr/local/lib/python3.10/site-packages/grpc/aio/_call.py\", line 236, in _raise_for_status\n    raise _create_rpc_error(await self.initial_metadata(), await\ngrpc.aio._call.AioRpcError: <AioRpcError of RPC that terminated with:\n\tstatus = StatusCode.ABORTED\n\tdetails = \"A second subscriber connected with the same SubscriberId and topic partition.\"\n\tdebug_error_string = \"UNKNOWN:Error received from peer ipv4:172.253.63.95:443 {created_time:\"2022-09-09T09:56:48.360743382+00:00\", grpc_status:10, grpc_message:\"A second subscriber connected with the same SubscriberId and topic partition.\"}\"\n>\n\nThe above exception was the direct cause of the following exception:\n\nTraceback (most recent call last):\n  File \"/usr/local/lib/python3.10/site-packages/google/cloud/pubsublite/internal/wire/retrying_connection.py\", line 100, in _run_loop\n    await self._loop_connection(\n  File \"/usr/local/lib/python3.10/site-packages/google/cloud/pubsublite/internal/wire/retrying_connection.py\", line 154, in _loop_connection\n    await self._read_queue.put(await read_task)\n  File \"/usr/local/lib/python3.10/site-packages/google/cloud/pubsublite/internal/wire/permanent_failable.py\", line 76, in await_unless_failed\n    raise self._failure_task.exception()\n  File \"/usr/local/lib/python3.10/site-packages/google/cloud/pubsublite/internal/wire/gapic_connection.py\", line 71, in read\n    raise self.error()\n  File \"/usr/local/lib/python3.10/site-packages/google/cloud/pubsublite/internal/wire/gapic_connection.py\", line 66, in read\n    return await self.await_unless_failed(response_it.__anext__())\n  File \"/usr/local/lib/python3.10/site-packages/google/cloud/pubsublite/internal/wire/permanent_failable.py\", line 75, in await_unless_failed\n    return await task\n  File \"/usr/local/lib/python3.10/site-packages/google/api_core/grpc_helpers_async.py\", line 105, in _wrapped_aiter\n    raise exceptions.from_grpc_error(rpc_error) from rpc_error\ngoogle.api_core.exceptions.Aborted: 409 A second subscriber connected with the same SubscriberId and topic partition. [reason: \"DUPLICATE_SUBSCRIBER_CONNECTIONS\"\ndomain: \"pubsublite.googleapis.com\"\n]"
}

格式化后的异常栈

Traceback (most recent call last):
  File "/usr/local/lib/python3.10/site-packages/google/api_core/grpc_helpers_async.py", line 102, in _wrapped_aiter
    async for response in self._call:  # pragma: no branch
  File "/usr/local/lib/python3.10/site-packages/grpc/aio/_call.py", line 326, in _fetch_stream_responses
    await self._raise_for_status()
  File "/usr/local/lib/python3.10/site-packages/grpc/aio/_call.py", line 236, in _raise_for_status
    raise _create_rpc_error(await self.initial_metadata(), await
grpc.aio._call.AioRpcError: <AioRpcError of RPC that terminated with:
        status = StatusCode.ABORTED
        details = "A second subscriber connected with the same SubscriberId and topic partition."
        debug_error_string = "UNKNOWN:Error received from peer ipv4:172.253.63.95:443 {created_time:"2022-09-09T09:56:48.360743382+00:00", grpc_status:10, grpc_message:"A second subscriber connected with the same SubscriberId and topic partition."}"
>

The above exception was the direct cause of the following exception:

Traceback (most recent call last):
  File "/usr/local/lib/python3.10/site-packages/google/cloud/pubsublite/internal/wire/retrying_connection.py", line 100, in _run_loop
    await self._loop_connection(
  File "/usr/local/lib/python3.10/site-packages/google/cloud/pubsublite/internal/wire/retrying_connection.py", line 154, in _loop_connection
    await self._read_queue.put(await read_task)
  File "/usr/local/lib/python3.10/site-packages/google/cloud/pubsublite/internal/wire/permanent_failable.py", line 76, in await_unless_failed
    raise self._failure_task.exception()
  File "/usr/local/lib/python3.10/site-packages/google/cloud/pubsublite/internal/wire/gapic_connection.py", line 71, in read
    raise self.error()
  File "/usr/local/lib/python3.10/site-packages/google/cloud/pubsublite/internal/wire/gapic_connection.py", line 66, in read
    return await self.await_unless_failed(response_it.__anext__())
  File "/usr/local/lib/python3.10/site-packages/google/cloud/pubsublite/internal/wire/permanent_failable.py", line 75, in await_unless_failed
    return await task
  File "/usr/local/lib/python3.10/site-packages/google/api_core/grpc_helpers_async.py", line 105, in _wrapped_aiter
    raise exceptions.from_grpc_error(rpc_error) from rpc_error
google.api_core.exceptions.Aborted: 409 A second subscriber connected with the same SubscriberId and topic partition. [reason: "DUPLICATE_SUBSCRIBER_CONNECTIONS"
domain: "pubsublite.googleapis.com"
]

代码实现

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())

解决方案

1. 强制每个客户端实例使用唯一SubscriberId

默认客户端自动生成的SubscriberId可能在重试、重启场景下重复,手动结合进程ID、UUID分配唯一标识:

import os
import uuid

# 在subscribe时添加subscriber_id参数
message_iterator = await async_subscriber_client.subscribe(
    subscription_path,
    per_partition_flow_control_settings=per_partition_flow_control_settings,
    subscriber_id=f"subscriber-{os.getpid()}-{uuid.uuid4()}"
)

2. 优雅处理连接中断与重试

遇到ABORTED错误时,用指数退避策略避免短时间重复连接:

from google.api_core.exceptions import Aborted
from tenacity import retry, stop_after_attempt, wait_exponential, retry_if_exception_type

@retry(
    stop=stop_after_attempt(5),
    wait=wait_exponential(multiplier=1, min=2, max=10),
    retry=retry_if_exception_type(Aborted)
)
async def async_receive_from_subscription(per_partition_count=100):
    per_partition_flow_control_settings = FlowControlSettings(
        messages_outstanding=per_partition_count,
        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,
            subscriber_id=f"subscriber-{os.getpid()}-{uuid.uuid4()}"
        )

        timed_iter = AsyncTimedIterable(message_iterator, 90)
        try:
            async for message in timed_iter:
                yield message
        except Aborted:
            raise

3. 优化超时后的资源清理

修改超时逻辑,确保客户端资源正常释放:

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:
                    # 超时后主动终止迭代,触发客户端资源清理
                    raise StopAsyncIteration

        self._factory = AsyncTimedIterator

4. 手动指定分区分配(可选)

避免自动分配的短暂竞争,为每个客户端固定分配无重叠的分区:

from google.cloud.pubsublite.types import Partition

# 每个客户端设置对应分区,比如第一个客户端处理0-1,第二个2-3...
message_iterator = await async_subscriber_client.subscribe(
    subscription_path,
    per_partition_flow_control_settings=per_partition_flow_control_settings,
    subscriber_id=f"subscriber-{os.getpid()}-{uuid.uuid4()}",
    partitions=[Partition(0), Partition(1)]
)

内容的提问来源于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.20 07:21:32