From 3684f4e2da87fcc6847a6a82c835ad16ace4e44b Mon Sep 17 00:00:00 2001 From: shellway-pc <413209390@qq.com> Date: Wed, 12 Aug 2026 22:13:30 +0800 Subject: [PATCH] =?UTF-8?q?=E5=A2=9E=E5=8A=A0=E4=BA=86=E5=85=88=E5=89=8D?= =?UTF-8?q?=E7=BC=BA=E5=A4=B1=E7=9A=84=E7=8E=B0=E9=87=91=E6=B5=81=E9=87=8F?= =?UTF-8?q?=E8=A1=A8=E7=9A=84=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 | 125 +++++++++++++++++++++++++++ quantitative_data/数据批量导入.ipynb | 59 +++++++++++-- 2 files changed, 177 insertions(+), 7 deletions(-) diff --git a/quantitative_data/importer.py b/quantitative_data/importer.py index a2fc43e..54d592a 100644 --- a/quantitative_data/importer.py +++ b/quantitative_data/importer.py @@ -933,6 +933,131 @@ def import_financial_statements( logger.info(" 财务数据导入完成") +def import_cashflow( + ts_code: Optional[str] = None, + period: Optional[str] = None, + start_date: Optional[str] = None, + end_date: Optional[str] = None, + conn=None, +) -> int: + """ + 导入现金流量表 (cashflow) + Tushare: cashflow_vip (VIP 接口) + + 调用方式 (互斥,按优先级生效): + 1. 按报告期全市场导入 (推荐,调用次数最少,一次拉取全市场某报告期): + import_cashflow(period="20181231") + 对应示例: df2 = pro.cashflow_vip(period='20181231', fields='') + 2. 按单只股票导入: + import_cashflow(ts_code="000001.SZ", start_date="2010-01-01", end_date="2025-12-31") + 3. 股票批量导入请配合 import_cashflow_batch 使用 + + 返回: 导入的记录数 + """ + pro = get_ts_pro() + own_conn = conn is None + if own_conn: + conn = get_pg_connection() + + log_desc = "" + try: + # ---- 构建 cashflow_vip 请求参数 ---- + kwargs = {} + if period: + # 按报告期 (全市场) 导入,period 形如 20181231 或 2018-12-31 + kwargs["period"] = str(period).replace("-", "") + log_desc = f"报告期 {kwargs['period']}" + 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( + " cashflow: 请指定 period (报告期, 如 '20241231') 或 ts_code (股票代码)" + ) + return 0 + + logger.info("=" * 60) + logger.info(f"[6.1] 导入现金流量表 (cashflow_vip): {log_desc}") + + def fetch(): + return pro.cashflow_vip(**kwargs) + + df = fetch_with_retry(fetch, max_retries=3) + if df is None or df.empty: + logger.warning(f" cashflow ({log_desc}): 未获取到数据") + return 0 + + df = normalize_columns(df) + + # 转换日期列 (YYYYMMDD -> DATE) + for col in ["ann_date", "f_ann_date", "end_date"]: + if col in df.columns: + df[col] = pd.to_datetime(df[col], format="%Y%m%d", errors="coerce") + + conflict_cols = ["ts_code", "end_date", "report_type"] + n = batch_insert("cashflow", df, conn, conflict_cols) + return n + + finally: + if own_conn: + conn.close() + + +def import_cashflow_batch( + stock_list: List[str], + start_date: Optional[str] = None, + end_date: Optional[str] = None, +) -> int: + """ + 按股票列表批量导入现金流量表 (cashflow) + - 逐只股票调用 cashflow_vip (VIP 接口) + - 适用于按股票维度补数据;全市场按报告期请用 import_cashflow(period=...) + + 返回: 累计导入的记录数 + """ + if start_date is None: + start_date = START_DATE + if end_date is None: + end_date = END_DATE + + total = len(stock_list) + logger.info("=" * 60) + logger.info( + f"[6.2] 批量导入现金流量表 (cashflow_vip): {start_date} ~ {end_date}, " + f"共 {total} 只股票" + ) + + conn = get_pg_connection() + success = 0 + for i, ts_code in enumerate(stock_list, 1): + try: + n = import_cashflow( + ts_code=ts_code, + start_date=start_date, + end_date=end_date, + conn=conn, + ) + success += n + except Exception as e: + logger.warning(f" [{ts_code}] 现金流量表导入失败: {e}") + conn.rollback() + + if i % 50 == 0 or i == total: + logger.info(f" 进度: {i}/{total}, 累计导入 {success} 条") + time.sleep(0.3) + + conn.close() + logger.info(f" 现金流量表批量导入完成, 共导入 {success} 条") + return success + + # ============================================================ # 7. 导入指数日线行情 # ============================================================ diff --git a/quantitative_data/数据批量导入.ipynb b/quantitative_data/数据批量导入.ipynb index ee3cb8a..a0be129 100644 --- a/quantitative_data/数据批量导入.ipynb +++ b/quantitative_data/数据批量导入.ipynb @@ -92,6 +92,8 @@ " import_adj_factor,\n", " import_adj_factor_batch,\n", " import_financial_statements,\n", + " import_cashflow,\n", + " import_cashflow_batch,\n", " import_index_daily,\n", " get_all_stock_codes,\n", " get_stock_codes_from_db,\n", @@ -519,13 +521,56 @@ ")" ] }, - { - "cell_type": "code", - "execution_count": null, - "metadata": {}, - "outputs": [], - "source": [ - "# 验证财务数据\n", + { + "cell_type": "markdown", + "metadata": {}, + "source": [ + "### 7.1 单独导入现金流量表 (cashflow_vip VIP 接口)\n", + "\n", + "`import_cashflow` 使用 Tushare VIP 接口 `cashflow_vip`,支持两种方式:\n", + "- 按报告期全市场导入: `import_cashflow(period=\"20181231\")`,对应示例 `df2 = pro.cashflow_vip(period='20181231', fields='')`\n", + "- 按股票导入: `import_cashflow(ts_code=\"000001.SZ\", start_date=..., end_date=...)`,或批量 `import_cashflow_batch(stock_list, ...)`" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# 方式A: 按报告期全市场导入 (推荐,一次拉取全市场某报告期数据,调用次数最少)\n", + "# n = import_cashflow(period=\"20241231\") # 对应 pro.cashflow_vip(period='20241231', fields='')\n", + "# print(f\"导入 20241231 报告期现金流量表: {n} 条\")\n", + "\n", + "# 方式B: 按股票列表批量导入 (适合补单只/部分股票数据)\n", + "if 'stock_list' not in dir():\n", + " stock_list = get_stock_codes_from_db()\n", + "\n", + "import_cashflow_batch(\n", + " stock_list,\n", + " start_date=\"2010-01-01\",\n", + " end_date=\"2025-12-31\",\n", + ")" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# 方式C: 导入单只股票的现金流量表\n", + "# n = import_cashflow(ts_code=\"000001.SZ\", start_date=\"2010-01-01\", end_date=\"2025-12-31\")\n", + "# print(f\"导入 000001.SZ 现金流量表: {n} 条\")" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "metadata": {}, + "outputs": [], + "source": [ + "# 验证财务数据\n", "conn = get_pg_connection()\n", "cursor = conn.cursor()\n", "for table in [\"income\", \"balancesheet\", \"cashflow\", \"fina_indicator\"]:\n",