feat(data): 双源5yr全市场部署 + baostock 15min + 断点续传
- raw_redownload 加断点续传(exists/skip已存在, 扩范围重下覆盖) - baostock_download: 15min双源(qfq+raw)下载器, NAS容器跑, 限速防封 - full_deploy_5yr.sh: 无人值守日线双源pipeline(qfq→raw→rsync→验证) - verify_dual_source: 容器内双源部署验证(fetch_day/iter_bars/除权日) 实测: 日线双源5yr 29600文件×2(rsync NAS), 浦发除权日 raw-5.9%/qfq-0.7%, 容器内 fetch_day/iter_bars/全市场抽检5只全通过. 15min沪深300 baostock限流121只.
This commit is contained in:
@@ -0,0 +1,166 @@
|
||||
#!/usr/bin/env python3
|
||||
"""baostock 15min 下载(NAS 容器跑,本机零负载)。
|
||||
|
||||
为何 baostock:新浪 stock_zh_a_minute 15min 历史仅半年;baostock 5.5 年 + 复权齐全
|
||||
(实测浦发 2021-01-04 起 21328 行;除权日 raw -5.9% vs qfq -0.7% 平滑)。
|
||||
|
||||
双源:qfq(adjustflag=2) + raw(adjustflag=3)。断点续传(skip 已存在,扩范围重下)。
|
||||
数据落 NAS 本地盘(容器挂载 /volume1/stock → /stock,最快,不经网络写入)。
|
||||
|
||||
约束:单线程 + SLEEP 限速(baostock 服务器温和限频,避免封)。
|
||||
|
||||
用法(NAS 容器,本机编排):
|
||||
/var/packages/Docker/target/usr/bin/docker run --rm -d --name dl15 \\
|
||||
--memory=512m -v /volume1/stock:/stock \\
|
||||
sanguo_vnpy_v2:with-sqlite \\
|
||||
sh -c "pip install baostock -q && python /stock/sanguo_vnpy/scripts/baostock_download.py --all --adjust both --start 2021-01-01"
|
||||
|
||||
env:
|
||||
DL_DIR 输出根(默认 /stock/A股数据/minute_kline)
|
||||
STOCK_LIST 全市场 csv(默认 NAS stock_info csv)
|
||||
SLEEP 每只每源间隔秒(默认 0.3)
|
||||
"""
|
||||
import argparse
|
||||
import csv
|
||||
import logging
|
||||
import os
|
||||
import sys
|
||||
import time
|
||||
|
||||
import pandas as pd
|
||||
import baostock as bs
|
||||
|
||||
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s")
|
||||
log = logging.getLogger("dl_15min")
|
||||
|
||||
DL_DIR = os.environ.get("DL_DIR", "/stock/A股数据/minute_kline")
|
||||
SLEEP = float(os.environ.get("SLEEP", "0.3"))
|
||||
DEFAULT_STOCK_LIST = "/stock/A股数据/stock_info/stock_basic_info_raw_20260326_113530.csv"
|
||||
|
||||
|
||||
def prefix(code: str) -> str:
|
||||
return "sh" if code.startswith(("60", "68", "51", "56", "58")) else "sz"
|
||||
|
||||
|
||||
def bs_symbol(code: str) -> str:
|
||||
return f"{prefix(code)}.{code}"
|
||||
|
||||
|
||||
def out_path(code: str, adjust: str) -> str:
|
||||
d = os.path.join(DL_DIR, f"15min_{adjust}")
|
||||
os.makedirs(d, exist_ok=True)
|
||||
return os.path.join(d, f"{prefix(code)}{code}_15min.parquet")
|
||||
|
||||
|
||||
def exists(code: str, adjust: str, start_year: int) -> bool:
|
||||
"""断点续传:文件存在且含 start_year 数据(扩范围 start 更早 → 重下覆盖)。"""
|
||||
f = out_path(code, adjust)
|
||||
if not os.path.exists(f):
|
||||
return False
|
||||
try:
|
||||
df = pd.read_parquet(f, columns=["date"])
|
||||
return str(start_year) in df["date"].astype(str).str[:4].unique()
|
||||
except Exception:
|
||||
return False
|
||||
|
||||
|
||||
def download_one(code: str, start: str, end: str, adjustflag: str):
|
||||
"""baostock 15min,返回 (df, None) 或 (None, err)。"""
|
||||
rs = bs.query_history_k_data_plus(
|
||||
bs_symbol(code),
|
||||
"date,time,open,high,low,close,volume,amount",
|
||||
start_date=start, end_date=end, frequency="15", adjustflag=adjustflag,
|
||||
)
|
||||
rows, fields = [], rs.fields
|
||||
while (rs.error_code == "0") & rs.next():
|
||||
rows.append(rs.get_row_data())
|
||||
if not rows:
|
||||
return None, rs.error_msg
|
||||
df = pd.DataFrame(rows, columns=fields)
|
||||
for c in ("open", "high", "low", "close", "volume", "amount"):
|
||||
df[c] = pd.to_numeric(df[c], errors="coerce")
|
||||
t = df["time"].astype(str)
|
||||
df["datetime"] = pd.to_datetime(
|
||||
df["date"] + " " + t.str[8:10] + ":" + t.str[10:12], errors="coerce"
|
||||
)
|
||||
df = df.dropna(subset=["close", "datetime"]).sort_values("datetime")
|
||||
return df, None
|
||||
|
||||
|
||||
def save_one(code, df, adjust):
|
||||
df.to_parquet(out_path(code, adjust))
|
||||
return len(df)
|
||||
|
||||
|
||||
def load_all_codes(stock_list):
|
||||
codes = []
|
||||
with open(stock_list, encoding="utf-8", errors="replace") as f:
|
||||
for row in csv.DictReader(f):
|
||||
for k in ("code", "symbol", "ts_code", "代码", "股票代码", "成分券代码"):
|
||||
if k in row and row[k]:
|
||||
c = str(row[k]).strip().split(".")[0]
|
||||
if c.isdigit() and len(c) == 6:
|
||||
codes.append(c)
|
||||
break
|
||||
return codes
|
||||
|
||||
|
||||
def main():
|
||||
ap = argparse.ArgumentParser()
|
||||
ap.add_argument("--symbols")
|
||||
ap.add_argument("--all", action="store_true")
|
||||
ap.add_argument("--start", default="2021-01-01")
|
||||
ap.add_argument("--end", default=None)
|
||||
ap.add_argument("--adjust", default="qfq", help="qfq / raw / both")
|
||||
ap.add_argument("--force", action="store_true")
|
||||
args = ap.parse_args()
|
||||
|
||||
end = args.end or time.strftime("%Y-%m-%d")
|
||||
start_year = int(args.start[:4])
|
||||
adjusts = ["qfq", "raw"] if args.adjust == "both" else [args.adjust]
|
||||
flag = {"qfq": "2", "raw": "3"}
|
||||
|
||||
if args.symbols:
|
||||
codes = [c.strip() for c in args.symbols.split(",") if c.strip()]
|
||||
elif args.all:
|
||||
sl = os.environ.get("STOCK_LIST", DEFAULT_STOCK_LIST)
|
||||
if not os.path.exists(sl):
|
||||
log.error("STOCK_LIST 不存在: %s", sl)
|
||||
sys.exit(1)
|
||||
codes = load_all_codes(sl)
|
||||
log.info("--all 读到 %d 只", len(codes))
|
||||
else:
|
||||
ap.error("需 --symbols 或 --all")
|
||||
|
||||
lg = bs.login()
|
||||
log.info("baostock login: %s %s", lg.error_code, lg.error_msg)
|
||||
|
||||
ok = fail = rows = skipped = 0
|
||||
for i, code in enumerate(codes, 1):
|
||||
for adj in adjusts:
|
||||
if not args.force and exists(code, adj, start_year):
|
||||
skipped += 1
|
||||
continue
|
||||
df, err = download_one(code, args.start, end, flag[adj])
|
||||
if df is None or df.empty:
|
||||
fail += 1
|
||||
if i <= 20 or i % 100 == 0:
|
||||
log.warning("[%d/%d] %s %s empty %s", i, len(codes), code, adj, err)
|
||||
else:
|
||||
try:
|
||||
n = save_one(code, df, adj)
|
||||
ok += 1
|
||||
rows += n
|
||||
if i <= 20 or i % 100 == 0:
|
||||
log.info("[%d/%d] %s %s ok %d", i, len(codes), code, adj, n)
|
||||
except Exception as e: # noqa: BLE001
|
||||
fail += 1
|
||||
log.error("save %s %s: %s", code, adj, e)
|
||||
time.sleep(SLEEP)
|
||||
bs.logout()
|
||||
log.info("=== 完成 ok=%d skip=%d fail=%d rows=%d dir=%s ===",
|
||||
ok, skipped, fail, rows, DL_DIR)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
Executable
+38
@@ -0,0 +1,38 @@
|
||||
#!/bin/bash
|
||||
# 无人值守日线双源 5yr 全市场部署 pipeline(本机串行,qfq→raw→rsync→验证)。
|
||||
# 不并发打新浪源(封IP):qfq 进程结束才下 raw。
|
||||
# 15min baostock 走 NAS 容器,与本机并行(不同源不同机)。
|
||||
set -uo pipefail
|
||||
cd /Users/chufeng/.openclaw/sanguo_projects/sanguo_vnpy_v2
|
||||
LOG=/tmp/full_deploy.log
|
||||
SL=data_cache/stock_info/stock_basic_info_raw_20260326_113530.csv
|
||||
log(){ echo "[$(date '+%m-%d %H:%M:%S')] $*" | tee -a "$LOG"; }
|
||||
|
||||
log "=== STAGE1 等 qfq 5yr 进程结束 ==="
|
||||
while pgrep -f "raw_redownload.py --all --adjust qfq" >/dev/null; do sleep 60; done
|
||||
log "qfq 完成: $(find data_cache/qfq -name '*.parquet'|wc -l) 文件"
|
||||
|
||||
log "=== STAGE2 raw 5yr 全市场(断点续传,单线程限速)==="
|
||||
STOCK_LIST="$SL" RAW_DIR=data_cache/raw SLEEP=1.0 \
|
||||
python3 scripts/data_platform/raw_redownload.py --all --start 2021-01-01 --adjust "" >> "$LOG" 2>&1
|
||||
log "raw 完成: $(find data_cache/raw -name '*.parquet'|wc -l) 文件"
|
||||
|
||||
log "=== STAGE3 rsync qfq 全量(--delete覆盖NAS旧5验证股) → NAS ==="
|
||||
rsync -az --delete \
|
||||
data_cache/qfq/ \
|
||||
sanguo-nas:/volume1/stock/A股数据/日线数据/qfq/ >> "$LOG" 2>&1
|
||||
log "qfq rsync done"
|
||||
|
||||
log "=== STAGE4 rsync raw 2021-2023 → NAS(不覆盖NAS已有24-26)==="
|
||||
for y in 2021 2022 2023; do
|
||||
rsync -az \
|
||||
"data_cache/raw/$y/" \
|
||||
"sanguo-nas:/volume1/stock/A股数据/日线数据/raw/$y/" >> "$LOG" 2>&1
|
||||
log " raw/$y done"
|
||||
done
|
||||
|
||||
log "=== STAGE5 验证 NAS 文件数 ==="
|
||||
ssh sanguo-nas 'for d in qfq raw; do for y in 2021 2022 2023 2024 2025 2026; do echo "$d/$y=$(ls /volume1/stock/A股数据/日线数据/$d/$y/*.parquet 2>/dev/null|wc -l)"; done; done' >> "$LOG" 2>&1
|
||||
|
||||
log "=== ALL DONE ==="
|
||||
echo "ALL_DONE_MARKER" >> "$LOG"
|
||||
@@ -79,6 +79,17 @@ def save_one(code: str, df) -> int:
|
||||
return n
|
||||
|
||||
|
||||
def exists(code: str, start_year: int) -> bool:
|
||||
"""symbol 在 start_year 是否已有 parquet(断点续传)。
|
||||
|
||||
同范围续跑 → skip;扩范围(start 更早)→ 新 start_year 不存在 → 重下全量覆盖。
|
||||
"""
|
||||
pref = prefix_for(code)
|
||||
return os.path.exists(
|
||||
os.path.join(RAW_DIR, str(start_year), f"{pref}{code}_daily.parquet")
|
||||
)
|
||||
|
||||
|
||||
def load_all_codes(stock_list: str) -> list[str]:
|
||||
"""从 stock_basic_info csv 读代码列表(容错列名)。"""
|
||||
codes = []
|
||||
@@ -101,6 +112,7 @@ def main():
|
||||
ap.add_argument("--start", default="2024-01-01")
|
||||
ap.add_argument("--end", default=None, help="默认今天")
|
||||
ap.add_argument("--adjust", default="", help="复权: '' raw / 'qfq' 前复权(双源用)")
|
||||
ap.add_argument("--force", action="store_true", help="强制重下(默认 skip 已存在=断点续传)")
|
||||
args = ap.parse_args()
|
||||
|
||||
end = args.end or time.strftime("%Y-%m-%d")
|
||||
@@ -121,8 +133,14 @@ def main():
|
||||
import warnings
|
||||
warnings.filterwarnings("ignore")
|
||||
|
||||
ok = fail = rows = 0
|
||||
start_year = int(args.start[:4])
|
||||
ok = fail = rows = skipped = 0
|
||||
for i, code in enumerate(codes, 1):
|
||||
if not args.force and exists(code, start_year):
|
||||
skipped += 1
|
||||
if skipped % 500 == 0:
|
||||
log.info("[%d/%d] ... skipped %d 已存在", i, len(codes), skipped)
|
||||
continue
|
||||
df, err = download_one(ak, code, args.start, end, args.adjust)
|
||||
if df is None:
|
||||
fail += 1
|
||||
@@ -138,7 +156,7 @@ def main():
|
||||
log.error("[%d/%d] %s SAVE FAIL %s", i, len(codes), code, e)
|
||||
time.sleep(SLEEP) # 限速(用户约束)
|
||||
|
||||
log.info("=== 完成: ok=%d fail=%d rows=%d,raw_dir=%s ===", ok, fail, rows, RAW_DIR)
|
||||
log.info("=== 完成: ok=%d skip=%d fail=%d rows=%d,raw_dir=%s ===", ok, skipped, fail, rows, RAW_DIR)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
|
||||
@@ -0,0 +1,64 @@
|
||||
#!/usr/bin/env python3
|
||||
"""双源(raw/qfq)部署验证:容器内读 NAS data_platform.yaml 配置的双源,
|
||||
确认 fetch_day(C-S3 实走接口)+ iter_bars 在 raw_dir/qfq_dir 正确路由,
|
||||
且除权日 raw 有缺口 / qfq 平滑(分红除权准确)。
|
||||
|
||||
用法(容器):
|
||||
docker run --rm -v /volume1/stock:/volume1/stock \
|
||||
-v /volume1/homes/admin/.sanguo_projects/sanguo_vnpy_v2:/app \
|
||||
--entrypoint sh sanguo_vnpy_v2:with-sqlite \
|
||||
-c "cd /app && python scripts/verify_dual_source.py"
|
||||
"""
|
||||
import sys
|
||||
sys.path.insert(0, "/app")
|
||||
from sanguo_data.config import load_config, find_config_path
|
||||
from sanguo_trader.data_source import fetch_day, iter_bars
|
||||
|
||||
cfg = load_config(find_config_path())
|
||||
print(f"config: {find_config_path()}")
|
||||
print(f"data_paths: raw_dir={cfg.data_paths.get('raw_dir')} qfq_dir={cfg.data_paths.get('qfq_dir')}")
|
||||
|
||||
OK = True
|
||||
|
||||
# 1) fetch_day(C-S3 实走接口)除权日双源
|
||||
print("\n=== fetch_day 除权日双源(浦发 2022-07-21 除权)===")
|
||||
res = {}
|
||||
for adj in ["qfq", "raw"]:
|
||||
b = fetch_day("600000", "2022-07-21", "d", adj, cfg)
|
||||
bp = fetch_day("600000", "2022-07-20", "d", adj, cfg)
|
||||
if b is None or bp is None:
|
||||
print(f" ❌ {adj}: fetch_day 返回 None(数据缺失或路由失败)")
|
||||
OK = False
|
||||
continue
|
||||
drop = ((b.close_price / bp.close_price) - 1) * 100
|
||||
res[adj] = drop
|
||||
print(f" {adj}: 07-20={bp.close_price:.4f} → 07-21={b.close_price:.4f} ({drop:+.1f}%)")
|
||||
|
||||
# 除权日 raw 应有明显缺口(< -3%),qfq 平滑(> -3%)
|
||||
if "raw" in res and "qfq" in res:
|
||||
if res["raw"] < -3 and res["qfq"] > -3:
|
||||
print(" ✅ 除权处理正确:raw 含缺口(撮合真实价)/ qfq 平滑(信号准)")
|
||||
else:
|
||||
print(f" ⚠️ 除权差异不明显 raw={res['raw']:.1f}% qfq={res['qfq']:.1f}%")
|
||||
|
||||
# 2) iter_bars cross-section(symbols=list)
|
||||
print("\n=== iter_bars cross-section(4 天)===")
|
||||
for adj in ["qfq", "raw"]:
|
||||
days = list(iter_bars(["600000"], "2022-07-19", "2022-07-22", "d", adj, cfg))
|
||||
if not days:
|
||||
print(f" ❌ {adj}: iter_bars 0 days")
|
||||
OK = False
|
||||
else:
|
||||
last = days[-1]
|
||||
print(f" {adj}: {len(days)} days, 末日={last[0]} close={last[1]['600000'].close_price:.4f}")
|
||||
|
||||
# 3) 全市场可读性抽检(5 只)
|
||||
print("\n=== 全市场抽检(5 只 fetch_day)===")
|
||||
for sym in ["600000", "000001", "000858", "300750", "688981"]:
|
||||
b = fetch_day(sym, "2026-07-07", "d", "qfq", cfg)
|
||||
print(f" {sym}: {'✅ close='+format(b.close_price,'.4f') if b else '❌ None'}")
|
||||
if b is None:
|
||||
OK = False
|
||||
|
||||
print("\n" + ("✅ 双源部署验证通过" if OK else "❌ 双源验证有失败项"))
|
||||
sys.exit(0 if OK else 1)
|
||||
Reference in New Issue
Block a user