ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

第17章:FastAPI异步编程与非阻塞 IO 实战

第17章:FastAPI异步编程与非阻塞 IO 实战 1. 项目背景业务场景聚合报价服务需要调用 3 个第三方 API物流运费、支付手续费、汇率换算然后计算出最终报价。小赵用最直观的方式实现app.get(/quote)defget_quote(product_id:int):shippingrequests.get(fhttps://api.shipping.com/calc?product{product_id})# 800msfeerequests.get(fhttps://api.payment.com/fee?product{product_id})# 600msraterequests.get(fhttps://api.forex.com/rate?fromUSDtoCNY)# 400mstotalshipping.json()[cost]fee.json()[fee]rate.json()[rate]return{total:total}接口响应时间800 600 400 1800ms。小赵想FastAPI 不是号称高性能吗怎么一个接口要 1.8 秒他尝试把def改成async defapp.get(/quote)asyncdefget_quote(product_id:int):# 加了 asyncshippingrequests.get(...)# 还是同步 requests...结果还是 1.8 秒而且并发 QPS 反而下降了。服务器 4 核 CPU100 个并发请求CPU 使用率只有 15%——因为所有协程都被阻塞在requests.get()上。痛点不掌握 Python 异步模型的核心原理FastAPI 的高并发能力完全是无效的伪异步async def里面调同步requests.get()——协程阻塞事件循环卡死这是最典型的 FastAPI 性能陷阱。串行等待3 个 API 顺序调用总耗时 最慢 API × 3。明明可以并发却串行执行。连接数爆炸每次请求新建一个 HTTP 连接三次握手 TLS 握手高并发下连接数超限。超时失控某个第三方 API 挂掉接口 hang 住 30 秒才报错——线程池沾满新请求排队等待。FastAPI 是 ASGI 框架它的高性能建立在async/await 非阻塞 IO之上。不理解这个模型就等于买了跑车但一直挂一档开。2. 项目设计场景小赵在监控面板上看到报价接口 P99 延迟 3.2 秒。大师走过来指着屏幕。小胖震惊“3.2 秒用户早关页面了。FastAPI 不是 Python 最快的框架吗这跟 Flask 有区别吗”小白“问题不在 FastAPI在小赵的代码。你看第 1 章我们讲过——async def里的同步阻塞 IOrequests.get()会卡住事件循环。但不止如此——他还串行调了 3 个 API。就像你去食堂打饭先排队打饭、再排队打菜、再排队打汤——为什么不三个窗口一起排”大师小白这个比喻好。今天我们把 Python 异步的三层概念讲透大家以后写 FastAPI 就不会踩坑第一层——协程是什么协程coroutine是一个可以在中途暂停和恢复的函数。Python 的async def定义协程await是暂停点。暂停时事件循环去执行其他协程。这就好比你在微波炉热饭的 3 分钟里顺便去洗了个水果——而不是干等着微波炉叮。技术映射Python 的asyncio是基于事件循环的单线程并发模型。await点 协程交出控制权。当你在async def里调同步阻塞函数如time.sleep(3)、requests.get()控制权交不出去——事件循环被卡住其他协程全部冻结。这叫做协程的协作式调度——你必须主动await。小赵“那我理解了——不能混用async def里必须用异步库。但httpx.AsyncClient为什么就比requests.get()好在 async 环境里”小白“requests.get()底层是同步 socket——socket.send()socket.recv()Python 线程在内核 I/O 上阻塞。而httpx.AsyncClient.get()是用asyncio的非阻塞 socket——当数据还没到达时它立刻交还事件循环控制权让其他协程继续执行。”大师“对。我再补一个容易忽略的细节——连接复用”# ❌ 串行 每次新建连接慢asyncdefbad():shippingawaithttpx.AsyncClient().get(url1)# 新建连接TCPTLSfeeawaithttpx.AsyncClient().get(url2)# 又新建连接rateawaithttpx.AsyncClient().get(url3)# 又新建连接# ✓ 串行 连接复用中asyncdefbetter():asyncwithhttpx.AsyncClient()asclient:shippingawaitclient.get(url1)feeawaitclient.get(url2)rateawaitclient.get(url3)# ✓✓ 并发 连接复用快asyncio.gather 同时发起三个请求asyncdefbest():asyncwithhttpx.AsyncClient()asclient:shipping,fee,rateawaitasyncio.gather(client.get(url1),client.get(url2),client.get(url3),)技术映射asyncio.gather()同时启动多个协程。总耗时 ≈ max(800ms, 600ms, 400ms) 800ms——比串行的 1800ms 快了 2.25 倍。httpx.AsyncClient内部维护一个连接池对同一 host 复用 TCP 连接省去三次握手和 TLS 握手。小胖“那如果 3 个 API 有依赖怎么办——第二个 API 的请求参数依赖第一个 API 的返回值”大师“那就是经典的’串行依赖’——没法并发。但可以优化把独立的部分并发依赖的部分串行。”# 假设报价需要运费汇率但汇率调用前需要先获取用户的国家代码asyncdefdependent():asyncwithhttpx.AsyncClient()asclient:# 并发运费和用户信息可以同时查shipping,user_infoawaitasyncio.gather(client.get(shipping_url),client.get(user_url),)# 串行汇率依赖用户的国家代码countryuser_info.json()[country]rateawaitclient.get(fhttps://api.forex.com/rate?country{country})returnshipping.json()[cost]rate.json()[rate]3. 项目实战——构建高性能报价服务环境准备pipinstallhttpx0.27.0 pytest-asyncio0.24.0分步实现步骤一搭建异步 HTTP 客户端目标连接复用 超时控制app/infrastructure/http_client.pyimporthttpxfromapp.core.configimportsettingsclassAsyncHTTPClient:异步 HTTP 客户端 —— 全局单例连接池复用_instance:httpx.AsyncClient|NoneNoneclassmethodasyncdefget_client(cls)-httpx.AsyncClient:ifcls._instanceisNone:cls._instancehttpx.AsyncClient(timeouthttpx.Timeout(connect5.0,# TCP 连接超时read10.0,# 读取响应超时write5.0,# 发送请求超时pool5.0,# 等待连接池可用连接超时),limitshttpx.Limits(max_keepalive_connections20,# 最大保活连接数max_connections50,# 总连接上限keepalive_expiry30,# 保活时间秒),)returncls._instanceclassmethodasyncdefclose(cls):ifcls._instance:awaitcls._instance.aclose()cls._instanceNone步骤二实现三种模式的报价服务目标直观对比性能差异app/domains/quote/service.pyimporttimeimportasyncioimporthttpxfromapp.infrastructure.http_clientimportAsyncHTTPClient# 模拟的第三方 API URL实际环境需替换SHIPPING_APIhttp://localhost:9001/shippingPAYMENT_APIhttp://localhost:9002/payment-feeFOREX_APIhttp://localhost:9003/forex-rateclassQuoteService:报价服务 —— 演示三种调用模式的性能差异# ═══════ 模式一同步串行最慢═══defquote_sync_serial(self,product_id:int)-dict:同步串行每个请求阻塞 0.5-1sstarttime.perf_counter()resp1httpx.get(f{SHIPPING_API}?product{product_id})# 阻塞resp2httpx.get(f{PAYMENT_API}?product{product_id})# 阻塞resp3httpx.get(FOREX_API)# 阻塞elapsedtime.perf_counter()-startreturn{mode:sync_serial,shipping:resp1.json().get(cost,0),fee:resp2.json().get(fee,0),rate:resp3.json().get(rate,0),elapsed_ms:round(elapsed*1000,2),}# ═══════ 模式二异步串行快于同步但未利用并发═══asyncdefquote_async_serial(self,product_id:int)-dict:异步串行非阻塞但顺序执行starttime.perf_counter()asyncwithhttpx.AsyncClient()asclient:resp1awaitclient.get(f{SHIPPING_API}?product{product_id})resp2awaitclient.get(f{PAYMENT_API}?product{product_id})resp3awaitclient.get(FOREX_API)elapsedtime.perf_counter()-startreturn{mode:async_serial,shipping:resp1.json().get(cost,0),fee:resp2.json().get(fee,0),rate:resp3.json().get(rate,0),elapsed_ms:round(elapsed*1000,2),}# ═══════ 模式三异步并发最快═══asyncdefquote_async_concurrent(self,product_id:int)-dict:异步并发三个请求同时发出总耗时 max(单个耗时)starttime.perf_counter()clientawaitAsyncHTTPClient.get_client()shipping_taskclient.get(f{SHIPPING_API}?product{product_id})payment_taskclient.get(f{PAYMENT_API}?product{product_id})forex_taskclient.get(FOREX_API)# asyncio.gather 同时执行三个协程resp1,resp2,resp3awaitasyncio.gather(shipping_task,payment_task,forex_task,# return_exceptionsTrue # 单个失败不影响其他)elapsedtime.perf_counter()-startreturn{mode:async_concurrent,shipping:resp1.json().get(cost,0),fee:resp2.json().get(fee,0),rate:resp3.json().get(rate,0),elapsed_ms:round(elapsed*1000,2),}步骤三增加并发控制目标使用 Semaphore 限制并发数classQuoteService:# ... 上面代码 ...# 信号量限制同时调用第三方 API 的并发数_semaphoreasyncio.Semaphore(10)asyncdefquote_with_limit(self,product_id:int)-dict:带并发限制的报价——防止打爆第三方 APIasyncwithself._semaphore:returnawaitself.quote_async_concurrent(product_id)步骤四创建报价 API 路由目标在接口中对比三种模式app/domains/quote/api.pyfromfastapiimportAPIRouter,Queryfromapp.domains.quote.serviceimportQuoteService routerAPIRouter(prefix/quote,tags[报价服务])quote_serviceQuoteService()router.get(/sync,summary同步串行报价慢)defquote_sync(product_id:intQuery(...,gt0)):def 端点 → 在线程池中执行不阻塞事件循环return{code:0,data:quote_service.quote_sync_serial(product_id)}router.get(/async-serial,summary异步串行报价)asyncdefquote_async_serial(product_id:intQuery(...,gt0)):return{code:0,data:awaitquote_service.quote_async_serial(product_id)}router.get(/async-concurrent,summary异步并发报价推荐)asyncdefquote_async_concurrent(product_id:intQuery(...,gt0)):return{code:0,data:awaitquote_service.quote_async_concurrent(product_id)}步骤五启动模拟服务并对比性能# 启动三个模拟的第三方 APIpython scripts/mock_apis.py# 起 3 个简单的 HTTP 服务每个 500-1000ms 延迟# 启动主服务uvicorn app.main:app--reload# ── 1. 同步串行 ──curl-shttp://localhost:8000/api/v1/quote/sync?product_id1|python-mjson.tool# elapsed_ms: 1850 ← 三个 API 延迟之和# ── 2. 异步串行 ──curl-shttp://localhost:8000/api/v1/quote/async-serial?product_id1|python-mjson.tool# elapsed_ms: 1800 ← 依然很慢虽然非阻塞但顺序执行# ── 3. 异步并发 ──curl-shttp://localhost:8000/api/v1/quote/async-concurrent?product_id1|python-mjson.tool# elapsed_ms: 620 ← 仅等于最慢的那个 API 延迟# ── 4. 并发压测比较 QPS ──# 同步模式 100 并发下 QPS ~50线程池耗尽# 异步并发模式 100 并发下 QPS ~800事件循环充分利用完整代码清单本章完整代码见column/code/chapter17/主要文件app/infrastructure/http_client.py异步 HTTP 客户端app/domains/quote/service.py三种模式的报价服务app/domains/quote/api.py报价 API 路由测试验证importpytestimportasynciofromapp.domains.quote.serviceimportQuoteServicepytest.mark.asyncioasyncdeftest_async_concurrent_is_parallel():验证 asyncio.gather 真正实现了并发总耗时 各任务之和serviceQuoteService()asyncdeffast_task():awaitasyncio.sleep(0.1)returnfastasyncdefslow_task():awaitasyncio.sleep(0.3)returnslow# 并发执行总耗时应接近 max(0.1, 0.3) 0.3sstartasyncio.get_event_loop().time()resultsawaitasyncio.gather(fast_task(),slow_task())elapsedasyncio.get_event_loop().time()-startassertelapsed0.35# 远小于 0.4串行之和assertresults[fast,slow]4. 项目总结优点 缺点对比模式async/await asyncio.gather多线程 (ThreadPoolExecutor)多进程Node.js 事件循环IO 并发优秀协程切换零开销中线程切换有开销低进程切换开销大优秀CPU 密集型差阻塞事件循环中受 GIL 限制优秀差编程模型async/await学习曲线中同步代码 线程池同步代码async/await内存占用极低一个协程 ~1KB高一个线程 ~8MB极高极低适用场景✓ 异步并发适用聚合多个下游 API 的 BFFBackend for Frontend接口需要同时查询多个数据库/缓存的只读接口WebSocket 长连接管理文件批量处理并发读写多个文件微服务间批量调用✗ 不适合异步CPU 密集型计算图片处理、加密解密——用def端点在独立线程池执行只有单一数据源的简单 CRUD——async 带来的收益不明显注意事项不要混用同步库async def函数内不要调time.sleep()、requests.get()、同步数据库驱动。用asyncio.sleep()、httpx.AsyncClient、asyncpg。asyncio.gather的 return_exceptions默认False——任一协程异常gather立即抛异常其他协程被取消。设return_exceptionsTrue让单个失败不影响整体。Semaphore 不是全局并发限制asyncio.Semaphore只限制当前事件循环内的并发。多 Worker 进程下需要 Redis 等外部计数器做全局限流。连接池耗尽表现大量httpx.PoolTimeout异常。调大max_connections或增加keepalive_expiry加速连接回收。常见踩坑经验案例一async def端点中的time.sleep()卡死事件循环现象100 并发请求只有一个请求在执行其余 99 个排队——QPS 只有 0.5。根因开发者在async def函数中调了time.sleep(2)事件循环被阻塞 2 秒。解决await asyncio.sleep(2)或改用def端点让线程池处理。案例二asyncio.gather中一个任务挂起导致所有任务超时现象3 个 API 并发调用其中一个超时 30s其余两个 200ms 就返回了一直被拦住。根因gather默认等待所有任务完成才返回。解决为每个任务单独设置 timeout —asyncio.wait_for(task, timeout5)或使用asyncio.as_completed()先返回先处理。案例三httpx.AsyncClient提前关闭现象服务启动正常运行几分钟后所有外部 API 调用报RuntimeError: Event loop is closed。根因在 Lifespan 中创建了AsyncClient但在某次异常中没有正确关闭。下次请求时复用了一个半关闭的 client。解决在app的 lifespan 事件中管理 client 的创建和关闭或每次请求创建新的AsyncClient性能略低但更安全。思考题初级修改报价服务新增一个超时兜底模式——如果某个 API 在 1 秒内未响应使用缓存中的上一次数据作为兜底stale-while-revalidate 策略。进阶如何使用asyncio.TaskGroupPython 3.11替代asyncio.gatherTaskGroup相比gather的优势是什么提示结构化并发。答案提示第 1 题使用asyncio.wait_for(task, timeout1)配合缓存。第 2 题TaskGroup是 Python 的结构化并发原语——如果组内任一任务抛异常所有子任务自动取消不会出现孤儿协程。第 37 章深入事件循环诊断与性能极限。延伸阅读与资源NumPy 从入门到生产落地全链路实战指南科学计算/向量化Redis 8 实战精讲从 CRUD 到源码构建高可用缓存系统Redis 实战修炼与原理进阶Python 3实战精进从脚本到高并发订单引擎python入门Rquests从菜鸟脚本到企业级SDK的网络实战圣经Milvus向量数据库实战修炼从 0 到 1精通向量检索与生产落地MongoDB 实战进阶与内核修炼后端工程师的 AI 转型第一课Ollama 与私有化大模型实战10倍开发者的 Dify 魔法书从零构建全栈 AI 应用后端工程师转型AI第一课-Ollama 与私有化大模型实战大型语言模型(LLM) vLLM 高性能推理落地实战Agent开发之LlamaIndex 实战修炼与源码进阶大语言模型Transformers 实战修炼与源码剖析
返回列表