diff --git a/quantitative_data/importer.py b/quantitative_data/importer.py index b087ff7..908b091 100644 --- a/quantitative_data/importer.py +++ b/quantitative_data/importer.py @@ -1566,6 +1566,70 @@ def get_missing_daily_dates( conn.close() +def check_daily_coverage( + start_date: Optional[str] = None, + end_date: Optional[str] = None, + conn=None, + tolerance: float = 0.05, +) -> List[Dict]: + """ + 检查 daily 表每日覆盖度(以 daily_basic 为基准,单条 SQL 聚合)。 + + 背景:resume_daily_by_date 只做「日期级」缺失检测(NOT EXISTS trade_date), + 若某交易日只有部分股票入库(如 2012-2013 沪市+创业板缺失),会误判为"已完整", + 导致覆盖度缺口永不补拉。本函数做「覆盖度级」校验: + + daily 当日股票数 < daily_basic 当日股票数 × (1 - tolerance) → 判定覆盖不足 + + 返回: [{"trade_date", "daily_cnt", "daily_basic_cnt", "coverage_pct"}, ...] + """ + if start_date is None: + start_date = START_DATE + if end_date is None: + end_date = END_DATE + + own_conn = conn is None + if own_conn: + conn = get_pg_connection() + + try: + cursor = conn.cursor() + cursor.execute( + """ + SELECT db.trade_date, + COALESCE(d.cnt, 0) AS daily_cnt, + db.cnt AS db_cnt, + ROUND(COALESCE(d.cnt, 0)::numeric / db.cnt * 100, 1) AS coverage_pct + FROM (SELECT trade_date, COUNT(DISTINCT ts_code) AS cnt + FROM daily_basic + WHERE trade_date BETWEEN %s AND %s + GROUP BY trade_date) db + LEFT JOIN (SELECT trade_date, COUNT(DISTINCT ts_code) AS cnt + FROM daily + WHERE trade_date BETWEEN %s AND %s + GROUP BY trade_date) d + ON d.trade_date = db.trade_date + WHERE COALESCE(d.cnt, 0) < db.cnt * (1 - %s) + ORDER BY db.trade_date + """, + (start_date, end_date, start_date, end_date, tolerance), + ) + rows = [ + { + "trade_date": r[0].strftime("%Y-%m-%d") if hasattr(r[0], "strftime") else str(r[0])[:10], + "daily_cnt": r[1], + "daily_basic_cnt": r[2], + "coverage_pct": float(r[3]), + } + for r in cursor.fetchall() + ] + cursor.close() + return rows + finally: + if own_conn: + conn.close() + + def resume_daily_by_date( start_date: Optional[str] = None, end_date: Optional[str] = None, diff --git a/quantitative_data/incremental_import.py b/quantitative_data/incremental_import.py index 816763b..876a5b7 100644 --- a/quantitative_data/incremental_import.py +++ b/quantitative_data/incremental_import.py @@ -35,6 +35,7 @@ from importer import ( import_adj_factor_batch, import_index_daily, import_financial_statements, + check_daily_coverage, ) logger = logging.getLogger("incremental") @@ -140,6 +141,27 @@ def run_daily(end_date, dry_run, limit): except Exception as e: logger.error(f" {table} 增量导入失败: {e}") + # 5. 覆盖度校验(以 daily_basic 为基准,检查 daily 是否缺部分股票) + # 防止"日期存在但覆盖不全"的缺口(如 2012-2013 沪市+创业板缺失)被增量逻辑跳过 + try: + coverage_start = "2010-01-01" # 全历史检查(单条 SQL 聚合,开销小) + logger.info(f"覆盖度校验: daily vs daily_basic ({coverage_start} ~ {end_date})") + if dry_run: + logger.info("[dry-run] 跳过覆盖度校验") + else: + partial = check_daily_coverage(coverage_start, end_date, conn=conn, tolerance=0.05) + if partial: + logger.warning(f"⚠ 发现 {len(partial)} 个交易日覆盖不足 (daily < daily_basic×95%):") + for p in partial[:10]: + logger.warning(f" {p['trade_date']}: daily={p['daily_cnt']} vs daily_basic={p['daily_basic_cnt']} ({p['coverage_pct']}%)") + if len(partial) > 10: + logger.warning(f" ... 其余 {len(partial)-10} 个交易日略") + logger.warning(" 请运行 repair_daily_backfill.py 修复历史缺口,或检查近期导入是否被中断") + else: + logger.info(" ✓ 覆盖度正常,无缺失交易日") + except Exception as e: + logger.error(f" 覆盖度校验失败: {e}") + conn.close() logger.info("每日增量导入完成") diff --git a/quantitative_data/repair_daily_backfill.py b/quantitative_data/repair_daily_backfill.py new file mode 100644 index 0000000..5d077ac --- /dev/null +++ b/quantitative_data/repair_daily_backfill.py @@ -0,0 +1,178 @@ +#!/usr/bin/env python3 +# -*- coding: utf-8 -*- +""" +修复 daily 表历史覆盖度缺口(2012-2013 沪市+创业板缺失) + +背景:daily 表 2012-2013 年只有深市主板/中小板(2012 年 1383 只、2013 年 677 只), + 缺全部沪市 + 创业板。原因:当年全量导入沪市请求失败/中断,但深市成功, + 导致断点续传的"日期级"缺失检测(NOT EXISTS trade_date)认为该日已有数据, + 永不补拉。 + +修复策略(覆盖度级校验): +1. 对指定日期范围,对比 daily 与 daily_basic 的当日去重股票数 + (daily_basic 同期数据完整,作为覆盖度基准) +2. daily 当日股票数 < daily_basic 当日股票数 × (1 - tolerance) 的日期 → 判定为"部分缺失" +3. 部分缺失的日期按全市场重新拉取(Tushare daily 按 trade_date 返回全市场) +4. 用 INSERT ON CONFLICT DO NOTHING 幂等写入,可重复执行 + +用法: + python3 repair_daily_backfill.py [--start 2012-01-01] [--end 2013-12-31] [--tolerance 0.05] [--dry-run] +""" +import os +import sys +import time +import argparse +import logging +from datetime import datetime + +HERE = os.path.dirname(os.path.abspath(__file__)) +sys.path.insert(0, HERE) + +# 显式加载 .env(config.py 会加载,但确保顺序正确) +try: + from dotenv import load_dotenv + env_path = os.path.join(HERE, ".env") + if os.path.exists(env_path): + load_dotenv(env_path, override=False) +except ImportError: + pass + +import importer +from importer import ( + get_pg_connection, get_ts_pro, batch_insert, + fetch_with_retry, _normalize_daily_df, +) + +logging.basicConfig( + level=logging.INFO, + format="%(asctime)s [%(levelname)s] %(message)s", + handlers=[ + logging.FileHandler(os.path.join(HERE, "repair_daily_backfill.log"), encoding="utf-8"), + logging.StreamHandler(), + ], +) +logger = logging.getLogger("repair_daily_backfill") + + +def get_coverage_ratio(conn, trade_date: str) -> tuple: + """返回 (daily_stocks, daily_basic_stocks)。无参考数据时 daily_basic_stocks=None""" + cur = conn.cursor() + try: + cur.execute("SELECT COUNT(DISTINCT ts_code) FROM daily WHERE trade_date=%s", (trade_date,)) + d_cnt = cur.fetchone()[0] + cur.execute("SELECT COUNT(DISTINCT ts_code) FROM daily_basic WHERE trade_date=%s", (trade_date,)) + db_cnt = cur.fetchone()[0] + return d_cnt, db_cnt + finally: + cur.close() + + +def find_partial_dates(conn, start_date: str, end_date: str, tolerance: float) -> list: + """ + 找出 daily 覆盖度不足的交易日(单条 SQL 聚合,避免逐日查询)。 + 返回 [(trade_date, daily_cnt, daily_basic_cnt), ...] + 判据:daily_basic 有数据且 daily 股票数 < daily_basic × (1 - tolerance) + """ + cur = conn.cursor() + try: + cur.execute( + """ + SELECT db.trade_date, + COALESCE(d.cnt, 0) AS daily_cnt, + db.cnt AS db_cnt + FROM (SELECT trade_date, COUNT(DISTINCT ts_code) AS cnt + FROM daily_basic + WHERE trade_date BETWEEN %s AND %s + GROUP BY trade_date) db + LEFT JOIN (SELECT trade_date, COUNT(DISTINCT ts_code) AS cnt + FROM daily + WHERE trade_date BETWEEN %s AND %s + GROUP BY trade_date) d + ON d.trade_date = db.trade_date + WHERE COALESCE(d.cnt, 0) < db.cnt * (1 - %s) + ORDER BY db.trade_date + """, + (start_date, end_date, start_date, end_date, tolerance), + ) + partial = [(r[0].strftime("%Y-%m-%d"), r[1], r[2]) for r in cur.fetchall()] + finally: + cur.close() + return partial + + +def backfill_date(pro, conn, td_str: str) -> bool: + """拉取单个交易日全市场 daily 数据并写入。成功返回 True""" + td_compact = td_str.replace("-", "") + + def fetch(): + return pro.daily(trade_date=td_compact) + + df = fetch_with_retry(fetch, max_retries=4, delay=3) + if df is None or df.empty: + logger.warning(f" [{td_str}] 返回空数据,跳过") + return False + + df = _normalize_daily_df(df) + # 幂等写入:已存在的行跳过(ON CONFLICT DO NOTHING) + batch_insert("daily", df, conn, ["ts_code", "trade_date"]) + return True + + +def main(): + parser = argparse.ArgumentParser(description="修复 daily 表历史覆盖度缺口") + parser.add_argument("--start", default="2012-01-01", help="起始日期 YYYY-MM-DD") + parser.add_argument("--end", default="2013-12-31", help="结束日期 YYYY-MM-DD") + parser.add_argument("--tolerance", type=float, default=0.05, + help="覆盖度容差,默认 0.05(daily 少于 daily_basic 的 95% 即判定缺失)") + parser.add_argument("--dry-run", action="store_true", help="只扫描不导入") + args = parser.parse_args() + + conn = get_pg_connection() + logger.info("=" * 60) + logger.info(f"[覆盖度扫描] {args.start} ~ {args.end} (容差 {args.tolerance:.0%})") + + partial = find_partial_dates(conn, args.start, args.end, args.tolerance) + if not partial: + logger.info(" 未发现覆盖度不足的交易日 ✅") + conn.close() + return + + logger.info(f" 发现 {len(partial)} 个覆盖度不足的交易日:") + # 按年份统计 + years = {} + for td, d_cnt, db_cnt in partial: + y = td[:4] + years.setdefault(y, []).append((td, d_cnt, db_cnt)) + for y in sorted(years): + lst = years[y] + logger.info(f" {y} 年: {len(lst)} 个交易日 | 样例 {lst[0][0]}(daily={lst[0][1]}/db={lst[0][2]})") + + if args.dry_run: + logger.info("[dry-run] 不执行导入,以上为待补拉清单") + conn.close() + return + + pro = get_ts_pro() + success = 0 + fail_list = [] + for i, (td, d_cnt, db_cnt) in enumerate(partial, 1): + ok = backfill_date(pro, conn, td) + if ok: + success += 1 + else: + fail_list.append(td) + if i % 20 == 0 or i == len(partial): + logger.info(f" 进度: {i}/{len(partial)} 成功={success} 失败={len(fail_list)}") + time.sleep(0.35) # 限流保护 + + conn.close() + logger.info(f" 补拉完成: 成功 {success}/{len(partial)}") + if fail_list: + logger.warning(f" 失败日期({len(fail_list)}): {fail_list[:20]}...") + # 写失败清单供重试 + with open(os.path.join(HERE, "repair_failed_dates.txt"), "w") as f: + f.write("\n".join(fail_list)) + + +if __name__ == "__main__": + main() diff --git a/quantitative_data/数据批量导入.ipynb b/quantitative_data/数据批量导入.ipynb index 34a381d..e509f2c 100644 --- a/quantitative_data/数据批量导入.ipynb +++ b/quantitative_data/数据批量导入.ipynb @@ -720,6 +720,68 @@ "display(df3[['ts_code', 'name', 'industry', 'roe', 'roa', 'eps', 'debt_to_assets']])\n", "conn.close()" ] + }, + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "---\n", + "## 数据完整性检验(覆盖度审计)\n", + "\n", + "> **背景**:2026-08-26 审计发现 daily 表 2012-2013 年存在覆盖度缺口(沪市+创业板整年缺失),\n", + "> 而断点续传的\"日期级\"检测(`NOT EXISTS trade_date`)无法发现这种\"日期存在但覆盖不全\"的问题。\n", + "> 本单元格做**覆盖度级**校验:以 `daily_basic` 当日股票数为基准,检查 `daily` 是否缺股。\n", + ">\n", + "> 运行 `repair_daily_backfill.py` 可自动补拉覆盖不足的交易日。\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# ============================================================\n", + "# 数据完整性检验:daily vs daily_basic 每日覆盖度审计\n", + "# 判据:daily 当日股票数 < daily_basic 当日股票数 × 95% → 覆盖不足\n", + "# ============================================================\n", + "import importlib\n", + "import importer\n", + "importlib.reload(importer)\n", + "from importer import get_pg_connection, check_daily_coverage, check_table_summary\n", + "\n", + "conn = get_pg_connection()\n", + "\n", + "# 1) 各表整体概览(行数 / 股票数 / 日期范围)\n", + "print('=' * 60)\n", + "print('各表整体概览:')\n", + "print('=' * 60)\n", + "check_table_summary(conn=conn)\n", + "\n", + "# 2) 全历史覆盖度审计(默认容差 5%)\n", + "print('\\n' + '=' * 60)\n", + "print('覆盖度审计: daily vs daily_basic (容差 5%)')\n", + "print('=' * 60)\n", + "from datetime import date\n", + "end_date = date.today().strftime('%Y-%m-%d')\n", + "partial = check_daily_coverage(start_date='2010-01-01', end_date=end_date, conn=conn, tolerance=0.05)\n", + "\n", + "if partial:\n", + " print(f'⚠ 发现 {len(partial)} 个交易日覆盖不足:')\n", + " # 按年份汇总\n", + " from collections import Counter\n", + " years = Counter(p['trade_date'][:4] for p in partial)\n", + " for y in sorted(years):\n", + " print(f' {y} 年: {years[y]} 个交易日覆盖不足')\n", + " print('\\n样例(前10条):')\n", + " for p in partial[:10]:\n", + " print(f\" {p['trade_date']}: daily={p['daily_cnt']} vs daily_basic={p['daily_basic_cnt']} ({p['coverage_pct']}%)\")\n", + " print('\\n→ 修复方法: 运行 python3 repair_daily_backfill.py 补拉')\n", + "else:\n", + " print('✅ 覆盖度正常,无缺失交易日')\n", + "\n", + "conn.close()\n" + ] } ], "metadata": {