Files
qlib/backend/app/api/jobs.py
T
Simon 195f5d41f4 perf(backend): 内存优化三项——全市场研究不再占满 8G
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 干净。
2026-09-06 22:12:44 +08:00

98 lines
3.1 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 API(Phase 4):提交 / 查询 / SSE 进度。
POST /api/jobs 创建 Job(BackgroundTasks 后台执行),立即返回 job_id
GET /api/jobs/{id} 状态 + 结果(成功时内嵌 result)
GET /api/jobs/{id}/events SSE 进度(queued→running→success|failed)
执行模式见 config job.mode:subprocess 时研究任务在独立子进程跑(内存隔离),
API worker 不被重任务拖垮(内存优化专项)。
"""
from __future__ import annotations
import asyncio
import json
from datetime import datetime
from fastapi import APIRouter, BackgroundTasks, HTTPException
from fastapi.responses import StreamingResponse
from app.api.deps import DbSession, JobRepoDep
from app.application.services.job_executor import new_id, run_job_background
from app.domain.entities.research import (
BacktestResult,
FactorTestReport,
JobRecord,
JobStatus,
ResearchSpec,
)
from app.infrastructure.persistence.sqlalchemy.repositories.jobs_impl import (
SqlAlchemyJobRepository,
)
from app.infrastructure.persistence.sqlalchemy.session import SessionLocal
router = APIRouter(prefix="/jobs", tags=["jobs"])
def _decode_result(kind: str, result_json: str | None):
if result_json is None:
return None
model = BacktestResult if kind == "backtest" else FactorTestReport
return model.model_validate_json(result_json)
def _job_view(job: JobRecord) -> dict:
view = job.model_dump()
view["result"] = _decode_result(job.kind, job.result_json)
view.pop("result_json", None)
view["spec"] = json.loads(job.spec_json)
return view
@router.post("", summary="创建异步研究 Job")
def create_job(
spec: ResearchSpec,
background: BackgroundTasks,
session: DbSession,
job_repo: JobRepoDep,
) -> dict:
job = JobRecord(
id=new_id("JOB"),
kind=spec.type,
spec_json=spec.model_dump_json(),
status=JobStatus.QUEUED,
created_at=datetime.now(),
)
job_repo.create(job)
session.commit()
background.add_task(run_job_background, job.id)
return {"job_id": job.id, "status": job.status}
@router.get("/{job_id}", summary="查询 Job 状态与结果")
def get_job(job_id: str, job_repo: JobRepoDep) -> dict:
job = job_repo.get(job_id)
if job is None:
raise HTTPException(status_code=404, detail=f"Job {job_id} 不存在")
return _job_view(job)
@router.get("/{job_id}/events", summary="Job 进度 SSE")
async def job_events(job_id: str) -> StreamingResponse:
async def gen():
while True:
with SessionLocal() as session:
job = SqlAlchemyJobRepository(session).get(job_id)
if job is None:
yield "event: error\ndata: job not found\n\n"
return
payload = json.dumps(
{"job_id": job.id, "status": job.status, "stage": job.stage}, ensure_ascii=False
)
yield f"data: {payload}\n\n"
if job.status in (JobStatus.SUCCESS, JobStatus.FAILED, JobStatus.CANCELLED):
return
await asyncio.sleep(0.4)
return StreamingResponse(gen(), media_type="text/event-stream")