"""Construye el dataset por equipo-partido (features RICAS) para los modelos de bloque defensivo y pasillos, procesando los preprocessed liga por liga EN STREAMING (baja → extrae 1 fila por (matchId,teamId) → borra el CSV). Resumable. Por equipo-partido extrae (todo desde los eventos; NO carga la columna qualifiers — centros/largos se derivan por geometría/distancia): - Resultado: goles a favor/en contra (para ELO y splits local/visita). - Estilo: % fast break / build up / juego directo / set piece / recovery / progressive; % ataques que terminan en centro; % pases cortos/largos/exitosos; pases progresivos; nº de ataques (llegan a último tercio); xT y xG por partido. - Formación más usada (id_formation) → one-hot en el entrenamiento. - Posición/presión: avg x/y, % campo rival, altura media de recuperación. - Distribución de BLOQUE DEFENSIVO enfrentado (alto/medio/bajo). - Distribución de PASILLOS de ataque (Izq/Centro/Der). - Peligro (xT) por (bloque × pasillo) en ATAQUE y en DEFENSA (concedido). El ELO, el rolling "hasta la fecha" y los splits local/visita se arman en el entrenamiento (necesitan orden cronológico entre partidos). Uso: python scripts/build_model_dataset.py --leagues "Spanish La Liga" ... python scripts/build_model_dataset.py --all """ from __future__ import annotations import argparse import sys import tempfile from pathlib import Path sys.path.insert(0, str(Path(__file__).resolve().parents[1] / "src")) import numpy as np import pandas as pd from racing_reports.datastore import DataStore OUT = Path(__file__).resolve().parents[1] / "vendor" / "data" / "modeling" / "matchteam_dataset.parquet" BLOCKS = {"Build Up against High Block": "H", "Build Up against Medium Block": "M", "Build Up against Low Block": "L"} SHOTS = {"Goal", "SavedShot", "MissedShots", "ShotOnPost"} LANES = ["Der", "Centro", "Izq"] # y<33.3 Der, <66.7 Centro, resto Izq (convención multitag) USECOLS = ["matchId", "teamId", "TeamName", "Competencia", "Temporada", "fecha", "home_team_id", "period_id", "sequenceId", "phaseLabel", "x", "y", "endX", "endY", "event_name", "xT", "xThreat", "xG_corr", "xG", "id_formation", "distanceTravelledByBall", "lastLineBroken", "isGoal", "outcome_type"] def _lane(y): return pd.cut(y, [-1, 33.333, 66.667, 1e9], labels=LANES) def _extract(df: pd.DataFrame, league: str, season: str) -> pd.DataFrame: df = df[[c for c in USECOLS if c in df.columns]].copy() if not {"matchId", "teamId", "x", "y", "event_name", "sequenceId", "phaseLabel"} <= set(df.columns): return pd.DataFrame() for c in ("x", "y", "endX", "endY", "period_id", "sequenceId", "distanceTravelledByBall"): if c in df.columns: df[c] = pd.to_numeric(df[c], errors="coerce") for c in ("TeamName", "fecha", "home_team_id", "id_formation", "lastLineBroken"): if c not in df.columns: df[c] = pd.NA df["matchId"] = df["matchId"].astype(str) df["teamId"] = df["teamId"].astype(str) df["xt"] = pd.to_numeric(df.get("xT"), errors="coerce") if df["xt"].isna().all() and "xThreat" in df.columns: df["xt"] = pd.to_numeric(df["xThreat"], errors="coerce") df["xt"] = df["xt"].fillna(0.0).clip(lower=0) df["xg"] = pd.to_numeric(df.get("xG_corr", df.get("xG")), errors="coerce").fillna(0.0) df["is_pass"] = df["event_name"].eq("Pass") df["is_shot"] = df["event_name"].isin(SHOTS) df["is_goal"] = df["event_name"].eq("Goal") df["blk"] = df["phaseLabel"].map(BLOCKS) df["lane"] = _lane(df["y"]) # centro (geom): pase desde zona ancha terminando en área central ex, ey = df.get("endX"), df.get("endY") df["is_cross"] = (df["is_pass"] & ex.notna() & ey.notna() & (ex >= 83) & ey.between(21, 79) & ((df["y"] < 21) | (df["y"] > 79) | (df["x"] >= 83))) # pase largo / corto por distancia recorrida del balón (fallback a |Δ|) dist = pd.to_numeric(df.get("distanceTravelledByBall"), errors="coerce") if dist.isna().all(): dist = ((ex - df["x"]) ** 2 + (ey - df["y"]) ** 2) ** 0.5 if ex is not None else pd.Series(np.nan, index=df.index) df["pass_long"] = df["is_pass"] & (dist >= 30) df["pass_short"] = df["is_pass"] & (dist < 15) df["pass_ok"] = df["is_pass"] & df.get("outcome_type", pd.Series("", index=df.index)).astype(str).str.contains("uccess", na=False) df["progressive"] = df["is_pass"] & df.get("lastLineBroken", pd.Series("", index=df.index)).astype(str).str.strip().isin(["last", "True", "true"]) df["opp_half"] = df["x"] > 50 # equipos / rival / local teams = df.dropna(subset=["teamId"]).groupby("matchId")["teamId"].agg(lambda s: list(dict.fromkeys(s))) opp = {} for mid, ts in teams.items(): if len(ts) == 2: opp[(mid, ts[0])] = ts[1]; opp[(mid, ts[1])] = ts[0] name = df.dropna(subset=["teamId"]).drop_duplicates("teamId").set_index("teamId")["TeamName"].to_dict() # ── secuencias (1 fila por matchId,period,sequenceId) ── seq = df.groupby(["matchId", "period_id", "sequenceId"], sort=False).agg( poss=("teamId", "first"), phase=("phaseLabel", "first"), blk=("blk", "first"), npass=("is_pass", "sum"), maxx=("x", "max"), has_cross=("is_cross", "any"), has_long=("pass_long", "any"), xt=("xt", "sum"), xg=("xg", "sum")).reset_index() seq["is_buildup"] = seq["phase"].astype(str).str.startswith("Build Up against") seq["is_counter"] = seq["phase"].eq("Counter Attack") seq["is_setpiece"] = seq["phase"].eq("Set Piece") seq["is_recovery"] = seq["phase"].eq("Recovery") seq["is_progr"] = seq["phase"].eq("Progressive Play") seq["is_direct"] = (seq["npass"] <= 4) & seq["has_long"] & ~seq["is_setpiece"] seq["reached_ft"] = seq["maxx"] >= 66 def _team_seq_feats(poss_team_col): g = seq.groupby(["matchId", poss_team_col]) f = g.agg(n_seq=("phase", "size"), sh_counter=("is_counter", "mean"), sh_buildup=("is_buildup", "mean"), sh_direct=("is_direct", "mean"), sh_setpiece=("is_setpiece", "mean"), sh_recovery=("is_recovery", "mean"), sh_progr=("is_progr", "mean"), cross_pct=("has_cross", "mean"), n_attacks=("reached_ft", "sum"), xt_total=("xt", "sum"), xg_total=("xg", "sum")).reset_index() f["xt_per_attack"] = f["xt_total"] / f["n_attacks"].clip(lower=1) return f.rename(columns={poss_team_col: "teamId"}) atk_feats = _team_seq_feats("poss") # ── bloque defensivo enfrentado (el rival construye → este equipo defiende) ── faced = (seq[seq["blk"].notna()].groupby(["matchId", "poss", "blk"]).size() .unstack("blk", fill_value=0).reset_index()) for b in ("H", "M", "L"): if b not in faced.columns: faced[b] = 0 faced["def_team"] = [opp.get((m, t)) for m, t in zip(faced["matchId"], faced["poss"])] defb = faced.dropna(subset=["def_team"]).rename(columns={"H": "def_H", "M": "def_M", "L": "def_L", "def_team": "teamId"})[["matchId", "teamId", "def_H", "def_M", "def_L"]] # ── pasillos de ataque (pases del equipo en campo rival) ── atk = df[df["is_pass"] & (df["x"] > 50) & df["y"].notna()] lanes = atk.groupby(["matchId", "teamId", "lane"], observed=True).size().unstack("lane", fill_value=0).reset_index() for ln in LANES: if ln not in lanes.columns: lanes[ln] = 0 lanes = lanes.rename(columns={ln: f"atk_{ln}" for ln in LANES}) # ── xT por (bloque × pasillo): ataque y defensa concedida ── # Mergeamos seq SOLO sobre los eventos con xT>0 (una fracción) para no inflar memoria. ev = df.loc[(df["xt"] > 0) & df["lane"].notna(), ["matchId", "period_id", "sequenceId", "lane", "xt"]].merge( seq[["matchId", "period_id", "sequenceId", "poss", "blk"]].rename(columns={"blk": "seq_blk"}), on=["matchId", "period_id", "sequenceId"], how="left") ev = ev[ev["seq_blk"].notna()].copy() # ataque: eventos del equipo en posesión, por (bloque, pasillo) off = (ev.groupby(["matchId", "poss", "seq_blk", "lane"], observed=True)["xt"].sum() .unstack(["seq_blk", "lane"], fill_value=0.0)) off.columns = [f"offxt_{b}_{l}" for b, l in off.columns] off = off.reset_index().rename(columns={"poss": "teamId"}) # defensa: el rival ataca → este equipo (def) concede; blk = el que impuso el defensor ev["def_team"] = [opp.get((m, t)) for m, t in zip(ev["matchId"], ev["poss"])] deff = (ev.dropna(subset=["def_team"]).groupby(["matchId", "def_team", "seq_blk", "lane"], observed=True)["xt"].sum() .unstack(["seq_blk", "lane"], fill_value=0.0)) deff.columns = [f"defxt_{b}_{l}" for b, l in deff.columns] deff = deff.reset_index().rename(columns={"def_team": "teamId"}) # ── resultado, formación, posición/presión ── df["formk"] = df["id_formation"].astype(str) base = df.groupby(["matchId", "teamId"]).agg( TeamName=("TeamName", "first"), fecha=("fecha", "first"), home_team_id=("home_team_id", "first"), n_events=("event_name", "size"), n_pass=("is_pass", "sum"), n_shots=("is_shot", "sum"), goals_for=("is_goal", "sum"), avg_x=("x", "mean"), avg_y=("y", "mean"), share_opp_half=("opp_half", "mean"), n_long=("pass_long", "sum"), n_short=("pass_short", "sum"), n_passok=("pass_ok", "sum"), n_progr=("progressive", "sum"), formation=("formk", lambda s: s.mode().iloc[0] if len(s.mode()) else "0")).reset_index() base["pct_long"] = base["n_long"] / base["n_pass"].clip(lower=1) base["pct_short"] = base["n_short"] / base["n_pass"].clip(lower=1) base["pct_pass_ok"] = base["n_passok"] / base["n_pass"].clip(lower=1) # altura media de recuperación (eventos defensivos) rec = df[df["event_name"].isin(["BallRecovery", "Interception", "Tackle"])].groupby(["matchId", "teamId"])["x"].mean().rename("recovery_height").reset_index() base["Competencia"] = league; base["Temporada"] = season out = base for t in (atk_feats, defb, lanes, off, deff, rec): out = out.merge(t, on=["matchId", "teamId"], how="left") out["rival_teamId"] = [opp.get((m, t)) for m, t in zip(out["matchId"], out["teamId"])] out["rival_name"] = out["rival_teamId"].map(name) out["is_home"] = out["teamId"].astype(str) == out["home_team_id"].astype(str) # goles en contra = goles del rival gf = out.set_index(["matchId", "teamId"])["goals_for"] out["goals_against"] = [gf.get((m, r), np.nan) if r is not None else np.nan for m, r in zip(out["matchId"], out["rival_teamId"])] fill0 = [c for c in out.columns if c.startswith(("def_", "atk_", "offxt_", "defxt_"))] out[fill0] = out[fill0].fillna(0) return out def _processed_keys() -> set: if not OUT.exists(): return set() d = pd.read_parquet(OUT, columns=["Competencia", "Temporada"]) return set(zip(d["Competencia"].astype(str), d["Temporada"].astype(str))) def _league_season_files(ds: DataStore): fs = ds._filesystem_client() root = ds.settings.azure_preprocessed_root.strip("/") out = [] for p in fs.get_paths(path=root, recursive=True): if getattr(p, "is_directory", False): continue name = getattr(p, "name", "") or "" fn = name.split("/")[-1] if fn.startswith("preprocessed_") and fn.endswith(".csv") and "etiquetado" not in fn: parts = name[len(root):].strip("/").split("/") if len(parts) >= 3: out.append((parts[0], parts[1], fn, getattr(p, "content_length", 0) or 0)) return out def main() -> None: ap = argparse.ArgumentParser() ap.add_argument("--leagues", nargs="*", default=None) ap.add_argument("--smallest", type=int, default=0) ap.add_argument("--all", action="store_true") ap.add_argument("--force", action="store_true") args = ap.parse_args() ds = DataStore() files = _league_season_files(ds) if args.leagues: files = [f for f in files if f[0] in set(args.leagues)] elif args.smallest: files = sorted(files, key=lambda f: f[3])[:args.smallest] elif not args.all: ap.error("Pasá --leagues, --smallest N o --all") done = set() if args.force else _processed_keys() OUT.parent.mkdir(parents=True, exist_ok=True) for lg, se, fn, sz in files: if (lg, se) in done: print(f"skip {lg}/{se}", flush=True); continue print(f">>> {lg}/{se} ({sz/1e6:.0f} MB) — bajando…", flush=True) fc = ds._filesystem_client().get_file_client(f"{ds.settings.azure_preprocessed_root}/{lg}/{se}/{fn}") with tempfile.TemporaryDirectory() as td: csv = Path(td) / fn with open(csv, "wb") as fh: fc.download_file(max_concurrency=8).readinto(fh) # descarga en paralelo df = pd.read_csv(csv, usecols=lambda c: c in USECOLS, low_memory=False) rows = _extract(df, lg, se) if rows.empty: print(f" {lg}/{se}: schema insuficiente, salteada", flush=True); continue existing = [pd.read_parquet(OUT)] if OUT.exists() else [] pd.concat(existing + [rows], ignore_index=True).to_parquet(OUT, index=False) print(f" {len(rows)} filas | {rows.shape[1]} cols | total {len(pd.read_parquet(OUT))}", flush=True) print("OK") if __name__ == "__main__": main()