Compare commits

..
2 Commits
3 changed files with 193 additions and 278 deletions
+181 -197
View File
@@ -766,6 +766,179 @@ def import_daily_basic_by_date(
logger.info(" 每日指标导入完成") logger.info(" 每日指标导入完成")
# ============================================================
# 4.1 导入资金流向 (moneyflow)
# ============================================================
def import_moneyflow(
ts_code: Optional[str] = None,
trade_date: Optional[str] = None,
start_date: Optional[str] = None,
end_date: Optional[str] = None,
conn=None,
) -> int:
"""
导入个股资金流向 (moneyflow)
Tushare: moneyflow
调用方式 (互斥,按优先级生效):
1. 按单日全市场导入 (推荐,一次拉取全市场):
import_moneyflow(trade_date="2026-08-14")
对应示例: pro.moneyflow(trade_date='20260814')
2. 按单只股票导入:
import_moneyflow(ts_code="000001.SZ", start_date="2026-01-01", end_date="2026-08-14")
对应示例: pro.moneyflow(ts_code='000001.SZ', start_date='20260101', end_date='20260814')
3. 按交易日批量全市场导入请配合 import_moneyflow_by_date 使用
返回: 导入的记录数
"""
pro = get_ts_pro()
own_conn = conn is None
if own_conn:
conn = get_pg_connection()
log_desc = ""
try:
# ---- 构建 moneyflow 请求参数 ----
kwargs = {}
if trade_date:
# 按单日 (全市场)
kwargs["trade_date"] = str(trade_date).replace("-", "")
log_desc = f"交易日 {kwargs['trade_date']}"
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(
" moneyflow: 请指定 trade_date (交易日, 如 '2026-08-14') 或 ts_code (股票代码)"
)
return 0
logger.info("=" * 60)
logger.info(f"[4.1] 导入资金流向 (moneyflow): {log_desc}")
def fetch():
return pro.moneyflow(**kwargs)
df = fetch_with_retry(fetch, max_retries=3)
if df is None or df.empty:
logger.warning(f" moneyflow ({log_desc}): 未获取到数据")
return 0
df = normalize_columns(df)
# 转换日期列 (YYYYMMDD -> DATE)
if "trade_date" in df.columns:
df["trade_date"] = pd.to_datetime(df["trade_date"], format="%Y%m%d", errors="coerce")
# 数值列安全转换 (NaN -> None)
numeric_cols = [
"buy_sm_vol", "buy_sm_amount", "sell_sm_vol", "sell_sm_amount",
"buy_md_vol", "buy_md_amount", "sell_md_vol", "sell_md_amount",
"buy_lg_vol", "buy_lg_amount", "sell_lg_vol", "sell_lg_amount",
"buy_elg_vol", "buy_elg_amount", "sell_elg_vol", "sell_elg_amount",
"net_mf_vol", "net_mf_amount",
]
for col in numeric_cols:
if col in df.columns:
df[col] = pd.to_numeric(df[col], errors="coerce")
conflict_cols = ["ts_code", "trade_date"]
return batch_insert("moneyflow", df, conn, conflict_cols)
finally:
if own_conn:
conn.close()
def import_moneyflow_by_date(
start_date: Optional[str] = None,
end_date: Optional[str] = None,
sleep_interval: float = 0.3,
):
"""
按交易日批量导入资金流向 (全市场)
Tushare moneyflow 接口可按交易日获取全市场数据,比较高效
用法:
import_moneyflow_by_date(start_date="2010-01-01", end_date="2025-12-31")
返回: 失败的交易日列表
"""
if start_date is None:
start_date = START_DATE
if end_date is None:
end_date = END_DATE
logger.info("=" * 60)
logger.info(f"[4.1] 导入资金流向 (moneyflow): {start_date} ~ {end_date}")
# 获取交易日列表
conn = get_pg_connection()
try:
cursor = conn.cursor()
cursor.execute(
"""
SELECT DISTINCT cal_date FROM trade_cal
WHERE is_open = 1
AND cal_date >= %s AND cal_date <= %s
ORDER BY cal_date
""",
(start_date, end_date),
)
trade_dates = [row[0].strftime("%Y%m%d") for row in cursor.fetchall()]
cursor.close()
finally:
conn.close()
total = len(trade_dates)
if total == 0:
logger.warning(f" 日期范围 {start_date} ~ {end_date} 内无交易日")
return []
logger.info(f" 共 {total} 个交易日")
conn = get_pg_connection()
success_count = 0
fail_list = []
for i, td in enumerate(trade_dates, 1):
try:
pro = get_ts_pro()
def fetch_moneyflow():
return pro.moneyflow(trade_date=td)
df = fetch_with_retry(fetch_moneyflow, max_retries=3)
if df is not None and not df.empty:
df = normalize_columns(df)
if "trade_date" in df.columns:
df["trade_date"] = pd.to_datetime(df["trade_date"], format="%Y%m%d", errors="coerce")
conflict_cols = ["ts_code", "trade_date"]
batch_insert("moneyflow", df, conn, conflict_cols)
success_count += 1
except Exception as e:
logger.warning(f" [{td}] 导入失败: {e}")
fail_list.append(td)
conn.rollback()
if i % 20 == 0 or i == total:
logger.info(f" 进度: {i}/{total} 成功={success_count} 失败={len(fail_list)}")
time.sleep(sleep_interval)
conn.close()
logger.info(f" 资金流向导入完成: 成功 {success_count}/{total}")
if fail_list:
logger.warning(f" 失败日期({len(fail_list)}): {fail_list[:20]}...")
return fail_list
# ============================================================ # ============================================================
# 5. 导入复权因子 # 5. 导入复权因子
# ============================================================ # ============================================================
@@ -933,200 +1106,6 @@ 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. 导入指数日线行情
# ============================================================ # ============================================================
@@ -1385,9 +1364,10 @@ def full_import(
3. 交易日历 3. 交易日历
4. 日线行情 4. 日线行情
5. 每日指标(估值) 5. 每日指标(估值)
6. 复权因子 6. 资金流向 (moneyflow)
7. 财务数据 (可选) 7. 复权因子
8. 指数日线行情 8. 财务数据 (可选)
9. 指数日线行情
参数: 参数:
- start_date, end_date: 数据范围 - start_date, end_date: 数据范围
@@ -1430,6 +1410,9 @@ def full_import(
# Step 4: 每日指标 (按日期导入) # Step 4: 每日指标 (按日期导入)
import_daily_basic_by_date(start_date, end_date) import_daily_basic_by_date(start_date, end_date)
# Step 4.1: 资金流向 (按交易日导入)
import_moneyflow_by_date(start_date, end_date)
# Step 5: 复权因子 # Step 5: 复权因子
import_adj_factor_batch(stock_codes, start_date, end_date) import_adj_factor_batch(stock_codes, start_date, end_date)
@@ -1500,6 +1483,7 @@ def check_table_summary(conn=None):
("trade_cal", (("trade_cal", "cal_date"),)), ("trade_cal", (("trade_cal", "cal_date"),)),
("daily", (("daily", "trade_date"),)), ("daily", (("daily", "trade_date"),)),
("daily_basic", (("daily_basic", "trade_date"),)), ("daily_basic", (("daily_basic", "trade_date"),)),
("moneyflow", (("moneyflow", "trade_date"),)),
("adj_factor", (("adj_factor", "trade_date"),)), ("adj_factor", (("adj_factor", "trade_date"),)),
("income", (("income", "end_date"),)), ("income", (("income", "end_date"),)),
("balancesheet", (("balancesheet", "end_date"),)), ("balancesheet", (("balancesheet", "end_date"),)),
+5 -1
View File
@@ -18,7 +18,6 @@
- 行情类表按 MAX(trade_date)+1 天 → 昨天 增量拉取 - 行情类表按 MAX(trade_date)+1 天 → 昨天 增量拉取
- daily / daily_basic 走按交易日全市场模式(快);adj_factor 走按股票批量 - daily / daily_basic 走按交易日全市场模式(快);adj_factor 走按股票批量
- 财务表按 MAX(end_date) 往前推 400 天 → 昨天(覆盖新公告的季度报告,UPSERT 幂等) - 财务表按 MAX(end_date) 往前推 400 天 → 昨天(覆盖新公告的季度报告,UPSERT 幂等)
- moneyflow 表暂未纳入(新版 importer 无对应导入函数,后续需要再补)
""" """
import argparse import argparse
import logging import logging
@@ -32,6 +31,7 @@ from importer import (
import_stock_basic, import_stock_basic,
import_daily_by_date, import_daily_by_date,
import_daily_basic_by_date, import_daily_basic_by_date,
import_moneyflow_by_date,
import_adj_factor_batch, import_adj_factor_batch,
import_index_daily, import_index_daily,
import_financial_statements, import_financial_statements,
@@ -49,6 +49,7 @@ logger.setLevel(logging.INFO)
TABLE_SPECS = { TABLE_SPECS = {
"daily": ("trade_date", ["ts_code", "trade_date"]), "daily": ("trade_date", ["ts_code", "trade_date"]),
"daily_basic": ("trade_date", ["ts_code", "trade_date"]), "daily_basic": ("trade_date", ["ts_code", "trade_date"]),
"moneyflow": ("trade_date", ["ts_code", "trade_date"]),
"adj_factor": ("trade_date", ["ts_code", "trade_date"]), "adj_factor": ("trade_date", ["ts_code", "trade_date"]),
"index_daily": ("trade_date", ["ts_code", "trade_date"]), "index_daily": ("trade_date", ["ts_code", "trade_date"]),
"income": ("end_date", ["ts_code", "end_date", "report_type"]), "income": ("end_date", ["ts_code", "end_date", "report_type"]),
@@ -105,6 +106,7 @@ def run_daily(end_date, dry_run, limit):
for table, date_col in [ for table, date_col in [
("daily", "trade_date"), ("daily", "trade_date"),
("daily_basic", "trade_date"), ("daily_basic", "trade_date"),
("moneyflow", "trade_date"),
("adj_factor", "trade_date"), ("adj_factor", "trade_date"),
("index_daily", "trade_date"), ("index_daily", "trade_date"),
]: ]:
@@ -129,6 +131,8 @@ def run_daily(end_date, dry_run, limit):
import_daily_by_date(start, end_date) import_daily_by_date(start, end_date)
elif table == "daily_basic": elif table == "daily_basic":
import_daily_basic_by_date(start, end_date) import_daily_basic_by_date(start, end_date)
elif table == "moneyflow":
import_moneyflow_by_date(start, end_date)
elif table == "adj_factor": elif table == "adj_factor":
import_adj_factor_batch(codes, start, end_date) import_adj_factor_batch(codes, start, end_date)
elif table == "index_daily": elif table == "index_daily":
@@ -96,9 +96,6 @@
" 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",
@@ -525,76 +522,6 @@
" end_date=\"2025-12-31\",\n", " end_date=\"2025-12-31\",\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": "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", "cell_type": "code",