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()
|