195 lines
7.2 KiB
Python
195 lines
7.2 KiB
Python
#!/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")
|
||
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"]),
|
||
"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:
|
||
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 == "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()
|