feat: 修复daily历史覆盖度缺口+增加覆盖度校验
- 新增 repair_daily_backfill.py:按daily_basic基准检测并补拉覆盖不足交易日 (解决2012-2013沪市+创业板整年缺失问题,断点续传日期级检测无法发现) - importer.py: 新增 check_daily_coverage() 覆盖度审计函数 - incremental_import.py: 增量导入后自动跑覆盖度校验并告警 - 数据批量导入.ipynb: 增加数据完整性检验单元格
This commit is contained in:
@@ -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,
|
||||
|
||||
@@ -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("每日增量导入完成")
|
||||
|
||||
|
||||
@@ -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()
|
||||
@@ -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": {
|
||||
|
||||
Reference in New Issue
Block a user