from __future__ import annotations import argparse import csv import json from collections import defaultdict from datetime import datetime from pathlib import Path from typing import Any PACKAGE_ROOT = Path(__file__).resolve().parents[1] MINUTE_ROOT = Path("E:/quant/2023_front_m") TARGET_DATES_CSV = PACKAGE_ROOT / "market_breadth_recalc.csv" OUT_CSV = PACKAGE_ROOT / "intraday_breadth_minute_pool.csv" SUMMARY_JSON = PACKAGE_ROOT / "intraday_breadth_minute_pool_summary.json" SUMMARY_MD = PACKAGE_ROOT / "intraday_breadth_minute_pool_summary.md" PROGRESS_JSON = PACKAGE_ROOT / "intraday_breadth_minute_pool_progress.json" TARGET_TIMES = [ "09:35:00", "09:40:00", "09:45:00", "09:50:00", "09:55:00", "10:00:00", "10:05:00", "10:10:00", "10:15:00", "10:20:00", "10:25:00", "10:30:00", "10:35:00", "10:40:00", ] def fnum(value: Any) -> float | None: try: text = str(value).strip() if not text: return None return float(text) except Exception: return None def date_key(value: str) -> str: return str(value).replace("-", "")[:8] def load_target_dates(limit: int) -> set[str]: rows = list(csv.DictReader(TARGET_DATES_CSV.open("r", newline="", encoding="utf-8-sig"))) dates = [date_key(row["signal_trade_date"]) for row in rows] if limit > 0: dates = dates[:limit] return set(dates) def minute_files(max_files: int) -> list[Path]: files: list[Path] = [] for market in ["SH", "SZ", "BJ"]: folder = MINUTE_ROOT / market if not folder.exists(): continue for path in sorted(folder.glob("price_*.csv")): if path.stat().st_size > 0: files.append(path) if max_files > 0 and len(files) >= max_files: return files return files def init_stats() -> dict[str, int]: return { "eligible_count": 0, "up_vs_prev_close_count": 0, "flat_vs_prev_close_count": 0, "down_vs_prev_close_count": 0, "up_vs_open_count": 0, "flat_vs_open_count": 0, "down_vs_open_count": 0, } def add_observation(agg: dict[tuple[str, str], dict[str, int]], d: str, t: str, close: float, prev_close: float, day_open: float) -> None: stats = agg[(d, t)] stats["eligible_count"] += 1 if close > prev_close: stats["up_vs_prev_close_count"] += 1 elif close < prev_close: stats["down_vs_prev_close_count"] += 1 else: stats["flat_vs_prev_close_count"] += 1 if close > day_open: stats["up_vs_open_count"] += 1 elif close < day_open: stats["down_vs_open_count"] += 1 else: stats["flat_vs_open_count"] += 1 def process_file(path: Path, target_dates: set[str], agg: dict[tuple[str, str], dict[str, int]]) -> tuple[int, int]: target_time_set = set(TARGET_TIMES) processed_target_days = 0 observations = 0 prev_date_last_close: float | None = None current_date = "" current_open: float | None = None current_time_close: dict[str, float] = {} def finish_day() -> None: nonlocal processed_target_days, observations if current_date in target_dates and prev_date_last_close is not None and current_open is not None: processed_target_days += 1 for t, close in current_time_close.items(): add_observation(agg, current_date, t, close, prev_date_last_close, current_open) observations += 1 with path.open("r", newline="", encoding="utf-8-sig") as f: reader = csv.DictReader(f) for row in reader: tag = row.get("timetag", "") if len(tag) < 17: continue d = tag[:8] t = tag[9:17] close = fnum(row.get("close")) if close is None: continue if d != current_date: if current_date: finish_day() prev_date_last_close = last_close current_date = d current_open = fnum(row.get("open")) current_time_close = {} last_close = close else: last_close = close if d in target_dates and t in target_time_set: current_time_close[t] = close if current_date: finish_day() return processed_target_days, observations def write_progress(payload: dict[str, Any]) -> None: PROGRESS_JSON.write_text(json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8") def write_manifest() -> None: rows = [] for path in sorted(PACKAGE_ROOT.rglob("*")): if path.is_file(): rows.append({"path": str(path.relative_to(PACKAGE_ROOT)).replace("\\", "/"), "bytes": path.stat().st_size}) with (PACKAGE_ROOT / "manifest.csv").open("w", newline="", encoding="utf-8-sig") as f: writer = csv.DictWriter(f, fieldnames=["path", "bytes"]) writer.writeheader() writer.writerows(rows) def main() -> None: parser = argparse.ArgumentParser() parser.add_argument("--date-limit", type=int, default=0) parser.add_argument("--max-files", type=int, default=0) parser.add_argument("--progress-every", type=int, default=100) args = parser.parse_args() target_dates = load_target_dates(args.date_limit) files = minute_files(args.max_files) agg: dict[tuple[str, str], dict[str, int]] = defaultdict(init_stats) total_target_days = 0 total_observations = 0 for idx, path in enumerate(files, start=1): try: days, observations = process_file(path, target_dates, agg) total_target_days += days total_observations += observations except Exception as exc: write_progress( { "status": "ERROR_BUT_CONTINUING", "file_index": idx, "file_count": len(files), "path": str(path), "error": str(exc), "generated_at": datetime.now().isoformat(timespec="seconds"), } ) if idx % args.progress_every == 0: write_progress( { "status": "RUNNING", "file_index": idx, "file_count": len(files), "target_dates": len(target_dates), "total_target_days_seen": total_target_days, "total_observations": total_observations, "generated_at": datetime.now().isoformat(timespec="seconds"), } ) rows = [] for d in sorted(target_dates): for t in TARGET_TIMES: stats = agg[(d, t)] rows.append( { "trade_date": f"{d[:4]}-{d[4:6]}-{d[6:8]}", "time": t, "universe_name": "local_nonempty_minute_file_pool", **stats, "price_basis": "minute_previous_trading_day_last_close", "notes": "Coverage pool only; not full A-share market breadth.", } ) with OUT_CSV.open("w", newline="", encoding="utf-8-sig") as f: writer = csv.DictWriter(f, fieldnames=list(rows[0].keys()) if rows else []) writer.writeheader() writer.writerows(rows) summary = { "generated_at": datetime.now().isoformat(timespec="seconds"), "target_dates": len(target_dates), "target_times": TARGET_TIMES, "minute_files_scanned": len(files), "total_target_days_seen": total_target_days, "total_observations": total_observations, "output_rows": len(rows), "boundaries": [ "This is breadth for the local non-empty minute-file coverage pool, not full A-share market breadth.", "Main up/down basis is each symbol's previous trading day's last minute close in the same minute file.", "The result carries eligible_count per timestamp because minute coverage is incomplete and uneven across markets.", ], } SUMMARY_JSON.write_text(json.dumps(summary, ensure_ascii=False, indent=2), encoding="utf-8") SUMMARY_MD.write_text( "\n".join( [ "# Intraday Breadth Minute Pool Summary", "", f"- generated_at: {summary['generated_at']}", f"- target_dates: {summary['target_dates']}", f"- minute_files_scanned: {summary['minute_files_scanned']}", f"- total_target_days_seen: {summary['total_target_days_seen']}", f"- total_observations: {summary['total_observations']}", f"- output_rows: {summary['output_rows']}", "", "## Boundaries", *[f"- {item}" for item in summary["boundaries"]], "", ] ), encoding="utf-8", ) write_progress({"status": "FINISHED", **summary}) write_manifest() print(json.dumps(summary, ensure_ascii=False, indent=2)) if __name__ == "__main__": main()