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

使用httpx+asyncio发送10万GET请求遇ConnectError问题求助

问题分析与解决方案

核心问题根源

  1. TCP端口耗尽与TIME_WAIT堆积:系统临时端口池容量有限(Linux默认约2.8万个),短时间内大量创建/销毁连接会导致端口被TIME_WAIT状态占用,无法及时复用,触发TCP Port numbers reused提示和httpx.ConnectError。
  2. 无超时配置风险:timeout=httpx.Timeout(None)会让连接无限等待,进一步占用端口和连接资源,加剧端口耗尽问题。
  3. 批量任务过载:一次性创建10万+请求任务,即使有连接池限制,也会导致底层网络栈处理压力过大,连接队列溢出。

针对性解决方案

1. 调整连接池参数

降低最大活跃连接数(建议200-500,根据系统和目标服务器承受能力调整),同时开启长连接复用,减少连接创建开销:

limits = httpx.Limits(
    max_connections=300,          # 降低并发连接数,避免端口耗尽
    max_keepalive_connections=100, # 保持长连接数量
    keepalive_expiry=30           # 长连接过期时间,单位秒
)

2. 添加合理超时配置

绝对禁止设置timeout=None,给连接、读取、池等待都设置超时,避免无效连接占用资源:

timeout = httpx.Timeout(
    connect=5,  # 连接超时5秒
    read=10,    # 读取响应超时10秒
    pool=30     # 连接池等待超时30秒
)

3. 分批次处理请求

不要一次性创建所有任务,改用分批执行,控制每批处理量,给系统和服务器喘息空间:

async def process_batch(client, items, url_template):
    tasks = [client.get(url_template.format(item), timeout=timeout) for item in items]
    # 捕获异常避免单个请求失败导致整批任务崩溃
    return await asyncio.gather(*tasks, return_exceptions=True)

async def search():
    guids = list(range(688001, 800000))
    batch_size = 500  # 每批处理500个请求,可根据实际调整
    timeout = httpx.Timeout(connect=5, read=10, pool=30)
    limits = httpx.Limits(max_connections=300, max_keepalive_connections=100)

    # 处理第一阶段请求
    blob_list = []
    async with httpx.AsyncClient(timeout=timeout, limits=limits) as client:
        for i in range(0, len(guids), batch_size):
            batch_guids = guids[i:i+batch_size]
            results = await process_batch(client, batch_guids, f"{SOME_URL}{{}}")
            # 过滤成功响应,处理异常
            for res, guid in zip(results, batch_guids):
                if isinstance(res, httpx.Response) and res.status_code == 200:
                    try:
                        blob_name = json.loads(res.text)["blob_name"]
                        blob_list.append((blob_name, guid))
                    except Exception as e:
                        print(f"解析响应失败 guid={guid}: {str(e)}")
                else:
                    print(f"请求失败 guid={guid}: {str(res)}")

    # 处理第二阶段请求
    game_results = []
    if blob_list:
        blobs, guids = zip(*blob_list)
        async with httpx.AsyncClient(timeout=timeout, limits=limits) as client:
            for i in range(0, len(blobs), batch_size):
                batch_blobs = blobs[i:i+batch_size]
                batch_guids = guids[i:i+batch_size]
                results = await process_batch(client, batch_blobs, f"{ANOTHER_URL}{{}}")
                for res, guid in zip(results, batch_guids):
                    if isinstance(res, httpx.Response) and res.status_code == 200:
                        game_results.append(guid)
                    else:
                        print(f"获取游戏结果失败 blob={batch_blobs[i]} guid={guid}: {str(res)}")

    print(game_results)

4. 系统层面优化(可选)

如果是Linux系统,可调整TCP参数减少TIME_WAIT影响(需root权限):

# 允许复用TIME_WAIT状态的端口
echo 1 > /proc/sys/net/ipv4/tcp_tw_reuse
# 加快TIME_WAIT端口回收,默认60秒改为30秒
echo 30 > /proc/sys/net/ipv4/tcp_fin_timeout

5. 异常重试增强(可选)

对失败的请求加入重试逻辑(注意确保接口幂等):

from tenacity import retry, stop_after_attempt, wait_exponential

@retry(stop=stop_after_attempt(3), wait=wait_exponential(multiplier=1, min=2, max=10))
async def fetch_with_retry(client, url):
    response = await client.get(url, timeout=timeout)
    response.raise_for_status()
    return response

# 在process_batch中替换为带重试的调用
tasks = [fetch_with_retry(client, url_template.format(item)) for item in items]

高效发送大量请求的最佳实践

  • 严格控制并发数:根据目标服务器QPS限制和自身系统资源,合理设置并发数,避免过载。
  • 强制使用长连接:开启HTTP keep-alive,减少TCP握手开销。
  • 分批次执行任务:避免一次性提交所有任务,分批执行降低系统压力。
  • 设置合理超时:杜绝无超时配置,避免无效连接占用资源。
  • 异常捕获与重试:对失败请求进行重试,同时记录失败日志方便后续排查。

内容的提问来源于stack exchange,提问作者Omer Dagry

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 15:35:51