Files
qlib/backend/app/application/services/job_executor.py
T
Simon 23972e7063 feat: 股息率案例口径 + 策略库与图表统一 + 回测存档完整化
汇总三轮未提交的开发(每轮均在本机 MariaDB + 真实浏览器上验证):

1) 股息率案例(全市场股息率最高 n 只,默认 20,每 m 月择股)
   - 新增日频估值表 daily_basic + 迁移;股息率因子(dv_ratio / dividend_yield / TTM)
   - 名称历史表 stock_name_history:剔除 ST 按**择股日当时名称**判定,消除
     「曾高股息后 ST」的股息陷阱(实测 3.70pp 偏差)
   - 区间择股/调仓双周期(m 择股 / y 调仓)、指数成分与白名单、停牌近似剔除
   - 复权因子口径核对(4,164,742 行、缺失 0.0%)、收盘价成交与涨跌停拦单
   - 案例实测:2020-01-01~2026-09-04 总收益 +24.86%(年化 3.52%、回撤 -28.58%)

2) 策略库与前端统一
   - strategy 表 + CRUD/PUT 原地更新 + `describe_strategy` 按 spec 真实推导
     「一句话说明 + 计算公式 + 执行步骤 + 注意事项」(与引擎实执行规则同源)
   - 任何出现股票代码处都成对显示名称且可点击进个股页
   - 全站图表基座统一 TradingView Lightweight Charts(ECharts 依赖、
     锁文件、组件与文档标注一并清除),买卖点标记只落在真实交易日上

3) 回测存档完整化(可往复查看)
   - 同步端点(POST /api/backtests、/api/factor-tests)此前完全不落库 → 现在同样归档,
     归档 id 经响应头 X-Experiment-Id 返回(不破坏 response_model)
   - data_version 首次真实写入(数据快照指纹:最新交易日 + 各表规模)
   - 个股收益曲线默认**全量保存**(此前硬截断 60 只);超出体积预算才裁剪,
     并写 archive_meta(机器可读)+ unimplemented(人可读)如实标注
   - 列表 kind/q 过滤 + X-Total-Count(此前 limit=50 静默截断)、DELETE 归档
   - 只读归档页 /experiments/{id}(Server Component,SSR 直出**选股条件**与
     **交易执行依据**);结果视图按 kind 分发(backtest/factor_test/selection),
     非回测归档不套用回测口径
   - 新增 CLI:prune_experiments(保留策略,默认 dry-run)、
     restore_experiment_from_job(从 Job 副本按原 id 重建被删的历史归档,默认 dry-run)

门禁:pytest 388 passed、ruff All checks passed、tsc 0 错误、图表单测 7 passed、
next build 成功、契约脚本 verify_strategy_workspace 59/59(含按 kind 逐类验证归档页)。
2026-09-20 07:31:04 +08:00

502 lines
19 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""Job 执行编排与 Experiment 归档(Phase 4)+ 内存隔离执行(内存优化专项)。
execute_job 运行状态机 queued→running→(success|failed),
成功时自动把 spec/result 存为 Experiment(含代码版本),实现「研究可复现」(AGENT.md §21)。
独立 Session 生命周期(不依赖请求 scope),可被 FastAPI BackgroundTasks 或测试直接调用。
执行模式(config job.mode):
- local —— 本进程内执行(开发 / 单测;无内存隔离)
- subprocess —— 独立子进程执行(python -m app.cli.run_job <job_id>),
子进程内施加 RLIMIT_AS 上限(config job.max_memory_gb),研究任务 OOM 只会
MemoryError 失败归档,不会拖垮 API worker;父进程在子进程异常退出时补记 failed。
"""
from __future__ import annotations
import json
import os
import subprocess
import sys
import threading
from collections.abc import Callable
from datetime import datetime
from pathlib import Path
from app.application.services.experiment_archive import archive_experiment, new_id
from app.domain.entities.research import (
JobRecord,
JobStatus,
ResearchSpec,
)
from app.quant.service import ResearchService
# backend/app/application/services/job_executor.py → parents[4] = 项目根
PROJECT_ROOT = Path(__file__).resolve().parents[4]
BACKEND_ROOT = PROJECT_ROOT / "backend"
# 本进程内正在运行的子进程 Job(并发上限保护宿主内存)
_active_jobs: dict[str, subprocess.Popen | None] = {}
_active_lock = threading.Lock()
def _execute_inner(
job_id: str,
*,
session_factory: Callable,
job_repo_factory: Callable,
experiment_repo_factory: Callable,
stock_repo_factory: Callable,
daily_repo_factory: Callable,
engine,
basic_repo_factory: Callable | None = None,
financial_repo_factory: Callable | None = None,
index_repo_factory: Callable | None = None,
name_repo_factory: Callable | None = None,
) -> None:
"""执行 Job。
可选工厂的语义(AGENT.md §24:不做静默降级):
- `basic_repo_factory`:daily_basic(dv_ratio 等每日指标)—— 因子/条件引用时必需
- 复权(qfq/hfq)折算在行情仓储 SQL 内完成,无需 adjust_factor 工厂
- `financial_repo_factory`:财务表 —— 条件引用 fundamental.* 时必需
- `index_repo_factory`:指数成分 —— universe.index_code 时必需
- `name_repo_factory`:名称变更历史 —— universe.exclude_st 时点口径;未注入则
回退最新名称(旧行为),结果如实标注残余偏差
未注入且 spec 需要时,由 Service 抛出明确错误(而非返回空结果)。
"""
with session_factory() as session:
job_repo = job_repo_factory(session)
experiment_repo = experiment_repo_factory(session)
job = job_repo.get(job_id)
if job is None:
return
job.status = JobStatus.RUNNING
job.started_at = datetime.now()
job_repo.update(job)
session.commit()
try:
is_selection = job.kind == "selection"
basic_repo = basic_repo_factory(session) if basic_repo_factory else None
index_repo = index_repo_factory(session) if index_repo_factory else None
name_repo = name_repo_factory(session) if name_repo_factory else None
if is_selection:
from app.application.services.selection_service import SelectionService
from app.domain.entities.selection import SelectionQuery
spec = SelectionQuery.model_validate_json(job.spec_json)
service = SelectionService(
stock_repo_factory(session),
daily_repo_factory(session),
financial_repo_factory(session) if financial_repo_factory else None,
index_repo=index_repo,
basic_repo=basic_repo,
name_repo=name_repo,
)
else:
spec = ResearchSpec.model_validate_json(job.spec_json)
service = ResearchService(
stock_repo_factory(session),
daily_repo_factory(session),
engine,
index_repo=index_repo,
basic_repo=basic_repo,
financial_repo=(
financial_repo_factory(session) if financial_repo_factory else None
),
name_repo=name_repo,
)
def _set_stage(name: str) -> None:
"""阶段上报(v3 §23):独立短会话写 job.stage 并 commit(子进程同样走 DB)。"""
try:
with session_factory() as st_sess:
st = job_repo_factory(st_sess).get(job_id)
if st is not None and st.status == JobStatus.RUNNING:
st.stage = name
job_repo_factory(st_sess).update(st)
st_sess.commit()
except Exception: # noqa: BLE001 —— 阶段上报失败不阻断执行
pass
if is_selection:
_set_stage("selection")
result = service.select(spec)
elif spec.type == "backtest":
result = service.run_backtest(spec, on_stage=_set_stage)
else:
result = service.run_factor_test(spec, on_stage=_set_stage)
experiment = archive_experiment(
session=session,
kind=job.kind,
spec_json=job.spec_json,
result=result,
job_id=job.id,
experiment_repo=experiment_repo,
)
# 完整结果**只在 experiment 存一份**(此前 job.result_json 与
# experiment.result_json 各存一份相同内容,完整存档后 2× 浪费)。
# GET /api/jobs/{id} 经 job.experiment_id 回读 experiment;
# 老记录(result_json 有值、experiment_id 为空)仍走 job 回退解码。
job.result_json = None
job.experiment_id = experiment.id
job.status = JobStatus.SUCCESS
job.error = None
except Exception as exc: # noqa: BLE001 —— 统一记为 failed 供前端展示
job.status = JobStatus.FAILED
job.error = f"{type(exc).__name__}: {exc}"
finally:
job.finished_at = datetime.now()
# 回读最后一次阶段上报(_set_stage 经独立会话写库),避免被本会话覆盖
try:
with session_factory() as last_sess:
last = job_repo_factory(last_sess).get(job_id)
if last is not None:
job.stage = last.stage
except Exception: # noqa: BLE001
pass
job_repo.update(job)
session.commit()
def execute_job(
job_id: str,
*,
session_factory: Callable,
job_repo_factory: Callable,
experiment_repo_factory: Callable,
stock_repo_factory: Callable,
daily_repo_factory: Callable,
engine,
basic_repo_factory: Callable | None = None,
financial_repo_factory: Callable | None = None,
index_repo_factory: Callable | None = None,
name_repo_factory: Callable | None = None,
) -> None:
"""入口包装:任何未预期异常都将 Job 标记 failed(防止卡在 queued/running)。"""
try:
_execute_inner(
job_id,
session_factory=session_factory,
job_repo_factory=job_repo_factory,
experiment_repo_factory=experiment_repo_factory,
stock_repo_factory=stock_repo_factory,
daily_repo_factory=daily_repo_factory,
engine=engine,
basic_repo_factory=basic_repo_factory,
financial_repo_factory=financial_repo_factory,
index_repo_factory=index_repo_factory,
name_repo_factory=name_repo_factory,
)
except Exception as exc: # noqa: BLE001
try:
with session_factory() as session:
repo = job_repo_factory(session)
job = repo.get(job_id)
if job is not None:
job.status = JobStatus.FAILED
job.error = f"内部错误: {type(exc).__name__}: {exc}"
job.finished_at = datetime.now()
repo.update(job)
session.commit()
except Exception: # noqa: BLE001
pass
def default_factories() -> dict:
"""后台执行所需的独立 Session / Repository / 引擎装配(跨请求生命周期)。"""
from app.infrastructure.persistence.sqlalchemy.repositories.index_impl import (
SqlAlchemyIndexConstituentRepository,
)
from app.infrastructure.persistence.sqlalchemy.repositories.jobs_impl import (
SqlAlchemyExperimentRepository,
SqlAlchemyJobRepository,
)
from app.infrastructure.persistence.sqlalchemy.repositories.market_impl import (
SqlAlchemyDailyBarRepository,
SqlAlchemyDailyBasicRepository,
SqlAlchemyFinancialRepository,
SqlAlchemyStockNameHistoryRepository,
SqlAlchemyStockRepository,
)
from app.infrastructure.persistence.sqlalchemy.session import SessionLocal
from app.quant.engine import LocalEngine
return {
"session_factory": SessionLocal,
"job_repo_factory": lambda s: SqlAlchemyJobRepository(s),
"experiment_repo_factory": lambda s: SqlAlchemyExperimentRepository(s),
"stock_repo_factory": lambda s: SqlAlchemyStockRepository(s),
"daily_repo_factory": lambda s: SqlAlchemyDailyBarRepository(s),
# 新增装配(daily_basic / adjust_factor / 财务 / 指数成分):
# 缺失时对应能力(dv_ratio 因子、qfq/hfq 复权、fundamental 条件、指数成份池)
# 会由 Service 明确报错,绝不静默降级
"basic_repo_factory": lambda s: SqlAlchemyDailyBasicRepository(s),
"financial_repo_factory": lambda s: SqlAlchemyFinancialRepository(s),
"index_repo_factory": lambda s: SqlAlchemyIndexConstituentRepository(s),
"name_repo_factory": lambda s: SqlAlchemyStockNameHistoryRepository(s),
"engine": LocalEngine(),
}
def submit_and_run(spec: ResearchSpec, *, factories: dict | None = None) -> JobRecord:
"""创建并执行一个 Job(复用 Job 状态机与 Experiment 归档),返回终态 Job。
按 job.mode 调度:subprocess 模式在独立进程内执行(内存隔离),否则本进程。
"""
facts = factories or default_factories()
session_factory = facts["session_factory"]
job_repo = facts["job_repo_factory"]
job = JobRecord(
id=new_id("JOB"),
kind=spec.type,
spec_json=json.dumps(spec.model_dump(mode="json"), ensure_ascii=False),
status=JobStatus.QUEUED,
created_at=datetime.now(),
)
with session_factory() as session:
job_repo(session).create(job)
session.commit()
if _job_mode() == "subprocess":
_run_in_subprocess(job.id, session_factory=session_factory, job_repo_factory=job_repo)
else:
execute_job(
job.id,
session_factory=session_factory,
job_repo_factory=job_repo,
experiment_repo_factory=facts["experiment_repo_factory"],
stock_repo_factory=facts["stock_repo_factory"],
daily_repo_factory=facts["daily_repo_factory"],
engine=facts["engine"],
basic_repo_factory=facts.get("basic_repo_factory"),
financial_repo_factory=facts.get("financial_repo_factory"),
index_repo_factory=facts.get("index_repo_factory"),
name_repo_factory=facts.get("name_repo_factory"),
)
with session_factory() as session:
done = job_repo(session).get(job.id)
# **仅内存**读透(不落库):完整结果只存 experiment 一份,但 submit_and_run 的
# 既有调用方(scripts/run_dividend_case.py、agent 工具)习惯从 job.result_json
# 取结果,这里按 experiment_id 回读一次填进返回对象,避免调用方静默拿到空结果。
# 数据库中的 job.result_json 仍然保持 NULL(P1:不重复存第二份)。
if done is not None and done.experiment_id and done.result_json is None:
exp = facts["experiment_repo_factory"](session).get(done.experiment_id)
if exp is not None:
done.result_json = exp.result_json
assert done is not None
return done
# ---------------------------------------------------------------------------
# 执行模式调度(config job.mode:local | subprocess)
# ---------------------------------------------------------------------------
def _job_mode() -> str:
from app.core.config import get_settings
return get_settings().job_mode
def _job_memory_limit_gb() -> int:
from app.core.config import get_settings
return get_settings().job_memory_limit_gb
def _job_max_concurrent() -> int:
from app.core.config import get_settings
return get_settings().job_max_concurrent
def _job_worker_cmd(job_id: str) -> list[str]:
"""研究子进程命令行(独立进程执行,cwd=backend 使 `-m app.cli.run_job` 可导入)。"""
return [sys.executable, "-m", "app.cli.run_job", job_id]
def run_job_background(
job_id: str,
*,
session_factory: Callable | None = None,
job_repo_factory: Callable | None = None,
) -> None:
"""API 后台任务入口:按 job.mode 调度 Job 执行。"""
if _job_mode() == "subprocess":
_run_in_subprocess(
job_id, session_factory=session_factory, job_repo_factory=job_repo_factory
)
else:
execute_job(job_id, **default_factories())
def _acquire_slot(job_id: str) -> bool:
with _active_lock:
if len(_active_jobs) >= max(_job_max_concurrent(), 1):
return False
_active_jobs[job_id] = None
return True
def _release_slot(job_id: str) -> None:
with _active_lock:
_active_jobs.pop(job_id, None)
def _mark_failed(
job_id: str,
error: str,
*,
session_factory: Callable | None = None,
job_repo_factory: Callable | None = None,
) -> None:
"""把非终态 Job 标记 failed(子进程异常退出 / 并发超限时兜底,防永久 running)。"""
if session_factory is None or job_repo_factory is None:
facts = default_factories()
session_factory = session_factory or facts["session_factory"]
job_repo_factory = job_repo_factory or facts["job_repo_factory"]
try:
with session_factory() as session:
repo = job_repo_factory(session)
job = repo.get(job_id)
if job is not None and job.status in (JobStatus.QUEUED, JobStatus.RUNNING):
job.status = JobStatus.FAILED
job.error = error
job.finished_at = datetime.now()
repo.update(job)
session.commit()
except Exception: # noqa: BLE001 —— 兜底标记失败自身异常不向上抛
pass
def _subprocess_log_target():
"""研究子进程 stderr 的落盘目标。
原先父进程用 DEVNULL,会吞掉子进程全部诊断(含 run_job 打印的内存上限设置失败
警告),使「如实提示未实现项」落空。改为追加写入 <项目根>/logs/job-subprocess.log。
记日志失败(无权限/磁盘满)绝不能导致 Job 失败 —— 此时退回 DEVNULL。
"""
try:
log_dir = BACKEND_ROOT.parent / "logs"
log_dir.mkdir(parents=True, exist_ok=True)
return (log_dir / "job-subprocess.log").open("a", encoding="utf-8")
except OSError:
return subprocess.DEVNULL
def _run_in_subprocess(
job_id: str,
*,
session_factory: Callable | None = None,
job_repo_factory: Callable | None = None,
) -> None:
"""在独立 python 进程执行 Job:子进程 RLIMIT 上限(防 OOM 整机),
父进程等待;子进程异常退出(如被杀)时把 Job 补记 failed。"""
if not _acquire_slot(job_id):
_mark_failed(
job_id,
"系统繁忙:并发研究任务已达上限,请稍后重试",
session_factory=session_factory,
job_repo_factory=job_repo_factory,
)
return
env = dict(os.environ)
env["QLIB_JOB_MEM_LIMIT_GB"] = str(_job_memory_limit_gb())
rc = -1
already_marked = False
errlog = _subprocess_log_target()
try:
proc = subprocess.Popen(
_job_worker_cmd(job_id),
cwd=BACKEND_ROOT,
env=env,
stdout=subprocess.DEVNULL,
stderr=errlog,
)
with _active_lock:
_active_jobs[job_id] = proc
rc = proc.wait()
except Exception as exc: # noqa: BLE001
_mark_failed(
job_id,
f"研究子进程启动失败: {type(exc).__name__}: {exc}",
session_factory=session_factory,
job_repo_factory=job_repo_factory,
)
already_marked = True
finally:
if errlog is not subprocess.DEVNULL:
errlog.close()
_release_slot(job_id)
if rc != 0 and not already_marked:
_mark_failed(
job_id,
f"研究子进程异常退出(code={rc}):任务未完成",
session_factory=session_factory,
job_repo_factory=job_repo_factory,
)
def terminate_active(job_id: str) -> bool:
"""终止该 Job 的活动子进程(如有)。返回是否找到并终止。"""
with _active_lock:
proc = _active_jobs.get(job_id)
if proc is None:
return False
import contextlib
with contextlib.suppress(Exception):
proc.terminate()
return True
def cancel_job(
job_id: str,
*,
session_factory: Callable | None = None,
job_repo_factory: Callable | None = None,
) -> bool:
"""取消 queued/running Job(置 CANCELLED 并终止子进程);不可取消返回 False。"""
facts = default_factories()
sf = session_factory or facts["session_factory"]
jr = job_repo_factory or facts["job_repo_factory"]
with sf() as session:
repo = jr(session)
job = repo.get(job_id)
if job is None or job.status not in (JobStatus.QUEUED, JobStatus.RUNNING):
return False
job.status = JobStatus.CANCELLED
job.stage = None
repo.update(job)
session.commit()
terminate_active(job_id)
return True
def mark_stale_jobs_failed(
*,
session_factory: Callable | None = None,
job_repo_factory: Callable | None = None,
reason: str | None = None,
) -> int:
"""服务启动调用:把上次进程异常退出遗留的 queued/running Job 标记 failed。
返回处理数量。防「进程被杀后 Job 永久 running」(AGENT.md §20 状态机闭环)。
"""
facts = default_factories()
sf = session_factory or facts["session_factory"]
jr = job_repo_factory or facts["job_repo_factory"]
counted = 0
with sf() as session:
repo = jr(session)
for status in (JobStatus.QUEUED, JobStatus.RUNNING):
for job in repo.list_by_status(status):
job.status = JobStatus.FAILED
job.error = reason or "服务重启:上次未完成任务被中断"
job.finished_at = datetime.now()
repo.update(job)
counted += 1
session.commit()
return counted