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本身不支持异步视图,但我们可以用asgiref的async_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) - 重试策略:
asyncio的retry装饰器,最多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.gather的return_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做链路追踪,否则异步代码出问题排查难度你懂的。我踩过的坑都在上面了,希望你能少走弯路。