Files
xiaoqiang 84912ef99c feat: 修复daily历史覆盖度缺口+增加覆盖度校验
- 新增 repair_daily_backfill.py:按daily_basic基准检测并补拉覆盖不足交易日
  (解决2012-2013沪市+创业板整年缺失问题,断点续传日期级检测无法发现)
- importer.py: 新增 check_daily_coverage() 覆盖度审计函数
- incremental_import.py: 增量导入后自动跑覆盖度校验并告警
- 数据批量导入.ipynb: 增加数据完整性检验单元格
2026-08-26 10:47:32 +00:00

221 lines
8.7 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""
量化数据增量导入脚本 (incremental_import.py)
=============================================
查询各表最大日期,只拉「最后日期 → 昨天」的数据,避免每天全量重导。
⚠️ 依赖仓库新版 importer.py(daily_vip / *_vip 接口,Tushare 需 5000+ 积分)
用法:
/opt/quant-venv/bin/python incremental_import.py # 每日增量(行情类)
/opt/quant-venv/bin/python incremental_import.py --weekly # 每日增量 + 财务表(建议每周一次)
/opt/quant-venv/bin/python incremental_import.py --dry-run # 只打印计划,不实际导入
/opt/quant-venv/bin/python incremental_import.py --end-date 2026-08-14 # 指定结束日期(默认昨天)
/opt/quant-venv/bin/python incremental_import.py --limit 100 # 只处理前 N 只股票(测试用)
说明:
- 行情类表按 MAX(trade_date)+1 天 → 昨天 增量拉取
- daily / daily_basic 走按交易日全市场模式(快);adj_factor 走按股票批量
- 财务表按 MAX(end_date) 往前推 400 天 → 昨天(覆盖新公告的季度报告,UPSERT 幂等)
"""
import argparse
import logging
import sys
from datetime import datetime, timedelta
from importer import (
get_pg_connection,
get_stock_codes_from_db,
import_trade_cal,
import_stock_basic,
import_daily_by_date,
import_daily_basic_by_date,
import_moneyflow_by_date,
import_adj_factor_batch,
import_index_daily,
import_financial_statements,
check_daily_coverage,
)
logger = logging.getLogger("incremental")
logger.propagate = False
if not logger.handlers:
h = logging.StreamHandler()
h.setFormatter(logging.Formatter("%(asctime)s [%(levelname)s] %(message)s"))
logger.addHandler(h)
logger.setLevel(logging.INFO)
# 各表 -> (日期列, 唯一约束列)
TABLE_SPECS = {
"daily": ("trade_date", ["ts_code", "trade_date"]),
"daily_basic": ("trade_date", ["ts_code", "trade_date"]),
"moneyflow": ("trade_date", ["ts_code", "trade_date"]),
"adj_factor": ("trade_date", ["ts_code", "trade_date"]),
"index_daily": ("trade_date", ["ts_code", "trade_date"]),
"income": ("end_date", ["ts_code", "end_date", "report_type"]),
"balancesheet": ("end_date", ["ts_code", "end_date", "report_type"]),
"cashflow": ("end_date", ["ts_code", "end_date", "report_type"]),
"fina_indicator": ("end_date", ["ts_code", "end_date"]),
"trade_cal": ("cal_date", ["exchange", "cal_date"]),
}
def get_max_date(conn, table, date_col):
"""查询表内最大日期(无数据返回 None)"""
cur = conn.cursor()
cur.execute(f'SELECT MAX({date_col}) FROM "{table}"')
row = cur.fetchone()
cur.close()
return row[0] if row and row[0] else None
def run_daily(end_date, dry_run, limit):
logger.info("=" * 60)
logger.info(f"每日增量导入 (截止 {end_date})")
conn = get_pg_connection()
# 1. 交易日历增量
max_cal = get_max_date(conn, "trade_cal", "cal_date")
if max_cal:
cal_start = (max_cal + timedelta(days=1)).strftime("%Y-%m-%d")
if cal_start <= end_date:
if dry_run:
logger.info(f"[dry-run] trade_cal: {cal_start} ~ {end_date}")
else:
import_trade_cal(cal_start, end_date)
else:
logger.info(f"trade_cal: 已最新 ({max_cal})")
else:
logger.warning("trade_cal 为空,先跑全量导入初始化")
conn.close()
return
# 2. 股票基本信息刷新(全量 UPSERT,量小)
if dry_run:
logger.info("[dry-run] stock_basic: 全量刷新")
else:
import_stock_basic()
# 3. 股票列表
codes = get_stock_codes_from_db(conn)
if limit:
codes = codes[:limit]
logger.info(f"股票数: {len(codes)}")
# 4. 行情类各表
for table, date_col in [
("daily", "trade_date"),
("daily_basic", "trade_date"),
("moneyflow", "trade_date"),
("adj_factor", "trade_date"),
("index_daily", "trade_date"),
]:
max_d = get_max_date(conn, table, date_col)
if max_d is None:
if table == "index_daily":
logger.warning(" index_daily: 表为空,自动全量初始化(6个指数,很快)")
if not dry_run:
import_index_daily(start_date=None, end_date=end_date)
else:
logger.warning(f" {table}: 表为空,请先执行全量导入 (full_import)")
continue
start = (max_d + timedelta(days=1)).strftime("%Y-%m-%d")
if start > end_date:
logger.info(f" {table}: 已最新 (max={max_d.strftime('%Y-%m-%d')})")
continue
logger.info(f" {table}: {start} ~ {end_date}")
if dry_run:
continue
try:
if table == "daily":
import_daily_by_date(start, end_date)
elif table == "daily_basic":
import_daily_basic_by_date(start, end_date)
elif table == "moneyflow":
import_moneyflow_by_date(start, end_date)
elif table == "adj_factor":
import_adj_factor_batch(codes, start, end_date)
elif table == "index_daily":
import_index_daily(start_date=start, end_date=end_date)
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("每日增量导入完成")
def run_financial(end_date, dry_run, limit):
logger.info("=" * 60)
logger.info(f"财务数据增量导入 (截止 {end_date})")
conn = get_pg_connection()
codes = get_stock_codes_from_db(conn)
if limit:
codes = codes[:limit]
if not codes:
logger.warning("stock_basic 为空,无法导入财务数据")
conn.close()
return
max_dates = {}
for t in ["income", "balancesheet", "cashflow", "fina_indicator"]:
md = get_max_date(conn, t, "end_date")
max_dates[t] = md
logger.info(f" {t}: max(end_date)={md}")
valid = [d for d in max_dates.values() if d]
if not valid:
logger.warning("财务表均为空,请先执行全量导入")
conn.close()
return
base = min(valid)
start = (base - timedelta(days=400)).strftime("%Y-%m-%d") # 往前推 400 天覆盖新公告
logger.info(f" 财务拉取范围: {start} ~ {end_date} (基准 max_end_date={base.strftime('%Y-%m-%d')})")
conn.close()
if dry_run:
return
import_financial_statements(codes, start, end_date)
logger.info("财务数据增量导入完成")
def main():
parser = argparse.ArgumentParser(description="量化数据增量导入")
parser.add_argument("--weekly", action="store_true", help="同时导入财务表(建议每周一次)")
parser.add_argument("--dry-run", action="store_true", help="只打印计划,不实际导入")
parser.add_argument("--end-date", default=None, help="结束日期 YYYY-MM-DD(默认昨天)")
parser.add_argument("--limit", type=int, default=0, help="只处理前 N 只股票(测试用)")
args = parser.parse_args()
end_date = args.end_date or (datetime.now() - timedelta(days=1)).strftime("%Y-%m-%d")
run_daily(end_date, args.dry_run, args.limit)
if args.weekly:
run_financial(end_date, args.dry_run, args.limit)
if __name__ == "__main__":
main()