增加了现金流量表初次导入的模块。

This commit is contained in:
2026-08-12 22:21:05 +08:00
parent 3684f4e2da
commit d468d3c330
2 changed files with 101 additions and 4 deletions
+69
View File
@@ -1058,6 +1058,75 @@ def import_cashflow_batch(
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. 导入指数日线行情
# ============================================================