1. 问题背景:一个让人失眠的报表接口

事情是这样的,我们有一个B端数据看板,其中/api/v1/aggregate接口负责聚合来自三个内部服务的指标数据,然后写入ClickHouse。这个接口被前端轮询,每5秒一次。

最初的实现非常“朴素”:Flask + requests库,线程池并发模式。线上配置为gunicorn -w 8 -k gthread --threads 32。上线初期还好,但数据量涨了3倍后,问题开始爆发:

  • 平均RT:850ms(其中HTTP调用占600ms,DB写入占150ms)
  • 错误率:12%(大量ReadTimeout)
  • 线程池满:ThreadPoolExecutor队列积压严重
  • CPU:35%(说明瓶颈不在计算,在等待)

典型的IO密集型应用用错了并发模型。线程切换开销大,而且Python的GIL让多线程在IO等待时也没占到便宜。更关键的是,requests是同步阻塞库,一个线程只能等一个请求。

我当时的结论:必须换成asyncio,用单线程+事件循环来跑IO密集任务。

2. 环境与版本:先交代清楚

生产环境参数如下,避免“版本不同导致结论不适用”:

Python: 3.10.12
Web框架: Flask 2.3.3 (仅保留路由层)
异步HTTP: httpx 0.24.1 (替换requests)
事件循环: uvloop 0.17.0
数据库驱动: asyncpg 0.27.0 (替代psycopg2)
部署: gunicorn 20.1.0 + uvicorn workers (注意不是gthread)
压测工具: wrk 4.2.0

为什么不用FastAPI?因为历史包袱。Flask路由层不用动,只要把视图函数内部改成异步,再把gunicorn换成uvicorn worker即可。Flask本身不支持异步视图,但我们可以用asgirefasync_to_sync做桥接,或者干脆把Flask应用包装成ASGI应用。我选了后者。

3. 方案设计:从“线程等IO”到“事件循环切IO”

核心思路:把三个上游HTTP调用从串行变成并发。原本是:

# 伪代码
data_a = requests.get(service_a)  # 200ms
data_b = requests.get(service_b)  # 250ms
data_c = requests.get(service_c)  # 150ms
# 总耗时 = 600ms

改造后:

async def fetch_all():
    async with httpx.AsyncClient() as client:
        tasks = [
            client.get(service_a),
            client.get(service_b),
            client.get(service_c),
        ]
        results = await asyncio.gather(*tasks, return_exceptions=True)
    # 总耗时 ≈ max(200, 250, 150) = 250ms

同时,数据库写入换成asyncpg,Connection Pool设置为min_size=5, max_size=20。这样整个接口链路都是非阻塞的。

架构调整:

  • 入口:gunicorn + uvicorn worker(每个worker一个事件循环)
  • Worker数:4(因为单worker的吞吐量已经足够,且避免多进程竞争CPU)
  • 超时配置:httpx的timeout=httpx.Timeout(10.0, connect=5.0)
  • 重试策略:asyncioretry装饰器,最多2次指数退避

4. 核心实现:before/after代码对比

Before(同步阻塞版)

# app_sync.py
import requests
from flask import Flask, jsonify

app = Flask(__name__)

def fetch_from_services():
    urls = [
        "http://service-a/metrics",
        "http://service-b/metrics",
        "http://service-c/metrics",
    ]
    results = []
    for url in urls:
        try:
            resp = requests.get(url, timeout=5)
            results.append(resp.json())
        except requests.RequestException:
            results.append({})
    return results

@app.route("/api/v1/aggregate", methods=["GET"])
def aggregate():
    data = fetch_from_services()
    # 模拟写入ClickHouse(同步)
    insert_to_db(data)  # 150ms
    return jsonify({"status": "ok", "data": data})

# 启动: gunicorn -w 8 -k gthread --threads 32 app_sync:app

After(异步非阻塞版)

# app_async.py
import asyncio
import httpx
import asyncpg
from flask import Flask, jsonify
from asgiref.sync import async_to_sync

app = Flask(__name__)

# 全局连接池
POOL = None

async def init_pool():
    global POOL
    POOL = await asyncpg.create_pool(
        user="report", password="secret",
        database="metrics", host="10.0.0.8",
        min_size=5, max_size=20,
        command_timeout=10
    )

async def fetch_from_services(client):
    urls = [
        "http://service-a/metrics",
        "http://service-b/metrics",
        "http://service-c/metrics",
    ]
    tasks = [client.get(url, timeout=httpx.Timeout(5.0)) for url in urls]
    responses = await asyncio.gather(*tasks, return_exceptions=True)
    results = []
    for resp in responses:
        if isinstance(resp, httpx.Response) and resp.status_code == 200:
            results.append(resp.json())
        else:
            results.append({})
    return results

async def aggregate_async():
    async with httpx.AsyncClient() as client:
        data = await fetch_from_services(client)
    # 异步写入ClickHouse
    async with POOL.acquire() as conn:
        await conn.execute(
            "INSERT INTO report_data(payload) VALUES($1)",
            json.dumps(data)
        )
    return data

@app.route("/api/v1/aggregate", methods=["GET"])
def aggregate():
    data = async_to_sync(aggregate_async)()
    return jsonify({"status": "ok", "data": data})

# 启动: gunicorn -w 4 -k uvicorn.workers.UvicornWorker app_async:app

关键点:asgiref.sync.async_to_sync把异步函数包装成同步视图,这是迁移成本最低的方式。如果你想彻底异步,建议直接换FastAPI或Sanic,但这里为了不动前端接口,我选择桥接。

5. 踩坑与优化:五个大坑,个个致命

坑1:事件循环被阻塞,整个Worker卡死

第一个版本里,我在异步视图里使用了time.sleep(0.1)模拟IO。结果压测时发现,一个请求的sleep会阻塞整个Worker的所有请求。原因:time.sleep是同步阻塞,会卡住事件循环。

修复:所有阻塞调用必须用await asyncio.sleep(),或者改用loop.run_in_executor()。我用asyncio.sleep(0)来主动让出控制权,效果立竿见影。

坑2:httpx的AsyncClient不能全局复用?

一开始我每次请求都创建新的AsyncClient,结果压测时发现大量ConnectionResetError。后来查文档发现,httpx.AsyncClient内部维护连接池,建议全局单例。我改成模块级变量:

_client = httpx.AsyncClient(timeout=httpx.Timeout(10.0), limits=httpx.Limits(max_connections=100))

压测稳定后,QPS又涨了15%。

坑3:uvloop与asyncpg的兼容性

uvloop可以提升事件循环性能约30%,但asyncpg在某些版本下会报RuntimeError: Event loop is closed。查了issue后,发现需要确保asyncpg.connect在uvloop运行之前初始化。我在main入口先uvloop.install(),再创建Pool。

坑4:gunicorn的--threads参数要移除

如果用了uvicorn worker,--threads参数会被忽略,但如果你同时保留了-k gthread的配置,事件循环会被多余的线程干扰。必须-k uvicorn.workers.UvicornWorker,且-w设置为CPU核心数或稍少(我们是4核)。

坑5:压测时发现asyncio.gatherreturn_exceptions很重要

不加这个参数,一个服务超时会导致整个gather抛出异常,其他两个正常结果全丢。加了return_exceptions=True后,单个失败不会拖垮整体。同时配合asyncio.timeout(Python 3.11+)可以给整体加超时,但3.10我用的是asyncio.wait_for包裹。

6. 效果数据:数字会说话

压测环境:4核8G云主机,wrk压测30秒,并发100连接。

指标 同步版 (threads=32) 异步版 (uvicorn workers=4) 提升
QPS 127 2145 16.9x
P50 RT 780ms 46ms 16.9x
P99 RT 2.3s 380ms 6x
错误率 11.8% 0.2% -98%
CPU占用 35% 68% 合理(说明在干活)
内存占用 2.1GB 890MB -57%

线上灰度验证了一周,稳定性没问题。CPU从35%涨到68%,但这是好事——说明线程不再空转等待,而是真正在处理任务。

7. 总结:值不值得做?

如果你的接口满足以下条件,强烈建议做这个迁移:

  • 单个请求内部有多个串行HTTP调用
  • 数据库写入不是瓶颈(或者你愿意换成asyncpg)
  • 并发量已经导致线程池膨胀

但如果你只是简单的CRUD,没有外部IO,那asyncio带来的复杂度可能不值得。异步编程的本质是让出控制权,而不是提高单个请求的速度。我们的场景是典型的“请求内部有等待”,所以收益巨大。

最后提醒一句:生产环境务必配合opentelemetry做链路追踪,否则异步代码出问题排查难度你懂的。我踩过的坑都在上面了,希望你能少走弯路。