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 干净。
29 lines
906 B
Python
29 lines
906 B
Python
"""Job / Experiment Repository Protocol(Phase 4)。"""
|
||
|
||
from __future__ import annotations
|
||
|
||
from typing import Protocol
|
||
|
||
from app.domain.entities.research import ExperimentRecord, JobRecord
|
||
|
||
|
||
class JobRepository(Protocol):
|
||
def create(self, job: JobRecord) -> JobRecord: ...
|
||
|
||
def get(self, job_id: str) -> JobRecord | None: ...
|
||
|
||
def update(self, job: JobRecord) -> None: ...
|
||
|
||
def list_recent(self, kind: str | None = None, limit: int = 20) -> list[JobRecord]: ...
|
||
|
||
def list_by_status(self, status: str, limit: int = 100) -> list[JobRecord]:
|
||
"""按状态查询(服务启动清理残留 queued/running 用)。"""
|
||
|
||
|
||
class ExperimentRepository(Protocol):
|
||
def save(self, experiment: ExperimentRecord) -> ExperimentRecord: ...
|
||
|
||
def get(self, experiment_id: str) -> ExperimentRecord | None: ...
|
||
|
||
def list_recent(self, limit: int = 50) -> list[ExperimentRecord]: ...
|