From 71444a8172dcd406fcff2fdb2a394ff08022b9c4 Mon Sep 17 00:00:00 2001 From: shellway-pc <413209390@qq.com> Date: Sat, 15 Aug 2026 23:11:14 +0800 Subject: [PATCH] =?UTF-8?q?fix:=E5=88=A0=E9=99=A4=E4=BA=867.1=E5=8D=95?= =?UTF-8?q?=E7=8B=AC=E5=AF=BC=E5=85=A5=E7=8E=B0=E9=87=91=E6=B5=81=E9=87=8F?= =?UTF-8?q?=E8=A1=A8=E6=A8=A1=E5=9D=97=EF=BC=8C=E5=85=88=E5=89=8D=E8=AF=AF?= =?UTF-8?q?=E5=B0=86moneyflow=EF=BC=88=E8=B5=84=E9=87=91=E6=B5=81=E5=90=91?= =?UTF-8?q?=E8=A1=A8=EF=BC=89=E5=BD=93=E6=88=90=E4=BA=86=E7=8E=B0=E9=87=91?= =?UTF-8?q?=E6=B5=81=E9=87=8F=E8=A1=A8=E3=80=82=E2=80=9C?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- quantitative_data/importer.py | 194 --------------------------- quantitative_data/数据批量导入.ipynb | 87 +----------- 2 files changed, 7 insertions(+), 274 deletions(-) diff --git a/quantitative_data/importer.py b/quantitative_data/importer.py index 9d46e09..b087ff7 100644 --- a/quantitative_data/importer.py +++ b/quantitative_data/importer.py @@ -1106,200 +1106,6 @@ 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 - - -def import_cashflow_initial( - start_date: str = "2010-01-01", - end_date: str = "2015-12-31", - sleep_interval: float = 0.3, -) -> int: - """ - 首次批量导入现金流量表 (cashflow) - 按报告期全市场循环导入 - - 自动生成 start_date ~ end_date 范围内的所有季度报告期 (0331/0630/0930/1231) - - 每个报告期调用一次 cashflow_vip(period=...) 一次性拉取全市场数据 - - 适用于首次初始化导入;后续增量/修补请用 import_cashflow(period=...) 或 import_cashflow_batch - - 示例: - import_cashflow_initial(start_date="2010-01-01", end_date="2015-12-31") - 内部循环调用: pro.cashflow_vip(period='20100331'), pro.cashflow_vip(period='20100630'), ..., - pro.cashflow_vip(period='20151231'),共 24 个报告期 - - 返回: 累计导入的记录数 - """ - # 解析年份范围 (支持 "2010-01-01" / "20100101" / "2010" 等格式) - start_compact = str(start_date).replace("-", "") - end_compact = str(end_date).replace("-", "") - start_year = int(start_compact[:4]) - end_year = int(end_compact[:4]) - - # 自动生成所有季度报告期 (YYYYMMDD) - periods = [] - for year in range(start_year, end_year + 1): - for month_day in ["0331", "0630", "0930", "1231"]: - period_str = f"{year}{month_day}" - # 过滤掉首尾年份中超出日期范围的报告期 - if period_str < start_compact or period_str > end_compact: - continue - periods.append(period_str) - - if not periods: - logger.warning(f" 日期范围 {start_date} ~ {end_date} 内无报告期") - return 0 - - total = len(periods) - logger.info("=" * 60) - logger.info( - f"[6.3] 首次批量导入现金流量表 (cashflow_vip 按报告期): " - f"{start_date} ~ {end_date}" - ) - logger.info(f" 共 {total} 个报告期: {periods[0]} ~ {periods[-1]}") - - conn = get_pg_connection() - success = 0 - fail_list = [] - for i, period in enumerate(periods, 1): - try: - n = import_cashflow(period=period, conn=conn) - success += n - except Exception as e: - logger.warning(f" [{period}] 现金流量表导入失败: {e}") - fail_list.append(period) - conn.rollback() - - if i % 5 == 0 or i == total: - logger.info(f" 进度: {i}/{total}, 累计导入 {success} 条") - time.sleep(sleep_interval) - - conn.close() - logger.info(f" 首次现金流量表批量导入完成: 成功 {success} 条") - if fail_list: - logger.warning(f" 失败报告期({len(fail_list)}): {fail_list}") - return success - - # ============================================================ # 7. 导入指数日线行情 # ============================================================ diff --git a/quantitative_data/数据批量导入.ipynb b/quantitative_data/数据批量导入.ipynb index 1f34490..34a381d 100644 --- a/quantitative_data/数据批量导入.ipynb +++ b/quantitative_data/数据批量导入.ipynb @@ -96,9 +96,6 @@ " import_adj_factor,\n", " import_adj_factor_batch,\n", " import_financial_statements,\n", - " import_cashflow,\n", - " import_cashflow_batch,\n", - " import_cashflow_initial,\n", " import_index_daily,\n", " get_all_stock_codes,\n", " get_stock_codes_from_db,\n", @@ -526,83 +523,13 @@ ")" ] }, - { - "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": "markdown", - "metadata": {}, - "source": [ - "#### 方式A2: 首次批量导入 (推荐初始化数据)\n", - "\n", - "`import_cashflow_initial` 自动生成日期范围内的所有季度报告期 (0331/0630/0930/1231),\n", - "每个报告期调用一次 `cashflow_vip(period=...)` 一次拉取全市场数据。\n", - "\n", - "以下示例一次导入 **2010年1月1日 ~ 2015年12月31日** 共 24 个报告期的所有上市公司现金流量表数据。" - ] - }, - { - "cell_type": "code", - "execution_count": null, - "metadata": {}, - "outputs": [], - "source": [ - "# 首次批量导入: 一次导入 2010-01-01 到 2015-12-31 的所有上市公司现金流量表\n", - "# 自动循环 24 个季度报告期 (20100331 ~ 20151231),每个报告期拉取全市场数据\n", - "import_cashflow_initial(\n", - " start_date=\"2010-01-01\",\n", - " end_date=\"2015-12-31\",\n", - " sleep_interval=0.3,\n", - ")" - ] - }, - { - "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", + { + "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",