diff --git a/backend/services/daily_collector.py b/backend/services/daily_collector.py index 6aed654..88e9e83 100644 --- a/backend/services/daily_collector.py +++ b/backend/services/daily_collector.py @@ -18,11 +18,15 @@ _CST = timezone(timedelta(hours=8)) CHECK_INTERVAL_SECONDS = 300 COLLECT_AFTER_TIME = dtime(15, 1) # 收盘后 15:01 开始允许采集(留 1 分钟等收盘数据稳定) CORE_STOCK_LIMIT = 100 # 题材领涨股涨幅前100 -TOP_THEME_LIMIT = 20 # 题材涨幅前20 +RANK_TOP_LIMIT = 50 # 每个榜单(涨幅/强度/热度/成交额)取前50 HOTMAP_TOP_N = 100 # 热点穿透采样题材数:涨幅榜+热度榜各取前100,合并(与热点穿透页一致) HOTMAP_CORE_LIMIT = 100 # 热点穿透核心股前100(按覆盖题材数降序) CORE_COVER_THRESHOLD = 2 # 热点穿透核心股门槛:覆盖题材数 ≥2 +# 题材榜单:sortField 1=涨幅(bf3) 3=强度(strengthValue) 4=热度排名(hotRank) 5=成交额(fex5) +# 每个榜单取前 RANK_TOP_LIMIT 个,按 themeCode 合并去重后入库 +THEME_RANK_LISTS = ((1, "bf3"), (3, "strengthValue"), (4, "hotRank"), (5, "fex5")) + # 每日缓存清理:交易日 9:31 清空全部缓存,保证开盘后数据全新 CLEANUP_TIME = dtime(9, 31) @@ -49,6 +53,7 @@ async def collect_daily(trade_date: str, dry_run: bool = False) -> dict: """采集指定交易日数据并入库。 核心股 = 题材领涨股涨幅前100 ∪ 热点穿透核心股(覆盖题材数≥2、按覆盖数降序前100),按股票代码去重。 + 题材 = 涨幅/强度/热度/成交额 4 榜各前50,按 themeCode 合并去重。 热点穿透构建失败时降级为仅领涨股前100,不影响当日采集。 Args: @@ -107,8 +112,15 @@ async def collect_daily(trade_date: str, dry_run: bool = False) -> dict: } f3_core = sorted(stock_map.values(), key=lambda x: -(x["f3"] or 0))[:CORE_STOCK_LIMIT] - # 4. 题材涨幅前20:bf3 降序 - top_themes = sorted(themes, key=lambda x: -(x.get("bf3") or 0))[:TOP_THEME_LIMIT] + # 4. 题材 4 榜(涨幅/强度/热度/成交额)各前50,按 themeCode 合并去重。 + # 涨幅榜列表与开头用于构建领涨股映射的 fetch_theme_list(1, False) 同缓存,直接复用。 + merged: dict[str, dict] = {} + rank_list = await fetch_theme_list(1, False) + for sort_field, _ in THEME_RANK_LISTS: + lst = rank_list if sort_field == 1 else await fetch_theme_list(sort_field, False) + for t in lst[:RANK_TOP_LIMIT]: + merged.setdefault(t["themeCode"], t) + top_themes = list(merged.values()) # 5. 合并去重:热点穿透核心股在前(覆盖数降序,rank 优先),随后补领涨股涨幅前100 core_stocks: list[dict] = [] @@ -145,7 +157,7 @@ async def collect_daily(trade_date: str, dry_run: bool = False) -> dict: theme_pairs.add((code, tc, tn)) if dry_run: - print(f"[collector] {trade_date} 核心股 {len(core_stocks)} 只(热点穿透 {len(hotmap_core)} + 领涨股 {len(f3_core)}),题材前20 {len(top_themes)} 只") + print(f"[collector] {trade_date} 核心股 {len(core_stocks)} 只(热点穿透 {len(hotmap_core)} + 领涨股 {len(f3_core)}),题材 {len(top_themes)} 只(4 榜前 {RANK_TOP_LIMIT} 合并)") return {"core_count": len(core_stocks), "theme_count": len(top_themes), "skipped": False} # 6. 入库(事务,UNIQUE 幂等) @@ -170,7 +182,7 @@ async def collect_daily(trade_date: str, dry_run: bool = False) -> dict: finally: conn.close() - print(f"[collector] {trade_date} 已采集:核心股 {len(core_stocks)} 只,题材前20 {len(top_themes)} 只") + print(f"[collector] {trade_date} 已采集:核心股 {len(core_stocks)} 只,题材 {len(top_themes)} 只(4 榜前 {RANK_TOP_LIMIT} 合并)") return {"core_count": len(core_stocks), "theme_count": len(top_themes), "skipped": False}