增加了先前缺失的现金流量表的导入模块。
This commit is contained in:
@@ -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. 导入指数日线行情
|
||||
# ============================================================
|
||||
|
||||
Reference in New Issue
Block a user