docs(data): 归档数据层验证产物 + 数据层总览README

- scripts/data_platform/_archive/legacy/: 归档20个独立探针/诊断/旧降级脚本(零引用验证)
- docs/archive/data/: 归档17个数据相关旧设计/plan/report(保留fusion spec作深读)
- docs/data-platform/README.md: 数据层单一权威记录(8节:架构/布局/源/管线/铁律/API/缺口/待办)
- 删除 _mootdx_depth_result.txt
- Phase2待办: 15m灌库链+旧回填import链(有测试/wrapper依赖,VPS schtask确认后归档)
This commit is contained in:
2026-07-29 10:11:38 +08:00
parent 1cc9126abb
commit c3e53fbef3
39 changed files with 436 additions and 0 deletions
@@ -1,854 +0,0 @@
# Plan 1: 数据层实施计划
> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking.
**Goal:** 移植 v1 数据层到 v2,适配 vnpy 4.4.0,提供统一 DataReader 输出 vnpy BarData,支撑后续回测/因子层。
**Architecture:** 继承 v1 的 validator/fallback/增量更新(已验证资产),重组为 4 个单一职责组件 + YAML 集中配置;修复 v1 的 BaoStock 超时坑;NAS 容器本地读 parquet/SQLite,不走 SMB。
**Tech Stack:** Python 3.10、pandas、pyarrowparquet)、vnpy 4.4.0 数据库模块、BaoStock、requests、PyYAML、pytest
## Global Constraints
- 不改造 vnpy 核心(只用其数据库模块读 DbBarData)
- 数据读 NAS `/volume1/stock/`(容器本地挂载,不走 SMB
- 继承 v1 资产:`~/.openclaw/sanguo_projects/sanguo_vnpy/data_platform/`
- 配置集中 YAML(v1 散落代码 → v2 `config/data_platform.yaml`
- TDD:每个 task 先写测试 → 失败 → 实现 → 通过 → commit
- v2 工作目录:`~/.openclaw/sanguo_projects/sanguo_vnpy_v2/`
---
## File Structure
| 文件 | 职责 | 来源 |
|------|------|------|
| `sanguo_data/__init__.py` | 包入口,导出公共接口 | 新建 |
| `sanguo_data/config.py` | YAML 配置加载 | 新建 |
| `sanguo_data/validator.py` | 7 条 fatal 校验 | copy v1 `data_platform/validator.py` |
| `sanguo_data/datafeed.py` | 多源接入 + fallback + BaoStock 超时 | 基于 v1 `data_platform/fallback.py` 改造 |
| `sanguo_data/datareader.py` | 统一读 parquet/SQLite → BarData | 新建 |
| `sanguo_data/datawriter.py` | 写 parquet + SQLite DbBarData | 基于 v1 `data_platform/import_vnpy_daily_fast.py` 改造 |
| `sanguo_data/scheduler.py` | 增量更新 + 断点续传 | 基于 v1 `data_platform/daily_all_update.py` 改造 |
| `config/data_platform.yaml` | 数据源/路径/限流配置 | 新建 |
| `tests/data/__init__.py` | 测试包 | 新建 |
| `tests/data/test_validator.py` | validator 测试 | 新建 |
| `tests/data/test_datareader.py` | DataReader 测试 | 新建 |
| `tests/data/test_datafeed.py` | DataFeed 测试(含 BaoStock 超时) | 新建 |
| `tests/data/test_config.py` | 配置加载测试 | 新建 |
| `tests/data/conftest.py` | 测试夹具(合成 BarData/parquet | 新建 |
---
## Task 1: 项目脚手架 + YAML 配置
**Files:**
- Create: `sanguo_data/__init__.py`, `sanguo_data/config.py`, `config/data_platform.yaml`
- Create: `tests/data/__init__.py`, `tests/data/test_config.py`
**Interfaces:**
- Produces: `load_config(path: str) -> DataConfig``DataConfig` 是 dataclass,含 `data_paths``data_sources``validation``performance` 字段
- [ ] **Step 1: 写失败测试**
```python
# tests/data/test_config.py
from sanguo_data.config import load_config, DataConfig
def test_load_config_returns_dataconfig(tmp_path):
yaml_content = """
data_paths:
daily_dir: /tmp/daily
minute_15_dir: /tmp/15min
vnpy_db: /tmp/quant.db
stock_list: /tmp/stock.csv
data_sources:
daily:
- name: eastmoney
enabled: true
interval: 4.0
validation:
price_positive: true
performance:
max_retries: 3
"""
p = tmp_path / "config.yaml"
p.write_text(yaml_content)
cfg = load_config(str(p))
assert isinstance(cfg, DataConfig)
assert cfg.data_paths["daily_dir"] == "/tmp/daily"
assert cfg.data_sources["daily"][0]["name"] == "eastmoney"
assert cfg.performance["max_retries"] == 3
```
- [ ] **Step 2: 运行测试确认失败**
Run: `pytest tests/data/test_config.py -v`
Expected: FAIL with "ModuleNotFoundError: sanguo_data.config"
- [ ] **Step 3: 实现 config.py**
```python
# sanguo_data/config.py
from dataclasses import dataclass
import yaml
@dataclass(frozen=True)
class DataConfig:
data_paths: dict
data_sources: dict
validation: dict
performance: dict
def load_config(path: str) -> DataConfig:
with open(path, "r", encoding="utf-8") as f:
raw = yaml.safe_load(f)
return DataConfig(
data_paths=raw.get("data_paths", {}),
data_sources=raw.get("data_sources", {}),
validation=raw.get("validation", {}),
performance=raw.get("performance", {}),
)
```
```python
# sanguo_data/__init__.py
from .config import DataConfig, load_config
__all__ = ["DataConfig", "load_config"]
```
- [ ] **Step 4: 创建 config/data_platform.yaml**
```yaml
# config/data_platform.yaml
data_paths:
daily_dir: /volume1/stock/A股数据/日线数据/daily
minute_15_dir: /volume1/stock/minute_kline/15min
vnpy_db: /volume1/stock/sanguo_vnpy/data/quant_trading.db
stock_list: /volume1/stock/A股数据/stock_info/stock_basic_info_raw_20260326_113530.csv
data_sources:
daily:
- name: eastmoney
enabled: true
interval: 4.0
- name: baostock
enabled: true
interval: 0.0
timeout: 30
- name: tencent
enabled: true
interval: 0.0
minute_15:
- name: eastmoney
enabled: true
interval: 4.0
validation:
price_positive: true
ohlc_consistency: true
no_future_dates: true
performance:
request_interval: 0.3
max_retries: 3
fail_window: 100
fail_threshold: 0.8
```
- [ ] **Step 5: 运行测试确认通过**
Run: `pytest tests/data/test_config.py -v`
Expected: PASS
- [ ] **Step 6: Commit**
```bash
git add sanguo_data/ config/data_platform.yaml tests/data/
git commit -m "feat(data): 脚手架 + YAML 配置加载"
```
---
## Task 2: Validatorcopy v1 + 测试)
**Files:**
- Create: `sanguo_data/validator.py`copy v1
- Create: `tests/data/test_validator.py`, `tests/data/conftest.py`
**Interfaces:**
- Produces: `validate_daily(df: pd.DataFrame) -> pd.DataFrame`(过滤非法行)
- [ ] **Step 1: copy v1 validator.py**
```bash
cp ~/.openclaw/sanguo_projects/sanguo_vnpy/data_platform/validator.py \
~/.openclaw/sanguo_projects/sanguo_vnpy_v2/sanguo_data/validator.py
```
读 copy 后的文件,确认 v1 原有校验函数名。如签名与下方测试不符,**以 v1 实际签名为准**调整。
- [ ] **Step 2: 写失败测试(合成数据)**
```python
# tests/data/conftest.py
import pandas as pd
import pytest
@pytest.fixture
def good_daily_df():
return pd.DataFrame({
"date": ["2026-01-01", "2026-01-02"],
"open": [10.0, 11.0], "high": [10.5, 11.5],
"low": [9.8, 10.8], "close": [10.2, 11.2],
"volume": [10000, 12000],
})
@pytest.fixture
def bad_daily_df():
return pd.DataFrame({
"date": ["2026-01-01"],
"open": [0.0], "high": [0.0], "low": [0.0], "close": [0.0],
"volume": [100],
})
```
```python
# tests/data/test_validator.py
from sanguo_data.validator import validate_daily
def test_validate_daily_keeps_good_rows(good_daily_df):
assert len(validate_daily(good_daily_df)) == 2
def test_validate_daily_drops_zero_price(bad_daily_df):
assert len(validate_daily(bad_daily_df)) == 0
```
- [ ] **Step 3: 运行测试**
Run: `pytest tests/data/test_validator.py -v`
Expected: FAIL(函数名不匹配)或 PASS(v1 直接可用)
- [ ] **Step 4: 适配导出接口**
如 v1 函数名/签名与测试不符,在 `validator.py` 末尾加薄适配(**不改 v1 校验逻辑**):
```python
# 适配层,不改 v1 校验逻辑
def validate_daily(df):
"""对外统一接口,委托 v1 校验规则"""
return _v1_validate(df) # 替换为 v1 实际函数名
```
- [ ] **Step 5: 运行确认通过**
Run: `pytest tests/data/test_validator.py -v`
Expected: PASS
- [ ] **Step 6: Commit**
```bash
git add sanguo_data/validator.py tests/data/test_validator.py tests/data/conftest.py
git commit -m "feat(data): 移植 v1 validator + 适配接口 + 测试"
```
---
## Task 3: DataReader — parquet 读取
**Files:**
- Create: `sanguo_data/datareader.py`
- Create: `tests/data/test_datareader.py`
**Interfaces:**
- Consumes: `DataConfig.data_paths["daily_dir"]`Task 1
- Produces: `read_parquet_daily(symbol: str, start: str, end: str, cfg: DataConfig) -> list[BarData]``_row_to_bar(symbol, row, interval) -> BarData`
- [ ] **Step 1: 写失败测试(合成 parquet)**
```python
# tests/data/test_datareader.py
import pandas as pd
from sanguo_data.config import DataConfig
from sanguo_data.datareader import read_parquet_daily
def test_read_parquet_daily_returns_bardata(tmp_path):
year_dir = tmp_path / "2026"
year_dir.mkdir()
df = pd.DataFrame({
"date": ["2026-01-05", "2026-01-06"],
"open": [10.0, 11.0], "high": [10.5, 11.5],
"low": [9.8, 10.8], "close": [10.2, 11.2],
"volume": [10000, 12000],
})
df.to_parquet(year_dir / "600000.parquet")
cfg = DataConfig(
data_paths={"daily_dir": str(tmp_path)},
data_sources={}, validation={}, performance={},
)
bars = read_parquet_daily("600000", "2026-01-01", "2026-12-31", cfg)
assert len(bars) == 2
assert bars[0].symbol == "600000"
assert bars[0].open_price == 10.0
```
- [ ] **Step 2: 运行确认失败**
Run: `pytest tests/data/test_datareader.py::test_read_parquet_daily_returns_bardata -v`
Expected: FAIL "ModuleNotFoundError"
- [ ] **Step 3: 实现 datareader.pyparquet 部分)**
```python
# sanguo_data/datareader.py
import pandas as pd
from datetime import datetime
from pathlib import Path
from vnpy.trader.object import BarData
from vnpy.trader.constant import Exchange, Interval
def read_parquet_daily(symbol: str, start: str, end: str, cfg) -> list[BarData]:
daily_dir = Path(cfg.data_paths["daily_dir"])
start_dt = datetime.strptime(start, "%Y-%m-%d")
end_dt = datetime.strptime(end, "%Y-%m-%d")
bars: list[BarData] = []
for year in range(start_dt.year, end_dt.year + 1):
f = daily_dir / str(year) / f"{symbol}.parquet"
if not f.exists():
continue
df = pd.read_parquet(f)
for _, row in df.iterrows():
d = pd.to_datetime(row["date"])
if start_dt <= d <= end_dt:
bars.append(_row_to_bar(symbol, row, Interval.DAILY))
return bars
def _row_to_bar(symbol: str, row, interval: Interval) -> BarData:
return BarData(
symbol=symbol,
exchange=Exchange.SSE, # Task 4 改为 guess_exchange
datetime=pd.to_datetime(row["date"]).to_pydatetime(),
interval=interval,
open_price=float(row["open"]),
high_price=float(row["high"]),
low_price=float(row["low"]),
close_price=float(row["close"]),
volume=float(row["volume"]),
gateway_name="DATA",
)
```
- [ ] **Step 4: 运行确认通过**
Run: `pytest tests/data/test_datareader.py::test_read_parquet_daily_returns_bardata -v`
Expected: PASS
- [ ] **Step 5: Commit**
```bash
git add sanguo_data/datareader.py tests/data/test_datareader.py
git commit -m "feat(data): DataReader parquet 读取 → BarData"
```
---
## Task 4: DataReader — SQLite DbBarData + 交易所判断
**Files:**
- Modify: `sanguo_data/datareader.py`(加 `read_db_daily` + `guess_exchange``_row_to_bar` 改用 `guess_exchange`
- Modify: `tests/data/test_datareader.py`
**Interfaces:**
- Produces: `read_db_daily(symbol, start, end, cfg) -> list[BarData]`(用 vnpy 4.4.0 数据库模块),`guess_exchange(symbol) -> Exchange`
- [ ] **Step 1: 写失败测试**
```python
# 追加到 tests/data/test_datareader.py
from sanguo_data.datareader import guess_exchange
def test_guess_exchange_sh():
assert guess_exchange("600000").value == "SSE"
def test_guess_exchange_sz():
assert guess_exchange("000001").value == "SZSE"
```
- [ ] **Step 2: 运行确认失败**
Run: `pytest tests/data/test_datareader.py::test_guess_exchange_sh -v`
Expected: FAIL "ImportError"
- [ ] **Step 3: 实现 guess_exchange + read_db_daily,并让 `_row_to_bar` 用 guess_exchange**
```python
# 追加到 sanguo_data/datareader.py;并把 _row_to_bar 的 exchange 改为 guess_exchange(symbol)
from vnpy.trader.database import get_database
def guess_exchange(symbol: str) -> Exchange:
"""按代码前缀判断交易所:6/68/5x→SSE0/3/15x→SZSE"""
if symbol.startswith(("60", "68", "51", "56", "58")):
return Exchange.SSE
if symbol.startswith(("00", "30", "15")):
return Exchange.SZSE
return Exchange.SSE
def read_db_daily(symbol: str, start: str, end: str, cfg) -> list[BarData]:
db = get_database()
start_dt = datetime.strptime(start, "%Y-%m-%d")
end_dt = datetime.strptime(end, "%Y-%m-%d")
return db.load_bar_data(
symbol=symbol,
exchange=guess_exchange(symbol),
interval=Interval.DAILY,
start=start_dt,
end=end_dt,
)
```
`_row_to_bar` 内的 `exchange=Exchange.SSE` 改为 `exchange=guess_exchange(symbol)`
> **Spike 检查点**`get_database()` 与 `load_bar_data` 签名需对照 vnpy 4.4.0 `vnpy/trader/database.py`。Task 8 spike 验证,若变更回头修正。
- [ ] **Step 4: 运行确认通过**
Run: `pytest tests/data/test_datareader.py -v`
Expected: PASS
- [ ] **Step 5: 真实 NAS 数据冒烟(手工)**
```bash
python -c "
from sanguo_data.config import load_config
from sanguo_data.datareader import read_db_daily
cfg = load_config('config/data_platform.yaml')
bars = read_db_daily('600000', '2026-01-01', '2026-06-30', cfg)
print(f'读取 {len(bars)} 条')
"
```
Expected: N > 0。若 0,检查 vnpy_db 路径与 4.4.0 接口。
- [ ] **Step 6: Commit**
```bash
git add sanguo_data/datareader.py tests/data/test_datareader.py
git commit -m "feat(data): DataReader SQLite + 交易所判断 + vnpy 4.4.0 spike 点"
```
---
## Task 5: DataFeed — 多源 fallback + BaoStock 超时(v1 卡死坑修复)
**Files:**
- Create: `sanguo_data/datafeed.py`(基于 v1 `fallback.py`
- Create: `tests/data/test_datafeed.py`
**Interfaces:**
- Produces: `fetch_daily(symbol, start, end, cfg) -> pd.DataFrame``fetch_with_fallback(symbol, start, end, sources) -> pd.DataFrame`
- [ ] **Step 1: 写失败测试(mock + 超时)**
```python
# tests/data/test_datafeed.py
import time
import pandas as pd
import pytest
from unittest.mock import patch
from sanguo_data.datafeed import fetch_with_fallback, _fetch_baostock_with_timeout
def test_fetch_with_fallback_uses_second_when_first_fails():
df_good = pd.DataFrame({"date": ["2026-01-01"], "open": [10.0]})
with patch("sanguo_data.datafeed._fetch_eastmoney", side_effect=Exception("limit")), \
patch("sanguo_data.datafeed._fetch_baostock", return_value=df_good):
out = fetch_with_fallback("600000", "2026-01-01", "2026-01-02", ["eastmoney", "baostock"])
assert len(out) == 1
def test_baostock_timeout_does_not_hang():
"""v1 卡死坑修复验证:超时必须返回,不能无限挂起"""
start = time.time()
with patch("sanguo_data.datafeed._fetch_baostock_raw", side_effect=lambda *a: time.sleep(60)):
with pytest.raises(TimeoutError):
_fetch_baostock_with_timeout("600000", "2026-01-01", "2026-01-02", timeout=2)
assert time.time() - start < 5
```
- [ ] **Step 2: 运行确认失败**
Run: `pytest tests/data/test_datafeed.py -v`
Expected: FAIL "ModuleNotFoundError"
- [ ] **Step 3: 实现 datafeed.pyv1 fallback 模式 + BaoStock 超时包装)**
```python
# sanguo_data/datafeed.py
import pandas as pd
from multiprocessing import Process, Queue
from sanguo_data.config import DataConfig
def fetch_with_fallback(symbol, start, end, sources: list[str]) -> pd.DataFrame:
fetchers = {
"eastmoney": _fetch_eastmoney,
"baostock": lambda s, a, b: _fetch_baostock_with_timeout(s, a, b, timeout=30),
"tencent": _fetch_tencent,
}
last_err = None
for name in sources:
try:
df = fetchers[name](symbol, start, end)
if df is not None and len(df) > 0:
return df
except Exception as e:
last_err = e
continue
raise RuntimeError(f"all sources failed: {last_err}")
def fetch_daily(symbol, start, end, cfg: DataConfig) -> pd.DataFrame:
sources = [s["name"] for s in cfg.data_sources.get("daily", []) if s.get("enabled", True)]
return fetch_with_fallback(symbol, start, end, sources)
def _fetch_baostock_with_timeout(symbol, start, end, timeout):
"""子进程隔离 BaoStock(修复 v1 无超时卡死坑)"""
q = Queue()
def worker():
try:
q.put(_fetch_baostock_raw(symbol, start, end))
except Exception as e:
q.put(e)
p = Process(target=worker)
p.start()
p.join(timeout)
if p.is_alive():
p.terminate(); p.join()
raise TimeoutError(f"baostock timeout after {timeout}s")
res = q.get()
if isinstance(res, Exception):
raise res
return res
def _fetch_baostock_raw(symbol, start, end):
"""从 v1 data_platform/fallback.py copy BaoStock 接入(baostock.query_history_k_data_plus"""
raise NotImplementedError("copy from v1 data_platform/fallback.py")
def _fetch_eastmoney(symbol, start, end):
raise NotImplementedError("copy from v1 data_platform/fallback.py")
def _fetch_tencent(symbol, start, end):
raise NotImplementedError("copy from v1 data_platform/fallback.py")
```
> **执行注意**`_fetch_baostock_raw` / `_fetch_eastmoney` / `_fetch_tencent` 从 v1 `data_platform/fallback.py` copy 接入逻辑。**超时包装是新增修复,不 copy。**
- [ ] **Step 4: 运行确认通过**
Run: `pytest tests/data/test_datafeed.py -v`
Expected: PASS
- [ ] **Step 5: Commit**
```bash
git add sanguo_data/datafeed.py tests/data/test_datafeed.py
git commit -m "feat(data): DataFeed 多源 fallback + BaoStock 超时修复"
```
---
## Task 6: DataWriter — 原子写 parquet + vnpy SQLite
**Files:**
- Create: `sanguo_data/datawriter.py`
- Create: `tests/data/test_datawriter.py`
**Interfaces:**
- Produces: `write_daily(symbol, df, cfg) -> None``atomic_write_parquet(path, df) -> None`
- [ ] **Step 1: 写失败测试**
```python
# tests/data/test_datawriter.py
import pandas as pd
from sanguo_data.config import DataConfig
from sanguo_data.datawriter import write_daily, atomic_write_parquet
def test_atomic_write_parquet(tmp_path):
f = tmp_path / "2026" / "600000.parquet"
df = pd.DataFrame({"date": ["2026-01-01"], "open": [10.0]})
atomic_write_parquet(str(f), df)
assert f.exists()
assert not list(tmp_path.glob("*.tmp"))
def test_write_daily_writes_parquet_and_db(tmp_path, monkeypatch):
cfg = DataConfig(
data_paths={"daily_dir": str(tmp_path / "daily"), "vnpy_db": str(tmp_path / "q.db")},
data_sources={}, validation={}, performance={},
)
df = pd.DataFrame({"date": ["2026-01-01"], "open": [10.0], "high": [10.0],
"low": [10.0], "close": [10.0], "volume": [100]})
called = {}
monkeypatch.setattr("sanguo_data.datawriter._save_to_vnpy_db", lambda bars, cfg: called.setdefault("bars", bars))
write_daily("600000", df, cfg)
assert (tmp_path / "daily" / "2026" / "600000.parquet").exists()
assert len(called["bars"]) == 1
```
- [ ] **Step 2: 运行确认失败**
Run: `pytest tests/data/test_datawriter.py -v`
Expected: FAIL
- [ ] **Step 3: 实现 datawriter.py**
```python
# sanguo_data/datawriter.py
import os
import pandas as pd
from pathlib import Path
from vnpy.trader.object import BarData
from vnpy.trader.constant import Interval
from sanguo_data.datareader import _row_to_bar
from sanguo_data.config import DataConfig
def atomic_write_parquet(path: str, df: pd.DataFrame) -> None:
p = Path(path)
p.parent.mkdir(parents=True, exist_ok=True)
tmp = str(p) + ".tmp"
df.to_parquet(tmp)
os.replace(tmp, str(p)) # 原子替换
def write_daily(symbol: str, df: pd.DataFrame, cfg: DataConfig) -> None:
# 1) parquet 增量合并(按年分区,去重保留最新)
for year, group in df.groupby(df["date"].str[:4]):
f = Path(cfg.data_paths["daily_dir"]) / year / f"{symbol}.parquet"
if f.exists():
old = pd.read_parquet(f)
combined = pd.concat([old, group]).drop_duplicates("date", keep="last")
else:
combined = group
atomic_write_parquet(str(f), combined)
# 2) vnpy SQLite
bars = [_row_to_bar(symbol, row, Interval.DAILY) for _, row in df.iterrows()]
_save_to_vnpy_db(bars, cfg)
def _save_to_vnpy_db(bars: list[BarData], cfg: DataConfig) -> None:
from vnpy.trader.database import get_database
db = get_database()
db.save_bar_data(bars) # spike 验证签名
```
- [ ] **Step 4: 运行确认通过**
Run: `pytest tests/data/test_datawriter.py -v`
Expected: PASS
- [ ] **Step 5: Commit**
```bash
git add sanguo_data/datawriter.py tests/data/test_datawriter.py
git commit -m "feat(data): DataWriter 原子写 parquet + vnpy SQLite"
```
---
## Task 7: UpdateScheduler — 增量更新 + 断点续传 + 熔断
**Files:**
- Create: `sanguo_data/scheduler.py`(基于 v1 `daily_all_update.py`
- Create: `tests/data/test_scheduler.py`
**Interfaces:**
- Consumes: fetch_dailyTask 5+ validate_dailyTask 2+ write_dailyTask 6
- Produces: `run_daily_update(cfg, symbols=None) -> UpdateReport`
- [ ] **Step 1: 写失败测试(断点续传)**
```python
# tests/data/test_scheduler.py
import json
import pandas as pd
from unittest.mock import patch
from sanguo_data.config import DataConfig
from sanguo_data.scheduler import run_daily_update, UpdateReport
def test_run_daily_update_skips_completed_on_resume(tmp_path):
progress_file = tmp_path / "progress.json"
progress_file.write_text('{"600000": "done"}')
cfg = DataConfig(
data_paths={"daily_dir": str(tmp_path), "vnpy_db": str(tmp_path / "q.db"),
"progress_file": str(progress_file)},
data_sources={"daily": [{"name": "eastmoney", "enabled": True}]},
validation={}, performance={},
)
with patch("sanguo_data.scheduler.fetch_daily", return_value=pd.DataFrame({
"date": ["2026-01-01"], "open": [10.0], "high": [10.0],
"low": [10.0], "close": [10.0], "volume": [100]})) as m_fetch, \
patch("sanguo_data.scheduler.write_daily") as m_write:
report = run_daily_update(cfg, symbols=["600000"])
assert m_fetch.call_count == 0 # 已 done,跳过
assert isinstance(report, UpdateReport)
assert report.skipped == 1
```
- [ ] **Step 2: 运行确认失败**
Run: `pytest tests/data/test_scheduler.py -v`
Expected: FAIL
- [ ] **Step 3: 实现 scheduler.py**
```python
# sanguo_data/scheduler.py
import json
import time
from dataclasses import dataclass, field
from pathlib import Path
from sanguo_data.config import DataConfig
from sanguo_data.datafeed import fetch_daily
from sanguo_data.validator import validate_daily
from sanguo_data.datawriter import write_daily
@dataclass
class UpdateReport:
total: int = 0
success: int = 0
failed: int = 0
skipped: int = 0
failures: list = field(default_factory=list)
def run_daily_update(cfg: DataConfig, symbols: list[str] | None = None) -> UpdateReport:
progress_path = Path(cfg.data_paths.get("progress_file", "progress.json"))
progress = json.loads(progress_path.read_text()) if progress_path.exists() else {}
symbols = symbols or _load_stock_list(cfg)
report = UpdateReport(total=len(symbols))
fail_window = cfg.performance.get("fail_window", 100)
fail_threshold = cfg.performance.get("fail_threshold", 0.8)
for sym in symbols:
if progress.get(sym) == "done":
report.skipped += 1
continue
try:
df = fetch_daily(sym, _last_date(sym, cfg), _today(), cfg)
df = validate_daily(df)
if len(df) > 0:
write_daily(sym, df, cfg)
progress[sym] = "done"
progress_path.write_text(json.dumps(progress, ensure_ascii=False))
report.success += 1
except Exception as e:
report.failed += 1
report.failures.append({"symbol": sym, "error": str(e)})
checked = report.success + report.failed
if checked >= fail_window and report.failed / max(checked, 1) > fail_threshold:
report.failures.append({"error": "FAIL_THRESHOLD_REACHED, abort"})
break
time.sleep(cfg.performance.get("request_interval", 0.3))
return report
def _load_stock_list(cfg: DataConfig) -> list[str]:
"""从 v1 data_platform/daily_all_update.py copy 全市场股票列表读取"""
raise NotImplementedError("copy from v1")
def _last_date(symbol: str, cfg: DataConfig) -> str:
return "2020-01-01" # 简化,实际读 parquet 最后日期
def _today() -> str:
return "2026-07-05"
```
- [ ] **Step 4: 运行确认通过**
Run: `pytest tests/data/test_scheduler.py -v`
Expected: PASS
- [ ] **Step 5: Commit**
```bash
git add sanguo_data/scheduler.py tests/data/test_scheduler.py
git commit -m "feat(data): UpdateScheduler 增量 + 断点续传 + 熔断"
```
---
## Task 8: vnpy 4.4.0 接口 spike + 端到端冒烟
**Files:**
- Create: `tests/data/test_spike_vnpy44.py`
- Modify: 按 spike 结果修正 Task 4/6 的 `get_database()` 调用
- [ ] **Step 1: 写 spike 测试(探查 vnpy 4.4.0 接口)**
```python
# tests/data/test_spike_vnpy44.py
"""Spike: 验证 vnpy 4.4.0 数据库接口。无 NAS 数据时可 skip。"""
import datetime
import pytest
from vnpy.trader.database import get_database
from vnpy.trader.constant import Interval, Exchange
def test_vnpy44_database_interface():
db = get_database()
assert hasattr(db, "load_bar_data")
assert hasattr(db, "save_bar_data")
bars = db.load_bar_data(
symbol="600000", exchange=Exchange.SSE, interval=Interval.DAILY,
start=datetime.datetime(2026, 1, 1), end=datetime.datetime(2026, 6, 30),
)
assert isinstance(bars, list)
```
- [ ] **Step 2: 容器内运行 spike**
```bash
docker exec sanguo_vnpy_v2 pytest tests/data/test_spike_vnpy44.py -v
```
Expected: PASS(接口符合假设)或 FAIL(4.4.0 接口变更 → 记录差异,修正 Task 4/6 调用)。
- [ ] **Step 3: 端到端冒烟**
```bash
python -c "
from sanguo_data.config import load_config
from sanguo_data.datareader import read_db_daily
from sanguo_data.scheduler import run_daily_update
cfg = load_config('config/data_platform.yaml')
bars = read_db_daily('600000', '2026-01-01', '2026-06-30', cfg)
print(f'读取 {len(bars)} 条')
report = run_daily_update(cfg, symbols=['600000'])
print(f'更新: {report.success} 成功, {report.failed} 失败')
"
```
Expected: 读取 N 条 + 更新成功。**这是 Plan 1 最终验收。**
- [ ] **Step 4: spike 发现差异则修正 Task 4/6,重测**
- [ ] **Step 5: Commit**
```bash
git add tests/data/test_spike_vnpy44.py
git commit -m "test(data): vnpy 4.4.0 接口 spike + 端到端冒烟"
```
---
## Self-Review(写完内联检查)
**1. Spec coverage(对照 design §2 数据层):**
- DataFeed(多源 fallback + BaoStock 超时)→ Task 5 ✅
- Validator7 条 fatal)→ Task 2 ✅
- DataWriter/Reader(双存储)→ Task 3/4/6 ✅
- UpdateScheduler(增量 + 断点续传 + 熔断)→ Task 7 ✅
- 配置集中 YAML → Task 1 ✅
- 早期 spikevnpy 4.4.0 接口)→ Task 4 + Task 8 ✅
- BaoStock 超时坑修复 → Task 5 ✅
**2. Placeholder 扫描:**
- `_fetch_baostock_raw` / `_fetch_eastmoney` / `_fetch_tencent` / `_load_stock_list``NotImplementedError("copy from v1 ...")`——这是**有意的执行指引**(指明从 v1 哪个文件 copy),不是 plan 占位。执行 subagent 按 v1 `fallback.py`/`daily_all_update.py` copy。
- 其余步骤均有完整代码/命令。
**3. 类型一致性:**
- `DataConfig`Task 1)在 Task 2-7 一致使用 ✅
- `BarData``_row_to_bar`Task 3)在 Task 4/6 复用 ✅
- `guess_exchange`Task 4)在 Task 6 间接复用 ✅
- `fetch_daily` / `validate_daily` / `write_daily` 跨 Task 一致 ✅
**4. 风险:**
- vnpy 4.4.0 `get_database()` / `load_bar_data` / `save_bar_data` 签名需 Task 8 spike 验证,若变更修正 Task 4/6。已在 plan 显式标注 spike 检查点。
@@ -1,221 +0,0 @@
# P0 数据补全实现计划(历史成份股 + ETF 全市场 + 退市 K 线)
> **For agentic workers:** 用 superpowers:subagent-driven-development 或 executing-plans 执行。Steps 用 `[ ]` 跟踪。
**Goal:** 补齐治幸存者偏差 + 策略核心缺口三类数据,落到 VPS 本地。
**Architecture:** 各源采集脚本 → staging parquet → 验证探针 → 合并主库;baostock 单登录守 48000/天;dbbardata 不动。
**Tech Stack:** python3.10 / akshare / baostock / xtquant(xtdata)/ pandas / pyarrow / sqlite3
---
## Global Constraints(所有 task 隐含)
- **baostock 单进程单登录**,不并发(防黑名单,日 ≤48000 query)
- **直连不走代理**:`$env:http_proxy=''; $env:https_proxy=''; $env:all_proxy=''`
- **dbbardata 不破坏**:只 INSERT OR REPLACE `daily_baostock_full` / 新表,不动 dbbardata 既有行
- **优先 baostock + miniQMT(xtdata)**
- **staging → 验证探针 → 合并主库**(用户铁律,不直接写主库)
- Windows VPS 49.232.102.198,`C:\Python310\python.exe -X utf8`,schtasks `/ru SYSTEM`
- 输出根:`C:\sanguo_vnpy_v2\data\`
---
## File Structure
| 文件 | 责任 |
|---|---|
| `scripts/data_platform/index_const_hist_download.py`(新) | 历史成份股采集(akshare 国证 + 新浪 + baostock 补时点) |
| `scripts/data_platform/build_daily_from_xtdata.py`(改 :40) | ETF universe 扩展(一次性全量) |
| `scripts/data_platform/daily_update_xtdata.py`(改 :114) | ETF 每日增量 universe |
| `scripts/data_platform/baostock_delisted_download.py`(新) | 退市股列表 + K 线采集 |
| `scripts/data_platform/import_delisted_to_db.py`(新) | 退市 K 线灌 `daily_baostock_full` |
| 各 `*_wrapper.ps1` + schtask | 部署 |
---
## Task 1: 历史成份股采集(治幸存者偏差)
**Files:** Create `scripts/data_platform/index_const_hist_download.py`;Output `data/index_const_hist/{code}.parquet`
**Interfaces:**
- Consumes: akshare `index_detail_hist_cni(symbol)` + `index_detail_hist_adjust_cni(symbol)`(国证源);新浪 `vII_HistoryComponent`(pandas.read_html, gb2312);baostock `query_hs300/zz500/sz50_stocks(date)`
- Produces: `data/index_const_hist/{code}.parquet`(列:`updateDate/index_code/code/code_name/adjust_type`);并集 = 曾经入选集
**指数清单:**
- 深证/国证(akshare 国证源):399001 / 399006 / 399101 / 399005 / 399330
- 中证(新浪):000852(中证1000)/ 932000(中证2000)/ 000300(交叉校验)/ 000016(上证50)
- baostock 已有(300/500/50 在 `bs_index_constituent`):Task1 补时点序列到同 schema
- [ ] **1.1 探针:akshare 国证源 hist 版**
```python
import akshare as ak
df = ak.index_detail_hist_cni(symbol="399101") # 历史样本(日期/样本代码/权重)
print(df.columns.tolist(), len(df), df.head(3))
adj = ak.index_detail_hist_adjust_cni(symbol="399101") # 调样记录(调整类型 OLD/+/-)
print(adj.columns.tolist(), len(adj))
```
预期:hist 有日期+样本+权重;adjust 有调整类型。**陷阱:必须 hist 版**(`index_detail_cni` 非 hist 版 2025-11-25 起只近期);`ak.index_stock_hist` 已下线别用。
- [ ] **1.2 探针:新浪中证历史成份**
```python
import pandas as pd
url = "http://vip.stock.finance.sina.com.cn/corp/go.php/vII_HistoryComponent/indexid/000852.phtml"
df = pd.read_html(url, encoding="gb2312")[0]
print(df.columns.tolist(), len(df), df.head(3))
```
预期:品种代码/品种名称/纳入日期/剔除日期(空=至今在列),含 *ST/退市股。
- [ ] **1.3 实现 `index_const_hist_download.py`**:三路采集 → 统一 schema(`updateDate/index_code/code/code_name/adjust_type`)→ 写 `data/index_const_hist/{code}.parquet`。串行 `time.sleep(0.8)`(akshare/新浪防封),单进程。环境变量 `BS_INDEX_HIST_OUT_DIR` 覆盖默认 Mac 路径(同 Day1 wrapper 模式)。
- [ ] **1.4 验证探针**:每指数 parquet 行数 + 抽样 3 行;**幸存者偏差校验** = 并集 `distinct code` 数 > 当前成份股数(证明含被踢股,例如 399101 并集 > 958 当前)。
- [ ] **1.5 wrapper + schtask**:`index_const_hist_wrapper.ps1`(设 OUT_DIR + utf8 + unset proxy + log);schtask `sanguo-index-hist` `/sc monthly /mo 2`(半年度调样后,6/12 月)`/ru SYSTEM`
- [ ] **1.6 commit**:`git add scripts/data_platform/index_const_hist_download.py scripts/data_platform/index_const_hist_wrapper.ps1 && git commit -m "feat(data): 历史成份股采集(治幸存者偏差,国证+新浪+baostock)"`
---
## Task 2: ETF 全市场日线
**Files:** Modify `scripts/data_platform/build_daily_from_xtdata.py:40` + `daily_update_xtdata.py:114`
**Interfaces:**
- Consumes: xtdata `get_stock_list_in_sector('沪深A股'/'沪深ETF'/'沪深基金')` + `get_market_data_ex(dividend_type='front')`
- Produces: 全市场 ETF(~1000 只)日线**前复权**,落 parquet/dbbardata(复用现有 xtdata 管线)
- [ ] **2.1 探针:ETF universe + 1 只 K 线**
```python
from xtquant import xtdata as xd
etf = xd.get_stock_list_in_sector('沪深ETF') or []
fund = xd.get_stock_list_in_sector('沪深基金') or []
a = xd.get_stock_list_in_sector('沪深A股') or []
u = list(set(a + etf + fund))
print(f"A={len(a)} ETF={len(etf)} fund={len(fund)} union={len(u)}")
r = xd.get_market_data_ex([], ['510300.SH'], period='1d',
start_time='20240101', end_time='20260721', dividend_type='front')
df = r.get('510300.SH')
print('510300 bars:', 0 if df is None else len(df), '| tail close:', None if df is None else df['close'].iloc[-1])
```
预期:ETF ~1000,union > A 股数;510300 前复权日线有值,close 非 NaN。
- [ ] **2.2 改 universe**:`build_daily_from_xtdata.py:40``daily_update_xtdata.py:114`
```python
u = xd.get_stock_list_in_sector("沪深A股") or []
```
改为
```python
u = list(set(
(xd.get_stock_list_in_sector("沪深A股") or []) +
(xd.get_stock_list_in_sector("沪深ETF") or []) +
(xd.get_stock_list_in_sector("沪深基金") or [])
))
```
保留 `dividend_type='front'`(前复权,§13 默认)。
- [ ] **2.3 全量下载 ETF**:跑改后的 `build_daily_from_xtdata.py`(走现有 xtdata 管线,**无限流**)→ parquet。
- [ ] **2.4 验证**:ETF 数 + 抽样(510300/513050/159919)+ 前复权 close 非 NaN + 日期范围。
- [ ] **2.5 schtask**:复用 `sanguo-daily-update`(universe 扩展后自动含 ETF,无需新 schtask)。
- [ ] **2.6 commit**:`git commit -m "feat(data): ETF 全市场日线(xtdata universe 扩展+前复权)"`
---
## Task 3: 退市股 K 线(反幸存者偏差核心)
**Files:** Create `scripts/data_platform/baostock_delisted_download.py` + `import_delisted_to_db.py`;Output → `daily_baostock_full`
**Interfaces:**
- Consumes: baostock `query_all_stock(day)` + `query_stock_basic(code)`(status + 退市日期)+ `query_history_k_data_plus(code, fields, adjustflag=3)`
- Produces: 退市股 K 线 INSERT OR REPLACE `daily_baostock_full`(18 列,复用 `parse_baostock_code`)
**范围:** 近 5 年退市(退市日期 ≥ 2021;守 48000/天;退市股分天跑)
- [ ] **3.1 探针:退市股列表字段**
```python
import baostock as bs, pandas as pd
bs.login()
rs = bs.query_all_stock(day="2026-07-18")
rows = []
while (rs.error_code == '0') & rs.next():
rows.append(rs.get_row_data())
df = pd.DataFrame(rows, columns=rs.fields)
print('query_all_stock fields:', rs.fields, '| rows:', len(df))
rs2 = bs.query_stock_basic(code="sh.600000")
b = []
while (rs2.error_code == '0') & rs2.next():
b.append(rs2.get_row_data())
print('query_stock_basic fields:', rs2.fields, '| sample:', b[0] if b else None)
bs.logout()
```
预期:`query_stock_basic``type`(1股)/`status`(1上市 0退市)/`outDate`(退市日期)。筛 `status=0 & outDate>='2021-01-01'`
- [ ] **3.2 实现 `baostock_delisted_download.py`**:
- 遍历全 code(或 `query_all_stock` 多日并集)→ `query_stock_basic``status=0 & outDate>='2021-01-01'` → 退市股列表
- 逐只 `query_history_k_data_plus(code, start_date='1990-01-01', end_date=outDate, fields=18字段, adjustflag=3)` → staging `data/delisted_kline/{code}.parquet`
- 单进程单登录,`time.sleep` 守预算,marker 断点续传(复用 Day1 模板),DAILY_LIMIT 计数器
- [ ] **3.3 `import_delisted_to_db.py`**:staging → INSERT OR REPLACE `daily_baostock_full`(复用 `parse_baostock_code` sh.600000→600000+SH + `executemany`,WAL + busy_timeout=60000,同 `import_baostock_to_db.py`)。**dbbardata 不碰**。
- [ ] **3.4 验证探针**:退市股数 + 抽样(某退市股 K 线行数 + max(date) ≤ 退市日)+ `daily_baostock_full` 行数增量 + distinct symbol 增量。
- [ ] **3.5 wrapper + schtask**:`baostock_delisted_wrapper.ps1`;schtask `sanguo-delisted` `/sc monthly /ru SYSTEM`(月度,守 48000,错开 day2b 02:00 + bs-daily-increment 17:00)。
- [ ] **3.6 commit**:`git commit -m "feat(data): 退市股 K 线采集(baostock,反幸存者偏差)"`
---
## Task 4: baostock 日增量 → daily_baostock_full(#7 daily_update_static)
> **串行约束**:本 task 与 Task3 都用 baostock 长会话,**必须串行**(Task3 probe → Task3 执行 → Task4),不可并发(防黑名单)。
**Files:** Create `scripts/data_platform/daily_update_static.py` + `daily_update_static_wrapper.ps1`
**背景:** 现有 `daily_update_xtdata.py` 只产 parquet 不灌 `daily_baostock_full`(已知 gap,memory `db-primary-parquet-fallback` 记录)。本 task 补 baostock 日线的**每日增量灌库**。
**Interfaces:**
- Consumes: baostock `query_stock_basic`(全 A,type=1 含退市,复用 `baostock_daily_fullmarket_download.py:fetch_all_stocks`)+ `query_history_k_data_plus`(LOOKBACK 窗口,adjustflag=3 raw,18 字段同 `BS_FIELDS`)
- Produces: staging `data/daily_baostock_increment/{YYYYMMDD}/{code}.{exc}_daily.parquet`(审计)→ 同进程 INSERT OR REPLACE `daily_baostock_full`(复用 `parse_baostock_code`+executemany+WAL+busy_timeout,同 `import_baostock_to_db.py`)
**设计(LOOKBACK 窗口 + 幂等,不同于全量 marker 模式):**
- **不用 marker 断点续传**(全量才需要;增量每日全量重拉最近 N 天)
- `LOOKBACK_DAYS=7`(覆盖周末/节假日;baostock 日终更新,17:00 跑时当日 bar 已就绪)
- 每只 1 query → 5537 query/run ≪ 48000/天 ✅(留足余量给 day2b/Task3)
- `sleep 0.4s × 5537 ≈ 37min`(17:00 schtask 可接受)
- `QUERY_COUNT` 计数器 + `DAILY_LIMIT=40000` 防御(复用全量脚本模式)
- **一脚本贯通**:download LOOKBACK → staging parquet(审计)→ in-memory df → executemany INSERT OR REPLACE(幂等,重复跑同一天安全,`drop_duplicates keep last` 不需要因 PK+OR REPLACE 天然去重)
**Steps:**
- [ ] **4.1 探针(可选,Day1 已实证 query_history_k_data_plus 可用)**:ssh VPS 跑 1 只近 7 天确认接口 + 当日 bar 就绪
- [ ] **4.2 写 `daily_update_static.py`**:自包含,结构
- `unset proxy` + `socket.setdefaulttimeout(30)`(同全量脚本,baostock 坑)
- `_login_once`/`_relogin`/`fetch_all_stocks`/`fetch_one_daily`/`parse_baostock_code` 复用(可 import 或复制;优先 from `baostock_daily_fullmarket_download import ...`,注意 `QUERY_COUNT` global 需在同进程)
- `LOOKBACK` 窗口:`start=today-7, end=today`
- 主循环:逐只 `fetch_one_daily` → staging parquet → 累积 df → 每 100 只 `executemany INSERT OR REPLACE`(WAL+busy_timeout=60000)
- `QUERY_COUNT`/`DAILY_LIMIT`/断路器/定期重登 复用
- 结束 verify:抽样 3 只 `max(date) ≈ today`、当日新增行数
- 环境变量 `BS_INCREMENT_OUT_DIR`/`DB_PATH` 覆盖默认(同 Day1 wrapper 模式适配 Win)
- [ ] **4.3 小样本**:`--limit 10` 跑 10 只,确认 staging 有行 + DB 抽样 max(date)≈today
- [ ] **4.4 全量跑**:5537 只,守预算
- [ ] **4.5 wrapper + schtask**:`daily_update_static_wrapper.ps1`(unset proxy+utf8+OUT_DIR+log);schtask `sanguo-bs-daily-increment` `/sc daily /st 17:00 /ru SYSTEM`(错开 daily-update 16:30 + day2b 02:00 + Task3 月度)
- [ ] **4.6 commit**:`git commit -m "feat(data): baostock 日增量灌库 daily_update_static(#7 gap 补)"`
---
## Self-Review
- **Spec 覆盖**:Task1→spec §4 成份股行 + §8 P0.1;Task2→§4 ETF 行 + §8 P0.2;Task3→§4 退市行 + §8 P0.3 ✅
- **Placeholder 扫描**:无 TBD/TODO;采集脚本给接口+探针+schema,实现者按骨架写完整(采集脚本完整代码由执行 agent 基于 接口/schema/陷阱 产出)✅
- **类型一致**:`index_const_hist` schema 各源统一;`daily_baostock_full` 18 列复用 `import_baostock_to_db.py``parse_baostock_code`+executemany ✅
- **陷阱纳入**:`ak.index_stock_hist` 下线(1.1 标注)/ csindex SPA 无历史(用国证+新浪)/ 新浪 gb2312(1.2)/ hist 版必须(1.1)✅
---
## Execution Handoff
计划存 `docs/superpowers/plans/2026-07-21-data-fusion-p0.md`。执行方式:
1. **Subagent-Driven**(推荐):每 Task 派 fresh agent + task 间 review
2. **Inline**:本 session 批量执行 + checkpoint
@@ -1,106 +0,0 @@
# 数据架构方案A迁移 + schtask 改造 实施计划
> **For agentic workers:** REQUIRED SUB-SKILL: superpowers:executing-plans。Steps use checkbox。
**Goal:** 落地 spec §14 方案A定稿 — DB 唯一表、每类数据唯一权威源、4 个新 schtask、迁移 5 单元,全程备份+staging+可回滚+审计。
**Architecture:** 以本地 DB 迁移为主(`daily_baostock_full`→dbbardata/parquet,无网络),schtask 改造(废弃旧 4 个新建 4 个)。每单元独立可回滚,按风险升序。
**Tech Stack:** Python3.10 / sqlite3(WAL) / pandas parquet / Windows schtasks / baostock+xtata+akshare
## Global Constraints
- baostock:单进程单登录,`DAILY_LIMIT=48000`,sleep 限速,login 探针 graceful skip,直连不走代理(unset proxy)
- xtata:单进程 download 不并发,无限流
- akshare:interval 4s 单线程,防东财封 IP
- 每迁移单元前:`sqlite3 .backup` 全库 + rsync 到 NAS `/volume1/stock/backup/` + WAL checkpoint
- 每单元:staging 隔离 → 验证探针 → 用户确认合并 → 旧 rename `_old` 保留 7 天
- 全程 nohup + 审计日志 `data/migration_logs/<unit>_<ts>.log`
- 不破坏 vnpy 回测:dbbardata schema 不动(只灌数据),`dbbardata` 12 列保持
## 文件结构
- 迁移脚本:`scripts/data_platform/migrate_*.py`(每单元一个)
- 验证脚本:`scripts/data_platform/verify_*.py`
- schtask wrapper:`scripts/data_platform/*_wrapper.ps1`
- 审计日志:`data/migration_logs/`
---
## 前置 Task 0:全库备份(所有单元前必做)
**Files:** `scripts/data_platform/backup_db.py`(新建,可复用)
- [ ] 写脚本:`sqlite3 .backup``quant_trading.db.bak_<YYYYMMDD>`(在线一致);WAL checkpoint;rsync 到 NAS
- [ ] 执行
- [ ] **verify**:`.bak` 存在 + 大小≈28GB + `PRAGMA integrity_check` ok
---
## 单元 1:存量垃圾清理(零风险)
**Files:** `scripts/data_platform/cleanup_staging.py`(新建)
- [ ] 写脚本:`--dry-run` 先列清单 → 删 `_staging_xtdata/`(14万)、`_xtdata.tar`(1.4G);移 `cta_*/dbg_*/smoke_*/trace_*``backtest_files/`
- [ ] dry-run 输出清单给用户确认
- [ ] 执行删除/移动
- [ ] **verify**:`data/` 根目录无散落 json/log;`backtest_files/` 收纳;`du -sh data/` 体积下降
- [ ] **回滚**:staging 可由 `build_daily_from_xtdata` 重建(已合并到 qfq/raw)
---
## 单元 2:config 统一 VPS 路径
**Files:** `config/data_platform.yaml`(VPS 实例)
- [ ] 核实 VPS 实际 config 路径(当前仓库版指 NAS /volume1,是容器版遗留)
- [ ] `daily_dir/raw_dir/qfq_dir/minute_15_dir``C:\sanguo_vnpy_v2\data\...`
- [ ] `daily_dir` 统一指 qfq(消除 `daily/` vs `qfq/` 分叉,`daily/`68文件归档)
- [ ] NAS config 保留 + 注释"备份用"
- [ ] **verify**:`datareader.read_parquet_daily` 抽样能读 + LocalParquetProvider 抽样
- [ ] **回滚**:yaml 改回
---
## 单元 3:成份股合并 → `constituent_unified`
**Files:** `scripts/data_platform/migrate_constituent.py` + `verify_constituent.py`
**Interfaces:**`bs_index_constituent`(baostock 300/500/50)+ `data/index_const_hist/*_union.parquet`(akshare cni 深证);写 `constituent_unified(date,index_code,code,code_name,source)`
- [ ] 写迁移脚本:按指数代码去重(300/500/50=baostock;深证 399xxx=akshare cni union;新浪 300/50 作校验丢弃);schema 映射 INSERT
- [ ] staging:先写 `constituent_unified_staging`
- [ ] **verify**:行数 / 指数覆盖 / 抽样某指数某日成份集 vs 源一致 / 无同指数同日重复
- [ ] 合并:rename staging → `constituent_unified`;`bs_index_constituent``_old`
- [ ] 7 天后删 `_old`
- [ ] **回滚**:rename `bs_index_constituent_old` 回来
---
## 单元 4:`daily_baostock_full` 拆分(本地 DB 迁移,无网络)
**Files:** `scripts/data_platform/migrate_daily_baostock.py` + `verify_daily_migration.py`
**Interfaces:**`daily_baostock_full`(含退市);写 `dbbardata('d')`(OHLCV 12 列)+ `data/valuation_baostock/<year>.parquet`
- [ ] 写迁移脚本:
- OHLCV:`daily_baostock_full` → dbbardata INSERT OR REPLACE(interval='d',exchange SH/SZ→SSE/SZSE,datetime=date)。含退市(治偏差)。ETF 不碰(已在 dbbardata)
- pe/pb:按年 group → `valuation_baostock/<year>.parquet` 宽表
- [ ] staging:先写 `dbbardata_staging_daily` 表 + parquet staging 目录,不动 dbbardata
- [ ] **verify**:
- 退市股(000005 等)在 dbbardata('d') 有了(治偏差验证)
- 在市股(600519)日线行数 / 抽样价格 vs daily_baostock_full 一致
- pe/pb parquet 按年覆盖 + 抽样值合理
- dbbardata 总行数变化合理(+退市日线)
- [ ] 合并:staging → dbbardata;`daily_baostock_full``_old`;valuation parquet → 正式目录
- [ ] 7 天后删 `_old`
- [ ] **回滚**:`daily_baostock_full_old` 还原 + dbbardata 从 `.bak` 恢复
---
## 单元 5:schtask 改造(废弃旧 4 个,新建 4 个)
**Files:** `scripts/data_platform/bs_eod.py`(日线+15min+pe/pb 拆)+ `xt_eod.py`(ETF+实时)+ `*_wrapper.ps1`
- [ ]`bs_eod.py`:基于 `daily_update_static.py` 扩展,+15min 增量,+pe/pb 拆 parquet;落 dbbardata('d'/'15m');`DAILY_LIMIT=48000`
- [ ]`xt_eod.py`:基于 `daily_update_xtdata.py`,universe 收窄 ETF/基金 + 个股当天实时;落 dbbardata('d')
- [ ] **verify**:`--limit 10` 小样本跑通 + 数据到当天
- [ ] 部署 schtask:废弃 `sanguo-daily-update`/`sanguo-bs-daily-increment`/`sanguo-index-hist`;新建 `sanguo-bs-eod`(18:05)/`sanguo-xt-eod`(18:40);`sanguo-bs-akshare` 调到 19:00;`sanguo-index`(月度 19:50)
- [ ] **verify**:`schtasks /query` + 首日运行结果码 + 数据抽查到当天
- [ ] **回滚**:重新注册旧 schtask
---
## 收尾:E2E 验证
- [ ] 回测 all_weather 一轮(读 dbbardata 日线含退市 + valuation parquet + constituent_unified)无回归
- [ ] LocalParquetProvider 接 constituent_unified + valuation_baostock 单测
- [ ] 更新 memory:`data-fusion-design-finalized`(标方案A落地)+ 新建 `data-arch-migration-done`
## 执行节奏
- 每单元独立提交 + 用户 review staging 再合并(单元 4/5 关键)
- 全程 VPS nohup 跑(Mac Mini 防休眠 caffeinate,长迁移)
- 顺序:0 → 1 → 2 → 3 → 4 → 5 → 收尾(严格风险升序)
@@ -1,168 +0,0 @@
# akshare 低频任务 schtask 部署 Plan (spec §14.5)
> **For agentic workers:** REQUIRED SUB-SKILL: superpowers:subagent-driven-development / executing-plans。本 plan 自包含(实测现状+脚本能力+VPS访问),fresh agent 可直接执行。
**Goal:** 部署 spec §14.5 akshare/index 低频 schtask — A 三表/估值增量 + B 成份股月度 + C 事件类 7 种。方案A 数据层(`dbbardata`/`constituent_unified`/`valuation_baostock`)的使用层配套,补全 `LocalUnifiedProvider` 依赖的静态数据源。
**Architecture:** 复用 `akshare_static_download.py`(16类/marker断点/4模式)+ `baostock_constituent_download.py`;改 `merge_constituent.py`/`migrate_constituent.py` 可重跑;按 akshare per-stock 全量慢 + 东财限流,拆多 schtask(日频/季频/月频)。
**Tech Stack:** Python 3.10, akshare, baostock, sqlite3, Windows schtasks
## Global Constraints(铁律)
- **akshare**: 单线程 0.8s sleep / 30s 超时 / 断路器(连30 failed exit) / marker 断点 / **unset proxy**`akshare_static_download.py` 已内置;东财限流严,**per-stock 全量慢,夜间跑**
- **baostock**: 单进程单登录 48000/天,不并发(防黑名单)
- **staging→验证→合并**(成份股 B,用户铁律:下载质量不可控,不直接写主库)
- **VPS**: ssh alias = `49.232.102.198`(IP 即 alias,User Administrator,key id_ed25519);`C:\Python310\python.exe -X utf8`;schtasks `/create /ru SYSTEM /rl HIGHEST /sc daily|monthly`;Windows ssh 引号地狱 → 脚本 scp + ssh python 跑最稳
- **dbbardata UNIQUE 不破坏**;constituent_unified 治偏差
- 直连不走代理(schtask wrapper 开头 `unset http_proxy https_proxy all_proxy`)
## 实测现状(2026-07-23 probe_akshare_status.py)
- **static/**(`C:\sanguo_vnpy_v2\data\static\`): balance/income/cashflow/valuation/financial_abstract 各 **5530 parquet,停 2026-07-22 18:25**(sanguo-bs-akshare disabled 前最后一次);provider fundamentals 依赖,要续更
- **events/**: **全 MISSING**(龙虎榜/北向/两融/解禁/大宗/可转债/研报从未采集)
- **constituent_unified**: 7110 行 = baostock 2938(000016:195/000300:940/000905:1803) + akshare_cni 1172(深证 399001:702/399005:145/399006:175/399330:150) + akshare_csindex 3000(000852:1000/932000:2000);**静态,月度更新 schtask 无**
- **schtask 状态**: sanguo-bs-akshare / sanguo-index / sanguo-index-hist **全无**(方案A `stop_all_data_schtasks.ps1` 清了)
- **akshare_static_download.py 16 类分 5 组**:
- PER_STOCK(5500股循环,unit=`{symbol}_{type}`): valuation/northbound/share_capital/balance/income/cashflow/financial_abstract
- TOP_HOLDERS(per-stock×period): top_holders
- PER_DATE(每交易日,unit=`{date}_{type}`): dragon_tiger/block_trade/margin_sse/restricted
- PER_PERIOD(报告期,unit=`{period}_{type}`): forecast/express
- ONE_SHOT: index_const/industry
- marker 是 **symbol 级**(非 period 级)→ 三表/估值要更新新数据必须 `--force`(否则 marker 跳过永不更新)
- **merge_constituent.py 不可重跑**: `ALTER TABLE bs_index_constituent RENAME TO bs_index_constituent_old` 只能一次(_old 已存在);无 DROP/REPLACE constituent_unified
- **existing wrappers**(参考模式): `bs_eod_wrapper.ps1` / `xt_eod_wrapper.ps1`(unset proxy + log + 调 python);`register_schtasks.ps1`(schtasks /create 模板)
## schtask 清单(方案A §14.5 适配,按频率拆)
| schtask | 频率 | 时间 | 脚本 | 内容 |
|---|---|---|---|---|
| `sanguo-ak-eod` | daily | 19:00 | ak_eod_wrapper.ps1 | valuation + financial_abstract `--force`(日频,5500×2×0.8s≈2.2h 夜间) |
| `sanguo-ak-quarter` | monthly(财报季 5/9/11 月+年报4月) | 周末 02:00 | ak_quarter_wrapper.ps1 | balance + income + cashflow + forecast + express `--force`(季频,5500×3×0.8≈2.2h) |
| `sanguo-ak-events` | daily | 19:30 | ak_events_wrapper.ps1 | dragon_tiger + block_trade + margin_sse + restricted `--start today --end today`(per-date 日频,4 unit 快) |
| `sanguo-ak-stock` | weekly | 周六 03:00 | ak_stock_wrapper.ps1 | northbound + share_capital + top_holders(per-stock 慢,周频) |
| `sanguo-index` | monthly | 19:50 | index_monthly_wrapper.ps1 | 成份股 3 源 + merge(B) |
## File Structure
- **Modify:** `scripts/data_platform/merge_constituent.py`(可重跑: DROP/REPLACE 替代 RENAME)
- **Modify:** `scripts/data_platform/migrate_constituent.py`(可重跑: staging 隔离 + DROP staging 重建)
- **Create:** `scripts/data_platform/ak_eod_wrapper.ps1` / `ak_quarter_wrapper.ps1` / `ak_events_wrapper.ps1` / `ak_stock_wrapper.ps1`(4 个 akshare wrapper)
- **Create:** `scripts/data_platform/index_monthly_wrapper.ps1`(B 成份股 3 源编排)
- **Create:** `scripts/data_platform/register_akshare_schtasks.ps1`(注册 5 schtask)
- **Create:** `scripts/data_platform/verify_akshare_e2e.py`(验证全部)
- **Test:** `tests/portfolio/test_merge_constituent_rerun.py`(B 改造 TDD)
---
## Task A: 三表/估值增量(2 schtask)
**Files:** Create 4 akshare wrapper + register; 复用 `akshare_static_download.py`(不改)。
- [ ] **A1: ak_eod_wrapper.ps1**(日频估值/财务摘要)
```powershell
# unset proxy + 调 akshare_static_download.py --types valuation,financial_abstract --force
$env:http_proxy=""; $env:https_proxy=""; $env:all_proxy=""
cd C:\sanguo_vnpy_v2
C:\Python310\python.exe -X utf8 scripts\data_platform\akshare_static_download.py `
--types valuation,financial_abstract --force `
*>> C:\sanguo_vnpy_v2\data\ak_eod.log
```
- [ ] **A2: ak_quarter_wrapper.ps1**(季频三表+预告/快报,财报季)— 同上 `--types balance,income,cashflow,forecast,express --force`
- [ ] **A3: 验证 A** — scp wrapper + 手动跑 `--limit 5` 确认 valuation parquet 更新今日:
```bash
scp scripts/data_platform/ak_eod_wrapper.ps1 49.232.102.198:'C:/sanguo_vnpy_v2/scripts/data_platform/'
ssh 49.232.102.198 'cd /d C:\sanguo_vnpy_v2 && C:\Python310\python.exe -X utf8 scripts\data_platform\akshare_static_download.py --types valuation --force --limit 3'
# 验 static/valuation/<code>_valuation.parquet mtime = 今日
```
- [ ] **A4: 注册 schtask**(register_akshare_schtasks.ps1 含 sanguo-ak-eod daily 19:00 + sanguo-ak-quarter monthly)
---
## Task B: 成份股月度(sanguo-index)— 代码改造 TDD
**Files:** Modify `merge_constituent.py` + `migrate_constituent.py`; Create `index_monthly_wrapper.ps1`; Test `test_merge_constituent_rerun.py`
- [ ] **B1: 写失败测试 — merge_constituent 可重跑**
```python
# tests/portfolio/test_merge_constituent_rerun.py
def test_merge_constituent_rerun_twice(tmp_path):
"""merge_constituent 跑两次不崩(第二次 DROP 重建,不 RENAME 已 _old 的表)。"""
db = tmp_path / "t.db"; c = sqlite3.connect(str(db))
# 造 constituent_unified + bs_index_constituent_old(已存在,_old 状态)
c.execute("CREATE TABLE constituent_unified(index_code,code,source,in_current,was_removed)")
c.execute("CREATE TABLE constituent_unified_staging(index_code,code,source,in_current,was_removed)")
c.execute("CREATE TABLE bs_index_constituent_old(code,date)") # _old 已存在
c.executemany("INSERT INTO constituent_unified_staging VALUES(?,?,?,?,?)",
[("000300","600519","baostock",1,0)])
c.commit(); c.close()
# 跑两次
import scripts.data_platform.merge_constituent as m # 或函数级 import
m.merge(str(db)) # 第一次: staging→unified(DROP 旧 unified 重建)
m.merge(str(db)) # 第二次: 不崩, unified 仍 1 行
c = sqlite3.connect(str(db))
assert c.execute("SELECT COUNT(*) FROM constituent_unified").fetchone()[0] == 1
```
- [ ] **B2: 改 merge_constituent.py 可重跑** — 把 `ALTER TABLE bs_index_constituent RENAME TO _old`(只能一次)改为:staging→`DROP TABLE IF EXISTS constituent_unified``CREATE constituent_unified AS SELECT FROM staging`。bs_index_constituent_old 已存在不碰。幂等。
- [ ] **B3: 改 migrate_constituent.py 可重跑** — staging 表 `DROP IF EXISTS constituent_unified_staging` 重建(每次重新聚合 baostock 988 时点 + akshare cni union + csindex),不依赖上次状态。
- [ ] **B4: index_monthly_wrapper.ps1**(3 源编排)
```powershell
$env:http_proxy=""; $env:https_proxy=""; $env:all_proxy=""
cd C:\sanguo_vnpy_v2
# 1. baostock 300/500/50 最新快照(单进程)
C:\Python310\python.exe -X utf8 scripts\data_platform\baostock_constituent_download.py --start 2026-01-01
# 2. akshare cni 深证 + csindex 中证(复用 P0 脚本 index_const_hist_download 或 akshare_static_download --types index_const)
C:\Python310\python.exe -X utf8 scripts\data_platform\akshare_constituent_download.py
# 3. migrate + merge(可重跑版)
C:\Python310\python.exe -X utf8 scripts\data_platform\migrate_constituent.py
C:\Python310\python.exe -X utf8 scripts\data_platform\merge_constituent.py
```
(注:akshare_constituent_download.py 若不存在,从 `akshare_static_download.py --types index_const` 或 P0 的 `index_const_hist_download.py` 复用;执行 agent 确认现有脚本)
- [ ] **B5: 验证 B** — 手动跑 wrapper,确认 `constituent_unified` 行数 ≥ 7110,source 含 baostock/akshare_cni/akshare_csindex,跑两次不崩。
- [ ] **B6: 注册 sanguo-index monthly 19:50**
---
## Task C: 事件类(per-date 日频 + per-stock 周频)
**Files:** Create `ak_events_wrapper.ps1` + `ak_stock_wrapper.ps1`(已在 A 的 register 注册)。
- [ ] **C1: ak_events_wrapper.ps1**(per-date 日频)— `--types dragon_tiger,block_trade,margin_sse,restricted --start {today} --end {today}`。每日 4 类×1 unit,快。落 `events/{type}/{date}_{type}.parquet`
- [ ] **C2: ak_stock_wrapper.ps1**(per-stock 周频慢)— `--types northbound,share_capital,top_holders --force`。5500×3 慢,周六 03:00。
- [ ] **C3: 可转债/研报**(spec §14.2 列但 akshare_static_download.py 16 类无对应 fetcher) — **评估**:若 akshare 有 `bond_zh_hs_cov_min`/`stock_research_info_em` 接口,加 fetcher;否则 N/A 标注使用说明。执行 agent 确认 akshare 接口可用性,不可用则跳过并在 verify 标注。
- [ ] **C4: 验证 C** — 手动跑 events wrapper `--start today --end today`,确认 `events/dragon_tiger/{today}_dragon_tiger.parquet` 生成。
---
## Task D: 注册 + E2E
- [ ] **D1: register_akshare_schtasks.ps1** — 注册 5 schtask(sanguo-ak-eod/ak-quarter/ak-events/ak-stock/index),`/create /ru SYSTEM /rl HIGHEST`,verify 段 `schtasks /query` 确认 Status=Ready。
- [ ] **D2: verify_akshare_e2e.py** — 验证全部:
- static/valuation 最新 mtime = 近日
- events/dragon_tiger 有 parquet
- constituent_unified ≥ 7110 + 3 source
- [ ] **D3: 更新 memory + 使用说明**`data-fusion-design-finalized` 的"剩余待办 akshare schtask"标完成;`docs/portfolio_local_unified_provider.md` 事件类从 N/A 更新。
---
## Self-Review
1. **spec §14.5 覆盖**: sanguo-bs-akshare(三表+事件)→ 拆 ak-eod/ak-quarter/ak-events/ak-stock(频率适配);sanguo-index 月度 → B。✓
2. **provider 依赖**: valuation/balance/income/cashflow/financial_abstract(LocalUnifiedProvider.get_fundamentals_df)→ A 覆盖。✓
3. **限流现实**: per-stock --force 全量慢,日频 valuation 2.2h / 季频三表 2.2h,夜间 + 周末 schtask。events per-date 日频快。✓
4. **B 可重跑**: merge DROP 重建幂等,migrate staging 隔离。TDD 跑两次不崩。✓
5. **风险**: akshare 东财封 IP(per-stock 全量)→ wrapper 内置断路器 + sleep;财报季触发 ak-quarter(`/sc monthly` 指定月或手动);可转债/研报接口待确认(C3)。
## Execution Handoff
Plan saved to `docs/superpowers/plans/2026-07-23-akshare-low-freq-schtask.md`。compact 后新 session 派 Sub Agent 执行(参考 LocalUnifiedProvider 模式:Task A→B→C→D,每 task 验证 + commit)。
## 关联文档/memory
- spec: `docs/superpowers/specs/2026-07-21-data-source-fusion-design.md` §14.5
- memory: `data-fusion-design-finalized`(方案A)/ `local-unified-provider-complete`(使用层)/ `vps-local-data-layout` / `baostock-concurrent-blacklist` / `schtasks-system-bat-gotchas`
- VPS 访问: ssh `49.232.102.198`,见 memory `windows-vps-access`
@@ -1,191 +0,0 @@
# 中证1000/2000 历史成份股补全 实施计划
> **For agentic workers:** REQUIRED SUB-SKILL: superpowers:executing-plans。Steps 用 checkbox `- [ ]` 跟踪。
**Goal:** 把中证1000(000852)/中证2000(932000)从"纯当前快照"补成"治幸存者偏差的全集"(含被踢出的股票),接入 `constituent_unified`,并部署定期更新 schtask。
**Architecture:** csindex 官方公告 JSON 接口(`queryAnnouncementByVo` + `queryAnnouncementById`)抓调整公告 → 解析附件 PDF/xlsx 的调入/调出名单 → 聚合成"曾经入选集"(全集型,非时点型)→ 入 `constituent_unified`,`in_current`=当前快照、`was_removed`=曾经入选−当前。
**Tech Stack:** Python3 + pandas + openpyxl + pdfplumber + sqlite3 + PowerShell schtask
## 诊断(已实证,2026-07-23)
现状 `constituent_unified`(VPS quant_trading.db):
- `000852`: total=1000, in_current=1000, **was_removed=0**(纯快照,未治偏差)
- `932000`: total=2000, in_current=2000, **was_removed=0**(纯快照)
- 对比 `000300`: total=940, in_current=300, was_removed=640(已治偏差)
**三处断点:**
1. **932000 launch xlsx 解析 bug**:`parse_csindex_announce.py:529``row[0]`(=指数代码 932000),应为 `row[3]`(证券代码)。→ 产出 distinct=1(2000 行全是 932000)。xlsx 实证 6 列:`指数代码/指数简称/指数英文简称/证券代码/证券中文简称/证券英文名称`
2. **000852 公告覆盖不全**:`filter_csi1000_notices`(162 行)用 `theme='指数调样'+title 含'中证1000'` 过滤,只拿 28 份(2018-07 起)。调查实证:列表 API payload 加 `indexCode:'000852'` 能拿 **96 条**(45 调样),可回溯到 **2014 发布期**(早期 HTML 表格,2018+ PDF/xlsx)。
3. **migrate 没接 announce_union**:`migrate_constituent.py:118-131` 只读 `_snapshot.parquet`,没读 `_announce_union.parquet`。→ 1220 个治偏差集白产了。路径也对不上(parse 在 Mac 产 announce_union,migrate 读 VPS HIST,没同步)。
## 关键简化
`constituent_unified` 是**全集型**(300/500/50 = baostock 988 时点聚合成 in_current/was_removed),**不是时点型**。所以:
- **不需要**反向回溯引擎(生效日边界、逐时点 asof join)
- 只要"曾经入选集"= 所有公告 add 记录 initial current 的 distinct code
- `in_current` = akshare 当前快照(权威),`was_removed` = 曾经入选 当前
调查 agent 提的"生效日≠公告日"等坑是**时点型**需求才需要,本计划(全集型)不涉及。
## Global Constraints(spec 铁律)
- baostock 单进程单登录不并发(本计划不碰 baostock,无冲突)
- 直连不走代理:`unset http_proxy https_proxy all_proxy`(脚本已内置)
- 单线程限速:csindex 接口 sleep 1.0~1.5s
- staging→验证→合并,不直接写主库(migrate 走 staging→merge 两步,已幂等)
- provider 读 VPS 本地,不调 online(本计划是采集层,可调 csindex)
- commit message 无 Co-Authored-By
---
### Task 1(#30):修 parse_csindex_announce.py 两处
**Files:**
- Modify: `scripts/data_platform/parse_csindex_announce.py:526-538`(932000 launch xlsx 列索引)
- Modify: `scripts/data_platform/parse_csindex_announce.py:120-182`(000852 列表搜索用 indexCode)
**改动 1a — 932000 launch xlsx 列索引(:526-538):**
现:`code = _norm_code(row[0])`, `name = str(row[1])`。改为按 header 定位列(稳健),或直接 `code=row[3]`, `name=row[4]`。推荐 header 定位:
```python
header = rows[0]
# 找"证券代码"和"证券中文简称"列(中英文混合 header)
code_idx = next((i for i,h in enumerate(header) if h and "证券代码" in str(h)), 3)
name_idx = next((i for i,h in enumerate(header) if h and "证券中文简称" in str(h)), 4)
for row in rows[1:]:
code = _norm_code(row[code_idx] if len(row)>code_idx else None)
name = str(row[name_idx]).strip() if len(row)>name_idx and row[name_idx] else ""
```
**改动 1b — 000852 列表搜索用 indexCode(:120-182):**
`fetch_all_notices` 拉全量再 `filter_csi1000_notices` title 过滤。改为:对 000852 用 `indexCode` payload 直接搜:
```python
payload = {"lang":"cn","classlist":[],"indexlist":[],
"indexCode":"000852", # ← 新增,直接按指数搜
"page":{"desc":"","key":"","page":page,"rows":100},
"related_topics":[],"typelist":[]}
```
保留旧 filter 作兜底(标题含中证1000+调整)。合并 indexCode 命中 已知 REGULAR/TEMP_IDS 去重。932000 走全局 `related_topics:["index_rebalance"]` + PDF grep "中证2000" section(parse_pdf_adjustments 已支持 target_section)。
**验证探针:**
```bash
python3 scripts/data_platform/parse_csindex_announce.py --only 1000
# 期望:filtered CSI 1000 公告 ≥ 40 条(原 28),date 范围早于 2018-07
python3 scripts/data_platform/parse_csindex_announce.py --only 2000
# 期望:932000_announce_union.parquet distinct codes ≈ 2000(原 bug=1)
```
- [ ] Step 1: 改 932000 launch xlsx 列索引(header 定位)
- [ ] Step 2: 改 000852 列表搜索(indexCode payload + 932000 related_topics)
- [ ] Step 3: Mac 重跑 `--only 1000` + `--only 2000`,验证探针
- [ ] Step 4: commit
---
### Task 2(#31):改 migrate_constituent.py 接 announce_union 聚合全集
**Files:**
- Modify: `scripts/data_platform/migrate_constituent.py:118-131`(加读 announce_union)
- Test: `tests/portfolio/test_migrate_announce_union.py`(新建,TDD)
**聚合逻辑(全集型):**
```python
# 读 000852_announce_union.parquet + 932000_announce_union.parquet
# announce_union schema: updateDate/index_code/code/code_name/adjust_type(add|remove|current|initial|current)/notice_id/source
# 全集聚合:
for idx in ['000852','932000']:
ann = read(f"{idx}_announce_union.parquet")
snap = read(f"{idx}_snapshot.parquet") # akshare 当前快照,权威 in_current
current_codes = set(snap['code']) # 当前在册
ever_codes = set(ann['code']) | current_codes # 曾经入选(所有 add/initial + current)
# 产出:ever_codes 每只一行
# in_current = code in current_codes
# was_removed = code not in current_codes(曾入选已踢)
# source = 'csindex_announce'
```
schema 对齐:`index_code/code/code_name/source/in_current/was_removed``code_name` 取 announce_union 或 snapshot 的(优先 snapshot 当前名)。
**合并进 staging:** 现有 `all_df = pd.concat([pool, df_deep, df_snap])`(:134)→ 把 000852/932000 的 announce_union 全集**替换** df_snap 里的 000852/932000 快照行(快照并入 announce 全集的 in_current),其他指数不动。
**TDD 测试(tests/portfolio/test_migrate_announce_union.py):**
- test announce_union 聚合:given announce(add A,B + remove C) + snapshot(current A,B,D),assert ever={A,B,C,D}, in_current={A,B,D}, was_removed={C}
- test 000852 distinct > 1000(治偏差证据)
- test 932000 distinct ≈ 2000(launch 修复)
- test 幂等(跑两次结果一致)
- [ ] Step 1: 写聚合测试(RED)
- [ ] Step 2: 改 migrate 加 announce_union 聚合(GREEN)
- [ ] Step 3: 测试通过
- [ ] Step 4: commit
---
### Task 3(#32):重跑→同步VPS→migrate→merge→验证
**Files:** 无新文件(运行现有 pipeline)
- [ ] Step 1: Mac 重跑 parse_csindex_announce.py --only both → 新 announce_union
- [ ] Step 2: scp 000852_announce_union.parquet + 932000_announce_union.parquet 到 VPS `C:\sanguo_vnpy_v2\data\index_const_hist\`
- [ ] Step 3: rsync 改后的 migrate_constituent.py 到 VPS
- [ ] Step 4: VPS 跑 migrate_constituent.py(SANGUO_DB 指向 quant_trading.db)→ merge_constituent.py
- [ ] Step 5: 验证(见下)
**验证标准(VPS 查 constituent_unified):**
```sql
SELECT index_code, COUNT(*), SUM(in_current), SUM(was_removed)
FROM constituent_unified WHERE index_code IN ('000852','932000') GROUP BY index_code;
```
- 000852: total > 1000(曾经入选 ~1200+), in_current=1000, **was_removed > 0**(治偏差)
- 932000: total ≈ 2000+, in_current=当前快照数, was_removed ≥ 0(launch current,中间调整无记录则 was_removed=0 可接受)
- 抽样:挑一只 known 被踢股(如 announce_union 里 remove 类型)→ constituent_unified 该 code was_removed=1
- 回归:300/500/50/深证 行数不变(没误伤)
---
### Task 4(#33):定期 schtask 方案+部署
**Files:**
- Create: `scripts/data_platform/csindex_constituent_wrapper.ps1`
- Create: `scripts/data_platform/register_csindex_schtasks.ps1`
**schtask 设计:**
- 名:`sanguo-csindex-constituent`
- 频率:**每月 16 号 + 6月/12月定调后额外**(中证1000 定期调整 6月/12月,临时调整不定期 → 月度抓足够,缓存增量)
- 时间:**20:30**(避开 baostock 18:05/xt 18:40/akshare 19:00-19:50 窗口)
- 流程:parse_csindex_announce.py --refresh-list(抓新公告)→ 同步 announce_union 已在本机 → migrate → merge
- 幂等:migrate/merge 已 DROP+CREATE 可重跑;parse 有 notice cache 增量
**wrapper ps1(仿 bs_eod_wrapper.ps1 风格):** unset proxy → Set-Location → timestamped log → python parse + migrate + merge → exit code
- [ ] Step 1: 写 wrapper ps1 + register ps1
- [ ] Step 2: VPS 部署 + schtasks /create /ru SYSTEM /rl HIGHEST
- [ ] Step 3: 手动触发一次验证(schtasks /run)
- [ ] Step 4: commit + 同步安装目录
---
### Task 5(#34):更新 memory
**Files:**
- Update: memory `data-fusion-design-finalized.md`(推翻 000852/932000 "永久 gap")
- Update: memory `static_data_gaps_design.md`(中证1000/2000 gap 关闭)
- Update: `MEMORY.md` 索引
**记:** csindex 公告 JSON 接口路推翻"永久 gap";000852 全集入库(曾经入选 1200+);932000 launch xlsx 列 bug 修复;全集型简化洞察(不需回溯引擎);定期 schtask;调查 agent 实证的 96 公告/45 调样/回溯到 2014。
- [ ] Step 1: 更新 3 个 memory 文件
- [ ] Step 2: MEMORY.md 索引行
---
## Self-Review
- spec 覆盖:① 调整补全→Task1-3 ② 定期抓取方案→Task4 ✓
- 全集型简化避免过度设计(调查 agent 的回溯引擎是 future 时点型需求,现不做)✓
- TDD:migrate 聚合逻辑先写测试 ✓
- 不破坏:300/500/50/深证 migrate 路径不动,只加 000852/932000 announce 段 ✓
- 约束:不走代理/单线程/staging→merge 幂等/不碰 baostock ✓
## 已知残留 gap(接受,不阻塞)
- 932000 中间调整(2023-08 launch 到 current 之间)csindex 无公告 → launch current 全集,中间被踢的不可补(2023 新指数,影响小)
- 000852 2014-2017 早期 HTML 表格解析格式松散,可能不全(扩 indexCode 搜索尽力补,实证 id=5/id=1585 等仍有表格)
@@ -1,601 +0,0 @@
# LocalUnifiedProvider Implementation Plan (spec §6 使用层)
> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking.
**Goal:** 实现 spec §6 使用层 `LocalUnifiedProvider`——读方案A 权威数据层(dbbardata/constituent_unified/valuation_baostock),零 online,治幸存者偏差,喂 all_weather 策略。
**Architecture:** 新建 `LocalUnifiedProvider(bullet_trade.DataProvider)`,内部按数据类路由方案A 权威表:日线读 `dbbardata('d')` raw + `bs_adjust_factor` 算前复权;成份股读 `constituent_unified` 并集(治偏差);估值读 `valuation_baostock` parquet + 市值读 static/valuation akshare parquet。Mac 测试用 `sqlite :memory:` + tmp parquet fixture,零 VPS 依赖。
**Tech Stack:** Python 3.10, pandas 2.3, sqlite3, pyarrow, pytest
## Global Constraints(spec + 用户铁律)
- **零 online**: provider 不 import baostock 调 online,纯读本地 DB/parquet(memory provider-local-data-only)。baostock 48000/天限频不波及使用层。
- **surgical**: 不改 `LocalParquetProvider`/`BaostockProvider`(旧链路保留,向后兼容)。
- **dbbardata 不破坏**: `UNIQUE(symbol,exchange,interval,datetime)`,只读不写。
- **复权**: dbbardata 存 raw,消费端按 `bs_adjust_factor.foreAdjustFactor` 算前复权(§14.7 最终目标,用户定不降级)。
- **constituent_unified 并集模型**: 表无 date 列,`get_index_stocks(date)` 返回 in_currentwas_removed 并集,date 参数无法精确时点过滤——治"纯当前幸存者"偏差,有轻微前视(使用说明标注)。
- **代码归一**: jq_code `600519.XSHG` ↔ dbbardata `symbol=600519, exchange=SSE`;`SSE→SH, SZSE→SZ`
## 实测 schema(VPS 2026-07-23 probe,执行 agent 必读)
DB = `C:\sanguo_vnpy_v2\data\quant_trading.db`(VPS) / Mac 测试用 fixture 路径。
**dbbardata('d')** — 唯一行情表,raw 真实价:
```
列: symbol TEXT, exchange TEXT(SSE/SZSE), datetime TEXT(YYYY-MM-DD HH:MM:SS),
interval TEXT('d'), volume REAL, turnover REAL, open_interest REAL,
open_price REAL, high_price REAL, low_price REAL, close_price REAL
样本: 600519 10056行 2001-08-27~2026-07-22; 000005退市 8146行~2024-04-26; 510300 ETF 3439行
```
**constituent_unified** — 成份股并集(无 date!):
```
列: index_code TEXT(如 '000300'), code TEXT(纯6位如 '000001'), code_name TEXT,
source TEXT('baostock'/'akshare'), in_current INT(0/1), was_removed INT(0/1)
分布: 000300=940(300当前+640被踢) 000905=1803 000016=195 000852=1000(全当前,历史不可补)
399001=702 399005=145 399006=175 399330=150 932000=2000(全当前)
```
**bs_adjust_factor** — 复权因子:
```
列: code TEXT('sh.600519'), dividOperateDate TEXT(YYYY-MM-DD),
foreAdjustFactor REAL, backAdjustFactor REAL, adjustFactor REAL
语义: foreAdjustFactor 按除权日分段,最新事件=1.0,递减往历史。qfq[t]=raw[t]*factor[date[t]]。
600519 有 12 事件: 2020-06-24=0.856267 ... 2026-06-26=1.0
```
**valuation_baostock/<year>.parquet** — baostock 估值(1990-2026 全年份):
```
列: symbol(6位), exchange(SH/SZ), date(YYYY-MM-DD), peTTM, psTTM, pcfNcfTTM, pbMRQ, turn, pctChg, isST
注: 无 market_cap/total_share 列! 市值从 static/valuation akshare 补。
```
**static/valuation/<code>_valuation.parquet** — akshare 估值(市值/股本来源,5530 文件):
```
中文列(见 LocalParquetProvider._VAL_COL_MAP): 总市值→total_market_cap, 流通市值→circ_market_cap,
总股本→total_share, PE(TTM)→pe_ttm, 市净率→pb ...
```
**static/{balance,income,cashflow}/<code>_<table>.parquet** — akshare 三表(balance 221列/income 170列):
```
通用列: SECUCODE, REPORT_DATE, REPORT_TYPE; balance 有 TOTAL_ASSETS/TOTAL_LIABILITIES/TOTAL_PARENT_EQUITY;
income 有 BASIC_EPS/OPERATE_INCOME/PARENT_NETPROFIT/OPERATE_INCOME_YOY
```
---
## File Structure
- **Create:** `sanguo_portfolio/providers/local_unified_provider.py` — LocalUnifiedProvider 类(~400行)
- **Modify:** `sanguo_portfolio/providers/__init__.py` — 导出 LocalUnifiedProvider
- **Modify:** `sanguo_portfolio/runner_backtest.py``build_provider``unified` 选项(choices + 分支)
- **Create:** `tests/portfolio/test_local_unified_provider.py` — DataProvider 契约单测(fixture: sqlite + tmp parquet)
- **Create:** `tests/portfolio/conftest.py` 追加 — `local_unified_provider` fixture(若需要,否则在测试文件内建)
- **Create:** `docs/portfolio_local_unified_provider.md` — 使用说明(架构/数据源/接口/复权/治偏差/Mac测试/部署)
---
## Task 0: 代码转换 + DB 连接辅助 + 复权因子构造
**Files:**
- Create: `sanguo_portfolio/providers/local_unified_provider.py`(本 task 建文件骨架 + 模块级辅助函数)
- Test: `tests/portfolio/test_local_unified_provider.py`
**Interfaces:**
- Produces: `jq_to_dbbardata(jq_code) -> (symbol, exchange)` / `dbbardata_to_jq(symbol, exchange) -> jq_code`; `_connect(cfg) -> sqlite3.Connection`; `_build_qfq_factor(code, conn, dates) -> pd.Series(factor indexed by date)`
- [ ] **Step 1: 写失败测试 — 代码转换**
```python
# tests/portfolio/test_local_unified_provider.py
from sanguo_portfolio.providers.local_unified_provider import (
jq_to_dbbardata, dbbardata_to_jq, LocalUnifiedProvider,
)
def test_jq_to_dbbardata_roundtrip():
assert jq_to_dbbardata("600519.XSHG") == ("600519", "SSE")
assert jq_to_dbbardata("000001.XSHE") == ("000001", "SZSE")
assert jq_to_dbbardata("600519") == ("600519", "SSE") # 纯6位推断
assert dbbardata_to_jq("600519", "SSE") == "600519.XSHG"
assert dbbardata_to_jq("000001", "SZSE") == "000001.XSHE"
```
- [ ] **Step 2: 跑测试确认 FAIL**`pytest tests/portfolio/test_local_unified_provider.py::test_jq_to_dbbardata_roundtrip -v`(ImportError)
- [ ] **Step 3: 实现模块骨架 + 代码转换**
```python
# sanguo_portfolio/providers/local_unified_provider.py
"""LocalUnifiedProvider: 读方案A 权威数据层, 零 online, 治幸存者偏差(spec §6)。
数据源(全本地 VPS C:\\sanguo_vnpy_v2\\data\\):
- 日线: dbbardata('d') raw + bs_adjust_factor 算前复权(§14.7)
- 成份股: constituent_unified 并集(治偏差,无 date 时点)
- 估值 pe/pb/ps/pcf: valuation_baostock/<year>.parquet(baostock 权威)
- 市值/股本: static/valuation akshare parquet(baostock valuation 无市值列)
- 三表: static/{balance,income,cashflow} akshare parquet
零 online: 不 import baostock 调 online。Mac 测试用 sqlite+parquet fixture。
"""
from __future__ import annotations
import logging, os, sqlite3
from datetime import datetime
from pathlib import Path
from typing import Any, Dict, List, Optional, Union
import pandas as pd
try:
from bullet_trade.data.providers.base import DataProvider # type: ignore
except ImportError:
class DataProvider: # type: ignore[no-redef]
name: str = "base"
logger = logging.getLogger(__name__)
_DEFAULT_DB = r"C:\sanguo_vnpy_v2\data\quant_trading.db"
_DEFAULT_DATA_DIR = r"C:\sanguo_vnpy_v2\data"
_JQ_SUFFIX_TO_EXC = {"XSHG": "SSE", "XSHE": "SZSE", "SH": "SSE", "SZ": "SZSE"}
_EXC_TO_JQ_SUFFIX = {"SSE": "XSHG", "SZSE": "XSHE"}
def jq_to_dbbardata(jq_code: str) -> tuple[str, str]:
"""600519.XSHG → ('600519', 'SSE')。纯6位按6开头=sh/0,3=sz 推断。"""
s = (jq_code or "").strip()
if "." not in s:
if len(s) == 6:
return s, ("SSE" if s.startswith("6") else "SZSE")
return s, "SSE"
code, suffix = s.split(".", 1)
return code, _JQ_SUFFIX_TO_EXC.get(suffix.upper(), "SSE")
def dbbardata_to_jq(symbol: str, exchange: str) -> str:
"""('600519','SSE') → '600519.XSHG'"""
jq_suffix = _EXC_TO_JQ_SUFFIX.get(str(exchange).upper(), "XSHG")
return f"{symbol}.{jq_suffix}"
# 复权因子代码转换: 600519.XSHG → 'sh.600519'(bs_adjust_factor.code 格式)
def _jq_to_bs_code(jq_code: str) -> str:
sym, exc = jq_to_dbbardata(jq_code)
prefix = "sh" if exc == "SSE" else "sz"
return f"{prefix}.{sym}"
```
- [ ] **Step 4: 跑测试确认 PASS**
- [ ] **Step 5: 写失败测试 — 复权因子构造**
```python
def test_build_qfq_factor(tmp_path):
# fixture: 2 除权事件, 最新=1.0
import sqlite3
db = tmp_path / "t.db"
c = sqlite3.connect(str(db))
c.execute("CREATE TABLE bs_adjust_factor(code TEXT, dividOperateDate TEXT, foreAdjustFactor REAL, backAdjustFactor REAL, adjustFactor REAL)")
c.executemany("INSERT INTO bs_adjust_factor VALUES(?,?,?,?,?)", [
("sh.600519", "2024-06-19", 0.90, 0, 0),
("sh.600519", "2025-06-19", 1.00, 0, 0),
])
c.commit(); c.close()
from sanguo_portfolio.providers.local_unified_provider import _build_qfq_factor
dates = pd.to_datetime(["2023-01-01", "2024-07-01", "2025-07-01"])
f = _build_qfq_factor("sh.600519", sqlite3.connect(str(db)), dates)
# 2023(早于最早事件)=0.90; 2024-07(between)=0.90; 2025-07(最新后)=1.00
assert abs(f.iloc[0] - 0.90) < 1e-6
assert abs(f.iloc[1] - 0.90) < 1e-6
assert abs(f.iloc[2] - 1.00) < 1e-6
```
- [ ] **Step 6: 实现 `_build_qfq_factor`** — asof join 逻辑(每个 date 找 ≤ 的最大 dividOperateDate 的 foreAdjustFactor;早于所有事件用最早;晚于所有用最新):
```python
def _build_qfq_factor(bs_code: str, conn: sqlite3.Connection,
dates: pd.Series) -> pd.Series:
"""构造每个 date 的前复权因子(asof)。qfq[t]=raw[t]*factor[t]。"""
rows = conn.execute(
"SELECT dividOperateDate, foreAdjustFactor FROM bs_adjust_factor "
"WHERE code=? ORDER BY dividOperateDate", (bs_code,)).fetchall()
if not rows:
return pd.Series([1.0] * len(dates), index=dates)
ev_dates = pd.to_datetime([r[0] for r in rows])
factors = [float(r[1]) for r in rows]
out = []
for d in pd.to_datetime(dates):
# 找 <= d 的最大事件; 全部 > d 用最早(第一个); 全部 <= d 用最后一个
mask = ev_dates <= d
out.append(factors[mask.argmax()] if mask.any() else factors[0])
# mask.argmax() 给第一个 True 的索引;但我们要"<= d 的最大事件"= 最后一个 True
# 修正:取最后一个 True
out = []
for d in pd.to_datetime(dates):
mask = ev_dates <= d
idx = int(np.where(mask)[0][-1]) if mask.any() else 0
out.append(factors[idx])
return pd.Series(out, index=pd.to_datetime(dates))
```
(注意:`np``import numpy as np`。实现时简化为单次循环取最后一个 True 索引。)
- [ ] **Step 7: 跑测试确认 PASS**
- [ ] **Step 8: Commit**`feat(portfolio): LocalUnifiedProvider 代码转换+复权因子(Task0)`
---
## Task 1: get_price(dbbardata raw + 前复权 + panel 长表)
**Files:** Modify `local_unified_provider.py``__init__` + `get_price`; Test 同文件。
**Interfaces:**
- Consumes: Task0 辅助函数 + `_connect`
- Produces: `LocalUnifiedProvider.get_price(security, start_date, end_date, frequency, fields, skip_paused, fq, count, panel, fill_paused) -> DataFrame`
策略契约(all_weather 实证):
- `get_price(hold_list, end_date, freq=daily, fields=[close,high_limit], count=1, panel=False)` — panel=False 长表需 time/code 列
- `get_price(stocks, freq=1d, fields=[close], count=n, panel=False)` — _trend_mean pivot(index=time,columns=code)
- `get_price(stock, freq=1m, fq="pre", count=1, panel=False)` — intraday(day 频率回测降级,1m 无数据返空)
- [ ] **Step 1: 写失败测试 — get_price daily 单股 + 复权**
```python
@pytest.fixture
def unified_provider(tmp_path):
"""造小样本 sqlite + parquet fixture。"""
db = tmp_path / "quant_trading.db"
c = sqlite3.connect(str(db))
c.execute("CREATE TABLE dbbardata(symbol,exchange,datetime,interval,volume,turnover,open_interest,open_price,high_price,low_price,close_price)")
rows = [
("600519","SSE","2024-06-18 00:00:00","d",1000,1e6,0,1000.0,1010.0,990.0,1000.0), # 除权前
("600519","SSE","2024-06-19 00:00:00","d",1000,1e6,0,900.0,910.0,890.0,900.0), # 除权日 raw 跳水
("600519","SSE","2024-06-20 00:00:00","d",1000,1e6,0,910.0,920.0,900.0,910.0),
]
c.executemany("INSERT INTO dbbardata VALUES(?,?,?,?,?,?,?,?,?,?,?)", rows)
c.execute("CREATE TABLE bs_adjust_factor(code,dividOperateDate,foreAdjustFactor,backAdjustFactor,adjustFactor)")
c.execute("INSERT INTO bs_adjust_factor VALUES('sh.600519','2024-06-19',0.9,0,0)") # 除权日 factor
c.commit(); c.close()
return LocalUnifiedProvider({"db_path": str(db), "data_dir": str(tmp_path)})
def test_get_price_raw_vs_qfq(unified_provider):
p = unified_provider
# raw: 除权日 900 跳水
df_raw = p.get_price("600519.XSHG", start_date="2024-06-18", end_date="2024-06-20", fq="raw")
assert len(df_raw) == 3
assert abs(df_raw.loc["2024-06-19", "close"] - 900.0) < 1e-6
# qfq: 06-18 = 1000*0.9 = 900; 06-19/20 = raw(factor=0.9 当 06-19 之后? 用最新段逻辑)
df_qfq = p.get_price("600519.XSHG", start_date="2024-06-18", end_date="2024-06-20", fq="qfq")
assert abs(df_qfq.loc["2024-06-18", "close"] - 900.0) < 1e-6 # 1000*0.9(早于事件用最早factor)
```
(复权断言:06-18 早于除权日 06-19 → 用 factor 0.9 → 1000*0.9=900;06-19/20 ≥ 事件日 → factor 取 06-19 的 0.9 → 900*0.9=810, 910*0.9=819。实现时按 `_build_qfq_factor` 语义校准断言。)
- [ ] **Step 2: 跑测试确认 FAIL**
- [ ] **Step 3: 实现 `__init__` + `get_price`**
```python
class LocalUnifiedProvider(DataProvider): # type: ignore[misc]
name: str = "sanguo_local_unified"
requires_live_data: bool = False
def __init__(self, config: Optional[Dict[str, Any]] = None) -> None:
cfg = config or {}
self.db_path: str = cfg.get("db_path", _DEFAULT_DB)
self.data_dir: str = cfg.get("data_dir", _DEFAULT_DATA_DIR)
self._conn: Optional[sqlite3.Connection] = None
self._val_bs_cache: Dict[int, pd.DataFrame] = {} # year -> valuation_baostock
def _connect(self) -> sqlite3.Connection:
if self._conn is None:
self._conn = sqlite3.connect(self.db_path, timeout=30)
self._conn.execute("PRAGMA busy_timeout = 30000")
return self._conn
def get_price(self, security, start_date=None, end_date=None, frequency="daily",
fields=None, skip_paused=False, fq="raw", count=None,
panel=True, fill_paused=True, **kwargs):
freq = str(frequency or "").lower()
if freq not in ("daily", "day", "1d", "d"):
return pd.DataFrame() # 1m/分钟 day 频率回测降级(数据层无 1m)
secs = [security] if isinstance(security, str) else list(security or [])
conn = self._connect()
start_str = self._to_date_str(start_date)
end_str = self._to_date_str(end_date) or datetime.now().strftime("%Y-%m-%d")
frames: Dict[str, pd.DataFrame] = {}
for jq_code in secs:
sym, exc = jq_to_dbbardata(jq_code)
q = "SELECT datetime, open_price, high_price, low_price, close_price, " \
"volume, turnover FROM dbbardata WHERE symbol=? AND exchange=? " \
"AND interval='d' AND datetime>=? AND datetime<=? ORDER BY datetime"
df = pd.read_sql(q, conn, params=(sym, exc, start_str + " 00:00:00", end_str + " 23:59:59"))
if df.empty:
frames[jq_code] = df; continue
df["datetime"] = pd.to_datetime(df["datetime"])
df = df.set_index("datetime")
df.index.name = None
if count:
df = df.tail(count)
# 复权
if fq in ("qfq", "pre", "前复权"):
factor = _build_qfq_factor(_jq_to_bs_code(jq_code), conn, df.index)
for col in ("open_price", "high_price", "low_price", "close_price"):
df[col] = df[col].values * factor.values
# 策略要 close/high_limit 字段名(jq 风格)
df = df.rename(columns={"open_price": "open", "high_price": "high",
"low_price": "low", "close_price": "close"})
# high_limit 不在 dbbardata, 留给 get_current_tick 语义;这里策略 prepare_stock_list 要 high_limit 列
# → 缺失列返 NaN(策略 hit = close==high_limit 不会命中,降级可接受)
if fields:
for f in fields:
if f not in df.columns:
df[f] = float("nan")
df = df[fields]
frames[jq_code] = df
if not frames or all(f.empty for f in frames.values()):
return pd.DataFrame()
if not panel:
parts = []
for jq_code, df in frames.items():
if df.empty:
continue
d = df.reset_index().rename(columns={"datetime": "time"})
d.insert(0, "code", jq_code)
parts.append(d)
return pd.concat(parts, ignore_index=True) if parts else pd.DataFrame()
if len(frames) == 1:
return next(iter(frames.values()))
return pd.concat(frames, axis=1)
```
- [ ] **Step 4: 跑测试确认 PASS**
- [ ] **Step 5: 写失败测试 — panel=False 多股长表 + count**
```python
def test_get_price_panel_false_multi(unified_provider):
df = unified_provider.get_price("600519.XSHG", end_date="2024-06-20", count=2, panel=False, fields=["close"])
assert "code" in df.columns and "time" in df.columns
assert len(df) == 2
```
- [ ] **Step 6: 实现(Step 3 已含 panel 分支),跑 PASS**
- [ ] **Step 7: Commit**`feat(portfolio): LocalUnifiedProvider get_price+前复权(Task1)`
---
## Task 2: get_index_stocks + get_constituent(constituent_unified 并集,治偏差)
**Files:** Modify `local_unified_provider.py`; Test 同文件。
**Interfaces:**
- Produces: `get_index_stocks(index_symbol, date) -> List[str]` + `get_constituent(index, date) -> List[str]`(语义别名)
- [ ] **Step 1: 写失败测试**
```python
def test_get_index_stocks_union(tmp_path):
db = tmp_path / "t.db"; c = sqlite3.connect(str(db))
c.execute("CREATE TABLE constituent_unified(index_code TEXT,code TEXT,code_name TEXT,source TEXT,in_current INT,was_removed INT)")
c.executemany("INSERT INTO constituent_unified VALUES(?,?,?,?,?,?)", [
("000300", "600519", "贵州茅台", "baostock", 1, 0),
("000300", "000001", "平安银行", "baostock", 1, 0),
("000300", "600811", "退市股", "baostock", 0, 1), # 被踢
])
c.commit(); c.close()
p = LocalUnifiedProvider({"db_path": str(db), "data_dir": str(tmp_path)})
stocks = p.get_index_stocks("000300.XSHG", "2020-01-01")
assert set(stocks) == {"600519.XSHG", "000001.XSHE", "600811.SH"} # 并集含被踢
# date 参数不报错(并集模型忽略)
assert p.get_constituent("000300", None) == stocks # 别名
```
- [ ] **Step 2: 跑测试确认 FAIL**
- [ ] **Step 3: 实现** — 查 constituent_unified,index_code 匹配(去 `.XXXX` 后缀),返回 in_current=1 OR was_removed=1 的并集,code→jq_code:
```python
def get_index_stocks(self, index_symbol, date=None) -> List[str]:
idx = index_symbol.split(".")[0] if "." in str(index_symbol) else str(index_symbol)
conn = self._connect()
rows = conn.execute(
"SELECT code FROM constituent_unified WHERE index_code=? "
"AND (in_current=1 OR was_removed=1)", (idx,)).fetchall()
out = []
for (code,) in rows:
code = str(code).strip()
if len(code) != 6:
continue
exc = "SSE" if code.startswith("6") else "SZSE"
out.append(dbbardata_to_jq(code, exc))
return out
def get_constituent(self, index, date=None) -> List[str]:
"""spec §6 语义别名 = get_index_stocks。"""
return self.get_index_stocks(index, date)
```
- [ ] **Step 4: 跑测试 PASS**
- [ ] **Step 5: Commit**`feat(portfolio): LocalUnifiedProvider 成份股并集治偏差(Task2)`
---
## Task 3: get_fundamentals_df(valuation_baostock + static akshare + 三表)
**Files:** Modify `local_unified_provider.py`; Test 同文件 + tmp parquet fixture。
**Interfaces:**
- Produces: `get_fundamentals_df(stocks, date) -> DataFrame` 列对齐 `_FUNDAMENTAL_COLUMNS`
数据源映射:
- `pe_ratio/pb_ratio/ps_ratio/pcf_ratio` ← valuation_baostock parquet(peTTM/pbMRQ/psTTM/pcfNcfTTM,baostock 权威)
- `market_cap/circulating_market_cap` ← static/valuation akshare parquet(total_market_cap/circ_market_cap,baston 无市值)
- 三表字段(eps/net_profit_margin/total_liability 等) ← static/{balance,income} akshare parquet(复用 LocalParquetProvider 读法)
- [ ] **Step 1: 写失败测试 — 估值字段从 valuation_baostock**
```python
def test_get_fundamentals_valuation(tmp_path):
# valuation_baostock/2024.parquet
vdir = tmp_path / "valuation_baostock"; vdir.mkdir()
pd.DataFrame({"symbol":["600519"],"exchange":["SH"],"date":["2024-09-30"],
"peTTM":[25.0],"psTTM":[15.0],"pcfNcfTTM":[20.0],"pbMRQ":[7.5],
"turn":[0.1],"pctChg":[1.0],"isST":[0]}).to_parquet(vdir/"2024.parquet")
# static/valuation akshare(市值)
sdir = tmp_path / "static" / "valuation"; sdir.mkdir(parents=True)
pd.DataFrame({"数据日期":["2024-09-30"],"总市值":[2e12],"流通市值":[2e12],"总股本":[1.256e9],
"PE(TTM)":[25],"市净率":[7.5]}).to_parquet(sdir/"600519.SH_valuation.parquet")
p = LocalUnifiedProvider({"db_path": str(tmp_path/"t.db"), "data_dir": str(tmp_path)})
df = p.get_fundamentals_df(["600519.XSHG"], date="2024-09-30")
assert abs(df.loc["600519.XSHG","pe_ratio"] - 25.0) < 1e-6 # baostock 权威
assert abs(df.loc["600519.XSHG","pb_ratio"] - 7.5) < 1e-6
assert abs(df.loc["600519.XSHG","market_cap"] - 2e4) < 1 # 2e12元→2e4亿
```
- [ ] **Step 2: 跑测试确认 FAIL**
- [ ] **Step 3: 实现** — 读 valuation_baostock parquet(year from date)+ static/valuation akshare;合并对齐 `_FUNDAMENTAL_COLUMNS`(复用 LocalParquetProvider 的 `_VAL_COL_MAP` / `to_yi` / 三表读法,import 复用):
```python
from .local_parquet_provider import (_VAL_COL_MAP, jq_to_file_code,
_to_float, _or_nan, _pct_to_decimal, _FUNDAMENTAL_COLUMNS)
from ..factors.valuation import to_yi
def get_fundamentals_df(self, stocks, date=None) -> pd.DataFrame:
if not stocks:
return pd.DataFrame(columns=_FUNDAMENTAL_COLUMNS)
date_str = self._to_date_str(date) or datetime.now().strftime("%Y-%m-%d")
rows = [self._build_fundamental_row(s, date_str) for s in stocks]
df = pd.DataFrame(rows, columns=_FUNDAMENTAL_COLUMNS)
if "code" in df.columns:
df = df.set_index("code", drop=False)
return df
def _read_valuation_baostock(self, year: int) -> pd.DataFrame:
if year in self._val_bs_cache:
return self._val_bs_cache[year]
p = os.path.join(self.data_dir, "valuation_baostock", f"{year}.parquet")
df = pd.read_parquet(p) if os.path.exists(p) else pd.DataFrame()
self._val_bs_cache[year] = df
return df
def _build_fundamental_row(self, jq_code, date_str) -> Dict[str, Any]:
sym, exc = jq_to_dbbardata(jq_code)
fc = jq_to_file_code(jq_code) # 600519.SH(static akshare 文件名)
row: Dict[str, Any] = {"code": jq_code}
# 1. pe/pb/ps/pcf ← valuation_baostock(baostock 权威)
year = int(date_str[:4])
vbs = self._read_valuation_baostock(year)
if not vbs.empty:
sub = vbs[(vbs["symbol"].astype(str) == sym) & (vbs["date"].astype(str) <= date_str)]
vrow = sub.iloc[-1] if not sub.empty else None
else:
vrow = None
def gbs(k):
return _to_float(vrow.get(k)) if vrow is not None else None
row["pe_ratio"] = _or_nan(gbs("peTTM"))
row["pb_ratio"] = _or_nan(gbs("pbMRQ"))
row["ps_ratio"] = _or_nan(gbs("psTTM"))
row["pcf_ratio"] = _or_nan(gbs("pcfNcfTTM"))
# 2. 市值/股本 + 三表 ← static akshare(复用 LocalParquetProvider 读法)
# 复用:直接实例化 LocalParquetProvider 读 static 部分,或内联读 static/valuation
ak_val = self._read_akshare_valuation(fc, date_str) # 返 renamed Series
mkt = _to_float(ak_val.get("total_market_cap")) if ak_val is not None else None
circ = _to_float(ak_val.get("circ_market_cap")) if ak_val is not None else None
row["market_cap"] = to_yi(mkt) if mkt else float("nan")
row["circulating_market_cap"] = to_yi(circ) if circ else float("nan")
# 3. 三表(income/balance)— 复用 LocalParquetProvider._read_quarter + 字段提取
# 简化:委托一个内部 LocalParquetProvider 实例读三表部分(eps/margin/liability)
lpp = self._get_lpp_helper()
inc = lpp._latest_row_before(lpp._read_quarter("income", fc), "REPORT_DATE", date_str)
bal = lpp._latest_row_before(lpp._read_quarter("balance", fc), "REPORT_DATE", date_str)
row["eps"] = _or_nan(_to_float(inc.get("BASIC_EPS")) if inc is not None else None)
# ... net_profit_margin/total_liability/roe 等(照 LocalParquetProvider._build_fundamental_row 逻辑)
return row
```
(实现时:`_get_lpp_helper()` 返一个复用的 `LocalParquetProvider(config)` 实例读 static 三表;`_read_akshare_valuation` 复用 LocalParquetProvider._read_valuation。DRY:不重写三表/akshare valuation 逻辑,委托 LocalParquetProvider。pe/pb 改 baostock 源覆盖 akshare 的。)
- [ ] **Step 4: 跑测试 PASS**
- [ ] **Step 5: 写测试 — 三表字段(eps/market_cap 全 _FUNDAMENTAL_COLUMNS 有值不 NaN)**
- [ ] **Step 6: 实现 + PASS**
- [ ] **Step 7: Commit**`feat(portfolio): LocalUnifiedProvider fundamentals baostock估值+akshare市值(Task3)`
---
## Task 4: 辅助方法(trade_days/all_securities/security_info/current_tick/split_dividend)
**Files:** Modify `local_unified_provider.py`; Test 同文件。
- [ ] **Step 1-2: 写失败测试 + FAIL**`get_trade_days(count=2)` 返 datetime list;`get_security_info` 返 display_name/start_date;`get_current_tick` 返 close+high_limit;`get_split_dividend` 返 bs_adjust_factor 事件;`get_all_securities` 返 dbbardata distinct symbol。
- [ ] **Step 3: 实现**:
- `get_trade_days`: 读 dbbardata 某 symbol(如 600519)distinct datetime,filter/count。
- `get_security_info`: dbbardata min/max datetime → start/end_date;display_name 从 constituent_unified code_name 或 code。
- `get_current_tick`: dbbardata 最近 close + valuation_baostock 最近 pctChg → high_limit=close×1.1(ST 0.05)。
- `get_split_dividend`: bs_adjust_factor → events(dividOperateDate + adjustFactor)。
- `get_all_securities`: dbbardata distinct symbol → DataFrame。
- [ ] **Step 4: 跑测试 PASS**
- [ ] **Step 5: Commit**`feat(portfolio): LocalUnifiedProvider 辅助方法(Task4)`
---
## Task 5: 接线(__init__ 导出 + runner build_provider 加 unified)
**Files:** Modify `sanguo_portfolio/providers/__init__.py`; Modify `sanguo_portfolio/runner_backtest.py`
- [ ] **Step 1: __init__.py 加导出**
```python
from .local_unified_provider import LocalUnifiedProvider
__all__ = ["SanguoMiniQmtProvider", "BaostockProvider", "LocalParquetProvider", "LocalUnifiedProvider"]
```
- [ ] **Step 2: runner_backtest build_provider 加 unified**
```python
# parse_args choices 加 "unified"; build_provider 加分支
p.add_argument("--provider", default="local", choices=["local", "baostock", "miniqmt", "unified"], ...)
# build_provider:
from .providers import LocalUnifiedProvider
if name == "unified":
return LocalUnifiedProvider(cfg)
```
- [ ] **Step 3: 跑 `pytest tests/portfolio/ -v` 全绿(回归)**
- [ ] **Step 4: Commit**`feat(portfolio): 接线 LocalUnifiedProvider 到 runner(Task5)`
---
## Task 6: 使用说明 + VPS E2E 验证
**Files:** Create `docs/portfolio_local_unified_provider.md`; VPS 跑 `python -m sanguo_portfolio.runner_backtest --provider unified --start 2024-01-01 --end 2024-03-31 --max-pool 20`
- [ ] **Step 1: 写使用说明** `docs/portfolio_local_unified_provider.md`(其他 session 直用)— 含:
- 一句话定位(读方案A权威层/零online/治偏差)
- 数据源映射表(每接口→哪张表/parquet)
- 接口清单(DataProvider 接口 + get_constituent)
- 复权说明(raw存储+消费端按bs_adjust_factor算qfq;fq参数 raw/qfq)
- **幸存者偏差说明**(constituent_unified 并集模型,治纯当前偏差,有轻微前视,date 参数忽略;中证1000/2000只快照永久gap)
- Mac 测试(fixture,零VPS依赖)
- 部署/运行(runner --provider unified;VPS 数据依赖 dbbardata/constituent_unified/valuation_baostock/static)
- 已知限制(high_limit 列 NaN→prepare_stock_list 涨停识别降级;1m 无数据;三表委托 LocalParquetProvider)
- 与旧 provider 关系(LocalParquetProvider/BaostockProvider 保留,unified 是方案A 后推荐)
- [ ] **Step 2: VPS E2E** — rsync 代码到 VPS,跑 `--provider unified --max-pool 20` 小样本回测,确认:
- get_price 读 dbbardata 出 K 线(含退市)
- get_index_stocks 出并集成份股
- get_fundamentals_df 出市值+pe/pb
- 回测不崩,有选股+指标输出
- [ ] **Step 3: Commit**`docs(portfolio): LocalUnifiedProvider 使用说明+VPS E2E(Task6)`
---
## Self-Review(plan 自检)
1. **Spec 覆盖**: spec §6 接口(get_daily/get_constituent/get_fundamentals/...)— get_constituent 别名✓;get_price 覆盖 get_daily+get_etf_daily(都读 dbbardata,ETF 也在);get_fundamentals_df ✓;其余 §6 方法(industry/longhubang/instrument)数据层未就绪(P1),使用说明标注 NotImplementedError。✓
2. **方案A §14 一致**: dbbardata 唯一行情✓;constituent_unified 治偏差✓;valuation_baostock pe/pb✓;raw+factor 复权✓;零online✓。
3. **类型一致**: `_build_qfq_factor(code, conn, dates) -> Series` 在 Task0/Task1 调用签名一致✓。
4. **占位扫描**: Task3 的 `_get_lpp_helper/_read_akshare_valuation` 标了"复用 LocalParquetProvider",实现 agent 须内联或委托,不留空✓。
5. **风险**: get_price 的 high_limit 列缺失(NaN)→策略 prepare_stock_list 涨停识别降级,使用说明标注(Task6)✓。
## Execution Handoff
Plan complete and saved to `docs/superpowers/plans/2026-07-23-local-unified-provider.md`.
@@ -1,77 +0,0 @@
# Phase 1 数据层 — 完成报告
**日期**2026-07-05
**状态**:✅ DONE(本地 16 passed + 端到端冒烟通过)
**分支**:已 merge to master`20e2437`
## 1. 目标
把 v1sanguo_vnpy)散落的数据资产(BaoStock parquet + vnpy SQLite db)收拢成统一的、可被 vnpy 原生调用的数据读取层;打通 NAS Docker 端到端链路。
**范围约束**:只做数据"读/写/调度/校验",不做回测/策略/Web(留 Phase 2+)。
## 2. 交付清单
### 代码模块(`sanguo_data/`
| 模块 | 职责 | 关键设计 |
|------|------|---------|
| `config.py` | 加载 `data_platform.yaml` | 统一配置入口 |
| `datareader.py` | 读日 K → `BarData` | `read_parquet_daily` + `read_db_daily`(配 vnpy SETTINGS+ `guess_exchange`(代码前缀判 SSE/SZSE |
| `datafeed.py` | 在线数据源抓取 | BaoStock + 东财/腾讯 fallback`_fetch_baostock_with_timeout` 子进程隔离(修 v1 卡死坑) |
| `datafeed/` | 数据源实现子包 | 可扩展 |
| `datawriter.py` | 数据落盘 | parquet 年分区 + vnpy sqlite 原子写 |
| `database/` | vnpy 数据库适配子包 | sqlite 路径配置 |
| `scheduler.py` | `UpdateScheduler` | 增量更新 + 断点续传 + 熔断 |
| `validator.py` | 数据完整性校验 | 缺口/异常值检测 |
### 测试(`tests/data/`16 passed
`test_config` / `test_datareader` / `test_datafeed` / `test_datawriter` / `test_scheduler` / `test_validator` / `test_spike_vnpy44`
## 3. 关键 Fix(收尾阶段)
| Gap | commit | 说明 |
|-----|--------|------|
| read_db_daily 没配 vnpy database 路径 | `339d85a` | 默认空 database.db 读不到 NAS 真实数据 → 加 `SETTINGS["database.database"]` 配置 |
| spike 测试被全局 SETTINGS 污染 | `8829501` | read_db_daily 改全局 SETTINGS → spike 用 monkeypatch mock load_bar_data 隔离 |
## 4. 部署产出(NAS Docker
- 容器 `sanguo_vnpy_v2` 重建加挂载 `/volume1/stock`RW
- 镜像 `sanguo_vnpy_v2:with-sqlite`docker commit 保 vnpy_sqlite,避免 NAS 弱 CPU pip 卡死)
- 外网链路 `vnpy.mysanguo.top` 恢复正常(未改动原代理/转发设计)
## 5. 端到端验证(真实数据)
`scripts/smoke_e2e.py` 在临时容器内读 NAS `quant_trading.db`
```
600000.SSE: 541 条
000001.SZSE: 563 条
300750.SZSE: 564 条
区间: 2024-01-01 ~ 2026-06-30
```
## 6. 风险与兜底
- **v1 db 被容器改**:当前 `quant_trading.db` 1.66→2.13GB(容器 RW 下 peewee create_tables 副作用)。
- **兜底**`/volume1/stock/sanguo_vnpy/data/` 含多个历史备份 `quant_trading_*.db.bak`1.5G 各,20260519~20260522 + pre_interval_migration),原始可恢复。数据完整性已验证。
## 7. 已知 Tech Debt(留后续 phase
| 项 | 影响 |
|----|------|
| `_load_stock_list` NotImplementedError | 部署前 copy v1,后续需补实现 |
| datawriter 一致性回滚 | final triage 待定 |
| validator row-by-row 全量校验性能 | 增量场景 OK,海量回测时需优化 |
| peewee 版本冲突告警 | empyrical-reloaded 要 `peewee<3.17.4`vnpy_sqlite 装 4.1.1 — **可能影响 Phase 2 因子层** |
## 8. Phase 2 衔接点
Phase 1 提供的下游可用接口:
- `read_db_daily(symbol, start, end, cfg) → list[BarData]` — vnpy 原生 BarData,可直接喂 vnpy.alpha / ctabacktester
- `read_parquet_daily(...)` — parquet 快速读取路径
- `UpdateScheduler` — 增量数据更新
Phase 2(因子/回测层)应基于这些接口构建,遵循 PRD ADR:vnpy.alpha 集成、因子层可插拔、multiprocessing+Ray 编排、不引 Qlib。