"""首页打板专题(tushare 同花顺版:limit_list_ths / limit_step / limit_cpt_list)。 整包 SWR(仿 index_global 列表层):进程内 state 新鲜直返 -> Redis 回填 -> 有旧值先返 + 后台刷新 -> 冷启动同步拉。盘中(交易日 09:15-15:30)数据源即有 当日快照(实测 quicksync 镜像盘中可取当日)-> fresh TTL 压到 5 分钟;其余时段 4 小时。 镜像坑(见 reference.py 注释):不传 fields;limit_list_ths 必须显式传 trade_date (缺省返回多日混包且 4000 行封顶);涨停/连扳池才有 tag/status/lu_desc/封单额, 炸板池只有价格与打开次数,跌停池几乎只有价格——行模型统一、字段可选。 """ from __future__ import annotations import asyncio import math import time from datetime import date, datetime, timedelta from .. import cache from ..config import settings from .sync_utils import call_retry, f_clean, get_pro_lazy, s_clean _BOARD_KEY = "limit_board:daily:v1" # Redis 整包缓存键(v1 起版) _REDIS_TTL = 3600 # 进程重启后的回填来源 _INTRADAY_TTL = 300.0 # 交易时段内的 fresh TTL(5 分钟准实时) _MAX_DATE_BACKTRACK = 5 # trade_date 定位回退天数(覆盖节假日/盘前) _BLOCKS_OUT = 12 # 最强板块输出条数 class LimitBoardError(RuntimeError): """三池全部拉不到(token/网络故障)——接口层转 503。""" def _yi(v) -> float | None: """元 -> 亿元(2 位小数)。""" f = f_clean(v) return None if f is None else round(f / 1e8, 2) def _fetch_pool_sync(pro, trade_date: str, limit_type: str): time.sleep(settings.screener_sync_interval) return call_retry(pro.limit_list_ths, trade_date=trade_date, limit_type=limit_type) def _pool_rows(df, mode: str) -> list[dict]: """行裁剪 + 单位换算。mode: up / broken / down。""" if df is None or df.empty: return [] rows: list[dict] = [] for _, r in df.iterrows(): row = { "ts_code": s_clean(r.get("ts_code")), "name": s_clean(r.get("name")), "price": f_clean(r.get("price")), "pct_chg": f_clean(r.get("pct_chg")), } if not row["ts_code"]: continue if mode == "up": row.update({ "tag": s_clean(r.get("tag")), "status": s_clean(r.get("status")), "lu_desc": s_clean(r.get("lu_desc")), "open_num": f_clean(r.get("open_num")), "limit_amount_yi": _yi(r.get("limit_amount")), # 封单额(亿) "turnover_yi": _yi(r.get("turnover")), # 成交额(亿) "first_lu_time": s_clean(r.get("first_lu_time")), "limit_up_suc_rate": f_clean(r.get("limit_up_suc_rate")), }) elif mode == "broken": row.update({ "open_num": f_clean(r.get("open_num")), "first_lu_time": s_clean(r.get("first_lu_time")), "last_lu_time": s_clean(r.get("last_lu_time")), }) rows.append(row) if mode == "up": # 封单额降序(打板看封单强度);封单额缺失(镜像个别行)沉底 rows.sort(key=lambda x: (x.get("limit_amount_yi") is None, -(x.get("limit_amount_yi") or 0))) return rows def _fetch_board_sync() -> dict: """定位交易日并拉三池 + 天梯 + 最强板块(同步网络 IO,需在 to_thread 里跑)。""" pro = get_pro_lazy() errors: list[str] = [] # trade_date 定位:今日起逐日回退,取第一个涨停池非空的日期 # (盘前/节假日当日为空;tushare 错误直接抛——定位失败无意义继续) trade_date: str | None = None up_rows: list[dict] = [] for i in range(_MAX_DATE_BACKTRACK): d = (date.today() - timedelta(days=i)).strftime("%Y%m%d") df = _fetch_pool_sync(pro, d, "涨停池") if df is not None and not df.empty: trade_date = d up_rows = _pool_rows(df, "up") break if trade_date is None: raise LimitBoardError(f"近 {_MAX_DATE_BACKTRACK} 天均无涨停池数据(节假日或数据源故障)") broken_rows: list[dict] = [] try: broken_rows = _pool_rows(_fetch_pool_sync(pro, trade_date, "炸板池"), "broken") except Exception as e: # noqa: BLE001 —— 单池失败不拖垮整包 errors.append(f"炸板池: {str(e)[:60]}") down_rows: list[dict] = [] try: down_rows = _pool_rows(_fetch_pool_sync(pro, trade_date, "跌停池"), "down") except Exception as e: # noqa: BLE001 errors.append(f"跌停池: {str(e)[:60]}") ladder: list[dict] = [] try: time.sleep(settings.screener_sync_interval) step = call_retry(pro.limit_step, trade_date=trade_date) if step is not None and not step.empty: for _, r in step.iterrows(): code = s_clean(r.get("ts_code")) n = f_clean(r.get("nums")) if code and n: ladder.append({"ts_code": code, "name": s_clean(r.get("name")), "nums": int(n)}) ladder.sort(key=lambda x: -x["nums"]) except Exception as e: # noqa: BLE001 errors.append(f"连板天梯: {str(e)[:60]}") blocks: list[dict] = [] try: time.sleep(settings.screener_sync_interval) cpt = call_retry(pro.limit_cpt_list, trade_date=trade_date) if cpt is not None and not cpt.empty: for _, r in cpt.head(_BLOCKS_OUT).iterrows(): blocks.append({ "name": s_clean(r.get("name")), "days": f_clean(r.get("days")), "up_stat": s_clean(r.get("up_stat")), "cons_nums": f_clean(r.get("cons_nums")), "up_nums": f_clean(r.get("up_nums")), "pct_chg": f_clean(r.get("pct_chg")), }) except Exception as e: # noqa: BLE001 errors.append(f"最强板块: {str(e)[:60]}") # 连板分布(limit_step 只含 2 板及以上;1 板 = 涨停池 tag 首板数) dist: dict[int, int] = {} for x in ladder: dist[x["nums"]] = dist.get(x["nums"], 0) + 1 summary = { "up_count": len(up_rows), "broken_count": len(broken_rows), "down_count": len(down_rows), "first_board_count": sum(1 for x in up_rows if x.get("tag") == "首板"), "max_ladder": ladder[0] if ladder else None, "ladder_dist": [{"nums": k, "count": v} for k, v in sorted(dist.items())], } return { "trade_date": f"{trade_date[:4]}-{trade_date[4:6]}-{trade_date[6:]}", "summary": summary, "up": up_rows, "broken": broken_rows, "down": down_rows, "ladder": ladder, "blocks": blocks, "errors": errors, } # ---------- 整包 SWR(进程内 -> Redis 回填 -> 旧值先返 + 后台刷新 -> 冷启动同步拉) ---------- _state: dict = {"payload": None} _refreshing = False _refresh_error: str | None = None _bg_tasks: set[asyncio.Task] = set() def _fresh_ttl(is_trading_day: bool | None) -> float: """交易时段 5 分钟(镜像盘中即有当日快照);其余 4 小时(盘后数据不变)。""" if is_trading_day is None: is_trading_day = datetime.now().weekday() < 5 # 判定失败回退 weekday 启发式 if is_trading_day: now = datetime.now() t = now.hour * 60 + now.minute if 9 * 60 + 15 <= t <= 15 * 60 + 30: return _INTRADAY_TTL return float(settings.market_eod_fresh_ttl) async def _refresh() -> dict: data = await asyncio.to_thread(_fetch_board_sync) payload = {**data, "updated_at": datetime.now().isoformat(), "fetched_ts": time.time()} _state["payload"] = payload await cache.cache_set(_BOARD_KEY, payload, ttl=_REDIS_TTL) return payload async def _refresh_wrapped() -> None: global _refresh_error, _refreshing try: await _refresh() _refresh_error = None except Exception as e: # noqa: BLE001 —— 后台刷新失败静默记错,下次并入 errors _refresh_error = f"打板专题后台刷新: {str(e)[:60]}" finally: _refreshing = False def _spawn_refresh() -> None: global _refreshing if _refreshing: return _refreshing = True task = asyncio.create_task(_refresh_wrapped()) _bg_tasks.add(task) task.add_done_callback(_bg_tasks.discard) async def fetch_limit_board(is_trading_day: bool | None) -> dict: """打板专题整包读取(SWR)。冷启动同步拉(5 次调用约 2-4s);此后盘中 5 分钟/盘后 4 小时。""" ttl = _fresh_ttl(is_trading_day) p = _state["payload"] if p is not None and time.time() - p["fetched_ts"] < ttl: return p if p is None: cached = await cache.cache_get(_BOARD_KEY) if cached: p = cached _state["payload"] = p if p is not None: _spawn_refresh() if _refresh_error and not p.get("errors"): p = {**p, "errors": [_refresh_error]} return p return await _refresh()