登录
推荐 文章 Go 技术 课程 下载 专题 AI
首页 >  文章 >  python教程

Python asyncio.Semaphore 如何限制并发:让批量请求在超时后可恢复

来源:17golang原创

时间:2026-08-29 12:14:13 339浏览 收藏

批量请求一多,最先暴露出来的通常不是业务错误,而是同时跑起来的协程数量失控:连接池被占满、对端开始返回 429,最后一个超时还会拖住整批结果。Python 的 asyncio.Semaphore 适合把这类并发限制放在任务入口,配合 asyncio.wait_for 给单个请求设上界,再用 asyncio.gather(return_exceptions=True) 把成功、超时和普通异常收回来。

限制并发的关键不是少创建协程,而是让每个协程在真正访问外部资源前取得 Semaphore 许可,并把超时处理放在单任务边界内。

要点速览
  • asyncio.Semaphore(3) 只允许三个任务同时进入受保护区。
  • 超时应包在单个请求周围,不能让一个失败取消整批任务。
  • gather(return_exceptions=True) 返回结果列表后,再按异常类型做统计与重试。
  • 最终验收要同时看成功数、超时数和仍在运行的任务数。

为什么直接 gather 会把批量请求推到失控

假设一次要检查 30 个对象,直接把 30 个 crawl_batch 协程交给 asyncio.gather,事件循环会尽快把它们都推进到网络调用附近。协程本身很轻,但连接、对端限流和本地解析并不会因为使用异步就无限扩容。

这里先别急着把批次数量硬编码成 3。并发上限应该对应资源边界:如果远端接口允许的稳定并发是 3,就把许可放在真正的外部调用前;解析结果、记录本地状态等不需要占用远端配额的部分,可以放在许可区外。

下面这张图只保留本文代码里真实出现的三个节点,展示任务从 crawl_batch 进入 asyncio.Semaphore 后才到达外部请求的路径。

Python asyncio.Semaphore 限制 crawl_batch 进入外部请求的并发路径

把 Semaphore 放在单个任务的资源边界

import asyncio

async def fetch_one(item, limit):
    async with limit:
        return await request_remote(item)

async def crawl_batch(items):
    limit = asyncio.Semaphore(3)
    tasks = [asyncio.create_task(fetch_one(item, limit)) for item in items]
    return await asyncio.gather(*tasks)

async with limit 会在进入代码块时取得一个许可,离开代码块时自动释放。即使 request_remote 抛出异常,异步上下文管理器也会执行退出逻辑,所以不会把 Semaphore 永久卡住。

常见误区是把整个 fetch_one 都放进一个更大的业务锁,连本地转换也一起排队。更细的边界通常更实用:只有外部请求占用许可,结果清洗和落盘不抢这三个名额。

让超时只影响当前请求

并发上限解决的是“同时有多少个”,不解决“一个请求卡多久”。把 asyncio.wait_for 放在单项请求外层,才能让慢任务及时释放自己的许可,并让批量任务继续处理其他项目。

async def fetch_one(item, limit):
    try:
        async with limit:
            return await asyncio.wait_for(
                request_remote(item),
                timeout=2.5,
            )
    except asyncio.TimeoutError:
        return {"item": item, "status": "timeout"}

超时发生时,wait_for 会取消它包住的等待操作;随后控制流进入 except asyncio.TimeoutError,返回一个可统计的状态。这个返回值比直接吞掉异常更容易在批处理结束后区分“成功”“超时”和“业务失败”。

用 gather 收敛成功、超时和普通异常

如果希望一项失败不影响整批任务,可以让每个任务把预期的超时转成状态对象,同时让意外异常由 gather(return_exceptions=True) 原样收集:

async def crawl_batch(items):
    limit = asyncio.Semaphore(3)
    tasks = [asyncio.create_task(fetch_one(item, limit)) for item in items]
    results = await asyncio.gather(*tasks, return_exceptions=True)

    success = [item for item in results if not isinstance(item, Exception)]
    failed = [item for item in results if isinstance(item, Exception)]
    return success, failed

这里的判断只处理“任务没有抛出 Python 异常”的结果;超时状态是一个普通字典,会进入 success 这一侧,因此生产代码里更建议再按 status 字段拆分,而不是把“有返回值”直接当成业务成功。

Python asyncio.wait_for 超时分支与 gather 结果收敛状态

回归检查要看三个可见结果

不要只打印“批处理完成”。在本地可以记录当前进入请求的数量,验收时至少确认三个结果:

检查项预期现象异常时先看哪里
并发峰值同时进入请求区不超过 3async with limit 是否包住了真正的请求
超时任务出现 status=timeout,其他任务仍有结果wait_for 的包裹范围和超时值
批次收敛gather 返回后没有遗留任务是否创建了未等待的 Task

如果仍然观察到连接数高于 3,优先检查是否还有另一条调用路径绕过了 fetch_one。如果只有超时数上升,不要立刻把并发调大;先看远端响应时间和连接复用情况。

常见问题

Semaphore 会限制任务创建数量吗?

不会。任务仍然可以先创建,Semaphore 限制的是进入受保护代码区的数量。

为什么不用 asyncio.wait 代替 gather?

wait 更适合需要分批取出已完成任务的场景;本文要收齐一批结果并保留异常,所以 gather(return_exceptions=True) 更直接。

超时后应该马上重试吗?

不一定。先确认对端是否限流、超时是否集中发生,再决定是否采用带退避的有限重试,并继续受同一个 Semaphore 约束。

迁移清单

  • 把外部请求封装到 fetch_one,只在请求边界取得 Semaphore 许可。
  • 给单请求设置明确的 asyncio.wait_for 超时,并把可预期超时转换为状态。
  • 使用 gather(return_exceptions=True) 收集整批结果,分别统计成功、超时和异常。
  • 用并发峰值、超时状态和任务收敛三项结果做回归检查。
声明:本文转载于:17golang原创 如有侵犯,请联系study_golang@163.com删除
相关阅读
更多>
最新阅读
更多>
课程推荐
更多>