114a69e997
根因: dbbardata UNIQUE(symbol,exchange,interval,datetime) 按字符串字面比较, 多写入路径混用 'YYYY-MM-DD' 与 'YYYY-MM-DD 00:00:00' -> 同一交易日双行, INSERT OR REPLACE 不去重 -> 回测交易日翻倍/pivot duplicate/信号异常。 方案A (统一纯日期, 详见 Main Agent 诊断): - 新增 scripts/data_platform/dbbardata_utils.py: normalize_daily_dt(s) 取前 10 字符, None/短串安全 - 4 个日线写入脚本写入前调 helper: - bs_eod.py (sanguo-bs-eod 个股日线 baostock) - migrate_daily_baostock.py (历史迁移) - xt_eod.py (sanguo-xt-eod ETF/基金 xtata) - import_vnpy_daily_fast.py (NAS 日线 parquet 导入, 加防御) - TDD: tests/data_platform/test_dbbardata_utils.py 9 cases 全过 - 回归: tests/data_platform + tests/portfolio 199 passed 12 skipped peewee DateTimeField formats 含 '%Y-%m-%d' (阶段0 VPS 实测确认), 读纯日期不崩, 方案A 前提成立。 15min 干净, 不动 (分钟必须带时分)。只改日线 interval='d'。 数据层根治, 不在 provider 适配兜底 (用户铁律)。
103 lines
4.2 KiB
Python
103 lines
4.2 KiB
Python
#!/usr/bin/env python3
|
|
# -*- coding: utf-8 -*-
|
|
"""migrate_daily_baostock.py — 单元4: daily_baostock_full 拆分 -> staging (本地DB, 无网络)。
|
|
|
|
- OHLCV -> dbbardata_staging_daily (interval='d', 含退市, 治回测幸存者偏差)
|
|
exchange SH/SZ -> SSE/SZSE; date -> 'YYYY-MM-DD 00:00:00'; amount->turnover
|
|
- pe/pb/turn/pctChg/isST -> data/valuation_baostock/<year>.parquet (按年宽表)
|
|
- 全量读 + groupby year (避免 17 次全表扫; 内存峰值~5GB, VPS 16GB OK)
|
|
- staging 验证后单独合并 (INSERT OR REPLACE dbbardata + daily_baostock_full->_old)
|
|
"""
|
|
import sqlite3
|
|
from pathlib import Path
|
|
|
|
import pandas as pd
|
|
|
|
from dbbardata_utils import normalize_daily_dt
|
|
|
|
DB = Path(r"C:\sanguo_vnpy_v2\data\quant_trading.db")
|
|
VAL_DIR = Path(r"C:\sanguo_vnpy_v2\data\valuation_baostock")
|
|
EXC_MAP = {"SH": "SSE", "SZ": "SZSE"}
|
|
|
|
|
|
def log(m):
|
|
print(m, flush=True)
|
|
|
|
|
|
c = sqlite3.connect(str(DB), timeout=120)
|
|
c.execute("PRAGMA busy_timeout = 120000")
|
|
c.execute("PRAGMA synchronous = NORMAL")
|
|
|
|
cols = [r[1] for r in c.execute("PRAGMA table_info(daily_baostock_full)")]
|
|
log(f"daily_baostock_full cols({len(cols)}): {cols}")
|
|
total = c.execute("SELECT COUNT(*) FROM daily_baostock_full").fetchone()[0]
|
|
mn, mx = c.execute("SELECT MIN(date), MAX(date) FROM daily_baostock_full").fetchone()
|
|
log(f" rows={total} date {mn}~{mx}")
|
|
|
|
# staging 表
|
|
c.execute("DROP TABLE IF EXISTS dbbardata_staging_daily")
|
|
c.execute("""CREATE TABLE dbbardata_staging_daily (
|
|
symbol TEXT, exchange TEXT, datetime TEXT, interval TEXT,
|
|
volume REAL, turnover REAL, open_interest REAL,
|
|
open_price REAL, high_price REAL, low_price REAL, close_price REAL)""")
|
|
c.execute("CREATE INDEX idx_staging_daily_sym ON dbbardata_staging_daily(symbol, datetime)")
|
|
c.commit()
|
|
|
|
VAL_DIR.mkdir(parents=True, exist_ok=True)
|
|
|
|
log("流式读 daily_baostock_full (chunksize=200000) -> staging + pe/pb 累积 (避 OOM) ...")
|
|
val_by_year = {}
|
|
db_total = val_total = 0
|
|
n_chunk = 0
|
|
for chunk in pd.read_sql("SELECT * FROM daily_baostock_full", c, chunksize=200000):
|
|
n_chunk += 1
|
|
odb = pd.DataFrame({
|
|
"symbol": chunk["symbol"].values,
|
|
"exchange": chunk["exchange"].map(EXC_MAP).values,
|
|
# datetime 归一纯日期 (dbbardata 双行根治方案A)
|
|
"datetime": chunk["date"].astype(str).map(normalize_daily_dt).values,
|
|
"interval": "d",
|
|
"volume": chunk["volume"].values,
|
|
"turnover": chunk["amount"].values,
|
|
"open_interest": 0.0,
|
|
"open_price": chunk["open"].values,
|
|
"high_price": chunk["high"].values,
|
|
"low_price": chunk["low"].values,
|
|
"close_price": chunk["close"].values,
|
|
})
|
|
c.executemany(
|
|
"INSERT INTO dbbardata_staging_daily VALUES (?,?,?,?,?,?,?,?,?,?,?)",
|
|
odb.itertuples(index=False, name=None))
|
|
c.commit()
|
|
db_total += len(odb)
|
|
chunk["_y"] = pd.to_datetime(chunk["date"]).dt.year
|
|
for yr, sub in chunk.groupby("_y"):
|
|
vdf = sub[["symbol", "exchange", "date", "peTTM", "psTTM", "pcfNcfTTM",
|
|
"pbMRQ", "turn", "pctChg", "isST"]]
|
|
val_by_year.setdefault(int(yr), []).append(vdf)
|
|
val_total += len(vdf)
|
|
if n_chunk % 10 == 0:
|
|
log(f" chunk#{n_chunk} db累计={db_total} val累计={val_total}")
|
|
log("写 valuation_baostock/<year>.parquet ...")
|
|
for yr in sorted(val_by_year):
|
|
df_y = pd.concat(val_by_year[yr], ignore_index=True).sort_values(["symbol", "date"])
|
|
df_y.to_parquet(VAL_DIR / f"{yr}.parquet", index=False)
|
|
log(f" {yr}: {len(df_y)} rows")
|
|
|
|
# 退市/在市抽样验证
|
|
log("\n[verify staging]")
|
|
for sym, label in [("000005", "退市"), ("000023", "退市"),
|
|
("600811", "退市"), ("600519", "在市"), ("000001", "在市")]:
|
|
r = c.execute(
|
|
"SELECT COUNT(*), MIN(datetime), MAX(datetime) "
|
|
"FROM dbbardata_staging_daily WHERE symbol=?", (sym,)).fetchone()
|
|
log(f" {sym}({label}): {r}")
|
|
log(f" staging distinct symbol: "
|
|
f"{c.execute('SELECT COUNT(DISTINCT symbol) FROM dbbardata_staging_daily').fetchone()[0]}")
|
|
|
|
c.close()
|
|
log(f"\nstaging dbbardata_staging_daily: {db_total} rows")
|
|
log(f"valuation_baostock: {val_total} rows, "
|
|
f"{len(list(VAL_DIR.glob('*.parquet')))} 年 parquet")
|
|
log("MIGRATE STAGING DONE (未合并主表, 验证 OK 后单独 merge)")
|