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 干净。
98 lines
3.1 KiB
Python
98 lines
3.1 KiB
Python
"""异步研究 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")
|