From a37254404524f0da5885a4febcf06098f796e172 Mon Sep 17 00:00:00 2001 From: claude_dev Date: Thu, 9 Jul 2026 19:22:36 +0800 Subject: [PATCH] =?UTF-8?q?feat(data):=20=E5=8F=8C=E6=BA=905yr=E5=85=A8?= =?UTF-8?q?=E5=B8=82=E5=9C=BA=E9=83=A8=E7=BD=B2=20+=20baostock=2015min=20+?= =?UTF-8?q?=20=E6=96=AD=E7=82=B9=E7=BB=AD=E4=BC=A0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 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只. --- scripts/data_platform/baostock_download.py | 166 +++++++++++++++++++++ scripts/data_platform/full_deploy_5yr.sh | 38 +++++ scripts/data_platform/raw_redownload.py | 22 ++- scripts/verify_dual_source.py | 64 ++++++++ 4 files changed, 288 insertions(+), 2 deletions(-) create mode 100644 scripts/data_platform/baostock_download.py create mode 100755 scripts/data_platform/full_deploy_5yr.sh create mode 100644 scripts/verify_dual_source.py diff --git a/scripts/data_platform/baostock_download.py b/scripts/data_platform/baostock_download.py new file mode 100644 index 0000000..370b18d --- /dev/null +++ b/scripts/data_platform/baostock_download.py @@ -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() diff --git a/scripts/data_platform/full_deploy_5yr.sh b/scripts/data_platform/full_deploy_5yr.sh new file mode 100755 index 0000000..7987d52 --- /dev/null +++ b/scripts/data_platform/full_deploy_5yr.sh @@ -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" diff --git a/scripts/data_platform/raw_redownload.py b/scripts/data_platform/raw_redownload.py index eb6a13e..fe78dfd 100644 --- a/scripts/data_platform/raw_redownload.py +++ b/scripts/data_platform/raw_redownload.py @@ -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__": diff --git a/scripts/verify_dual_source.py b/scripts/verify_dual_source.py new file mode 100644 index 0000000..d87fd81 --- /dev/null +++ b/scripts/verify_dual_source.py @@ -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)