From 51079bd8dca34566abbe7bc4bca6e745a63a1042 Mon Sep 17 00:00:00 2001 From: shellway-pc <413209390@qq.com> Date: Sat, 15 Aug 2026 23:02:26 +0800 Subject: [PATCH] =?UTF-8?q?=E5=A2=9E=E5=8A=A0=E4=BA=86moneyflow=E8=A1=A8?= =?UTF-8?q?=E6=A0=BC=E5=AF=BC=E5=85=A5=E6=A8=A1=E5=9D=97=E3=80=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- quantitative_data/importer.py | 184 +++++++++++++++++++++++- quantitative_data/incremental_import.py | 6 +- 2 files changed, 186 insertions(+), 4 deletions(-) diff --git a/quantitative_data/importer.py b/quantitative_data/importer.py index 5545806..9d46e09 100644 --- a/quantitative_data/importer.py +++ b/quantitative_data/importer.py @@ -766,6 +766,179 @@ def import_daily_basic_by_date( logger.info(" 每日指标导入完成") +# ============================================================ +# 4.1 导入资金流向 (moneyflow) +# ============================================================ + +def import_moneyflow( + ts_code: Optional[str] = None, + trade_date: Optional[str] = None, + start_date: Optional[str] = None, + end_date: Optional[str] = None, + conn=None, +) -> int: + """ + 导入个股资金流向 (moneyflow) + Tushare: moneyflow + + 调用方式 (互斥,按优先级生效): + 1. 按单日全市场导入 (推荐,一次拉取全市场): + import_moneyflow(trade_date="2026-08-14") + 对应示例: pro.moneyflow(trade_date='20260814') + 2. 按单只股票导入: + import_moneyflow(ts_code="000001.SZ", start_date="2026-01-01", end_date="2026-08-14") + 对应示例: pro.moneyflow(ts_code='000001.SZ', start_date='20260101', end_date='20260814') + 3. 按交易日批量全市场导入请配合 import_moneyflow_by_date 使用 + + 返回: 导入的记录数 + """ + pro = get_ts_pro() + own_conn = conn is None + if own_conn: + conn = get_pg_connection() + + log_desc = "" + try: + # ---- 构建 moneyflow 请求参数 ---- + kwargs = {} + if trade_date: + # 按单日 (全市场) + kwargs["trade_date"] = str(trade_date).replace("-", "") + log_desc = f"交易日 {kwargs['trade_date']}" + elif ts_code: + # 按单只股票 (日期范围) + if start_date is None: + start_date = START_DATE + if end_date is None: + end_date = END_DATE + kwargs["ts_code"] = ts_code + kwargs["start_date"] = start_date.replace("-", "") + kwargs["end_date"] = end_date.replace("-", "") + log_desc = f"{ts_code} ({start_date} ~ {end_date})" + else: + logger.warning( + " moneyflow: 请指定 trade_date (交易日, 如 '2026-08-14') 或 ts_code (股票代码)" + ) + return 0 + + logger.info("=" * 60) + logger.info(f"[4.1] 导入资金流向 (moneyflow): {log_desc}") + + def fetch(): + return pro.moneyflow(**kwargs) + + df = fetch_with_retry(fetch, max_retries=3) + if df is None or df.empty: + logger.warning(f" moneyflow ({log_desc}): 未获取到数据") + return 0 + + df = normalize_columns(df) + + # 转换日期列 (YYYYMMDD -> DATE) + if "trade_date" in df.columns: + df["trade_date"] = pd.to_datetime(df["trade_date"], format="%Y%m%d", errors="coerce") + + # 数值列安全转换 (NaN -> None) + numeric_cols = [ + "buy_sm_vol", "buy_sm_amount", "sell_sm_vol", "sell_sm_amount", + "buy_md_vol", "buy_md_amount", "sell_md_vol", "sell_md_amount", + "buy_lg_vol", "buy_lg_amount", "sell_lg_vol", "sell_lg_amount", + "buy_elg_vol", "buy_elg_amount", "sell_elg_vol", "sell_elg_amount", + "net_mf_vol", "net_mf_amount", + ] + for col in numeric_cols: + if col in df.columns: + df[col] = pd.to_numeric(df[col], errors="coerce") + + conflict_cols = ["ts_code", "trade_date"] + return batch_insert("moneyflow", df, conn, conflict_cols) + + finally: + if own_conn: + conn.close() + + +def import_moneyflow_by_date( + start_date: Optional[str] = None, + end_date: Optional[str] = None, + sleep_interval: float = 0.3, +): + """ + 按交易日批量导入资金流向 (全市场) + Tushare moneyflow 接口可按交易日获取全市场数据,比较高效 + + 用法: + import_moneyflow_by_date(start_date="2010-01-01", end_date="2025-12-31") + + 返回: 失败的交易日列表 + """ + if start_date is None: + start_date = START_DATE + if end_date is None: + end_date = END_DATE + + logger.info("=" * 60) + logger.info(f"[4.1] 导入资金流向 (moneyflow): {start_date} ~ {end_date}") + + # 获取交易日列表 + conn = get_pg_connection() + try: + cursor = conn.cursor() + cursor.execute( + """ + SELECT DISTINCT cal_date FROM trade_cal + WHERE is_open = 1 + AND cal_date >= %s AND cal_date <= %s + ORDER BY cal_date + """, + (start_date, end_date), + ) + trade_dates = [row[0].strftime("%Y%m%d") for row in cursor.fetchall()] + cursor.close() + finally: + conn.close() + + total = len(trade_dates) + if total == 0: + logger.warning(f" 日期范围 {start_date} ~ {end_date} 内无交易日") + return [] + + logger.info(f" 共 {total} 个交易日") + + conn = get_pg_connection() + success_count = 0 + fail_list = [] + for i, td in enumerate(trade_dates, 1): + try: + pro = get_ts_pro() + + def fetch_moneyflow(): + return pro.moneyflow(trade_date=td) + + df = fetch_with_retry(fetch_moneyflow, max_retries=3) + if df is not None and not df.empty: + df = normalize_columns(df) + if "trade_date" in df.columns: + df["trade_date"] = pd.to_datetime(df["trade_date"], format="%Y%m%d", errors="coerce") + conflict_cols = ["ts_code", "trade_date"] + batch_insert("moneyflow", df, conn, conflict_cols) + success_count += 1 + except Exception as e: + logger.warning(f" [{td}] 导入失败: {e}") + fail_list.append(td) + conn.rollback() + + if i % 20 == 0 or i == total: + logger.info(f" 进度: {i}/{total} 成功={success_count} 失败={len(fail_list)}") + time.sleep(sleep_interval) + + conn.close() + logger.info(f" 资金流向导入完成: 成功 {success_count}/{total}") + if fail_list: + logger.warning(f" 失败日期({len(fail_list)}): {fail_list[:20]}...") + return fail_list + + # ============================================================ # 5. 导入复权因子 # ============================================================ @@ -1385,9 +1558,10 @@ def full_import( 3. 交易日历 4. 日线行情 5. 每日指标(估值) - 6. 复权因子 - 7. 财务数据 (可选) - 8. 指数日线行情 + 6. 资金流向 (moneyflow) + 7. 复权因子 + 8. 财务数据 (可选) + 9. 指数日线行情 参数: - start_date, end_date: 数据范围 @@ -1430,6 +1604,9 @@ def full_import( # Step 4: 每日指标 (按日期导入) import_daily_basic_by_date(start_date, end_date) + # Step 4.1: 资金流向 (按交易日导入) + import_moneyflow_by_date(start_date, end_date) + # Step 5: 复权因子 import_adj_factor_batch(stock_codes, start_date, end_date) @@ -1500,6 +1677,7 @@ def check_table_summary(conn=None): ("trade_cal", (("trade_cal", "cal_date"),)), ("daily", (("daily", "trade_date"),)), ("daily_basic", (("daily_basic", "trade_date"),)), + ("moneyflow", (("moneyflow", "trade_date"),)), ("adj_factor", (("adj_factor", "trade_date"),)), ("income", (("income", "end_date"),)), ("balancesheet", (("balancesheet", "end_date"),)), diff --git a/quantitative_data/incremental_import.py b/quantitative_data/incremental_import.py index 49a5e87..816763b 100644 --- a/quantitative_data/incremental_import.py +++ b/quantitative_data/incremental_import.py @@ -18,7 +18,6 @@ - 行情类表按 MAX(trade_date)+1 天 → 昨天 增量拉取 - daily / daily_basic 走按交易日全市场模式(快);adj_factor 走按股票批量 - 财务表按 MAX(end_date) 往前推 400 天 → 昨天(覆盖新公告的季度报告,UPSERT 幂等) - - moneyflow 表暂未纳入(新版 importer 无对应导入函数,后续需要再补) """ import argparse import logging @@ -32,6 +31,7 @@ from importer import ( 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, @@ -49,6 +49,7 @@ 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"]), @@ -105,6 +106,7 @@ def run_daily(end_date, dry_run, limit): for table, date_col in [ ("daily", "trade_date"), ("daily_basic", "trade_date"), + ("moneyflow", "trade_date"), ("adj_factor", "trade_date"), ("index_daily", "trade_date"), ]: @@ -129,6 +131,8 @@ def run_daily(end_date, dry_run, limit): 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":