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

