1. 问题背景:同步调用如何拖垮你的Web服务

先交代一下业务场景:这是一个典型的BFF(Backend For Frontend)订单聚合接口,前端需要一次拿到订单状态、物流轨迹、商品快照三个维度的数据。最初的实现是纯Flask同步代码,按顺序调用三个内部HTTP服务:

# 重构前:同步串行调用
@app.route('/api/order/detail')
def order_detail():
    order_id = request.args.get('order_id')
    # 1. 查订单状态 - 平均120ms
    order_info = requests.get(f'http://order-svc/orders/{order_id}').json()
    # 2. 查物流轨迹 - 平均130ms  
    logistics = requests.get(f'http://logistics-svc/tracks/{order_id}', timeout=2).json()
    # 3. 查商品快照 - 平均110ms
    snapshot = requests.get(f'http://product-svc/snapshots/{order_id}').json()
    # 业务聚合逻辑
    return jsonify({...})

三个上游接口平均耗时分别是120ms、130ms、110ms,串行执行总耗时约360ms,加上JSON序列化和业务逻辑,接口平均响应380ms。这还不是最糟糕的——当并发上来,requests库的同步阻塞导致线程池(Flask默认线程池)被占满,每个线程都在等IO,CPU利用率不到30%,但P99延迟飙升到1.2秒。运维给的监控图显示,数据库连接池也被打满(因为线程等待时持有数据库连接)。

当时的第一反应是“加机器”,但细想不对——加机器只能线性扩展并发量,无法降低单请求延迟,而且成本不低。既然瓶颈在IO等待,正确的解法是异步化,让等待期间去处理其他任务。

2. 环境与版本:别用错库

先明确技术栈。我用的是:

  • Python 3.10.12(内置asyncio,无需额外安装)
  • Flask 2.3.3(同步框架,配合asyncio需要特殊处理,后面讲)
  • httpx 0.25.0(支持异步的HTTP客户端,比aiohttp更顺手,API和requests几乎一致)
  • gunicorn 20.1.0 + gevent worker(生产部署)

为什么不用aiohttp?因为httpx的API设计更接近requests,重构时改动最小,而且支持HTTP/2(虽然没有用到)。为什么不用FastAPI?因为现有代码是Flask,迁移成本高,而且Flask+asyncio也能达到目的。

3. 方案设计:从线程模型到事件循环

核心思路很简单:把三个串行的HTTP调用改成并发发起,利用asyncio的事件循环在等待IO时切换协程。但有几个关键决策点:

第一,怎么让Flask支持异步? Flask是WSGI同步框架,不能直接await。有三种方案:
- 使用asgirefsync_to_async包装(Django生态的,不适用)
- 用asyncio.run()在同步视图里跑一个协程(简单,但每个请求会创建新事件循环,性能差)
- 用loop.run_in_executorasyncio.ensure_future结合全局事件循环(推荐)

我选了第三种方案:在模块加载时创建一个全局事件循环,每个请求通过asyncio.run_coroutine_threadsafe提交协程到该循环。这样事件循环常驻,避免反复创建销毁的开销。

第二,并发度怎么控制? 三个请求全部并发,如果上游服务扛不住怎么办?需要加信号量限制最大并发数。这里设为10,意思是每个请求最多同时发10个HTTP请求(实际只有3个),主要防止上游雪崩时本服务被打爆。

第三,超时怎么处理? 同步版本用的是requests.get(timeout=2),异步版本用asyncio.wait_for包一层,同样设2秒。如果某个上游超时,不能影响其他两个请求的结果。

4. 核心实现:异步重构代码对比

先看重构后的完整代码。注意httpx.AsyncClient要复用,不能每次请求都实例化(开销大)。

# 重构后:asyncio并发调用
import asyncio
import httpx
from flask import Flask, jsonify, request

app = Flask(__name__)
# 全局事件循环和异步客户端
loop = asyncio.new_event_loop()
asyncio.set_event_loop(loop)
client = httpx.AsyncClient(timeout=2.0, limits=httpx.Limits(max_connections=50))
# 信号量限制并发
semaphore = asyncio.Semaphore(10)

async def fetch_url(url: str) -> dict:
    """带信号量和超时控制的异步请求"""
    async with semaphore:
        try:
            resp = await client.get(url)
            resp.raise_for_status()
            return resp.json()
        except Exception as e:
            # 记录日志,返回空字典保证聚合逻辑不崩溃
            app.logger.error(f"fetch {url} failed: {e}")
            return {}

async def fetch_order_data(order_id: str) -> dict:
    """并发获取三个上游数据"""
    tasks = [
        fetch_url(f'http://order-svc/orders/{order_id}'),
        fetch_url(f'http://logistics-svc/tracks/{order_id}'),
        fetch_url(f'http://product-svc/snapshots/{order_id}'),
    ]
    results = await asyncio.gather(*tasks, return_exceptions=False)
    return {
        'order_info': results[0],
        'logistics': results[1],
        'snapshot': results[2],
    }

@app.route('/api/order/detail')
def order_detail():
    order_id = request.args.get('order_id')
    # 提交协程到全局事件循环,等待结果
    future = asyncio.run_coroutine_threadsafe(fetch_order_data(order_id), loop)
    data = future.result(timeout=3)  # 等待最多3秒
    # 业务聚合逻辑(同步代码,不涉及IO)
    return jsonify({...})

关键点解析:

  1. asyncio.new_event_loop()创建全局事件循环,不要在视图函数里asyncio.run()——那样每次请求都创建新循环,开销抵消异步收益。
  2. future.result(timeout=3)是阻塞调用,会阻塞Flask的worker线程。但因为协程内部是异步非阻塞的,三个HTTP请求的总耗时才120ms左右(取最大那个),所以worker线程只阻塞很短时间。这比同步版本让线程空等360ms高效得多。
  3. semaphore是全局的,控制并发请求数。如果某瞬间有100个请求同时进来,最多只有10个HTTP请求在飞行,其余排队等待。

如果用的是FastAPI,代码会更简洁:

# 顺带提一下FastAPI版本(仅供参考)
from fastapi import FastAPI
import httpx

app = FastAPI()
client = httpx.AsyncClient(timeout=2.0)

@app.get("/api/order/detail")
async def order_detail(order_id: str):
    async with asyncio.Semaphore(10):  # 每次请求内限制
        async with client as c:  # 注意:FastAPI推荐每次请求独立client?不,应该复用
        ...

不过这是后话,本次重构以Flask为主。

5. 踩坑与优化:三个意想不到的坑

坑1:协程泄漏导致内存增长
第一次上线后,内存稳步上涨,两天后OOM。排查发现是httpx.AsyncClient没有正确关闭。在同步版本里,每个请求用requests.get(),连接用完自动释放。而异步版本如果每个请求都async with client,会频繁创建和销毁连接池,反而性能差。正确做法是全局单例,但必须保证进程退出时正确释放。我在app.before_first_request里初始化,在app.teardown_appcontext里不关闭(因为多个请求共享),只在进程退出时通过atexit关闭。

坑2:EventLoop被阻塞卡死
有段时间发现接口偶尔全部超时,日志显示asyncio.run_coroutine_threadsafe一直等待。原因是某个第三方库(非异步)在协程里被同步调用了,比如time.sleep(0.5)redis.get()同步操作。这会阻塞整个事件循环,所有协程都无法切换。解决方案:协程里严禁用同步阻塞API,必须全部换成await asyncio.sleep()或异步客户端。最终用grep -r "time.sleep"扫了一遍,改掉了两个隐藏的同步调用。

坑3:gunicorn worker模型不匹配
默认gunicorn用sync worker,每个worker单线程,即使协程并发,一个worker也只能同时处理一个请求。必须用geventuvicorn worker。我用的是gunicorn -k gevent -w 4,gevent worker会为每个请求分配一个greenlet,配合asyncio的事件循环,才能发挥并发效果。如果继续用sync worker,异步化只降低单请求延迟,不提升并发吞吐。

优化1:连接复用
httpx.AsyncClient默认开启HTTP连接复用,但需要显式设置limits=httpx.Limits(max_connections=50, max_keepalive_connections=20)。如果不设置,默认连接池太小(10),高并发时会有连接等待。

优化2:超时分级
三个上游的超时时间不一样:订单服务必须2秒内返回,物流可以容忍3秒,商品快照1.5秒。不能用统一的timeout=2。我给fetch_url增加timeout参数,内部用asyncio.wait_for覆盖:

async def fetch_url(url: str, timeout: float = 2.0) -> dict:
    async with semaphore:
        try:
            resp = await client.get(url, timeout=timeout)
            return resp.json()
        except Exception:
            return {}

6. 效果数据:压测结果对比

locust压测,环境:4核8G云服务器,上游服务用mock_server模拟(每个接口固定延迟120ms)。压测参数:200并发用户,持续5分钟。

指标 同步版本 asyncio版本 提升
平均响应时间 380ms 47ms 88% ↓
P95延迟 520ms 78ms 85% ↓
P99延迟 1.2s 128ms 89% ↓
QPS(吞吐量) 80 req/s 620 req/s 7.75x ↑
CPU利用率 30% 55% 更充分利用
数据库连接占用 200(打满) 45(正常) 78% ↓

QPS提升7倍多,主要原因是原来一个请求占一个线程(Flask默认线程池),线程等待IO时无法处理新请求。异步化后,一个worker线程能同时处理多个请求的协程,线程数不需要增加。

额外收益: 数据库连接占用下降明显。因为每个请求持有数据库连接的时间从380ms降到47ms(业务聚合逻辑仍然访问数据库),连接释放更快,池子不用开那么大。

7. 总结:异步化的适用边界

这次重构效果显著,但并非所有场景都适合。总结几点适用条件:

  1. IO密集且多依赖调用:如果接口只查一次数据库(IO耗时20ms),异步化收益不大,反而增加复杂度。
  2. 下游服务支持并发:如果上游接口本身有严格限流,并发调用会触发429,需要降级。
  3. 团队熟悉asyncio:异步代码的调试比同步难得多,尤其协程泄漏和死锁问题,需要经验。

最后给个建议:如果新项目直接用FastAPI(原生异步),别用Flask硬凑。但存量Flask项目,通过全局事件循环+run_coroutine_threadsafe的方式,可以低成本的获得异步收益,不必重写整个服务。


附:完整代码已上传Gist,见评论区链接。有问题欢迎交流,尤其是协程泄漏的排查过程比较曲折,如果有人感兴趣可以单独写一篇。