commit 54efc3c361c744b8662131a0fe824bde3c92f3e8 Author: gavindiaz Date: Tue Jul 14 07:31:13 2026 +0800 first commit diff --git a/README.md b/README.md new file mode 100644 index 0000000..b4ec042 --- /dev/null +++ b/README.md @@ -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 线重试 / 路径稳健化 | + +--- + +## 十一、许可 + +仅供学习和研究使用。加密货币交易有风险,请自行评估。 \ No newline at end of file diff --git a/config/config.yaml b/config/config.yaml new file mode 100644 index 0000000..ef60ce5 --- /dev/null +++ b/config/config.yaml @@ -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" \ No newline at end of file diff --git a/main.py b/main.py new file mode 100644 index 0000000..26d53a0 --- /dev/null +++ b/main.py @@ -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()) \ No newline at end of file diff --git a/output/btc_results.json b/output/btc_results.json new file mode 100644 index 0000000..5150cad --- /dev/null +++ b/output/btc_results.json @@ -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": [] +} \ No newline at end of file diff --git a/output/btc_usdt_results.json b/output/btc_usdt_results.json new file mode 100644 index 0000000..16240ac --- /dev/null +++ b/output/btc_usdt_results.json @@ -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": [] +} \ No newline at end of file diff --git a/output/universe_results.json b/output/universe_results.json new file mode 100644 index 0000000..7b3a25f --- /dev/null +++ b/output/universe_results.json @@ -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": [] +} \ No newline at end of file diff --git a/output/usdc_results.json b/output/usdc_results.json new file mode 100644 index 0000000..5477fef --- /dev/null +++ b/output/usdc_results.json @@ -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": [] +} \ No newline at end of file diff --git a/requirements.txt b/requirements.txt new file mode 100644 index 0000000..c48946f --- /dev/null +++ b/requirements.txt @@ -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 diff --git a/universe_selector/__init__.py b/universe_selector/__init__.py new file mode 100644 index 0000000..5350f1e --- /dev/null +++ b/universe_selector/__init__.py @@ -0,0 +1 @@ +from .engine import UniverseEngine diff --git a/universe_selector/__pycache__/__init__.cpython-312.pyc b/universe_selector/__pycache__/__init__.cpython-312.pyc new file mode 100644 index 0000000..222867e Binary files /dev/null and b/universe_selector/__pycache__/__init__.cpython-312.pyc differ diff --git a/universe_selector/__pycache__/engine.cpython-312.pyc b/universe_selector/__pycache__/engine.cpython-312.pyc new file mode 100644 index 0000000..976828a Binary files /dev/null and b/universe_selector/__pycache__/engine.cpython-312.pyc differ diff --git a/universe_selector/__pycache__/fetcher.cpython-312.pyc b/universe_selector/__pycache__/fetcher.cpython-312.pyc new file mode 100644 index 0000000..24f0a80 Binary files /dev/null and b/universe_selector/__pycache__/fetcher.cpython-312.pyc differ diff --git a/universe_selector/__pycache__/filters.cpython-312.pyc b/universe_selector/__pycache__/filters.cpython-312.pyc new file mode 100644 index 0000000..5663ff0 Binary files /dev/null and b/universe_selector/__pycache__/filters.cpython-312.pyc differ diff --git a/universe_selector/__pycache__/utils.cpython-312.pyc b/universe_selector/__pycache__/utils.cpython-312.pyc new file mode 100644 index 0000000..e113113 Binary files /dev/null and b/universe_selector/__pycache__/utils.cpython-312.pyc differ diff --git a/universe_selector/engine.py b/universe_selector/engine.py new file mode 100644 index 0000000..4c25aa6 --- /dev/null +++ b/universe_selector/engine.py @@ -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}") \ No newline at end of file diff --git a/universe_selector/fetcher.py b/universe_selector/fetcher.py new file mode 100644 index 0000000..7f2bc40 --- /dev/null +++ b/universe_selector/fetcher.py @@ -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 session(ThreadedResolver 绕过 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 通常为 True;onchainDate/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 \ No newline at end of file diff --git a/universe_selector/filters.py b/universe_selector/filters.py new file mode 100644 index 0000000..1360cbe --- /dev/null +++ b/universe_selector/filters.py @@ -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 diff --git a/universe_selector/utils.py b/universe_selector/utils.py new file mode 100644 index 0000000..efdc385 --- /dev/null +++ b/universe_selector/utils.py @@ -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 diff --git a/use_engine.log b/use_engine.log new file mode 100644 index 0000000..a85c7a6 --- /dev/null +++ b/use_engine.log @@ -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: +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: +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: +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