Files
quanxiel/quantitative_data/incremental_import.py
T

199 lines
7.3 KiB
Python
Raw 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,
)
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}")
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()