Spaces:
Running
Running
Worker Pool 饱和问题修复方案
问题背景
在 HuggingFace Space 部署的 stock_data_api_service 中,以下端点频繁出现 timeout / cooldown / worker_pool_saturated 错误:
stock_search(/api/v1/search)boards_flow(/api/v1/boards/flow、/boards/concepts/flow、/boards/industries/flow)sector_concept_hot(/api/v1/boards/concepts/flow)
根本原因
HuggingFace Space 的 Worker 限制
- Free tier:1-2 个并发 worker
- Pro tier:10+ 个并发 worker
- 每个慢请求(15-60s)会占用一个 worker,导致 worker pool 饱和
无全局并发控制
source_runner.py中的call_with_timeout每次都创建新的ThreadPoolExecutor- 没有全局的并发限制和请求排队机制
- 多个慢请求同时到达,立即超出 HF 限制
网络请求慢
- eastmoney/AKShare 接口响应时间不稳定(15-60s+)
- 当前
SOURCE_TIMEOUT_SECONDS = 15太短,但增加超时会占用更多 worker
修复方案
1. 引入全局源调用池(GlobalSourcePool)
修改文件: app/services/source_runner.py
class GlobalSourcePool:
"""全局源调用池,控制并发数,防止 worker pool 饱和"""
def __init__(self, max_concurrent: int = 2, max_workers: int = 4):
self._max_concurrent = max_concurrent
self._semaphore = threading.Semaphore(max_concurrent)
self._executor = concurrent.futures.ThreadPoolExecutor(
max_workers=max_workers,
thread_name_prefix="source_pool"
)
self._stats = {
"active": 0,
"queued": 0,
"completed": 0,
"failed": 0,
"timeout": 0,
}
关键特性:
- ✅ 使用
Semaphore控制最大并发数 - ✅ 线程池复用,不再每次创建新的 ThreadPoolExecutor
- ✅ 统计指标(active, queued, completed, failed, timeout)
- ✅ 懒初始化,全局单例
2. 添加环境变量配置
修改文件: app/core/config.py
# 源调用池配置
max_concurrent_sources: int = _int_env("MAX_CONCURRENT_SOURCES", 2) # HF Free: 2, Pro: 10
source_pool_workers: int = _int_env("SOURCE_POOL_WORKERS", 4) # 线程池大小
配置项说明:
| 配置项 | 默认值 | HF Free | HF Pro | 说明 |
|---|---|---|---|---|
MAX_CONCURRENT_SOURCES |
2 | 2 | 10 | 最大并发源调用数 |
SOURCE_POOL_WORKERS |
4 | 4 | 10 | 线程池大小 |
SOURCE_TIMEOUT_SECONDS |
15 | 15-30 | 15 | 单个请求超时 |
3. 添加监控端点
修改文件: app/api/routes.py
@router.get("/admin/source-pool-stats")
def source_pool_stats():
"""获取源调用池的统计信息"""
stats = get_source_pool_stats()
return {
"pool_stats": stats,
"recommendations": {
"worker_pool_saturated": stats.get("active", 0) >= stats.get("max_concurrent", 2),
"suggestion": "如果 active >= max_concurrent,说明 worker pool 已饱和",
},
}
@router.get("/admin/health-detailed")
def health_detailed():
"""详细的健康检查"""
pool_stats = get_source_pool_stats()
return {
"ok": True,
"version": "v24-concurrent-control",
"source_pool": pool_stats,
"config": {...},
}
部署指南
步骤 1:更新环境变量
在 HuggingFace Space 的 Settings > Repository secrets 中添加:
MAX_CONCURRENT_SOURCES=2 # HF Free tier 保持 2
SOURCE_POOL_WORKERS=4 # 线程池大小
SOURCE_TIMEOUT_SECONDS=15 # 可以适当增加到 30(如果响应慢)
步骤 2:提交代码
cd stock_data_api_service
git add -A
git commit -m "fix: 引入全局源调用池,防止 HF worker pool 饱和
- 添加 GlobalSourcePool 类,使用 Semaphore 控制并发数
- 配置 MAX_CONCURRENT_SOURCES 和 SOURCE_POOL_WORKERS 环境变量
- 添加 /admin/source-pool-stats 和 /admin/health-detailed 监控端点
- 复用线程池,避免每次请求创建新的 ThreadPoolExecutor
- 问题:board_flow/sector_concept_hot 端点因无并发控制导致 worker_pool_saturated"
步骤 3:验证修复
监控端点检查:
curl -H "X-API-Key: YOUR_KEY" https://your-space.hf.space/api/v1/admin/health-detailed并发测试:
# 同时发起 5 个并发请求(HF Free tier 只有 2 个 worker) for i in {1..5}; do curl -H "X-API-Key: YOUR_KEY" \ "https://your-space.hf.space/api/v1/boards/flow?category=industry&limit=10" & done wait查看统计:
curl -H "X-API-Key: YOUR_KEY" https://your-space.hf.space/api/v1/admin/source-pool-stats
调优指南
场景 1:Free tier(1-2 workers)
配置:
MAX_CONCURRENT_SOURCES=2
SOURCE_POOL_WORKERS=2
SOURCE_TIMEOUT_SECONDS=20
特点:
- 最大 2 个并发请求,其他排队等待
- 线程池大小 2,避免线程切换开销
- 超时 20s,给上游足够时间响应
场景 2:Pro tier(10+ workers)
配置:
MAX_CONCURRENT_SOURCES=10
SOURCE_POOL_WORKERS=10
SOURCE_TIMEOUT_SECONDS=15
特点:
- 最大 10 个并发请求
- 线程池大小 10
- 超时 15s
场景 3:上游响应慢(60s+)
配置:
MAX_CONCURRENT_SOURCES=2
SOURCE_POOL_WORKERS=2
SOURCE_TIMEOUT_SECONDS=30
特点:
- 减少并发数,避免 worker 被长时间占用
- 增加超时时间,给上游更多时间
性能指标
修复前
- ❌
worker_pool_saturated错误:多个并发请求时立即饱和 - ❌ 请求排队无限制,所有请求竞争
- ❌ 每次请求创建新的线程池,资源浪费
修复后
- ✅ 最多 N 个并发请求(N = MAX_CONCURRENT_SOURCES)
- ✅ 超出的请求排队等待,而不是立即饱和
- ✅ 线程池复用,资源效率提升
- ✅ 统计指标可用于监控和调优
常见问题
Q1:为什么 boards_flow 和 sector_concept_hot 容易饱和?
A1:这些端点的特征:
- 调用 eastmoney/AKShare 的网络接口
- 响应时间不稳定(15-60s)
- 没有本地缓存(不像
/stocks/{code}/daily有长期缓存) - 并发调用时立即占用所有 worker
Q2:如何判断 worker pool 是否饱和?
A2:调用 /admin/source-pool-stats,查看:
{
"pool_stats": {
"active": 2, // 当前活跃请求数
"queued": 5, // 排队等待的请求数
"max_concurrent": 2, // 最大并发数
"completed": 100, // 已完成请求数
"timeout": 10, // 超时请求数
}
}
如果 active >= max_concurrent,说明 worker pool 已饱和。
Q3:应该增加 MAX_CONCURRENT_SOURCES 还是增加 SOURCE_TIMEOUT_SECONDS?
A3:取决于你的 HF tier 和上游响应时间:
| 情况 | 建议 |
|---|---|
| HF Free tier + 上游响应快(<15s) | 保持 MAX=2, TIMEOUT=15 |
| HF Free tier + 上游响应慢(>30s) | MAX=2, TIMEOUT=30 |
| HF Pro tier + 上游响应快(<15s) | MAX=10, TIMEOUT=15 |
| HF Pro tier + 上游响应慢(>30s) | MAX=5, TIMEOUT=30 |
关键原则:
MAX_CONCURRENT_SOURCES×SOURCE_TIMEOUT_SECONDS≈ HF Worker 数 × 每个 worker 的时间- Free tier:2 workers × 30s = 最多 60s 的并发时间
- Pro tier:10 workers × 15s = 最多 150s 的并发时间
监控和报警
指标
- pool_stats.active - 当前活跃请求数(>max_concurrent 时饱和)
- pool_stats.queued - 排队请求数(>0 说明有排队)
- pool_stats.timeout - 超时请求数(持续增加说明上游响应慢)
- pool_stats.failed - 失败请求数(>0 需要排查)
报警规则
active >= max_concurrent持续 5 分钟 → 考虑增加MAX_CONCURRENT_SOURCEStimeout > 10持续 5 分钟 → 考虑增加SOURCE_TIMEOUT_SECONDSqueued > 0持续 10 分钟 → 考虑优化上游或增加 worker
参考资料
- HuggingFace Spaces Worker Pool: https://huggingface.co/docs/hub/en/spaces-overview
- FastAPI Concurrent Requests: https://fastapi.tiangolo.com/async/
- Python ThreadPoolExecutor: https://docs.python.org/3/library/concurrent.futures.html