Compare commits
3
Commits
9697b59d4f
...
3e4fab99d3
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
3e4fab99d3 | ||
|
|
d468d3c330 | ||
|
|
3684f4e2da |
@@ -933,6 +933,200 @@ def import_financial_statements(
|
|||||||
logger.info(" 财务数据导入完成")
|
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. 导入指数日线行情
|
# 7. 导入指数日线行情
|
||||||
# ============================================================
|
# ============================================================
|
||||||
|
|||||||
@@ -77,7 +77,11 @@
|
|||||||
"metadata": {},
|
"metadata": {},
|
||||||
"outputs": [],
|
"outputs": [],
|
||||||
"source": [
|
"source": [
|
||||||
"# 导入核心模块\n",
|
"# 导入核心模块 (强制重新加载,确保使用最新代码)\n",
|
||||||
|
"import importlib\n",
|
||||||
|
"import importer\n",
|
||||||
|
"importlib.reload(importer)\n",
|
||||||
|
"\n",
|
||||||
"from importer import (\n",
|
"from importer import (\n",
|
||||||
" get_pg_connection,\n",
|
" get_pg_connection,\n",
|
||||||
" get_ts_pro,\n",
|
" get_ts_pro,\n",
|
||||||
@@ -92,6 +96,9 @@
|
|||||||
" import_adj_factor,\n",
|
" import_adj_factor,\n",
|
||||||
" import_adj_factor_batch,\n",
|
" import_adj_factor_batch,\n",
|
||||||
" import_financial_statements,\n",
|
" import_financial_statements,\n",
|
||||||
|
" import_cashflow,\n",
|
||||||
|
" import_cashflow_batch,\n",
|
||||||
|
" import_cashflow_initial,\n",
|
||||||
" import_index_daily,\n",
|
" import_index_daily,\n",
|
||||||
" get_all_stock_codes,\n",
|
" get_all_stock_codes,\n",
|
||||||
" get_stock_codes_from_db,\n",
|
" get_stock_codes_from_db,\n",
|
||||||
@@ -519,13 +526,83 @@
|
|||||||
")"
|
")"
|
||||||
]
|
]
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
"cell_type": "code",
|
"cell_type": "markdown",
|
||||||
"execution_count": null,
|
"metadata": {},
|
||||||
"metadata": {},
|
"source": [
|
||||||
"outputs": [],
|
"### 7.1 单独导入现金流量表 (cashflow_vip VIP 接口)\n",
|
||||||
"source": [
|
"\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",
|
||||||
"conn = get_pg_connection()\n",
|
"conn = get_pg_connection()\n",
|
||||||
"cursor = conn.cursor()\n",
|
"cursor = conn.cursor()\n",
|
||||||
"for table in [\"income\", \"balancesheet\", \"cashflow\", \"fina_indicator\"]:\n",
|
"for table in [\"income\", \"balancesheet\", \"cashflow\", \"fina_indicator\"]:\n",
|
||||||
|
|||||||
Reference in New Issue
Block a user