mirror of
https://ghfast.top/https://github.com/aeroxw/tick-stock-panel.git
synced 2026-09-12 23:44:16 +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>
97 lines
3.7 KiB
Python
97 lines
3.7 KiB
Python
"""TickFlow capability rate-limit helpers.
|
|
|
|
This module centralizes the small pieces of batch/rpm resolution used by
|
|
TickFlow-backed services. It intentionally does not manage custom data sources.
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import threading
|
|
import time
|
|
from dataclasses import dataclass
|
|
from typing import TypeVar
|
|
|
|
from app.tickflow.capabilities import Cap, CapabilitySet
|
|
|
|
T = TypeVar("T")
|
|
|
|
# 进程级共享限速器: 原先每个调用方各自本地 sleep(60/rpm), 并发同步 (kline/index/
|
|
# depth/watchlist/custom) 时聚合请求速率会成倍超过单能力 rpm → 429。
|
|
# 这里用一张按 rpm 分桶的「下一个可用时刻」表 (Lock 守护), 所有调用方按同一时间轴
|
|
# 排队, 使跨调用方的聚合发包间隔 >= 60/rpm。以 rpm 为键 (调用方签名只带 rpm, 不带 cap;
|
|
# rpm 是各能力速率的代理); 恰好同 rpm 的不同能力会共享一队, 偏保守但绝不超速。
|
|
_slot_lock = threading.Lock()
|
|
_next_slot: dict[int, float] = {}
|
|
|
|
|
|
def _reserve_slot(rpm: int, interval: float) -> float:
|
|
"""在共享时间轴上为一次请求预约一个发包槽, 返回需等待的秒数 (>=0)。
|
|
|
|
interval = 60/rpm。now 早于该 rpm 桶的 next_slot 时排到 next_slot, 否则排到 now;
|
|
随后把该桶 next_slot 后移 interval。持锁仅做时间账目, 不在锁内 sleep。
|
|
"""
|
|
key = rpm if rpm and rpm > 0 else -1
|
|
with _slot_lock:
|
|
now = time.monotonic()
|
|
scheduled = max(now, _next_slot.get(key, now))
|
|
_next_slot[key] = scheduled + interval
|
|
return scheduled - now
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class ResolvedLimit:
|
|
batch: int | None
|
|
rpm: int | None
|
|
|
|
|
|
def resolve_limit(
|
|
capset: CapabilitySet,
|
|
cap: Cap,
|
|
*,
|
|
default_batch: int | None = None,
|
|
default_rpm: int | None = None,
|
|
default_rpm_when_unset: bool = True,
|
|
) -> ResolvedLimit:
|
|
"""Return a capability's batch/rpm with caller-provided fallbacks."""
|
|
lim = capset.limits(cap)
|
|
if lim is None:
|
|
return ResolvedLimit(batch=default_batch, rpm=default_rpm)
|
|
return ResolvedLimit(
|
|
batch=lim.batch if lim.batch else default_batch,
|
|
rpm=lim.rpm if lim.rpm else (default_rpm if default_rpm_when_unset else None),
|
|
)
|
|
|
|
|
|
def batch_interval(rpm: int | None, *, default: float = 0.0) -> float:
|
|
"""Return the existing uniform batch interval formula: 60 / rpm."""
|
|
return 60.0 / rpm if rpm and rpm > 0 else default
|
|
|
|
|
|
def chunked(items: list[T], batch_size: int | None) -> list[list[T]]:
|
|
"""Split items by batch size, preserving the existing None-as-one-batch behavior."""
|
|
if batch_size is None:
|
|
return [items]
|
|
return [items[i:i + batch_size] for i in range(0, len(items), batch_size)]
|
|
|
|
|
|
def sleep_between_batches(index: int, rpm: int | None, *, default_interval: float = 0.0) -> None:
|
|
"""Sleep before every batch after the first, using the existing interval formula.
|
|
|
|
内部改用进程级共享限速器 (_reserve_slot): 保持「首批不 sleep, 后续每批间隔 60/rpm」
|
|
的单调用方观感, 同时让并发调用方按同一时间轴排队, 聚合速率不再超过单能力 rpm。
|
|
"""
|
|
interval = batch_interval(rpm, default=default_interval)
|
|
if interval <= 0:
|
|
return
|
|
if index <= 0:
|
|
# 首批不 sleep, 但登记一个占位槽, 让后续/并发调用方在同一时间轴上排队
|
|
_reserve_slot(rpm or -1, interval)
|
|
return
|
|
wait = _reserve_slot(rpm or -1, interval)
|
|
if wait > 0:
|
|
time.sleep(wait)
|
|
|
|
|
|
def min_batch(preferred: int, limit: ResolvedLimit) -> int:
|
|
"""Clamp a user-preferred batch size by a resolved capability batch limit."""
|
|
return min(preferred, limit.batch) if limit.batch else preferred
|