Files
sanguo_vnpy_v2/scripts/data_platform/raw_redownload.py
T
claude_dev a372544045 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只.
2026-07-09 19:22:36 +08:00

164 lines
6.1 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
#!/usr/bin/env python3
"""raw 日线重下(task #79):akshare 新浪源 adjust="" 真实价。
根因:daily_dir 历史数据是 mixed-adjusthfq bulk + akshare raw tail 拼接),
3-30 类单日 -94% 假跌。本脚本从头重下**单一 raw** 到 raw_dir,与 daily_dir 同构
{prefix}{symbol}_daily.parquet,按 year 分目录),datareader 直接读。
**约束(用户反馈,见 memory/feedback-data-download-constraints**
- 直连不走代理(unset proxy env + NO_PROXY=*
- 单线程 + SLEEP 间隔限速(不并发猛打,防数据源封 IP)
用法:
# 验证单只/几只
python3 raw_redownload.py --symbols 600000,000001 --start 2024-01-01
# 全市场(读 STOCK_LIST csv,后台跑)
python3 raw_redownload.py --all --start 2024-01-01
env:
RAW_DIR 本地 raw 根(默认 /tmp/stock_dl/A股数据/日线数据/raw
SLEEP 每只请求间隔秒(默认 1.0,限速防封)
STOCK_LIST stock_basic_info csv 路径(--all 读,默认拉到本地的 csv)
"""
import argparse
import csv
import logging
import os
import sys
import time
# 直连:进程级 unset 代理(用户约束)
for _k in ["HTTP_PROXY", "HTTPS_PROXY", "http_proxy", "https_proxy", "ALL_PROXY", "all_proxy"]:
os.environ.pop(_k, None)
os.environ["NO_PROXY"] = "*"
os.environ["no_proxy"] = "*"
import pandas as pd # noqa: E402
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s")
log = logging.getLogger("raw_redl")
RAW_DIR = os.environ.get("RAW_DIR", "/tmp/stock_dl/A股数据/日线数据/raw")
SLEEP = float(os.environ.get("SLEEP", "1.0"))
DEFAULT_STOCK_LIST = "/tmp/stock_dl/A股数据/stock_info/stock_basic_info_raw_20260326_113530.csv"
def prefix_for(code: str) -> str:
"""sh/sz 前缀(与 datareader.guess_exchange 一致)。"""
return "sh" if code.startswith(("60", "68", "51", "56", "58")) else "sz"
def download_one(ak, code: str, start: str, end: str, adjust: str = ""):
"""新浪源拉日线(adjust="" raw / "qfq" 前复权),返回 (df, None) 或 (None, err)。"""
sym = f"{prefix_for(code)}{code}"
try:
df = ak.stock_zh_a_daily(
symbol=sym,
start_date=start.replace("-", ""),
end_date=end.replace("-", ""),
adjust=adjust,
)
except Exception as e: # noqa: BLE001
return None, f"{type(e).__name__}: {str(e)[:100]}"
if df is None or df.empty:
return None, "empty"
df["date"] = pd.to_datetime(df["date"])
df["year"] = df["date"].dt.year
return df, None
def save_one(code: str, df) -> int:
"""按 year 分组写 parquet(与 daily_dir 同构),返回写入行数。"""
n = 0
for year, g in df.groupby("year"):
ydir = os.path.join(RAW_DIR, str(int(year)))
os.makedirs(ydir, exist_ok=True)
out = os.path.join(ydir, f"{prefix_for(code)}{code}_daily.parquet")
g.drop(columns=["year"]).to_parquet(out)
n += len(g)
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 = []
with open(stock_list, encoding="utf-8", errors="replace") as f:
reader = csv.DictReader(f)
for row in reader:
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", help="逗号分隔代码,如 600000,000001")
ap.add_argument("--all", action="store_true", help="全市场(读 STOCK_LIST csv")
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")
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(先跑 run_daily_update.sh 拉取)", sl)
sys.exit(1)
codes = load_all_codes(sl)
log.info("--all 从 %s 读到 %d", sl, len(codes))
else:
ap.error("需指定 --symbols 或 --all")
import akshare as ak
import warnings
warnings.filterwarnings("ignore")
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
log.warning("[%d/%d] %s FAIL %s", i, len(codes), code, err)
else:
try:
n = save_one(code, df)
ok += 1
rows += n
log.info("[%d/%d] %s ok %d rows", i, len(codes), code, n)
except Exception as e: # noqa: BLE001
fail += 1
log.error("[%d/%d] %s SAVE FAIL %s", i, len(codes), code, e)
time.sleep(SLEEP) # 限速(用户约束)
log.info("=== 完成: ok=%d skip=%d fail=%d rows=%draw_dir=%s ===", ok, skipped, fail, rows, RAW_DIR)
if __name__ == "__main__":
main()