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

Zeep AsyncClient异步功能未达预期,求问题排查

问题: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、绑定地址)是全局共享的。多个并发任务同时修改这些值,会导致请求串用参数、状态竞争,甚至被迫变成串行执行——这直接抵消了异步的优势,大数据量时这个问题被放大。

具体修复点&优化建议

  1. 每个任务独立创建AsyncClient和Transport
    不要共享这些实例,每个任务单独初始化,避免状态互相干扰。

  2. 把同步数据处理移出事件循环
    zeep.helpers.serialize_object和json.dumps是同步阻塞操作,处理2MiB数据时会卡住事件循环。用asyncio.to_thread把这些操作放到线程池执行,不影响其他任务并发。

  3. 修正Blob客户端的上下文管理
    你在循环里创建blob_client并通过async with托管,但随即把它传给任务,会导致上下文提前退出、客户端被关闭。应该在任务内部创建并管理Blob客户端的上下文。

  4. 调整并发数与连接池匹配
    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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 03:13:16