mirror of
https://ghfast.top/https://github.com/aeroxw/tick-stock-panel.git
synced 2026-09-12 19:04:15 +08:00
* 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>
294 lines
11 KiB
Python
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()
|