并发与异步执行
FastAPI 在一个 asyncio 事件循环上处理所有请求。在该循环上进行的阻塞调用——同步的 SQLAlchemy 查询、pandas 计算、yfinance 下载——会使进程中所有其他请求和 WebSocket 停顿,直到该调用返回。
run_sync()
concurrency.py 提供了从异步代码中执行阻塞代码的唯一方式:
from concurrency import run_sync
result = await run_sync(database.get_asset_context, ticker=ticker)
run_sync 在工作线程中运行函数(asyncio.to_thread),传播上下文变量,并接受关键字参数。tests/test_async_boundary.py 检查 main.py、路由器、用例和功能包只通过它来调用阻塞代码。
event loop ──► async route ──► await run_sync(blocking_call) ──► worker thread
│ │
└── keeps serving other requests and WebSockets ◄───────────────┘
已经是异步的部分
- 历史查询、新闻、监控和实时 worker 使用异步 SQLAlchemy 引擎(asyncpg),由
DATABASE_URL派生。 - 数据提供方通过
BaseFinancialProvider使用httpx.AsyncClient;同一请求中的各 数据提供方并行运行(asyncio.gather),每个都有各自的并发请求数限制。 - 缓存使用异步 Redis 客户端,
CacheService除外(同步,通过run_sync调用)。
后台任务
| 任务 | 运行方式 |
|---|---|
| 实时数据流 | 线程池中的 TradingView 客户端,同时最多 TV_MAX_CONNECTIONS 个;tick 交回事件循环处理 |
| 每日金丝雀 | APScheduler AsyncIOScheduler,UTC 时间 CANARY_RUN_HOUR |
| 使用日志 | 由中间件放入队列,由后台任务批量写入;响应从不等待它 |
| 新闻刷新、按需金丝雀运行 | FastAPI 后台任务 |
贡献者指南
- 使用
async def编写路由。 - 使用
await run_sync(service.method, ...)调用同步服务;切勿在协程中直接调用。 - 直接 await 异步服务(Redis asyncio、
httpx、asyncpg)。 - 保持只有一个 Gunicorn worker:数据流、WebSocket 客户端和金丝雀都存在于进程内存中。