diff --git a/quantitative_data/incremental_import.py b/quantitative_data/incremental_import.py new file mode 100644 index 0000000..fa5f5f3 --- /dev/null +++ b/quantitative_data/incremental_import.py @@ -0,0 +1,188 @@ +#!/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 幂等) + - moneyflow 表暂未纳入(新版 importer 无对应导入函数,后续需要再补) +""" +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_adj_factor_batch, + import_index_daily, + import_financial_statements, +) + +logger = logging.getLogger("incremental") +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"]), + "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"), + ("adj_factor", "trade_date"), + ("index_daily", "trade_date"), + ]: + max_d = get_max_date(conn, table, date_col) + if max_d is None: + 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 == "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}") + + 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()