mirror of
https://ghfast.top/https://github.com/aeroxw/tick-stock-panel.git
synced 2026-09-12 19:04:15 +08:00
* feat(screener): 选股引擎支持 ETF - 12 个内置策略打 asset_types 白名单 + strategy_supports_asset;涨停类 (连板/断板反包)仅股票,其余 10 个技术类对 ETF 开放 - ScreenerService(repo, asset_type) 分流取数,ETF 复用 kline_etf_enriched, 跳过股票专用历史缓存与涨停信号;进程级 _history_cache key 含 asset_type - API /run、/run_preset 透传 asset_type;/strategies 按资产过滤; 股票专有策略在 ETF 下返回空 - 新增 enriched_dirname(asset_type) 共享 helper;get_enriched_latest_asset 增 refresh 参数(供轮询线程避免冷缓存同步重算) - 前端「策略」页加 股票/ETF 切换,ETF 走实时单跑(空日期→用 ETF 自身最新日); QK.screenerStrategies 按 asset_type keyed - 测试:test_screener_etf.py Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> * feat(backtest): 回测支持 ETF(个股/因子/策略组合) - 三条回测路径 + 共用 BacktestEngine 面板加载按 asset_type 路由到 kline_etf_enriched(复用 enriched_dirname);PanelCache key 隔离资产; ETF 跳过股票专用 get_enriched_range 缓存 - 面板 compute_all/名称 JOIN 按 asset_type 取维表(get_instruments_asset), 修复 ETF 策略回测用错股票维表致名称为空/涨停信号算错 - BacktestConfig/FactorConfig/StrategyBacktestConfig 增 asset_type - 三个回测 API + SSE stream 透传 asset_type;_make_job_key 纳入 asset_type (修复 stream 与 cancel job_key 不对齐致取消失效的回归) - 前端策略组合页/因子页加 股票/ETF 切换,标的搜索与策略列表跟随资产; assetType 持久化 - 测试:test_backtest_etf.py(含 job_key 一致性回归);既有回测测试替身同步 Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> * feat(monitor): 监控规则支持 ETF - engine.evaluate(df, asset_type) 按规则 asset_type 分轮评估;quote_service 增开 ETF 评估轮(用 ETF enriched 快照),股票轮不受影响、不重置其策略结果 - ETF 评估轮独立 try(异常不丢弃已算出的股票告警)+ refresh=False(不在轮询 线程触发 ETF 冷缓存同步重算) - ETF 版历史加载器(main.py 注入)+ 按规则 asset_type 选加载器 - _strategy_pools 按 (sid, asset_type) 键,避免同策略股票/ETF 规则互相覆盖 - name_map 仅在有 ETF 规则时补 ETF 维表, setdefault 保股票名优先 - RuleModel/normalize 增 asset_type(默认 stock,持久化往返) - 前端 RuleEditor 加 股票/ETF 选择,策略列表与标的搜索跟随资产 - 测试:test_monitor_etf.py Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> * feat(etf): 前端 API 绑定透传 asset_type + 文档 - api.ts: screener/backtest 绑定加 assetType 参数,MonitorRule 类型加 asset_type - docs/features.md: 标注选股/回测/监控的 ETF 支持范围与前提 Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> * fix(reliability): 管道并发/原子写/能力探测/监控告警多处加固 后端可靠性专项修复(均带回归测试, backend 全套 64 passed): 并发与数据完整性: - 盘后管道单飞: JobStore.create() 去重纳入 pending∨running, 关闭"两次快速点击" 并发双跑窗口; 新增 _heavy_run_lock 执行槽挡住 reap 后僵尸线程并发写 parquet - adj_factor/minute 全部改走原子写(tmp+replace), 消除 kill/断电致 all.parquet 损坏 - 分块拉取失败聚合 WARNING 可见化(不再静默当成功); 复权失败标的会保持旧价已提示 能力探测: - 周期重探(60min)热更新 app.state.capabilities, 付费 Key 过期/续费无需重启即可见 - 瞬时探测失败(超时/连接/5xx, 按 _is_transient 判定)不降级、保留旧付费档; 真 401/无权限仍正常降级回落 free-api 监控告警: - 评估仅在连续竞价(9:30-11:30/13:00-15:00)+ 快照当日新鲜度下进行, 避开集合竞价/ 收盘后陈旧价与节假日误告警 - scope=sector fail-closed(validate 拒绝新建 + _apply_scope 返回空), 修复板块规则 对全市场刷屏 - 飞书 webhook 加退避重试并移到独立线程池 fire-and-forget, 不再阻塞行情轮询线程 单标的新鲜度: 新增 repo.symbols_lagging() 检测掉队标的并 WARNING + 计入 job 结果 Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
102 lines
3.9 KiB
Python
102 lines
3.9 KiB
Python
"""盘后管道 API — 异步触发 + 进度跟踪。"""
|
||
from __future__ import annotations
|
||
|
||
import asyncio
|
||
import concurrent.futures as _cf
|
||
import logging
|
||
|
||
from fastapi import APIRouter, HTTPException, Request
|
||
|
||
from app.jobs import daily_pipeline
|
||
from app.services.pipeline_jobs import job_store, release_run_slot, try_acquire_run_slot
|
||
from app.api.data import invalidate_storage_cache
|
||
|
||
# 长时间任务专用线程池(隔离于 FastAPI 默认线程池,防止阻塞请求处理)
|
||
_long_task_executor = _cf.ThreadPoolExecutor(max_workers=2, thread_name_prefix="long-task")
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
router = APIRouter(prefix="/api/pipeline", tags=["pipeline"])
|
||
|
||
|
||
@router.post("/run")
|
||
async def run_now(request: Request) -> dict:
|
||
"""异步触发盘后管道,立即返回 job_id。客户端轮询 /jobs/{id} 拿进度。
|
||
|
||
若已有任务在跑,**返回该任务 id 而不是开新任务**(防止并发拉数据撞限流)。
|
||
但如果该任务已运行超过 10 分钟 (可能因 reload 卡死), 强制标记为失败后重新创建。
|
||
"""
|
||
repo = request.app.state.repo
|
||
capset = request.app.state.capabilities
|
||
|
||
# 检测卡死的 running job (如 reload 后孤儿 task / 网络读无限阻塞)。
|
||
# reap_stale 会在 /run 和 /jobs/{id} 轮询端点都调用,保证卡死后能自愈。
|
||
job_store.reap_stale()
|
||
|
||
# 单飞: 复用任何活跃 (pending∨running) 任务, is_new=False 时不再调度新任务
|
||
job_id, is_new = job_store.create()
|
||
if not is_new:
|
||
return {"job_id": job_id, "reused": True}
|
||
|
||
# 在 executor 里跑同步任务(pipeline 内部都是阻塞 IO + CPU)
|
||
async def task() -> None:
|
||
# 重任务执行槽: 防僵尸并发(reap 后线程仍活时新任务不得并行写 parquet)
|
||
if not try_acquire_run_slot():
|
||
job_store.fail(job_id, "已有数据任务在运行(或上一次任务卡死未结束),请稍后再试")
|
||
return
|
||
try:
|
||
job_store.start(job_id)
|
||
loop = asyncio.get_event_loop()
|
||
|
||
def progress(stage: str, pct: int, msg: str, stage_pct: int | None = None,
|
||
skip_log: bool = False) -> None:
|
||
job_store.progress(job_id, stage, pct, msg, stage_pct=stage_pct, skip_log=skip_log)
|
||
|
||
result = await loop.run_in_executor(
|
||
_long_task_executor,
|
||
lambda: daily_pipeline.run_now(repo, capset, on_progress=progress),
|
||
)
|
||
job_store.succeed(job_id, result)
|
||
invalidate_storage_cache()
|
||
repo.refresh_cache() # 刷新 Polars 缓存
|
||
except Exception as e: # noqa: BLE001
|
||
logger.exception("pipeline failed")
|
||
job_store.fail(job_id, str(e))
|
||
invalidate_storage_cache()
|
||
finally:
|
||
release_run_slot()
|
||
|
||
asyncio.create_task(task())
|
||
return {"job_id": job_id, "reused": False}
|
||
|
||
|
||
@router.get("/jobs/{job_id}")
|
||
def get_job(job_id: str) -> dict:
|
||
# 每次轮询都检查卡死 job — 前端每秒轮询,STALE_JOB_TIMEOUT_S(10min)后必定自愈,
|
||
# 无需用户再次手动点「同步」。
|
||
job_store.reap_stale()
|
||
j = job_store.get(job_id)
|
||
if not j:
|
||
raise HTTPException(status_code=404, detail="job not found")
|
||
return j
|
||
|
||
|
||
@router.post("/jobs/{job_id}/cancel")
|
||
def cancel_job(job_id: str) -> dict:
|
||
"""手动取消一个 running 的 job。"""
|
||
j = job_store.get(job_id)
|
||
if not j:
|
||
raise HTTPException(status_code=404, detail="job not found")
|
||
if j["status"] not in ("running", "pending"):
|
||
raise HTTPException(status_code=400, detail=f"job status is {j['status']}, cannot cancel")
|
||
job_store.fail(job_id, "用户手动取消")
|
||
return {"cancelled": job_id}
|
||
|
||
|
||
@router.get("/jobs")
|
||
def list_jobs(limit: int = 20) -> dict:
|
||
return {
|
||
"active_id": job_store.active_id(),
|
||
"jobs": job_store.list_recent(limit=limit),
|
||
}
|