Files
tick-stock-panel/backend/app/services/ext_pull.py
T
Jinfeng SunandClaude Opus 4.8 9aa96edbd7 改进: 并发韧性 + 数据性能 + 死代码清理 + 前端 UX + ST/Sharpe 修复 (#78)
* fix(concurrency): 共享缓存/任务表加锁, 全局限速, depth 原子写, 认证热路径缓存

修复多线程下的竞态与阻塞:
- overview/strategy_cache/PanelCache/StrategyMonitor._watching 四处共享状态加锁,
  消除 "dict/OrderedDict mutated" 与丢更新/半写读取
- strategy_cache/depth parquet 改临时文件 + os.replace 原子写
- rate_limits 改进程级共享时间轴限速, 并发同步不再聚合超过单能力 rpm;
  scheduler 令牌账目与 sleep 分离, sleep 不再独占锁串行化其他请求
- auth.is_configured() 内存缓存, 认证中间件不再每请求读盘阻塞事件循环
- api/backtest 任务清理/取消全程持 _jobs_lock, 并用 Semaphore(2) 限并发重回测

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* perf(data): limit_ladder 去 N+1 全市场重算, 指标裁剪, factor 向量化

- limit_ladder 前一日 consecutive 改窄读单日 parquet 存储列 (谓词/投影下推),
  替代 range(1,10) 逐日 _load_enriched_for_date 全市场指标重算 (最坏 9x)
- compute_indicators 新增可选 needed 裁剪 (默认 None 行为逐位不变, 已对照验证),
  factor 只算所需因子列
- factor._calc_period_return 用 Polars join 替代 Python 逐行 price_map 循环,
  _add_groups 去 map_elements 改纯表达式 (输出逐位一致)
- screener ext value_map 按 parquet mtime 记忆化, 免每请求磁盘重读
  (DuckDB 过滤仍用隔离 :memory: 连接, 不扩大注入面)

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* refactor(backend): 报表存储去重, 删死代码, DuckDB 视图重建收敛, 管道失败如实标记

- 三份近乎逐字复制的 *_reports.py 收敛到共享 JsonReportStore (原子写 + 锁),
  各模块公有 API/id 格式/上限/落盘 schema 完全保持不变
- 删除 ext_pull.py 中字节相同的死 _run_loop (Python 只绑第二个) 及无用 import
- 13 张 DuckDB 视图重建收敛为唯一权威 repository.rebuild_views(),
  daily_pipeline 与 /api/data/clear 改为调用 (修好 clear 路径漏挂视图的漂移)
- daily_pipeline 累积 stage_errors 并在末尾抛出, 部分失败不再误报成功;
  free/None 模式的能力门控跳过不计入失败

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* feat(frontend): SSE 连接态, 路由代码分割, 查询失效修复, 三态与无障碍

- 实时行情 SSE: 连接态 store + 指数退避 + 断线徽标/toast (避免静默丢告警);
  回测 SSE 断线有界重连 + 可重试, 不再永久卡住进度条
- router 全部 React.lazy + Suspense, vite manualChunks 拆图表库
  (echarts 变独立 1MB 懒加载 chunk, 首屏包显著减小)
- 修 Data 清库后其它页显示旧数据 (改回广域失效); 修 Watchlist kline 失效键
  永不匹配; query key 收敛到 QK 工厂 (新增 strategyDetail)
- Monitor/Analysis/StockAnalysis/ExtPages/CustomSignals 补 loading 门控与
  error/empty 三态区分
- 新增共享 Modal 原语 (焦点陷阱/ESC/焦点还原/aria), 改造 3 个高频弹窗;
  Toast/AlertToast 加 aria-live 与键盘可达; Watchlist/LimitUpLadder 卡片 memo

修复本轮 review 发现的缺陷:
- Modal 焦点 effect 依赖 onClose 致每次输入抢焦点 → 改 ref 只装一次
- StrategySettingsDialog 删除确认框被 Modal 面板裁剪 → 移出作兄弟节点

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

* fix(quant): 修正 ST 板块限价套错 与 因子 Sharpe 年化频率

两个不报错但会算错数的领域 bug:

1. ST 5% 涨跌停限幅被无条件套到创业板/科创板 ST 股:
   注册制改革后 创业板(300/301)、科创板(688/689) 的风险警示股仍执行 20%,
   北交所 30%, 只有主板 ST 才是 5%。原代码 _is_st 先判且覆盖板块限幅, 导致
   创业板/科创板 ST 的涨停价按 5% 计算 → +5% 被误报涨停、真 +20% 涨停被漏报,
   污染 signal_limit_up / consecutive_limit_ups / 连板梯队 / near_limit_up。
   修正: ST 5% 仅在 ~(创业板|科创板|北交所) 时生效 (EOD + 盘中两条路径 + near_limit_up)。

2. 因子回测 Sharpe 一律乘 √252, 但 group_nav 每点是一个调仓周期收益:
   月频调仓下是月收益, 乘 √252 会把 Sharpe 高估 √(252/12) ≈ 4.6x (周频 ≈2.2x),
   使无效因子显示成明星因子, 废掉"先筛无效指标"的用途。
   修正: 年化系数按 config.rebalance 取 √252/√52/√12。

新增 tests/test_st_limit_and_sharpe.py (5 例) 覆盖两处修正。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-08 18:10:19 +08:00

294 lines
11 KiB
Python

"""扩展数据定时拉取引擎 — 从外部 API 拉取数据写入 Parquet。"""
from __future__ import annotations
import asyncio
import json
import logging
import threading
from datetime import date, datetime, timezone
from functools import reduce
from typing import Any
import httpx
from app.services.ext_data import (
ExtConfig,
ExtConfigStore,
rows_to_parquet,
)
logger = logging.getLogger(__name__)
# ---------------------------------------------------------------------------
# 响应解析
# ---------------------------------------------------------------------------
def _extract_rows(data: Any, path: str) -> list[dict]:
"""按 dot-path 从 JSON 响应中提取行数组。
例: path="data.list" → response["data"]["list"]
如果 path 为空,直接将 data 视为数组。
"""
if not path:
if isinstance(data, list):
return data
raise ValueError("response_path 为空但响应不是数组")
keys = path.split(".")
current = data
for key in keys:
if isinstance(current, dict):
if key not in current:
raise ValueError(f"响应中不存在路径 '{path}',缺失键 '{key}'")
current = current[key]
elif isinstance(current, list):
try:
current = current[int(key)]
except (ValueError, IndexError) as e:
raise ValueError(f"响应路径 '{path}' 解析失败: {e}") from e
else:
raise ValueError(f"响应路径 '{path}' 中间值不是 dict/list: {type(current)}")
if not isinstance(current, list):
raise ValueError(f"路径 '{path}' 指向的不是数组,而是 {type(current)}")
return current
def _apply_field_map(rows: list[dict], field_map: dict[str, str]) -> list[dict]:
"""将外部字段名映射为内部配置字段名。field_map: {外部名: 内部名}。"""
if not field_map:
return rows
mapped = []
for row in rows:
new_row: dict = {}
for k, v in row.items():
mapped_key = field_map.get(k, k)
new_row[mapped_key] = v
mapped.append(new_row)
return mapped
# ---------------------------------------------------------------------------
# 拉取执行
# ---------------------------------------------------------------------------
def _apply_preset_flatten(config_id: str, rows: list[dict]) -> list[dict]:
"""对内置预设 (概念/行业) 应用结构转换, 与 fetch_preset 保持一致。
延迟导入避免与 ext_presets 形成循环依赖。
非预设 id 原样返回。
"""
if config_id not in ("ext_gn_ths", "ext_hy_ths"):
return rows
from app.services.ext_presets import _flatten_concept_rows, _flatten_industry_rows
flatten = _flatten_concept_rows if config_id == "ext_gn_ths" else _flatten_industry_rows
return flatten(rows)
async def fetch_and_ingest(
config: ExtConfig,
data_dir,
) -> tuple[int, str]:
"""执行一次拉取: 请求外部 API → 解析响应 → 写入 Parquet。
Returns:
(rows_written, date_str)
"""
pull = config.pull
if not pull or not pull.url:
raise ValueError("拉取未配置或 URL 为空")
async with httpx.AsyncClient(timeout=30) as client:
headers = pull.headers or {}
kwargs: dict[str, Any] = {"headers": headers}
if pull.method.upper() == "POST" and pull.body:
kwargs["content"] = pull.body
if "content-type" not in {k.lower() for k in headers}:
kwargs["headers"]["Content-Type"] = "application/json"
resp = await client.request(pull.method.upper(), pull.url, **kwargs)
resp.raise_for_status()
# 解析 JSON
try:
data = resp.json()
except Exception as e:
raise ValueError(f"响应不是有效 JSON: {e}") from e
# 提取行
rows = _extract_rows(data, pull.response_path)
if not rows:
raise ValueError("提取到的行数为 0")
# 内置预设 (概念/行业): 应用结构转换, 让产出 schema 与分析页一致。
# 否则 raw 接口列 (concepts/industries 数组、name) 会直接覆盖正确的 part.parquet,
# 导致分析页因找不到维度字段 (所属概念/所属同花顺行业) 而"数据消失"。
# 见 ext_presets._flatten_* —— 手动拉取 / 定时拉取都必须走同一套转换。
rows = _apply_preset_flatten(config.id, rows)
# 字段映射
rows = _apply_field_map(rows, pull.field_map)
# 校验可关联标的的字段:直接 symbol/code,或配置里声明的映射源列。
row_keys = set(rows[0]) if rows else set()
mapped_cols = {
m.get("col")
for m in (config.symbol_map or {}, config.code_map or {})
if m.get("type") == "mapped" and m.get("col")
}
if rows and not ({"symbol", "code"} & row_keys or mapped_cols & row_keys):
raise ValueError("数据行中缺少 symbol/code 字段,请配置字段映射或标的映射")
# 写入
snap = date.today()
n = rows_to_parquet(rows, config, data_dir, snapshot_date=snap)
return n, snap.isoformat()
# ---------------------------------------------------------------------------
# 调度器
# ---------------------------------------------------------------------------
class PullScheduler:
"""后台调度器:为每个启用了 pull 的 ExtConfig 维护定时任务。
线程安全说明:
refresh()/stop() 可能从主事件循环 (lifespan startup) 或同步路由的
worker 线程 (configure_pull 是 def 而非 async def, FastAPI 丢进线程池)
调用。worker 线程里没有 running loop, 直接 asyncio.create_task 会抛
"no running event loop"。因此对 task 的增删一律通过
call_soon_threadsafe 提交到主循环执行 —— 同一套代码两种调用场景都安全。
"""
def __init__(self) -> None:
self._tasks: dict[str, asyncio.Task] = {}
self._running = False
self._loop: asyncio.AbstractEventLoop | None = None
def start(self, data_dir) -> None:
"""启动调度(在 lifespan startup 调用,主事件循环内)。"""
self._running = True
self._data_dir = data_dir
try:
self._loop = asyncio.get_running_loop()
except RuntimeError:
self._loop = None
logger.info("PullScheduler started")
def _submit(self, fn, *args) -> None:
"""把一个 callable 提交到主事件循环执行 (线程安全)。
startup 在主循环内调用时 fn 立即排队; worker 线程调用时跨线程排队。
两者都通过 call_soon_threadsafe, 保证 _tasks 字典的读写只在主循环里发生。
"""
loop = self._loop
if loop is None or loop.is_closed():
raise RuntimeError(
"PullScheduler: 事件循环不可用 (start() 未在事件循环中调用?)"
)
loop.call_soon_threadsafe(fn, *args)
def stop(self) -> None:
"""停止所有任务 (从 shutdown 调用)。"""
self._running = False
for task in self._tasks.values():
task.cancel()
self._tasks.clear()
logger.info("PullScheduler stopped")
def refresh(self, data_dir) -> None:
"""重新加载配置,更新调度任务(增/删/改)。线程安全。"""
self._data_dir = data_dir
store = ExtConfigStore(data_dir)
configs = store.load_all()
active_ids: set[str] = set()
new_configs: list[ExtConfig] = []
for config in configs:
if not config.pull or not config.pull.enabled or not config.pull.url:
continue
active_ids.add(config.id)
if config.id not in self._tasks:
new_configs.append(config)
# 需要移除的 id (快照当前 task 字典的键, 避免遍历时改字典)
remove_ids = [cid for cid in list(self._tasks) if cid not in active_ids]
# 所有对 _tasks 的修改都提交到主循环里执行, 保证线程安全
def _apply() -> None:
for config in new_configs:
if config.id not in self._tasks: # 二次校验, 防重复
self._tasks[config.id] = self._loop.create_task(
self._run_loop(config)
)
logger.info(
"PullScheduler: scheduled %s (every %d min)",
config.id, config.pull.schedule_minutes,
)
for cid in remove_ids:
task = self._tasks.pop(cid, None)
if task is not None:
task.cancel()
logger.info("PullScheduler: removed %s", cid)
self._submit(_apply)
async def _run_loop(self, config: ExtConfig) -> None:
"""单个配置的定时拉取循环。
策略: 启用后立即执行一次, 之后按 interval 循环。
每次循环重读最新配置 (fresh), interval 取自 fresh.pull.schedule_minutes,
这样用户中途修改间隔也能立即生效 (无需重启)。
"""
try:
while self._running:
# 每轮重读最新配置 — 用户可能修改了 url / interval / enabled
store = ExtConfigStore(self._data_dir)
fresh = store.get(config.id)
if not fresh or not fresh.pull or not fresh.pull.enabled:
break
pull = fresh.pull
# 先执行一次 (启用即拉取, 让用户立刻看到生效)
try:
n, d = await fetch_and_ingest(fresh, self._data_dir)
fresh.pull.last_run = datetime.now(timezone.utc).isoformat()
fresh.pull.last_status = "success"
fresh.pull.last_message = f"{n} rows @ {d}"
fresh.pull.last_rows = n
store.upsert(fresh)
logger.info("PullScheduler: %s success, %d rows", config.id, n)
except Exception as e:
fresh2 = store.get(config.id)
if fresh2 and fresh2.pull:
fresh2.pull.last_run = datetime.now(timezone.utc).isoformat()
fresh2.pull.last_status = "error"
fresh2.pull.last_message = str(e)[:200]
store.upsert(fresh2)
logger.warning("PullScheduler: %s error: %s", config.id, e)
# 间隔取自最新配置 (每次重新读取, 修复改间隔不生效)
interval = max(pull.schedule_minutes * 60, 60) # 至少 60s
# 预告下次运行时间, 供前端展示
next_dt = datetime.now(timezone.utc).timestamp() + interval
latest = store.get(config.id)
if latest and latest.pull:
latest.pull.next_run = datetime.fromtimestamp(
next_dt, tz=timezone.utc
).isoformat()
store.upsert(latest)
await asyncio.sleep(interval)
if not self._running:
break
except asyncio.CancelledError:
pass
# 全局单例
pull_scheduler = PullScheduler()