Zeep AsyncClient异步功能未达预期,求问题排查
我尝试用Zeep的AsyncClient执行异步数据拉取,但没实现预期的异步效果。我知道AsyncClient加载WSDL是同步的,但这解释不了当前的问题。试了多种代码配置都不行,怀疑对用法有根本性误解。
简化后的核心代码如下:
import azure.storage.blob.aio as blob import asyncio import datetime import time import zeep from zeep.transports import AsyncTransport import zeep.helpers import httpx import logging import json logging.getLogger('zeep.wsdl.bindings.soap').setLevel(logging.ERROR) async def gather_with_concurrency(n, *coros): semaphore = asyncio.Semaphore(n) async def sem_coro(coro): async with semaphore: return await coro return await asyncio.gather(*(sem_coro(c) for c in coros)) async def get_location_data(access_token,locations,businessdate,aclient,transport,blob_client): transport.client.headers.update({'AccessToken': access_token,'LocationToken':locations[1]}) aclient.service._binding_options.update({'address':'actual_address_for_calls'}) req_type = aclient.get_type('OperationName') #Include parameters in request body req_data = req_type(BusinessDate=businessdate) data = await aclient.service.OperationName(req_data) ### Add extra points to data for better tracking ### json_data = zeep.helpers.serialize_object(data) await blob_client.upload_blob(json.dumps(json_data,default=str),overwrite=True) print(f'Finished getting data for {locations[0]} at {time.time()-ref}') return json_data async def main(): ### Get credentials needed to establish connections and grab location tokens/access token ### tasks = [] businessdate = datetime.date(2023,3,7) wsdl_client = httpx.Client() async_client = httpx.AsyncClient() transport= AsyncTransport(client=async_client,wsdl_client=wsdl_client) aclient = zeep.AsyncClient("Link.to.WSDL",transport=transport) print('Getting Location Data') for location in location_tokens: blob_client = blob_service_client.get_blob_client(container='container_name',blob=f'blob_name') async with blob_client: tasks.append(get_location_data(access_token,location,businessdate,aclient,transport,blob_client)) location_data = await gather_with_concurrency(30,*tasks) ref = time.time() loop = asyncio.ProactorEventLoop() asyncio.set_event_loop(loop) try: asyncio.run(main()) finally: loop.run_until_complete(loop.shutdown_asyncgens()) loop.close()
该代码处理小数据量(约20KiB)时正常,但处理大数据量(约2MiB)时运行时间问题显著。请问使用AsyncClient时我遗漏了什么?
核心问题:共享AsyncClient/Transport实例导致状态冲突
你在所有任务里复用同一个aclient和transport,但这两个对象的状态(headers、绑定地址)是全局共享的。多个并发任务同时修改这些值,会导致请求串用参数、状态竞争,甚至被迫变成串行执行——这直接抵消了异步的优势,大数据量时这个问题被放大。
具体修复点&优化建议
每个任务独立创建AsyncClient和Transport
不要共享这些实例,每个任务单独初始化,避免状态互相干扰。把同步数据处理移出事件循环
zeep.helpers.serialize_object和json.dumps是同步阻塞操作,处理2MiB数据时会卡住事件循环。用asyncio.to_thread把这些操作放到线程池执行,不影响其他任务并发。修正Blob客户端的上下文管理
你在循环里创建blob_client并通过async with托管,但随即把它传给任务,会导致上下文提前退出、客户端被关闭。应该在任务内部创建并管理Blob客户端的上下文。调整并发数与连接池匹配
30的并发量可能超过服务端限制或客户端连接池上限,反而导致排队。先调低到10左右,同时给httpx.AsyncClient设置匹配的连接池大小。
修正后的代码示例
import azure.storage.blob.aio as blob import asyncio import datetime import time import zeep from zeep.transports import AsyncTransport import zeep.helpers import httpx import logging import json logging.getLogger('zeep.wsdl.bindings.soap').setLevel(logging.ERROR) async def gather_with_concurrency(n, *coros): semaphore = asyncio.Semaphore(n) async def sem_coro(coro): async with semaphore: return await coro return await asyncio.gather(*(sem_coro(c) for c in coros)) async def get_location_data(access_token, location, businessdate, blob_service_client, wsdl_url, target_address): # 每个任务独立创建客户端和传输层 async_client = httpx.AsyncClient(limits=httpx.Limits(max_connections=10)) wsdl_client = httpx.Client() transport = AsyncTransport(client=async_client, wsdl_client=wsdl_client) aclient = zeep.AsyncClient(wsdl_url, transport=transport) # 仅当前任务生效的header和地址设置 transport.client.headers.update({ 'AccessToken': access_token, 'LocationToken': location[1] }) aclient.service._binding_options.update({'address': target_address}) req_type = aclient.get_type('OperationName') req_data = req_type(BusinessDate=businessdate) data = await aclient.service.OperationName(req_data) # 同步操作移到线程池,不阻塞事件循环 json_data = await asyncio.to_thread(zeep.helpers.serialize_object, data) json_str = await asyncio.to_thread(json.dumps, json_data, default=str) # 任务内部管理Blob客户端上下文 blob_client = blob_service_client.get_blob_client( container='container_name', blob=f'blob_name_{location[0]}' # 确保Blob名称唯一 ) async with blob_client: await blob_client.upload_blob(json_str, overwrite=True) print(f'Finished getting data for {location[0]} at {time.time()-ref}') # 手动关闭客户端释放资源 await async_client.aclose() wsdl_client.close() return json_data async def main(): # 替换为实际的凭证和位置数据 access_token = "your_access_token" location_tokens = [("loc1", "token1"), ("loc2", "token2")] blob_service_client = blob.BlobServiceClient.from_connection_string("your_connection_string") businessdate = datetime.date(2023,3,7) wsdl_url = "Link.to.WSDL" target_address = "actual_address_for_calls" print('Getting Location Data') tasks = [] for location in location_tokens: tasks.append(get_location_data(access_token, location, businessdate, blob_service_client, wsdl_url, target_address)) location_data = await gather_with_concurrency(10, *tasks) ref = time.time() # 无需手动设置ProactorEventLoop,asyncio.run在Windows会自动处理 try: asyncio.run(main()) finally: pass
内容的提问来源于stack exchange,提问作者Daniel Woods

