#!/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()