增加了moneyflow表格导入模块。
This commit is contained in:
@@ -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. 导入复权因子
|
||||||
# ============================================================
|
# ============================================================
|
||||||
@@ -1385,9 +1558,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 +1604,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 +1677,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"),)),
|
||||||
|
|||||||
@@ -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":
|
||||||
|
|||||||
Reference in New Issue
Block a user