first commit

This commit is contained in:
2026-07-14 07:31:13 +08:00
commit 54efc3c361
19 changed files with 973 additions and 0 deletions
+350
View File
@@ -0,0 +1,350 @@
# Universe Selection Engine (USE)
一个基于币安实时数据的加密货币标的池筛选引擎,支持 **多报价币扫描** 与三种主流量化策略:**趋势跟踪**、**均值回归**、**统计套利**。
> 所有阈值、报价币、代理地址、K 线周期等参数均通过 `config/config.yaml` 管理,**零硬编码**。
---
## 一、项目特点
-**多报价币混扫**USDT / USDC / BUSD / BTC 任选组合,一次扫描多个计价市场
-**按报价币独立阈值**:BTC 交叉盘流动性差,可单独放宽成交额门槛
-**自适应稳定币黑名单**:扫 USDC 时不会误杀 USDC 标的
-**上市天数过滤**:剔除刚上线的新币
-**K 线拉取带重试**:网络抖动不会丢失标的
-**异步并发 + 限频**:使用 `asyncio.Semaphore` 防 429 封禁
-**三种策略并行**:趋势 / 均值回归 / 统计套利
---
## 二、快速开始
### 2.1 环境要求
- Python >= 3.10
- Windows / Linux / macOS 全平台
- **代理工具**(国内访问币安 API 必需)
### 2.2 安装
```bash
cd crypto_use
pip install -r requirements.txt
```
依赖:`ccxt`, `pandas`, `numpy`, `statsmodels`, `pyyaml`, `scipy`
### 2.3 代理配置
> ⚠️ **国内直连币安 API 会被阻断,必须配置代理**。
编辑 `config/config.yaml`,修改 `exchange.proxy`
```yaml
exchange:
proxy: "http://127.0.0.1:7890" # 改成你的代理地址
```
如不需要代理,填 `null`。系统会判断:有代理则注入自定义 aiohttp session`ThreadedResolver` 绕过 aiodns DNS 失败问题),无代理则走 ccxt 默认行为。
### 2.4 运行
```bash
python main.py
```
输出结果保存到 `output/universe_results.json`。**任何目录下运行均可**,项目根路径自动锁定 `main.py` 所在位置。
---
## 三、输出说明
```json
{
"timestamp": "2026-07-14 07:24:52",
"timeframe": "1h",
"quote_currencies": ["USDT", "USDC"],
"trend_universe": ["BTC/USDT", "ETH/USDT", "BNB/USDT", ...],
"mr_universe": [],
"starb_pairs": []
}
```
| 字段 | 类型 | 含义 |
|------|------|------|
| `timestamp` | string | 运行时间 |
| `timeframe` | string | 使用的 K 线周期 |
| `quote_currencies` | list | 本次扫描的报价币列表 |
| `trend_universe` | list | 趋势跟踪策略标的池 |
| `mr_universe` | list | 均值回归策略标的池 |
| `starb_pairs` | list | 统计套利配对池,格式 `[["ETH/USDT", "SOL/USDT"], ...]` |
---
## 四、配置参数详解
所有参数集中在 `config/config.yaml`,共 8 个顶级字段。
### 4.1 `exchange` 交易所设置
```yaml
exchange:
name: binance # 交易所名称(ccxt 支持列表)
market: spot # spot=现货, future=U本位合约
proxy: "http://127.0.0.1:7890" # 代理地址,留空或 null 直连
```
### 4.2 `data` 数据设置
```yaml
data:
timeframe: 1h # K线周期:1m, 5m, 15m, 1h, 4h, 1d
lookback_candles: 500 # 拉取历史K线根数
max_concurrency: 15 # K 线并发拉取上限(防 429
```
| 参数 | 调小 | 调大 |
|------|------|------|
| `timeframe` | 1m/5m → 捕捉短期波动,适合日内 | 1d → 过滤噪声,适合长线 |
| `lookback_candles` | 100 → 计算快但统计不显著 | 1000+ → 更稳健但增加 API 压力 |
| `max_concurrency` | 5~8 → 更安全 | 20 → 更快但易触发 429 |
### 4.3 `base_filter` 基础流动性过滤
```yaml
base_filter:
# 报价币列表 —— 引擎会扫描所有这些计价市场,写几个就扫几个
# 支持混配:['USDT', 'USDC'] 会同时扫 USDT 对和 USDC 对
quote_currencies:
- USDT
- USDC
# 单一计价模式(兼容老配置),如果 quote_currencies 未设置则使用此值
quote_currency: USDT
min_quote_volume_24h: 20000000 # 默认成交额门槛(按对应报价币计算)
# 按报价币覆盖成交额门槛 —— 低流动性报价币可单独放宽
# 未列出的报价币沿用 min_quote_volume_24h
volume_thresholds:
BTC: 50 # BTC 交叉盘流动性远低于 USDT 对
BUSD: 5000000 # BUSD 对同样冷门
USDC: 5000000 # USDC 对略低于 USDT
max_spread_pct: 0.05 # 最大买卖价差 (%)
min_listing_days: 90 # 最小上市天数(0 = 不限制)
# 报价币对应的稳定币黑名单(不同计价场景需要不同定义)
stablecoin_bases:
USDT: [USDC, FDUSD, DAI, TUSD, BUSD, EUR, TRY]
USDC: [USDT, FDUSD, DAI, TUSD, BUSD, EUR, TRY]
BUSD: [USDT, USDC, FDUSD, DAI, TUSD, EUR, TRY]
BTC: [USDT, USDC, FDUSD, DAI, TUSD, BUSD, EUR, TRY, WBTC]
USD: [USDT, USDC, FDUSD, DAI, TUSD, BUSD, EUR, TRY]
exclude_stablecoins: true # 是否启用稳定币过滤
```
**核心设计**:扫 USDT 时,`USDC/FDUSD/DAI...` 这些 base 会被当作稳定币排除;扫 USDC 时不会误杀 USDC 标的;扫 BTC 时还会排除 `WBTC`(防止 WBTC/BTC 这种没意义的循环对)。
### 4.4 `trend_filter` 趋势跟踪过滤
```yaml
trend_filter:
min_annualized_vol: 0.40 # 最小年化波动率(40%
min_hurst: 0.55 # 最小 Hurst 指数(>0.5 表示趋势持续)
```
### 4.5 `mean_reversion_filter` 均值回归过滤
```yaml
mean_reversion_filter:
max_adf_pvalue: 0.05 # ADF 检验最大 p-value<0.05 表示平稳)
max_hurst: 0.45 # 最大 Hurst 指数(<0.5 表示反持续)
max_kurtosis: 10.0 # 最大峰度(剔除极端肥尾/黑天鹅币)
```
### 4.6 `starb_filter` 统计套利过滤
```yaml
starb_filter:
max_coint_pvalue: 0.05 # 协整检验最大 p-value
# 板块分类 —— 键为板块名,值为标的列表
# 标的格式必须为 BASE/QUOTE,必须在 quote_currencies 里有对应 quote
sectors:
L1: ["ETH/USDT", "SOL/USDT", "AVAX/USDT", "ADA/USDT", "DOT/USDT", "NEAR/USDT"]
DeFi: ["UNI/USDT", "AAVE/USDT", "MKR/USDT", "SNX/USDT", "COMP/USDT", "LDO/USDT"]
Meme: ["DOGE/USDT", "SHIB/USDT", "PEPE/USDT", "WIF/USDT", "FLOKI/USDT", "BONK/USDT"]
AI: ["RNDR/USDT", "FET/USDT", "AGIX/USDT", "TAO/USDT", "ARKM/USDT"]
# 交叉盘示例:
# BTC_Pairs: ["ETH/BTC", "SOL/BTC", "AVAX/BTC"]
```
> ⚠️ 板块分类需要**根据市场叙事动态更新**,不是一成不变的。
### 4.7 `fetcher` K 线拉取设置
```yaml
fetcher:
max_retries: 3 # 单个标的最大重试次数
retry_delay: 1.0 # 重试基础间隔(秒),实际为 retry_delay × 第几次
```
### 4.8 `output` 输出设置
```yaml
output:
save_path: "./output/universe_results.json" # 相对路径相对项目根,也可写绝对路径
```
---
## 五、参数调优指南
### 5.1 标的池太少?→ 放宽阈值
| 场景 | 操作 |
|------|------|
| 没有任何标的过基础过滤 | 降低 `min_quote_volume_24h` 或在 `volume_thresholds` 给冷门 quote 设小门槛 |
| 趋势池为空 | 降低 `min_hurst`(如 0.55 → 0.52)或降低 `min_annualized_vol`0.40 → 0.25 |
| 均值回归池为空 | 降低 `max_kurtosis`10.0 → 15.0)或提高 `max_hurst`0.45 → 0.50 |
| 套利配对为空 | 提高 `max_coint_pvalue`(0.05 → 0.10)或扩大板块覆盖 |
### 5.2 标的池太多?→ 收紧阈值
| 场景 | 操作 |
|------|------|
| 趋势池全是"死币" | 提高 `min_hurst` 到 0.60+ |
| 均值回归池出现单边下跌币 | 降低 `max_kurtosis` 到 5.0 以下 |
| 套利配对太多 | 降低 `max_coint_pvalue` 到 0.01 |
### 5.3 多报价币策略
| 需求 | 配置示例 |
|------|----------|
| 同时扫 USDT + USDC 主流市场 | `quote_currencies: [USDT, USDC]` |
| 只扫 BTC 交叉盘(量化做 BTC 对统计套利) | `quote_currencies: [BTC]` + `volume_thresholds.BTC: 50` |
| 全量扫描(不推荐,慢) | `quote_currencies: [USDT, USDC, BUSD, BTC, USD]` + 对应 `volume_thresholds` |
---
## 六、与实盘对接
### 6.1 读取结果
```python
import json
from pathlib import Path
ROOT = Path(__file__).resolve().parent
with open(ROOT / "output" / "universe_results.json", "r", encoding="utf-8") as f:
universe = json.load(f)
trend_symbols = universe["trend_universe"] # 趋势策略使用
mr_symbols = universe["mr_universe"] # 均值回归策略使用
starb_pairs = universe["starb_pairs"] # 统计套利策略使用
quote_currencies = universe["quote_currencies"] # 本次扫描的报价币
```
### 6.2 定时调度
**建议执行频率**:每天凌晨 00:00 UTC 运行一次(流动性/统计量在日内变化不大)。
**Linux cron 示例**
```bash
0 0 * * * cd /path/to/crypto_use && /usr/bin/python main.py >> /var/log/use_cron.log 2>&1
```
**Windows 任务计划程序**:新建基本任务 → 触发器"每天 00:00" → 操作"启动程序" `python main.py`,起始于 `crypto_use` 目录。
---
## 七、项目结构
```
crypto_use/
├── config/
│ └── config.yaml # 全局配置(零硬编码的唯一来源)
├── universe_selector/
│ ├── __init__.py
│ ├── fetcher.py # 异步数据拉取(ccxt + 限频 + 重试)
│ ├── filters.py # 三种策略过滤器(Hurst/ADF/Coint
│ ├── engine.py # 流水线控制器
│ └── utils.py # 统计算法辅助函数
├── output/
│ └── universe_results.json # 筛选结果(每次运行覆盖)
├── main.py # 入口脚本
├── use_engine.log # 运行日志(每次追加)
├── requirements.txt # 依赖包
└── README.md # 本文件
```
---
## 八、常见问题
**Q: 运行报 DNS 错误?**
A: 确认代理 7890 端口已启动,且 `config.yaml``exchange.proxy` 已正确填写。本项目在有代理时会自动注入 `ThreadedResolver`,绕过 aiodns 在代理环境下的 DNS 失败问题。
**Q: 报 429 限频?**
A: 降低 `data.max_concurrency` 到 8~10;或降低 `fetcher.max_retries` 减少重试风暴。
**Q: 某些标的 K 线拉不到?**
A: 项目对每个标的自动重试 3 次(指数退避),最终失败会在日志中输出 ERROR 但不中断流程。如大量失败,检查网络或降低并发。
**Q: 均值回归池一直是空的?**
A: 这是正常现象。`max_hurst < 0.45 + ADF p < 0.05 + kurtosis < 10` 三条件同时满足的标的在币安全市场现货里极少。降低 `max_kurtosis` 到 15+ 或提高 `max_hurst` 到 0.50 试试。
**Q: 统计套利配对一直是 0**
A: 在 `starb_filter.sectors` 里扩大板块覆盖,或提高 `max_coint_pvalue` 到 0.10。
**Q: 想扫描 BTC 交叉盘但出不了标的?**
A: 在 `base_filter.volume_thresholds.BTC` 设置小门槛(如 50 BTC)。BTC 交叉盘流动性远低于 USDT 对,必须单独配置门槛。
**Q: 想看某币是因什么原因被过滤掉的?**
A: 当前版本未提供 verbose 调试模式。如需要可临时将 `fetcher.py``get_liquid_symbols` 的循环里加 `print(f"[skip] {symbol}: 原因")`
**Q: 支持合约市场吗?**
A: 支持。把 `exchange.market` 改成 `future` 即可。但注意币安合约 API 的 ticker 字段与现货略有差异,可能需要根据实际情况微调。
**Q: 在哪个目录运行 main.py**
A: 任意目录均可。项目根路径通过 `Path(__file__).resolve().parent` 自动锁定,配置/输出/日志路径与 cwd 无关。
---
## 九、扩展开发
### 9.1 新增一种策略过滤器
`universe_selector/filters.py``StrategyFilters` 类里添加新方法,例如:
```python
def filter_breakout(self, data_dict: dict[str, pd.DataFrame]) -> list[str]:
"""突破策略:检测近期价格是否突破 N 日高点"""
universe = []
for sym, df in data_dict.items():
close = df['close'].values
high_20d = close[:-1][-20*24:].max() # 前 20 天最高
if close[-1] > high_20d:
universe.append(sym)
return universe
```
然后在 `engine.py``run_pipeline` 中调用:
```python
results['breakout_universe'] = self.filters.filter_breakout(data_dict)
```
### 9.2 接入实盘自动下单
引擎输出 `output/universe_results.json` 后,下游策略可订阅文件变化(`watchdog` 库)或定时轮询,将新标的池与持仓比对,自动调整下单白名单。
---
## 十、版本
| 版本 | 日期 | 变更 |
|------|------|------|
| 1.0 | 2026-07 | 初始版本(USDT only |
| 1.1 | 2026-07 | 多报价币 / 按 quote 独立门槛 / 上市天数 / K 线重试 / 路径稳健化 |
---
## 十一、许可
仅供学习和研究使用。加密货币交易有风险,请自行评估。
+64
View File
@@ -0,0 +1,64 @@
exchange:
name: binance # 交易所名称 (ccxt 支持的名称)
market: spot # spot (现货) 或 future (U本位合约)
proxy: "http://127.0.0.1:7890" # 代理地址(国内必需),如不需要填 null
data:
timeframe: 1h # K线周期 (1m, 5m, 15m, 1h, 4h, 1d)
lookback_candles: 500 # 拉取的历史K线数量
max_concurrency: 15 # 并发拉取上限(防 429
base_filter:
# 报价币列表 —— 引擎会拉取所有这些报价对,写几个就扫几个
# 支持混配:['USDT', 'USDC'] 会同时扫两种计价
quote_currencies:
- USDT
- USDC
# 单一计价模式(兼容老配置),如果 quote_currencies 未设置则使用此值
quote_currency: USDT
min_quote_volume_24h: 20000000 # 24小时最小成交额默认值 (按对应报价币计算)
# 按报价币覆盖成交额门槛 —— 低流动性报价币(如 BTC/BUSD 交叉盘)可单独放宽
# 未列出的报价币沿用 min_quote_volume_24h
volume_thresholds:
BTC: 50 # BTC 交叉盘流动性远低于 USDT 对
BUSD: 5000000 # BUSD 对同样冷门
USDC: 5000000 # USDC 对略低于 USDT
max_spread_pct: 0.05 # 最大买卖价差 (%)
min_listing_days: 90 # 最小上市天数(0 = 不限制)
# 报价币对应的稳定币黑名单(这些 base 资产会被排除)
# 不同报价币场景下需要不同的稳定币定义
stablecoin_bases:
USDT: [USDC, FDUSD, DAI, TUSD, BUSD, EUR, TRY]
USDC: [USDT, FDUSD, DAI, TUSD, BUSD, EUR, TRY]
BUSD: [USDT, USDC, FDUSD, DAI, TUSD, EUR, TRY]
BTC: [USDT, USDC, FDUSD, DAI, TUSD, BUSD, EUR, TRY, WBTC]
USD: [USDT, USDC, FDUSD, DAI, TUSD, BUSD, EUR, TRY]
exclude_stablecoins: true # 是否排除稳定币 base
trend_filter:
min_annualized_vol: 0.40 # 最小年化波动率 (40%)
min_hurst: 0.55 # 最小 Hurst 指数 (趋势持续性)
mean_reversion_filter:
max_adf_pvalue: 0.05 # ADF 检验最大 p-value (平稳性)
max_hurst: 0.45 # 最大 Hurst 指数 (反持续性)
max_kurtosis: 10.0 # 最大峰度 (剔除极端肥尾/黑天鹅币)
starb_filter:
max_coint_pvalue: 0.05 # 协整检验最大 p-value
# 统计套利板块分类 —— 键为板块名,值为标的列表
# 标的格式必须为 BASE/QUOTE,必须在 quote_currencies 里有对应 quote
sectors:
L1: ["ETH/USDT", "SOL/USDT", "AVAX/USDT", "ADA/USDT", "DOT/USDT", "NEAR/USDT"]
DeFi: ["UNI/USDT", "AAVE/USDT", "MKR/USDT", "SNX/USDT", "COMP/USDT", "LDO/USDT"]
Meme: ["DOGE/USDT", "SHIB/USDT", "PEPE/USDT", "WIF/USDT", "FLOKI/USDT", "BONK/USDT"]
AI: ["RNDR/USDT", "FET/USDT", "AGIX/USDT", "TAO/USDT", "ARKM/USDT"]
# 计价币示例(注释掉或删除不需要的板块):
# BTC_Pairs: ["ETH/BTC", "SOL/BTC", "AVAX/BTC"]
fetcher:
max_retries: 3 # 单个标的最大重试次数
retry_delay: 1.0 # 重试间隔(秒)
output:
save_path: "./output/universe_results.json"
+61
View File
@@ -0,0 +1,61 @@
import asyncio
import logging
import sys
from pathlib import Path
import yaml
from universe_selector import UniverseEngine
# 定位项目根目录(main.py 所在目录),保证配置/输出路径与执行位置无关
PROJECT_ROOT = Path(__file__).resolve().parent
LOG_FILE = PROJECT_ROOT / "use_engine.log"
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s - %(name)s - %(levelname)s - %(message)s',
handlers=[
logging.StreamHandler(sys.stdout),
logging.FileHandler(LOG_FILE, encoding='utf-8'),
],
)
def load_config(path: Path | None = None) -> dict:
cfg_path = path or PROJECT_ROOT / "config" / "config.yaml"
with open(cfg_path, 'r', encoding='utf-8') as f:
return yaml.safe_load(f)
def resolve_output_path(raw: str) -> Path:
"""支持相对路径(相对项目根)和绝对路径。"""
p = Path(raw)
if p.is_absolute():
return p
return PROJECT_ROOT / p
async def main():
logging.info("=== 启动 Universe Selection Engine ===")
config = load_config()
# 输出路径统一解析为绝对路径,避免 cwd 影响
if 'output' in config and 'save_path' in config['output']:
config['output']['save_path'] = str(resolve_output_path(config['output']['save_path']))
engine = UniverseEngine(config)
results = await engine.run_pipeline()
if results:
print("\n" + "=" * 50)
print(f"运行时间: {results['timestamp']}")
print(f"趋势标的数: {len(results['trend_universe'])}")
print(f"均值回归标的数: {len(results['mr_universe'])}")
print(f"套利配对数: {len(results['starb_pairs'])}")
print("=" * 50 + "\n")
if __name__ == "__main__":
if sys.platform == 'win32':
asyncio.set_event_loop_policy(asyncio.WindowsSelectorEventLoopPolicy())
asyncio.run(main())
+13
View File
@@ -0,0 +1,13 @@
{
"timestamp": "2026-07-14 07:19:36",
"timeframe": "1h",
"quote_currencies": [
"BTC"
],
"trend_universe": [
"TCT/BTC",
"JASMY/BTC"
],
"mr_universe": [],
"starb_pairs": []
}
+24
View File
@@ -0,0 +1,24 @@
{
"timestamp": "2026-07-14 07:25:26",
"timeframe": "1h",
"quote_currencies": [
"BTC",
"USDT"
],
"trend_universe": [
"BTC/USDT",
"ETH/USDT",
"BNB/USDT",
"XRP/USDT",
"ZEC/USDT",
"DOGE/USDT",
"SOL/USDT",
"DEXE/USDT",
"WLD/USDT",
"ALLO/USDT",
"MUB/USDT",
"SPCXB/USDT"
],
"mr_universe": [],
"starb_pairs": []
}
+29
View File
@@ -0,0 +1,29 @@
{
"timestamp": "2026-07-14 07:24:52",
"timeframe": "1h",
"quote_currencies": [
"USDT",
"USDC"
],
"trend_universe": [
"BTC/USDT",
"ETH/USDT",
"BNB/USDT",
"XRP/USDT",
"BNB/USDC",
"BTC/USDC",
"ETH/USDC",
"XRP/USDC",
"ZEC/USDT",
"ZEC/USDC",
"DOGE/USDT",
"SOL/USDT",
"SOL/USDC",
"WLD/USDT",
"ALLO/USDT",
"MUB/USDT",
"SPCXB/USDT"
],
"mr_universe": [],
"starb_pairs": []
}
+13
View File
@@ -0,0 +1,13 @@
{
"timestamp": "2026-07-14 07:19:50",
"timeframe": "1h",
"quote_currencies": [
"USDC"
],
"trend_universe": [
"BTC/USDC",
"ETH/USDC"
],
"mr_universe": [],
"starb_pairs": []
}
+6
View File
@@ -0,0 +1,6 @@
ccxt>=4.0.0
pandas>=2.0.0
numpy>=1.24.0
statsmodels>=0.14.0
pyyaml>=6.0
scipy>=1.10.0
+1
View File
@@ -0,0 +1 @@
from .engine import UniverseEngine
Binary file not shown.
Binary file not shown.
+59
View File
@@ -0,0 +1,59 @@
import json
import logging
from datetime import datetime
from pathlib import Path
from .fetcher import BinanceFetcher
from .filters import StrategyFilters
logger = logging.getLogger(__name__)
class UniverseEngine:
def __init__(self, config: dict):
self.config = config
self.fetcher = BinanceFetcher(config)
self.filters = StrategyFilters(config)
async def run_pipeline(self):
try:
# 1. 基础池
liquid_symbols = await self.fetcher.get_liquid_symbols()
if not liquid_symbols:
logger.error("基础池为空,请检查网络/代理或放宽 base_filter 阈值。")
return {}
# 2. 获取数据
data_dict = await self.fetcher.fetch_klines_batch(liquid_symbols)
if not data_dict:
logger.error("K 线拉取全部失败,请检查网络或降低并发。")
return {}
# 3. 执行策略过滤
results = {
'timestamp': datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
'timeframe': self.config['data']['timeframe'],
'quote_currencies': self.fetcher.quote_currencies,
'trend_universe': self.filters.filter_trend(data_dict),
'mr_universe': self.filters.filter_mean_reversion(data_dict),
'starb_pairs': self.filters.filter_starb(data_dict),
}
# 4. 保存结果
self._save_results(results)
return results
finally:
await self.fetcher.close()
def _save_results(self, results: dict):
save_path = Path(
self.config.get('output', {}).get('save_path', './output/universe_results.json')
)
save_path.parent.mkdir(parents=True, exist_ok=True)
# tuple 转 list 方便 JSON 序列化
results['starb_pairs'] = [list(p) for p in results['starb_pairs']]
with open(save_path, 'w', encoding='utf-8') as f:
json.dump(results, f, indent=4, ensure_ascii=False)
logger.info(f"标的池结果已保存至: {save_path}")
+192
View File
@@ -0,0 +1,192 @@
import asyncio
import time
from datetime import datetime, timezone
from typing import Iterable
import aiohttp
import aiohttp.resolver
import ccxt.async_support as ccxt
import pandas as pd
import logging
logger = logging.getLogger(__name__)
class BinanceFetcher:
def __init__(self, config: dict):
self.config = config
self.quote_currencies = self._resolve_quote_currencies(config)
self.exchange = self._build_exchange(config)
self._markets_meta: dict[str, dict] = {} # symbol -> market dict(含 onchainDate 等)
@staticmethod
def _resolve_quote_currencies(config: dict) -> list[str]:
"""优先取 quote_currencies 列表,回退到 quote_currency 单值。"""
base = config.get('base_filter', {})
quotes = base.get('quote_currencies')
if quotes:
return [q.upper() for q in quotes]
single = base.get('quote_currency')
return [single.upper()] if single else ['USDT']
@staticmethod
def _build_exchange(config: dict):
exchange_class = getattr(ccxt, config['exchange']['name'])
kwargs = {
'enableRateLimit': True,
'options': {'defaultType': config['exchange']['market']},
}
proxy = config['exchange'].get('proxy')
if proxy:
kwargs['httpsProxy'] = proxy
exchange = exchange_class(kwargs)
# 若需代理,注入自定义 aiohttp sessionThreadedResolver 绕过 aiodns DNS 失败)
if proxy:
resolver = aiohttp.resolver.ThreadedResolver()
connector = aiohttp.TCPConnector(resolver=resolver, enable_cleanup_closed=True)
exchange._custom_session = aiohttp.ClientSession(connector=connector)
exchange._own_session = False # 阻止 ccxt 在 open() 中重建 session
exchange.session = exchange._custom_session
return exchange
async def close(self):
await self.exchange.close()
custom = getattr(self.exchange, '_custom_session', None)
if custom and not custom.closed:
await custom.close()
# ------------------------------------------------------------------
# Phase 1: 流动性 + 上市天数过滤
# ------------------------------------------------------------------
async def get_liquid_symbols(self) -> list[str]:
"""按所有配置的报价币拉取并过滤。"""
logger.info("正在加载交易所市场信息...")
await self.exchange.load_markets()
self._markets_meta = dict(self.exchange.markets)
tickers = await self.exchange.fetch_tickers()
base_conf = self.config['base_filter']
min_listing_days = int(base_conf.get('min_listing_days', 0))
exclude_stablecoins = bool(base_conf.get('exclude_stablecoins', True))
stablecoin_map = base_conf.get('stablecoin_bases', {})
# 通用稳定币列表(用于报价币不在 stablecoin_map 时的回退)
default_stable = {'USDT', 'USDC', 'FDUSD', 'DAI', 'TUSD', 'BUSD', 'EUR', 'TRY', 'WBTC'}
# 按报价币解析成交额门槛 —— 未配置时沿用默认值
default_volume = float(base_conf['min_quote_volume_24h'])
volume_thresholds: dict[str, float] = {
q.upper(): float(v) for q, v in (base_conf.get('volume_thresholds') or {}).items()
}
for q in self.quote_currencies:
t = volume_thresholds.get(q, default_volume)
logger.info(f" 报价币 {q}: 成交额门槛 = {t:,.0f}")
valid: list[str] = []
for symbol, ticker in tickers.items():
market = self._markets_meta.get(symbol, {})
quote = market.get('quote')
base = market.get('base')
if quote not in self.quote_currencies:
continue
# 1. 成交额(按 quote 独立阈值)
qv = ticker.get('quoteVolume') or 0
threshold = volume_thresholds.get(quote.upper(), default_volume)
if qv < threshold:
continue
# 2. 价差
bid, ask = ticker.get('bid'), ticker.get('ask')
if bid and ask and bid > 0:
spread = (ask - bid) / bid * 100
if spread > base_conf['max_spread_pct']:
continue
# 3. 上市天数(币安 market.active 通常为 TrueonchainDate/listingDate 不一定存在)
if min_listing_days > 0 and not self._passes_listing_filter(market, min_listing_days):
continue
# 4. 稳定币过滤
if exclude_stablecoins:
banned = set(stablecoin_map.get(quote, default_stable))
if base in banned:
continue
valid.append(symbol)
logger.info(
f"基础流动性过滤完成:报价币={self.quote_currencies},剩余 {len(valid)} 个标的。"
)
return valid
@staticmethod
def _passes_listing_filter(market: dict, min_days: int) -> bool:
"""检查市场是否满足上市天数要求。"""
# 币安 markets 通常不直接返回 listing date;这里使用 active 字段做软校验
if not market.get('active', True):
return False
# 部分交易所/合约类型会返回 info.createdAt / info.listDate
info = market.get('info', {}) if isinstance(market, dict) else {}
ts_ms = info.get('onboardDate') or info.get('listDate') or info.get('createdAt')
if ts_ms is None:
# 拿不到日期,按通过处理(避免误杀)
return True
try:
ts_ms = int(ts_ms)
except (TypeError, ValueError):
return True
onboard = datetime.fromtimestamp(ts_ms / 1000, tz=timezone.utc)
elapsed = (datetime.now(tz=timezone.utc) - onboard).days
return elapsed >= min_days
# ------------------------------------------------------------------
# Phase 2: 批量 K 线拉取(带重试)
# ------------------------------------------------------------------
async def fetch_klines_batch(self, symbols: list[str]) -> dict[str, pd.DataFrame]:
timeframe = self.config['data']['timeframe']
limit = self.config['data']['lookback_candles']
concurrency = int(self.config['data'].get('max_concurrency', 15))
fetcher_conf = self.config.get('fetcher', {})
max_retries = int(fetcher_conf.get('max_retries', 3))
retry_delay = float(fetcher_conf.get('retry_delay', 1.0))
sem = asyncio.Semaphore(concurrency)
async def fetch_single(sym):
async with sem:
return sym, await self._fetch_with_retry(sym, timeframe, limit, max_retries, retry_delay)
logger.info(f"开始异步拉取 {len(symbols)} 个标的的 {timeframe} K线...")
results = await asyncio.gather(*(fetch_single(s) for s in symbols))
return {sym: df for sym, df in results if df is not None}
async def _fetch_with_retry(
self, symbol: str, timeframe: str, limit: int, max_retries: int, retry_delay: float
) -> pd.DataFrame | None:
last_err = None
for attempt in range(1, max_retries + 1):
try:
ohlcv = await self.exchange.fetch_ohlcv(symbol, timeframe=timeframe, limit=limit)
if not ohlcv:
raise ValueError("empty ohlcv response")
df = pd.DataFrame(
ohlcv, columns=['timestamp', 'open', 'high', 'low', 'close', 'volume']
)
df['timestamp'] = pd.to_datetime(df['timestamp'], unit='ms')
df.set_index('timestamp', inplace=True)
return df
except Exception as e:
last_err = e
logger.warning(f"[{symbol}] 第 {attempt}/{max_retries} 次拉取失败: {e}")
if attempt < max_retries:
await asyncio.sleep(retry_delay * attempt)
logger.error(f"[{symbol}] 拉取失败,已重试 {max_retries} 次,最终放弃。last_err={last_err}")
return None
+71
View File
@@ -0,0 +1,71 @@
import numpy as np
import pandas as pd
from . import utils
import logging
logger = logging.getLogger(__name__)
class StrategyFilters:
def __init__(self, config: dict):
self.config = config
# 根据 timeframe 计算一年有多少个 period (用于波动率年化)
tf = config['data']['timeframe']
self.periods_per_year = {
'1m': 525600, '5m': 105120, '15m': 35040,
'1h': 8760, '4h': 2190, '1d': 365
}.get(tf, 8760)
def filter_trend(self, data_dict: dict[str, pd.DataFrame]) -> list[str]:
conf = self.config['trend_filter']
universe = []
for sym, df in data_dict.items():
returns = np.diff(np.log(df['close'].values))
vol = utils.calculate_annualized_vol(returns, self.periods_per_year)
hurst = utils.calculate_hurst(df['close'].values)
if vol >= conf['min_annualized_vol'] and hurst >= conf['min_hurst']:
universe.append(sym)
logger.info(f"趋势跟踪池筛选完成: {len(universe)} 个标的")
return universe
def filter_mean_reversion(self, data_dict: dict[str, pd.DataFrame]) -> list[str]:
conf = self.config['mean_reversion_filter']
universe = []
for sym, df in data_dict.items():
close = df['close'].values
returns = np.diff(np.log(close))
hurst = utils.calculate_hurst(close)
# 峰度计算 (剔除肥尾)
kurtosis = pd.Series(returns).kurtosis()
is_stationary = utils.check_adf_stationarity(close, conf['max_adf_pvalue'])
if (is_stationary and
hurst <= conf['max_hurst'] and
kurtosis <= conf['max_kurtosis']):
universe.append(sym)
logger.info(f"均值回归池筛选完成: {len(universe)} 个标的")
return universe
def filter_starb(self, data_dict: dict[str, pd.DataFrame]) -> list[tuple[str, str]]:
conf = self.config['starb_filter']
pairs = []
sectors = conf.get('sectors', {})
for sector, symbols in sectors.items():
# 过滤出有数据的标的
valid_syms = [s for s in symbols if s in data_dict]
# 两两组合计算协整
for i in range(len(valid_syms)):
for j in range(i+1, len(valid_syms)):
sym1, sym2 = valid_syms[i], valid_syms[j]
ts1 = data_dict[sym1]['close'].values
ts2 = data_dict[sym2]['close'].values
if utils.check_cointegration(ts1, ts2, conf['max_coint_pvalue']):
pairs.append((sym1, sym2))
logger.info(f"统计套利配对筛选完成: {len(pairs)}")
return pairs
+48
View File
@@ -0,0 +1,48 @@
import numpy as np
import statsmodels.api as sm
from statsmodels.tsa.stattools import adfuller, coint
def calculate_hurst(ts: np.ndarray) -> float:
"""使用 R/S 分析计算 Hurst 指数"""
try:
lags = range(2, min(100, len(ts) // 2))
tau = [np.std(np.subtract(ts[lag:], ts[:-lag])) for lag in lags]
# 避免 log(0)
tau = [t for t in tau if t > 0]
lags = list(lags)[:len(tau)]
if len(lags) < 2: return 0.5
poly = np.polyfit(np.log(lags), np.log(tau), 1)
return poly[0] * 2.0
except Exception:
return 0.5 # 计算失败返回随机游走假设
def calculate_annualized_vol(returns: np.ndarray, periods_per_year: int) -> float:
"""计算年化波动率"""
if len(returns) < 2: return 0.0
return np.std(returns) * np.sqrt(periods_per_year)
def check_adf_stationarity(ts: np.ndarray, max_pvalue: float) -> bool:
"""ADF 平稳性检验"""
try:
# 剔除 NaN
ts = ts[~np.isnan(ts)]
if len(ts) < 20: return False
result = adfuller(ts, autolag='AIC')
return result[1] <= max_pvalue
except Exception:
return False
def check_cointegration(ts1: np.ndarray, ts2: np.ndarray, max_pvalue: float) -> bool:
"""Engle-Granger 协整检验"""
try:
# 对齐长度并剔除 NaN
min_len = min(len(ts1), len(ts2))
ts1, ts2 = ts1[-min_len:], ts2[-min_len:]
mask = ~np.isnan(ts1) & ~np.isnan(ts2)
ts1, ts2 = ts1[mask], ts2[mask]
if len(ts1) < 50: return False
score, pvalue, _ = coint(ts1, ts2)
return pvalue <= max_pvalue
except Exception:
return False
+42
View File
@@ -0,0 +1,42 @@
2026-07-14 06:58:42,217 - root - INFO - === 启动 Universe Selection Engine ===
2026-07-14 06:58:42,226 - universe_selector.fetcher - INFO - 正在拉取全市场 24hr Ticker 和 BookTicker...
2026-07-14 06:58:42,831 - ccxt.base.exchange - WARNING - binance requires to release all resources with an explicit call to the .close() coroutine. If you are using the exchange instance with async coroutines, add `await exchange.close()` to your code into a place when you're done with the exchange and don't need the exchange instance anymore (at the end of your async coroutine).
2026-07-14 06:58:42,832 - asyncio - ERROR - Unclosed client session
client_session: <aiohttp.client.ClientSession object at 0x0000019ADB2E3920>
2026-07-14 06:59:23,317 - root - INFO - === 启动 Universe Selection Engine ===
2026-07-14 06:59:23,324 - universe_selector.fetcher - INFO - 正在拉取全市场 24hr Ticker 和 BookTicker...
2026-07-14 06:59:23,888 - ccxt.base.exchange - WARNING - binance requires to release all resources with an explicit call to the .close() coroutine. If you are using the exchange instance with async coroutines, add `await exchange.close()` to your code into a place when you're done with the exchange and don't need the exchange instance anymore (at the end of your async coroutine).
2026-07-14 06:59:23,889 - asyncio - ERROR - Unclosed client session
client_session: <aiohttp.client.ClientSession object at 0x0000029348E53260>
2026-07-14 07:00:39,766 - root - INFO - === 启动 Universe Selection Engine ===
2026-07-14 07:00:39,785 - universe_selector.fetcher - INFO - 正在拉取全市场 24hr Ticker 和 BookTicker...
2026-07-14 07:00:40,337 - asyncio - ERROR - Unclosed client session
client_session: <aiohttp.client.ClientSession object at 0x0000024E9DA8AEA0>
2026-07-14 07:03:41,494 - root - INFO - === 启动 Universe Selection Engine ===
2026-07-14 07:03:41,501 - universe_selector.fetcher - INFO - 正在拉取全市场 24hr Ticker 和 BookTicker...
2026-07-14 07:07:06,043 - root - INFO - === 启动 Universe Selection Engine ===
2026-07-14 07:07:06,050 - universe_selector.fetcher - INFO - 正在拉取全市场 24hr Ticker 和 BookTicker...
2026-07-14 07:07:08,847 - universe_selector.fetcher - INFO - 基础流动性过滤完成,剩余 16 个标的。
2026-07-14 07:07:08,864 - universe_selector.fetcher - INFO - 开始异步拉取 16 个标的的 1h K线...
2026-07-14 07:07:10,544 - universe_selector.filters - INFO - 趋势跟踪池筛选完成: 13 个标的
2026-07-14 07:07:10,976 - universe_selector.filters - INFO - 均值回归池筛选完成: 0 个标的
2026-07-14 07:07:11,006 - universe_selector.filters - INFO - 统计套利配对筛选完成: 0 对
2026-07-14 07:07:11,008 - universe_selector.engine - INFO - 标的池结果已保存至: ./output/universe_results.json
2026-07-14 07:19:10,946 - root - INFO - === 启动 Universe Selection Engine ===
2026-07-14 07:19:10,958 - universe_selector.fetcher - INFO - 正在加载交易所市场信息...
2026-07-14 07:19:13,460 - universe_selector.fetcher - INFO - 基础流动性过滤完成:报价币=['USDT', 'USDC'],剩余 18 个标的。
2026-07-14 07:19:13,465 - universe_selector.fetcher - INFO - 开始异步拉取 18 个标的的 1h K线...
2026-07-14 07:19:15,119 - universe_selector.filters - INFO - 趋势跟踪池筛选完成: 14 个标的
2026-07-14 07:19:15,459 - universe_selector.filters - INFO - 均值回归池筛选完成: 0 个标的
2026-07-14 07:19:15,479 - universe_selector.filters - INFO - 统计套利配对筛选完成: 0 对
2026-07-14 07:19:15,481 - universe_selector.engine - INFO - 标的池结果已保存至: C:\Users\Administrator\Desktop\新建文件夹 (2)\crypto_use\output\universe_results.json
2026-07-14 07:24:47,904 - root - INFO - === 启动 Universe Selection Engine ===
2026-07-14 07:24:47,916 - universe_selector.fetcher - INFO - 正在加载交易所市场信息...
2026-07-14 07:24:50,643 - universe_selector.fetcher - INFO - 报价币 USDT: 成交额门槛 = 20,000,000
2026-07-14 07:24:50,644 - universe_selector.fetcher - INFO - 报价币 USDC: 成交额门槛 = 5,000,000
2026-07-14 07:24:50,648 - universe_selector.fetcher - INFO - 基础流动性过滤完成:报价币=['USDT', 'USDC'],剩余 21 个标的。
2026-07-14 07:24:50,653 - universe_selector.fetcher - INFO - 开始异步拉取 21 个标的的 1h K线...
2026-07-14 07:24:52,363 - universe_selector.filters - INFO - 趋势跟踪池筛选完成: 17 个标的
2026-07-14 07:24:52,989 - universe_selector.filters - INFO - 均值回归池筛选完成: 0 个标的
2026-07-14 07:24:53,025 - universe_selector.filters - INFO - 统计套利配对筛选完成: 0 对
2026-07-14 07:24:53,026 - universe_selector.engine - INFO - 标的池结果已保存至: C:\Users\Administrator\Desktop\新建文件夹 (2)\crypto_use\output\universe_results.json