1) 数据装配流式+列裁剪:Repository 新增 stream_range_many_columns(只 SELECT 所需列、SQL 侧转 REAL、yield_per 分批),引擎按 required_columns 取数 (LocalEngine 仅 close+因子字段),消除 ORM/Decimal 全量物化; 2) 研究 Job 独立子进程执行(job.mode=subprocess):python -m app.cli.run_job 在子进程内设 RLIMIT_AS 上限,OOM 归档 failed 而非拖垮 API worker; 子进程异常退出由父进程补记 failed;并发上限 2; 3) 服务启动清理:残留 queued/running Job 标记 failed(防永久 running)。 实测同款全市场回测:uvicorn worker RSS 稳定 ~220MB,任务峰值内存由 4.1GB+ 降至 ~470MB,24s 完成并归档(此前 43s 未完成即 OOM)。 新增/更新测试 96 passed,ruff 干净。
92 lines
3.8 KiB
Python
92 lines
3.8 KiB
Python
"""QlibEngine —— 基于 Qlib 数据管线的研究引擎(QuantEngine 实现)。
|
||
|
||
v1 能力(本阶段已打通并测试):
|
||
1. provider.build_qlib_dataset:把本地行情按 qlib 二进制格式落盘(data/qlib)
|
||
2. dataset.ensure_qlib_init + D.features:从 QlibDataset 读取行情(真实 qlib 通路)
|
||
3. 在 Qlib 读取的行情面板上执行 TopK 因子回测 → 标准 BacktestResult
|
||
(与 LocalEngine 记账规则一致:无未来函数、成本/涨跌停/停牌近似 + unimplemented 标注)
|
||
|
||
说明:
|
||
- 因子分(score)在源行情上按 app.quant.factors 计算(注册表口径一致);
|
||
qlib 侧负责数据读取供给(未来可由 Qlib Dataset 直接产出特征)。
|
||
- Alpha158 特征集 + LightGBM 预测信号的模型增强(walk-forward 训练/预测)为下一步
|
||
TODO,详见 docs/ROADMAP.md §2 与 docs/QLIB_VERIFICATION.md。
|
||
- 默认研究引擎仍是 LocalEngine(app/quant/engine.py);切换只需注入本类。
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
from pathlib import Path
|
||
|
||
import pandas as pd
|
||
|
||
from app.domain.entities.research import BacktestResult, FactorTestReport, ResearchSpec
|
||
from app.quant.engine import QuantEngine
|
||
from app.quant.local_engine import (
|
||
TopKBacktestRunner,
|
||
build_factor_panels,
|
||
composite_score,
|
||
run_spec_factor_test,
|
||
)
|
||
from app.quant.qlib_adapter.dataset import ensure_qlib_init, load_close_panel
|
||
from app.quant.qlib_adapter.provider import build_qlib_dataset
|
||
|
||
_ENGINE_NOTE = "QlibEngine v1:本地行情→QlibDataset(bin)→D.features 读取→TopK 因子回测"
|
||
|
||
|
||
class QlibEngine(QuantEngine):
|
||
"""基于 Qlib 数据管线的引擎;因子评估复用共享实现,回测从 QlibDataset 读取行情。"""
|
||
|
||
name = "qlib"
|
||
|
||
# qlib 落盘需要 OHLCV(vwap/factor 由本地合成,不来自行情表)
|
||
_DUMP_COLUMNS = {"open", "high", "low", "close", "volume", "amount"}
|
||
|
||
def __init__(self, qlib_dir: Path | None = None) -> None:
|
||
# 默认落盘到 data/qlib(与 storage.qlib_dir 一致);可注入临时目录便于测试
|
||
self.qlib_dir = qlib_dir or _default_qlib_dir()
|
||
|
||
def required_columns(self, spec: ResearchSpec) -> set[str]:
|
||
# 回测需把全字段落盘成 qlib 数据集;因子测试走共享实现,只需 close+因子字段
|
||
if spec.type == "backtest":
|
||
return set(self._DUMP_COLUMNS)
|
||
from app.quant.engine import factor_required_columns
|
||
|
||
return factor_required_columns(spec)
|
||
|
||
def run_factor_test(
|
||
self, daily: pd.DataFrame, spec: ResearchSpec, horizon_days: int = 21
|
||
) -> FactorTestReport:
|
||
# 因子评估与数据无关,直接复用共享实现(与 LocalEngine 同口径)
|
||
report, _panels = run_spec_factor_test(daily, spec, horizon_days)
|
||
return report
|
||
|
||
def run_backtest(self, daily: pd.DataFrame, spec: ResearchSpec) -> BacktestResult:
|
||
panels = build_factor_panels(daily, spec.factors)
|
||
score = composite_score(panels)
|
||
|
||
self.qlib_dir.mkdir(parents=True, exist_ok=True)
|
||
uri = build_qlib_dataset(daily, self.qlib_dir)
|
||
ensure_qlib_init(uri)
|
||
|
||
symbols = [str(s) for s in daily["symbol"].unique()]
|
||
start, end = spec.period
|
||
close = load_close_panel(uri, symbols, start, end)
|
||
if close.empty:
|
||
raise RuntimeError(
|
||
"QlibDataset 读取为空:请检查 build_qlib_dataset 落盘与 provider_uri"
|
||
)
|
||
close = close.sort_index()
|
||
|
||
result = TopKBacktestRunner(spec, score, close).run()
|
||
result.config_snapshot = spec.model_dump(mode="json")
|
||
note = _ENGINE_NOTE
|
||
result.unimplemented = [note, *result.unimplemented]
|
||
return result
|
||
|
||
|
||
def _default_qlib_dir() -> Path:
|
||
from app.core.config import PROJECT_ROOT
|
||
|
||
return PROJECT_ROOT / "data" / "qlib"
|