merge(HEL-193): 以 .11 线上基线 013ed29 整合行情历史补档 8e94c7b(HEL-190/HEL-199)
- 合入 agent/agent/ebb4e3e0638a:真实交易日历补最近60日快照 + sys.path 引导返工 - architecture-inventory 于合并后重新生成 Co-authored-by: multica-agent <github@multica.ai>
This commit is contained in:
@@ -0,0 +1,202 @@
|
||||
"""Auditable recent-trading-day snapshot backfill helpers.
|
||||
|
||||
Planning and backup stay free of provider imports so feature boundary tests remain green.
|
||||
The service layer supplies open trading dates from the live calendar and executes sync.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import sqlite3
|
||||
from datetime import date, datetime, timedelta
|
||||
from pathlib import Path
|
||||
from typing import Any, Iterable
|
||||
|
||||
|
||||
MAX_RANGE_TRADING_DAYS = 15
|
||||
MAX_RECENT_TRADING_DAYS = 60
|
||||
DEFAULT_RECENT_TRADING_DAYS = 60
|
||||
|
||||
# Tables touched by a successful historical dashboard sync. User / token / model
|
||||
# tables must never appear here.
|
||||
SNAPSHOT_BACKFILL_WRITE_TABLES = frozenset(
|
||||
{
|
||||
"dashboard_snapshots",
|
||||
"data_snapshots",
|
||||
"sync_runs",
|
||||
}
|
||||
)
|
||||
|
||||
|
||||
def clamp_recent_lookback(lookback: int) -> int:
|
||||
value = int(lookback)
|
||||
if value < 1:
|
||||
raise ValueError("回补交易日数量至少为 1。")
|
||||
if value > MAX_RECENT_TRADING_DAYS:
|
||||
raise ValueError(f"单次最多回补最近 {MAX_RECENT_TRADING_DAYS} 个交易日。")
|
||||
return value
|
||||
|
||||
|
||||
def calendar_window_start(end_date: str, lookback: int) -> str:
|
||||
"""Natural-day lower bound large enough to cover lookback open sessions."""
|
||||
end = datetime.strptime(end_date, "%Y%m%d").date()
|
||||
span = max(40, int(lookback * 2) + 20)
|
||||
return (end - timedelta(days=span)).strftime("%Y%m%d")
|
||||
|
||||
|
||||
def select_open_trade_dates(
|
||||
calendar_rows: Iterable[dict[str, Any]],
|
||||
end_date: str,
|
||||
lookback: int,
|
||||
) -> list[str]:
|
||||
"""Pick the last ``lookback`` open SSE sessions on or before ``end_date``."""
|
||||
lookback = clamp_recent_lookback(lookback)
|
||||
end = normalize_compact_date(end_date)
|
||||
open_dates = sorted(
|
||||
{
|
||||
normalize_compact_date(str(row.get("cal_date") or ""))
|
||||
for row in calendar_rows
|
||||
if int(row.get("is_open") or 0) == 1 and row.get("cal_date")
|
||||
}
|
||||
)
|
||||
open_dates = [item for item in open_dates if item <= end]
|
||||
if not open_dates:
|
||||
raise ValueError("交易日历未返回可用交易日,请检查行情 Token。")
|
||||
return open_dates[-lookback:]
|
||||
|
||||
|
||||
def select_open_trade_dates_in_range(
|
||||
calendar_rows: Iterable[dict[str, Any]],
|
||||
start_date: str,
|
||||
end_date: str,
|
||||
*,
|
||||
maximum: int = MAX_RANGE_TRADING_DAYS,
|
||||
) -> tuple[list[str], list[str]]:
|
||||
"""Return (open_dates, skipped_non_trading_days) inside an inclusive range."""
|
||||
start = normalize_compact_date(start_date)
|
||||
end = normalize_compact_date(end_date)
|
||||
if start > end:
|
||||
raise ValueError("开始日期不能晚于结束日期。")
|
||||
open_set = {
|
||||
normalize_compact_date(str(row.get("cal_date") or ""))
|
||||
for row in calendar_rows
|
||||
if int(row.get("is_open") or 0) == 1 and row.get("cal_date")
|
||||
}
|
||||
open_dates: list[str] = []
|
||||
skipped: list[str] = []
|
||||
cursor = datetime.strptime(start, "%Y%m%d").date()
|
||||
last = datetime.strptime(end, "%Y%m%d").date()
|
||||
while cursor <= last:
|
||||
compact = cursor.strftime("%Y%m%d")
|
||||
if compact in open_set:
|
||||
open_dates.append(compact)
|
||||
else:
|
||||
skipped.append(compact)
|
||||
cursor += timedelta(days=1)
|
||||
if len(open_dates) > maximum:
|
||||
raise ValueError(f"单次最多回补 {maximum} 个交易日。")
|
||||
return open_dates, skipped
|
||||
|
||||
|
||||
def classify_snapshot_coverage(
|
||||
trade_dates: list[str],
|
||||
existing_dates: Iterable[str],
|
||||
) -> dict[str, Any]:
|
||||
present_set = {
|
||||
normalize_compact_date(item)
|
||||
for item in existing_dates
|
||||
if item
|
||||
}
|
||||
present = [item for item in trade_dates if item in present_set]
|
||||
missing = [item for item in trade_dates if item not in present_set]
|
||||
return {
|
||||
"trade_dates": list(trade_dates),
|
||||
"present": present,
|
||||
"missing": missing,
|
||||
"present_count": len(present),
|
||||
"missing_count": len(missing),
|
||||
}
|
||||
|
||||
|
||||
def create_sqlite_backup(
|
||||
source_path: Path,
|
||||
backup_dir: Path,
|
||||
*,
|
||||
label: str = "pre-backfill",
|
||||
stamped_at: datetime | None = None,
|
||||
) -> Path:
|
||||
"""Create a timestamped SQLite backup via the native backup API."""
|
||||
source = Path(source_path)
|
||||
if not source.exists():
|
||||
raise FileNotFoundError(f"数据库不存在:{source}")
|
||||
stamp = (stamped_at or datetime.now().astimezone()).strftime("%Y%m%d-%H%M%S")
|
||||
safe_label = "".join(ch if ch.isalnum() or ch in "-_" else "-" for ch in label).strip("-") or "backup"
|
||||
backup_dir = Path(backup_dir)
|
||||
backup_dir.mkdir(parents=True, exist_ok=True)
|
||||
target = backup_dir / f"review-{safe_label}-{stamp}.db"
|
||||
source_conn = sqlite3.connect(f"file:{source}?mode=ro", uri=True)
|
||||
try:
|
||||
target_conn = sqlite3.connect(target)
|
||||
try:
|
||||
source_conn.backup(target_conn)
|
||||
target_conn.commit()
|
||||
finally:
|
||||
target_conn.close()
|
||||
finally:
|
||||
source_conn.close()
|
||||
return target
|
||||
|
||||
|
||||
def display_date(compact: str) -> str:
|
||||
value = normalize_compact_date(compact)
|
||||
return f"{value[:4]}-{value[4:6]}-{value[6:8]}"
|
||||
|
||||
|
||||
def normalize_compact_date(value: str) -> str:
|
||||
compact = str(value or "").replace("-", "").strip()
|
||||
if len(compact) != 8 or not compact.isdigit():
|
||||
raise ValueError("日期格式应为 YYYY-MM-DD。")
|
||||
datetime.strptime(compact, "%Y%m%d")
|
||||
return compact
|
||||
|
||||
|
||||
def build_backfill_audit(
|
||||
*,
|
||||
mode: str,
|
||||
end_date: str,
|
||||
lookback: int | None,
|
||||
coverage: dict[str, Any],
|
||||
skipped_non_trading_days: list[str] | None = None,
|
||||
backup_path: str | None = None,
|
||||
dry_run: bool = False,
|
||||
results: list[dict[str, Any]] | None = None,
|
||||
) -> dict[str, Any]:
|
||||
results = list(results or [])
|
||||
succeeded = [row for row in results if row.get("status") == "success"]
|
||||
skipped = [row for row in results if row.get("status") == "skipped"]
|
||||
failed = [row for row in results if row.get("status") == "failed"]
|
||||
return {
|
||||
"ok": not failed,
|
||||
"mode": mode,
|
||||
"dry_run": dry_run,
|
||||
"end_date": display_date(end_date),
|
||||
"lookback": lookback,
|
||||
"backup_path": backup_path,
|
||||
"write_tables": sorted(SNAPSHOT_BACKFILL_WRITE_TABLES),
|
||||
"trade_dates": [display_date(item) for item in coverage.get("trade_dates") or []],
|
||||
"present": [display_date(item) for item in coverage.get("present") or []],
|
||||
"missing": [display_date(item) for item in coverage.get("missing") or []],
|
||||
"skipped_non_trading_days": [
|
||||
display_date(item) for item in (skipped_non_trading_days or [])
|
||||
],
|
||||
"present_count": int(coverage.get("present_count") or 0),
|
||||
"missing_count": int(coverage.get("missing_count") or 0),
|
||||
"results": results,
|
||||
"succeeded_count": len(succeeded),
|
||||
"skipped_count": len(skipped),
|
||||
"failed_count": len(failed),
|
||||
"created_dates": [
|
||||
str(row.get("trade_date") or "")
|
||||
for row in succeeded
|
||||
if row.get("action") == "created"
|
||||
],
|
||||
}
|
||||
@@ -227,6 +227,31 @@ class MarketRepositoryMixin:
|
||||
result.append(payload)
|
||||
return result
|
||||
|
||||
def list_snapshot_trade_dates(
|
||||
self,
|
||||
start_date: str = "",
|
||||
end_date: str = "",
|
||||
) -> list[str]:
|
||||
clauses: list[str] = []
|
||||
parameters: list[Any] = []
|
||||
if start_date:
|
||||
clauses.append("trade_date >= ?")
|
||||
parameters.append(start_date)
|
||||
if end_date:
|
||||
clauses.append("trade_date <= ?")
|
||||
parameters.append(end_date)
|
||||
where = f"WHERE {' AND '.join(clauses)}" if clauses else ""
|
||||
with self.connect() as connection:
|
||||
rows = connection.execute(
|
||||
f"""
|
||||
SELECT trade_date FROM dashboard_snapshots
|
||||
{where}
|
||||
ORDER BY trade_date
|
||||
""",
|
||||
parameters,
|
||||
).fetchall()
|
||||
return [str(row["trade_date"]) for row in rows]
|
||||
|
||||
def start_sync(self, trade_date: str, source: str) -> int:
|
||||
started_at = datetime.now().astimezone().isoformat(timespec="seconds")
|
||||
with self.connect() as connection:
|
||||
|
||||
@@ -3,9 +3,11 @@ from __future__ import annotations
|
||||
import copy
|
||||
import re
|
||||
from datetime import date, datetime, time as dt_time, timedelta
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
|
||||
from backend.bootstrap.config import (
|
||||
DATA_DIR,
|
||||
normalize_date,
|
||||
tushare_code,
|
||||
validate_stock_code,
|
||||
@@ -13,6 +15,17 @@ from backend.bootstrap.config import (
|
||||
)
|
||||
from backend.data.providers.ifind_client import IfindError
|
||||
from backend.data.providers.tushare_client import TushareClient, TushareError
|
||||
from backend.features.market.backfill_history import (
|
||||
DEFAULT_RECENT_TRADING_DAYS,
|
||||
MAX_RANGE_TRADING_DAYS,
|
||||
build_backfill_audit,
|
||||
calendar_window_start,
|
||||
classify_snapshot_coverage,
|
||||
create_sqlite_backup,
|
||||
display_date,
|
||||
select_open_trade_dates,
|
||||
select_open_trade_dates_in_range,
|
||||
)
|
||||
from backend.features.market.charts import ChartDataError
|
||||
from backend.features.market.insights import MarketInsightsService
|
||||
from backend.features.sentiment.engine import SENTIMENT_ENGINE_VERSION
|
||||
@@ -890,31 +903,226 @@ class MarketServiceMixin:
|
||||
"intraday": intraday_points,
|
||||
}
|
||||
|
||||
def backfill(self, start_date: str, end_date: str) -> list[dict[str, Any]]:
|
||||
start = datetime.strptime(normalize_date(start_date), "%Y%m%d").date()
|
||||
end = datetime.strptime(normalize_date(end_date), "%Y%m%d").date()
|
||||
if start > end:
|
||||
raise ValueError("开始日期不能晚于结束日期。")
|
||||
weekdays = []
|
||||
current = start
|
||||
while current <= end:
|
||||
if current.weekday() < 5:
|
||||
weekdays.append(current)
|
||||
current += timedelta(days=1)
|
||||
if len(weekdays) > 15:
|
||||
raise ValueError("单次最多回补 15 个工作日。")
|
||||
results = []
|
||||
for day in weekdays:
|
||||
dashboard = self.sync_dashboard(day.strftime("%Y%m%d"))
|
||||
def backfill(
|
||||
self,
|
||||
start_date: str = "",
|
||||
end_date: str = "",
|
||||
*,
|
||||
lookback: int | None = None,
|
||||
dry_run: bool = False,
|
||||
force: bool = False,
|
||||
create_backup: bool = True,
|
||||
) -> dict[str, Any]:
|
||||
"""Backfill dashboard snapshots for real trading days only.
|
||||
|
||||
- Date-range mode keeps the admin UI contract (max 15 open sessions).
|
||||
- Recent mode fills the last N open sessions (default/max 60).
|
||||
Weekends and holidays are reported as skipped non-trading days, not errors.
|
||||
"""
|
||||
if not self.configured:
|
||||
raise ValueError("公共行情尚未配置,无法回补历史快照。")
|
||||
normalized_end = normalize_date(end_date or date.today().isoformat())
|
||||
if lookback is not None or not (start_date and end_date):
|
||||
target_lookback = (
|
||||
DEFAULT_RECENT_TRADING_DAYS if lookback is None else int(lookback)
|
||||
)
|
||||
return self.backfill_recent_trading_days(
|
||||
end_date=normalized_end,
|
||||
lookback=target_lookback,
|
||||
dry_run=dry_run,
|
||||
force=force,
|
||||
create_backup=create_backup,
|
||||
)
|
||||
return self._backfill_date_range(
|
||||
start_date=normalize_date(start_date),
|
||||
end_date=normalized_end,
|
||||
dry_run=dry_run,
|
||||
force=force,
|
||||
create_backup=create_backup,
|
||||
)
|
||||
|
||||
def backfill_recent_trading_days(
|
||||
self,
|
||||
end_date: str = "",
|
||||
lookback: int = DEFAULT_RECENT_TRADING_DAYS,
|
||||
*,
|
||||
dry_run: bool = False,
|
||||
force: bool = False,
|
||||
create_backup: bool = True,
|
||||
) -> dict[str, Any]:
|
||||
normalized_end = normalize_date(end_date or date.today().isoformat())
|
||||
trade_dates = self._load_recent_open_trade_dates(normalized_end, lookback)
|
||||
existing = self.database.list_snapshot_trade_dates(
|
||||
trade_dates[0], trade_dates[-1]
|
||||
)
|
||||
coverage = classify_snapshot_coverage(trade_dates, existing)
|
||||
return self._execute_snapshot_backfill(
|
||||
mode="recent",
|
||||
end_date=normalized_end,
|
||||
lookback=lookback,
|
||||
coverage=coverage,
|
||||
skipped_non_trading_days=[],
|
||||
dry_run=dry_run,
|
||||
force=force,
|
||||
create_backup=create_backup,
|
||||
)
|
||||
|
||||
def _backfill_date_range(
|
||||
self,
|
||||
start_date: str,
|
||||
end_date: str,
|
||||
*,
|
||||
dry_run: bool = False,
|
||||
force: bool = False,
|
||||
create_backup: bool = True,
|
||||
) -> dict[str, Any]:
|
||||
window_start = calendar_window_start(end_date, MAX_RANGE_TRADING_DAYS)
|
||||
calendar_rows = self._tushare_client().query(
|
||||
"trade_cal",
|
||||
{
|
||||
"exchange": "SSE",
|
||||
"start_date": min(window_start, start_date),
|
||||
"end_date": end_date,
|
||||
},
|
||||
"cal_date,is_open,pretrade_date",
|
||||
)
|
||||
trade_dates, skipped = select_open_trade_dates_in_range(
|
||||
calendar_rows,
|
||||
start_date,
|
||||
end_date,
|
||||
maximum=MAX_RANGE_TRADING_DAYS,
|
||||
)
|
||||
if not trade_dates:
|
||||
raise ValueError("选定区间内没有交易日,周末或节假日无需回补。")
|
||||
existing = self.database.list_snapshot_trade_dates(trade_dates[0], trade_dates[-1])
|
||||
coverage = classify_snapshot_coverage(trade_dates, existing)
|
||||
return self._execute_snapshot_backfill(
|
||||
mode="range",
|
||||
end_date=end_date,
|
||||
lookback=None,
|
||||
coverage=coverage,
|
||||
skipped_non_trading_days=skipped,
|
||||
dry_run=dry_run,
|
||||
force=force,
|
||||
create_backup=create_backup,
|
||||
)
|
||||
|
||||
def _load_recent_open_trade_dates(self, end_date: str, lookback: int) -> list[str]:
|
||||
start_date = calendar_window_start(end_date, lookback)
|
||||
calendar_rows = self._tushare_client().query(
|
||||
"trade_cal",
|
||||
{
|
||||
"exchange": "SSE",
|
||||
"start_date": start_date,
|
||||
"end_date": end_date,
|
||||
},
|
||||
"cal_date,is_open,pretrade_date",
|
||||
)
|
||||
return select_open_trade_dates(calendar_rows, end_date, lookback)
|
||||
|
||||
def _execute_snapshot_backfill(
|
||||
self,
|
||||
*,
|
||||
mode: str,
|
||||
end_date: str,
|
||||
lookback: int | None,
|
||||
coverage: dict[str, Any],
|
||||
skipped_non_trading_days: list[str],
|
||||
dry_run: bool,
|
||||
force: bool,
|
||||
create_backup: bool,
|
||||
) -> dict[str, Any]:
|
||||
targets = list(coverage["trade_dates"] if force else coverage["missing"])
|
||||
backup_path: str | None = None
|
||||
if create_backup and not dry_run and targets:
|
||||
backup = create_sqlite_backup(
|
||||
Path(self.database.path),
|
||||
DATA_DIR / "backups",
|
||||
label=f"pre-{mode}-backfill",
|
||||
)
|
||||
backup_path = str(backup)
|
||||
|
||||
results: list[dict[str, Any]] = []
|
||||
if dry_run:
|
||||
for trade_date in coverage["trade_dates"]:
|
||||
exists = trade_date in coverage["present"]
|
||||
if exists and not force:
|
||||
status = "skipped"
|
||||
action = "exists"
|
||||
else:
|
||||
status = "planned"
|
||||
action = "refresh" if exists else "create"
|
||||
results.append(
|
||||
{
|
||||
"requested_date": display_date(trade_date),
|
||||
"trade_date": display_date(trade_date),
|
||||
"status": status,
|
||||
"action": action,
|
||||
}
|
||||
)
|
||||
return build_backfill_audit(
|
||||
mode=mode,
|
||||
end_date=end_date,
|
||||
lookback=lookback,
|
||||
coverage=coverage,
|
||||
skipped_non_trading_days=skipped_non_trading_days,
|
||||
backup_path=backup_path,
|
||||
dry_run=True,
|
||||
results=results,
|
||||
)
|
||||
|
||||
present_before = set(coverage["present"])
|
||||
for trade_date in targets:
|
||||
existed = trade_date in present_before
|
||||
try:
|
||||
dashboard = self.sync_dashboard(trade_date)
|
||||
actual = normalize_date(
|
||||
str(dashboard.get("meta", {}).get("trade_date") or trade_date)
|
||||
)
|
||||
results.append(
|
||||
{
|
||||
"requested_date": display_date(trade_date),
|
||||
"trade_date": display_date(actual),
|
||||
"status": "success",
|
||||
"action": "refreshed" if existed else "created",
|
||||
"source": dashboard.get("meta", {}).get("source"),
|
||||
"records": self._record_count(dashboard),
|
||||
}
|
||||
)
|
||||
except Exception as exc:
|
||||
results.append(
|
||||
{
|
||||
"requested_date": display_date(trade_date),
|
||||
"trade_date": display_date(trade_date),
|
||||
"status": "failed",
|
||||
"action": "refresh" if existed else "create",
|
||||
"error": str(exc),
|
||||
}
|
||||
)
|
||||
|
||||
for trade_date in coverage["present"]:
|
||||
if force:
|
||||
continue
|
||||
results.append(
|
||||
{
|
||||
"requested_date": day.isoformat(),
|
||||
"trade_date": dashboard["meta"]["trade_date"],
|
||||
"source": dashboard["meta"]["source"],
|
||||
"records": self._record_count(dashboard),
|
||||
"requested_date": display_date(trade_date),
|
||||
"trade_date": display_date(trade_date),
|
||||
"status": "skipped",
|
||||
"action": "exists",
|
||||
}
|
||||
)
|
||||
return results
|
||||
|
||||
results.sort(key=lambda row: str(row.get("requested_date") or ""))
|
||||
return build_backfill_audit(
|
||||
mode=mode,
|
||||
end_date=end_date,
|
||||
lookback=lookback,
|
||||
coverage=coverage,
|
||||
skipped_non_trading_days=skipped_non_trading_days,
|
||||
backup_path=backup_path,
|
||||
dry_run=False,
|
||||
results=results,
|
||||
)
|
||||
|
||||
def _stock_identity(self, code: str, trade_date: str) -> tuple[str, str]:
|
||||
snapshot = self.database.get_snapshot(trade_date) or {}
|
||||
|
||||
@@ -29,11 +29,17 @@ class SystemRoutesMixin:
|
||||
def backfill_data(self) -> None:
|
||||
try:
|
||||
body = self.read_json_body()
|
||||
results = self.application_service.backfill(
|
||||
lookback_raw = body.get("lookback")
|
||||
lookback = int(lookback_raw) if lookback_raw not in (None, "") else None
|
||||
audit = self.application_service.backfill(
|
||||
str(body.get("start_date") or ""),
|
||||
str(body.get("end_date") or ""),
|
||||
lookback=lookback,
|
||||
dry_run=bool(body.get("dry_run")),
|
||||
force=bool(body.get("force")),
|
||||
create_backup=body.get("create_backup", True) is not False,
|
||||
)
|
||||
self.send_json({"ok": True, "results": results})
|
||||
self.send_json({"ok": True, **audit, "results": audit.get("results") or []})
|
||||
except ValueError as exc:
|
||||
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
|
||||
except Exception as exc:
|
||||
|
||||
Reference in New Issue
Block a user