From 40559a35d79fabcf95f9e26e5a04d1f1d531e3ab Mon Sep 17 00:00:00 2001 From: elmartinj Date: Wed, 8 Apr 2026 17:22:26 -0600 Subject: [PATCH 1/9] Add initial CENACE pipeline integration --- src/cenace_pipeline.py | 69 +++++++++++++++ src/data/cenace/__init__.py | 0 src/data/cenace/aggregate/core.py | 48 +++++++++++ src/data/cenace/config.py | 15 ++++ src/data/cenace/transform/core.py | 0 src/data/cenace/utils/cenace_data.py | 88 +++++++++++++++++++ src/dataset_registry.py | 25 ++++++ src/evaluation/cenace/__init__.py | 0 src/evaluation/cenace/core.py | 92 ++++++++++++++++++++ src/evaluation/cenace/metrics.py | 47 +++++++++++ src/forecast/cenace/__init__.py | 0 src/forecast/cenace/core.py | 93 ++++++++++++++++++++ src/forecast/cenace/models.py | 122 +++++++++++++++++++++++++++ src/pipeline.py | 32 +++++++ 14 files changed, 631 insertions(+) create mode 100644 src/cenace_pipeline.py create mode 100644 src/data/cenace/__init__.py create mode 100644 src/data/cenace/aggregate/core.py create mode 100644 src/data/cenace/config.py create mode 100644 src/data/cenace/transform/core.py create mode 100644 src/data/cenace/utils/cenace_data.py create mode 100644 src/dataset_registry.py create mode 100644 src/evaluation/cenace/__init__.py create mode 100644 src/evaluation/cenace/core.py create mode 100644 src/evaluation/cenace/metrics.py create mode 100644 src/forecast/cenace/__init__.py create mode 100644 src/forecast/cenace/core.py create mode 100644 src/forecast/cenace/models.py create mode 100644 src/pipeline.py diff --git a/src/cenace_pipeline.py b/src/cenace_pipeline.py new file mode 100644 index 0000000..e437e7c --- /dev/null +++ b/src/cenace_pipeline.py @@ -0,0 +1,69 @@ +from __future__ import annotations + +import argparse + +import pandas as pd + +from src.data.cenace.aggregate.core import build_hourly_partitions +from src.evaluation.cenace.core import run_evaluation +from src.forecast.cenace.core import run_forecast + + +def run_cenace_pipeline( + cutoff: str, + model: str, + h: int = 24, + max_window_size: int = 48, + skip_aggregate: bool = False, +) -> tuple[str, str]: + try: + cutoff_ts = pd.Timestamp(cutoff) + except Exception as exc: + raise ValueError(f"Invalid cutoff timestamp: {cutoff}") from exc + + if not skip_aggregate: + n_written = build_hourly_partitions() + print(f"Aggregated {n_written} partitions") + + forecast_path = run_forecast( + cutoff=cutoff_ts, + model=model, + h=h, + max_window_size=max_window_size, + ) + print(f"Forecasts saved to: {forecast_path}") + + eval_path = run_evaluation( + cutoff=cutoff_ts, + model=model, + h=h, + max_window_size=max_window_size, + ) + print(f"Metrics saved to: {eval_path}") + + return str(forecast_path), str(eval_path) + + +def parse_args() -> argparse.Namespace: + parser = argparse.ArgumentParser() + parser.add_argument("--cutoff", required=True) + parser.add_argument("--model", required=True) + parser.add_argument("--h", type=int, default=24) + parser.add_argument("--max-window-size", type=int, default=48) + parser.add_argument("--skip-aggregate", action="store_true") + return parser.parse_args() + + +def main() -> None: + args = parse_args() + run_cenace_pipeline( + cutoff=args.cutoff, + model=args.model, + h=args.h, + max_window_size=args.max_window_size, + skip_aggregate=args.skip_aggregate, + ) + + +if __name__ == "__main__": + main() diff --git a/src/data/cenace/__init__.py b/src/data/cenace/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/src/data/cenace/aggregate/core.py b/src/data/cenace/aggregate/core.py new file mode 100644 index 0000000..99433f8 --- /dev/null +++ b/src/data/cenace/aggregate/core.py @@ -0,0 +1,48 @@ +from __future__ import annotations + +import pandas as pd + +from src.data.cenace.config import PROCESSED_CSV, PROCESSED_EVENTS_HOURLY_DIR + +INPUT_CSV = PROCESSED_CSV +OUTPUT_ROOT = PROCESSED_EVENTS_HOURLY_DIR + + +def build_hourly_partitions() -> int: + df = pd.read_csv(INPUT_CSV) + + df["ds"] = pd.to_datetime(df["ds"], errors="coerce") + df["y"] = pd.to_numeric(df["y"], errors="coerce") + + df = df.dropna(subset=["unique_id", "ds", "y"]).copy() + df = df.sort_values(["unique_id", "ds"]).drop_duplicates(["unique_id", "ds"]) + + df["year"] = df["ds"].dt.year + df["month"] = df["ds"].dt.month + df["day"] = df["ds"].dt.day + + OUTPUT_ROOT.mkdir(parents=True, exist_ok=True) + + n_written = 0 + for (year, month, day), part in df.groupby(["year", "month", "day"], sort=True): + part_dir = ( + OUTPUT_ROOT / f"year={year:04d}" / f"month={month:02d}" / f"day={day:02d}" + ) + part_dir.mkdir(parents=True, exist_ok=True) + + out_path = part_dir / "series.parquet" + part[["unique_id", "ds", "y"]].to_parquet(out_path, index=False) + + print(f"Saved: {out_path}") + n_written += 1 + + return n_written + + +def main() -> None: + n_written = build_hourly_partitions() + print(f"\nDone. Wrote {n_written} daily partitions.") + + +if __name__ == "__main__": + main() diff --git a/src/data/cenace/config.py b/src/data/cenace/config.py new file mode 100644 index 0000000..baf3893 --- /dev/null +++ b/src/data/cenace/config.py @@ -0,0 +1,15 @@ +from __future__ import annotations + +from pathlib import Path + +ROOT = Path(__file__).resolve().parents[3] + +DATA_ROOT = ROOT / "data" / "cenace" + +TMP_DIR = DATA_ROOT / "tmp" +PROCESSED_DIR = DATA_ROOT / "processed" +PROCESSED_CSV = PROCESSED_DIR / "cenace.csv" + +PROCESSED_EVENTS_HOURLY_DIR = DATA_ROOT / "processed-events" / "hourly" +FORECASTS_HOURLY_DIR = DATA_ROOT / "forecasts" / "hourly" +EVALUATIONS_HOURLY_DIR = DATA_ROOT / "evaluations" / "hourly" diff --git a/src/data/cenace/transform/core.py b/src/data/cenace/transform/core.py new file mode 100644 index 0000000..e69de29 diff --git a/src/data/cenace/utils/cenace_data.py b/src/data/cenace/utils/cenace_data.py new file mode 100644 index 0000000..5142720 --- /dev/null +++ b/src/data/cenace/utils/cenace_data.py @@ -0,0 +1,88 @@ +from __future__ import annotations + +from dataclasses import dataclass +from pathlib import Path + +import duckdb +import pandas as pd + + +@dataclass +class CENACEData: + base_path: Path + freq: str = "hourly" + h: int = 24 + max_window_size: int = 24 * 90 + + def __post_init__(self) -> None: + self.base_path = Path(self.base_path) + + def _date_to_partition(self, d: pd.Timestamp) -> Path: + return ( + self.base_path + / f"year={d.year:04d}" + / f"month={d.month:02d}" + / f"day={d.day:02d}" + / "series.parquet" + ) + + def _paths_for_range(self, start: pd.Timestamp, end: pd.Timestamp) -> list[str]: + days = pd.date_range(start.normalize(), end.normalize(), freq="D") + paths = [self._date_to_partition(d) for d in days] + existing = [str(p) for p in paths if p.exists()] + if not existing: + raise FileNotFoundError( + f"No parquet files found between {start} and \ + {end} under {self.base_path}" + ) + return existing + + def get_df( + self, + cutoff: str | pd.Timestamp, + max_window_size: int | None = None, + sort: bool = True, + ) -> pd.DataFrame: + cutoff = pd.Timestamp(cutoff) + window = max_window_size or self.max_window_size + start = cutoff - pd.Timedelta(hours=window - 1) + + paths = self._paths_for_range(start, cutoff) + + query = f""" + SELECT unique_id, ds, y + FROM read_parquet({paths}) + WHERE ds >= TIMESTAMP '{start}' + AND ds <= TIMESTAMP '{cutoff}' + """ + + df = duckdb.sql(query).df() + df["ds"] = pd.to_datetime(df["ds"]) + + if sort: + df = df.sort_values(["unique_id", "ds"]).reset_index(drop=True) + + return df + + def get_actuals( + self, cutoff: str | pd.Timestamp, h: int | None = None + ) -> pd.DataFrame: + cutoff = pd.Timestamp(cutoff) + horizon = h or self.h + + start = cutoff + pd.Timedelta(hours=1) + end = cutoff + pd.Timedelta(hours=horizon) + + paths = self._paths_for_range(start, end) + + query = f""" + SELECT unique_id, ds, y + FROM read_parquet({paths}) + WHERE ds >= TIMESTAMP '{start}' + AND ds <= TIMESTAMP '{end}' + """ + + df = duckdb.sql(query).df() + df["ds"] = pd.to_datetime(df["ds"]) + df = df.sort_values(["unique_id", "ds"]).reset_index(drop=True) + return df diff --git a/src/dataset_registry.py b/src/dataset_registry.py new file mode 100644 index 0000000..87e674c --- /dev/null +++ b/src/dataset_registry.py @@ -0,0 +1,25 @@ +from __future__ import annotations + +from collections.abc import Callable +from typing import Any + +from src.cenace_pipeline import run_cenace_pipeline + +DatasetRunner = Callable[..., tuple[str, str]] + +DATASET_REGISTRY: dict[str, DatasetRunner] = { + "cenace": run_cenace_pipeline, +} + +PIPELINE_DATASET_CHOICES = sorted(DATASET_REGISTRY) + + +def run_dataset_pipeline(dataset: str, **kwargs: Any) -> tuple[str, str]: + try: + runner = DATASET_REGISTRY[dataset] + except KeyError as exc: + raise ValueError( + f"Unsupported dataset: {dataset}. " f"Available: {PIPELINE_DATASET_CHOICES}" + ) from exc + + return runner(**kwargs) diff --git a/src/evaluation/cenace/__init__.py b/src/evaluation/cenace/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/src/evaluation/cenace/core.py b/src/evaluation/cenace/core.py new file mode 100644 index 0000000..b6207e9 --- /dev/null +++ b/src/evaluation/cenace/core.py @@ -0,0 +1,92 @@ +from __future__ import annotations + +import argparse +from pathlib import Path + +import pandas as pd + +from src.data.cenace.config import ( + EVALUATIONS_HOURLY_DIR, + FORECASTS_HOURLY_DIR, + PROCESSED_EVENTS_HOURLY_DIR, +) +from src.data.cenace.utils.cenace_data import CENACEData +from src.evaluation.cenace.metrics import evaluate_forecasts + + +def cutoff_partition(root: Path, cutoff: pd.Timestamp) -> Path: + return ( + root + / f"year={cutoff.year:04d}" + / f"month={cutoff.month:02d}" + / f"day={cutoff.day:02d}" + ) + + +def run_evaluation( + cutoff: str | pd.Timestamp, + model: str, + h: int = 24, + max_window_size: int = 48, +) -> Path: + cutoff = pd.Timestamp(cutoff) + + data = CENACEData( + base_path=PROCESSED_EVENTS_HOURLY_DIR, + freq="hourly", + h=h, + max_window_size=max_window_size, + ) + + forecast_path = ( + FORECASTS_HOURLY_DIR + / model + / f"year={cutoff.year:04d}" + / f"month={cutoff.month:02d}" + / f"day={cutoff.day:02d}" + / "forecasts.parquet" + ) + + actuals = data.get_actuals(cutoff, h=h) + forecasts = pd.read_parquet(forecast_path) + + merged = forecasts.merge(actuals, on=["unique_id", "ds"], how="inner") + if merged.empty: + raise ValueError("Merged forecasts/actuals is empty") + + metrics = evaluate_forecasts(merged) + + eval_root = EVALUATIONS_HOURLY_DIR / model + out_dir = cutoff_partition(eval_root, cutoff) + out_dir.mkdir(parents=True, exist_ok=True) + out_path = out_dir / "metrics.parquet" + + metrics.to_parquet(out_path, index=False) + return out_path + + +def parse_args() -> argparse.Namespace: + parser = argparse.ArgumentParser() + parser.add_argument("--cutoff", required=True) + parser.add_argument("--model", required=True) + parser.add_argument("--h", type=int, default=24) + parser.add_argument("--max-window-size", type=int, default=48) + return parser.parse_args() + + +def main() -> None: + args = parse_args() + out_path = run_evaluation( + cutoff=args.cutoff, + model=args.model, + h=args.h, + max_window_size=args.max_window_size, + ) + metrics = pd.read_parquet(out_path) + print(f"Saved metrics: {out_path}") + print(metrics.head()) + print(metrics.shape) + + +if __name__ == "__main__": + main() diff --git a/src/evaluation/cenace/metrics.py b/src/evaluation/cenace/metrics.py new file mode 100644 index 0000000..435fe70 --- /dev/null +++ b/src/evaluation/cenace/metrics.py @@ -0,0 +1,47 @@ +from __future__ import annotations + +import pandas as pd + + +def mae(y_true: pd.Series, y_pred: pd.Series) -> float: + return (y_true - y_pred).abs().mean() + + +def rmse(y_true: pd.Series, y_pred: pd.Series) -> float: + return ((y_true - y_pred) ** 2).mean() ** 0.5 + + +def smape(y_true: pd.Series, y_pred: pd.Series) -> float: + denom = (y_true.abs() + y_pred.abs()) / 2 + out = (y_true - y_pred).abs() / denom + out = out.where(denom != 0, 0.0) + return 100 * out.mean() + + +def evaluate_forecasts(merged: pd.DataFrame) -> pd.DataFrame: + per_uid = ( + merged.groupby("unique_id", as_index=False) + .apply( + lambda g: pd.Series( + { + "mae": mae(g["y"], g["y_hat"]), + "rmse": rmse(g["y"], g["y_hat"]), + "smape": smape(g["y"], g["y_hat"]), + } + ) + ) + .reset_index(drop=True) + ) + + overall = pd.DataFrame( + [ + { + "unique_id": "__overall__", + "mae": mae(merged["y"], merged["y_hat"]), + "rmse": rmse(merged["y"], merged["y_hat"]), + "smape": smape(merged["y"], merged["y_hat"]), + } + ] + ) + + return pd.concat([per_uid, overall], ignore_index=True) diff --git a/src/forecast/cenace/__init__.py b/src/forecast/cenace/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/src/forecast/cenace/core.py b/src/forecast/cenace/core.py new file mode 100644 index 0000000..1bc8b72 --- /dev/null +++ b/src/forecast/cenace/core.py @@ -0,0 +1,93 @@ +from __future__ import annotations + +import argparse +from pathlib import Path + +import pandas as pd + +from src.data.cenace.config import FORECASTS_HOURLY_DIR, PROCESSED_EVENTS_HOURLY_DIR +from src.data.cenace.utils.cenace_data import CENACEData +from src.forecast.cenace.models import MODEL_REGISTRY + + +def cutoff_partition(root: Path, cutoff: pd.Timestamp) -> Path: + return ( + root + / f"year={cutoff.year:04d}" + / f"month={cutoff.month:02d}" + / f"day={cutoff.day:02d}" + ) + + +def parse_args() -> argparse.Namespace: + parser = argparse.ArgumentParser() + parser.add_argument("--cutoff", required=True) + parser.add_argument("--model", required=True, choices=sorted(MODEL_REGISTRY)) + parser.add_argument("--h", type=int, default=24) + parser.add_argument("--max-window-size", type=int, default=48) + return parser.parse_args() + + +def run_forecast( + cutoff: str | pd.Timestamp, + model: str, + h: int = 24, + max_window_size: int = 48, +) -> Path: + cutoff = pd.Timestamp(cutoff) + + data = CENACEData( + base_path=PROCESSED_EVENTS_HOURLY_DIR, + freq="hourly", + h=h, + max_window_size=max_window_size, + ) + + train = data.get_df(cutoff, max_window_size=max_window_size) + + if model not in MODEL_REGISTRY: + raise ValueError( + f"Unknown CENACE model: {model}. " f"Available: {sorted(MODEL_REGISTRY)}" + ) + + model_fn = MODEL_REGISTRY[model] + forecasts = model_fn(train, cutoff=cutoff, h=h) + + forecast_root = FORECASTS_HOURLY_DIR / model + out_dir = cutoff_partition(forecast_root, cutoff) + out_dir.mkdir(parents=True, exist_ok=True) + out_path = out_dir / "forecasts.parquet" + + forecasts.to_parquet(out_path, index=False) + return out_path + + +def main() -> None: + args = parse_args() + cutoff = pd.Timestamp(args.cutoff) + + data = CENACEData( + base_path=PROCESSED_EVENTS_HOURLY_DIR, + freq="hourly", + h=args.h, + max_window_size=args.max_window_size, + ) + + train = data.get_df(cutoff, max_window_size=args.max_window_size) + model_fn = MODEL_REGISTRY[args.model] + forecasts = model_fn(train, cutoff=cutoff, h=args.h) + + forecast_root = FORECASTS_HOURLY_DIR / args.model + out_dir = cutoff_partition(forecast_root, cutoff) + out_dir.mkdir(parents=True, exist_ok=True) + out_path = out_dir / "forecasts.parquet" + + forecasts.to_parquet(out_path, index=False) + + print(f"Saved forecasts: {out_path}") + print(forecasts.head()) + print(forecasts.shape) + + +if __name__ == "__main__": + main() diff --git a/src/forecast/cenace/models.py b/src/forecast/cenace/models.py new file mode 100644 index 0000000..3a11d21 --- /dev/null +++ b/src/forecast/cenace/models.py @@ -0,0 +1,122 @@ +from __future__ import annotations + +import pandas as pd + +from src.forecast.forecast import generate_forecast + + +def naive_last_value(train_df: pd.DataFrame, cutoff: str, h: int = 24) -> pd.DataFrame: + cutoff = pd.Timestamp(cutoff) + + last_values = ( + train_df.sort_values(["unique_id", "ds"]) + .groupby("unique_id", as_index=False) + .tail(1)[["unique_id", "y"]] + .rename(columns={"y": "y_hat"}) + ) + + future_ds = pd.date_range( + cutoff + pd.Timedelta(hours=1), + periods=h, + freq="h", + ) + + out = [] + for ds in future_ds: + tmp = last_values.copy() + tmp["ds"] = ds + out.append(tmp) + + fcst = pd.concat(out, ignore_index=True) + return ( + fcst[["unique_id", "ds", "y_hat"]] + .sort_values(["unique_id", "ds"]) + .reset_index(drop=True) + ) + + +def seasonal_naive_24(train_df: pd.DataFrame, cutoff: str, h: int = 24) -> pd.DataFrame: + cutoff = pd.Timestamp(cutoff) + + expected_start = cutoff - pd.Timedelta(hours=23) + last_day = train_df.loc[ + (train_df["ds"] >= expected_start) & (train_df["ds"] <= cutoff), + ["unique_id", "ds", "y"], + ].copy() + + counts = last_day.groupby("unique_id")["ds"].count() + bad_ids = counts[counts != 24] + if not bad_ids.empty: + raise ValueError( + "seasonal_naive_24 needs exactly 24 hourly observations in the last day " + f"for every series. Bad series count: {len(bad_ids)}" + ) + + last_day["hour_ahead"] = ( + (last_day["ds"] - expected_start) / pd.Timedelta(hours=1) + ).astype(int) + profile = last_day[["unique_id", "hour_ahead", "y"]].rename(columns={"y": "y_hat"}) + + unique_ids = profile["unique_id"].drop_duplicates().sort_values().tolist() + future_ds = pd.date_range(cutoff + pd.Timedelta(hours=1), periods=h, freq="h") + + future_index = pd.DataFrame( + [(uid, ds, i % 24) for i, ds in enumerate(future_ds) for uid in unique_ids], + columns=["unique_id", "ds", "hour_ahead"], + ) + + fcst = ( + future_index.merge(profile, on=["unique_id", "hour_ahead"], how="left")[ + ["unique_id", "ds", "y_hat"] + ] + .sort_values(["unique_id", "ds"]) + .reset_index(drop=True) + ) + + if fcst["y_hat"].isna().any(): + raise ValueError( + "seasonal_naive_24" + " produced missing forecasts after profile merge." + ) + + return fcst + + +def auto_arima(train_df: pd.DataFrame, cutoff: str, h: int = 24) -> pd.DataFrame: + cutoff = pd.Timestamp(cutoff) + + raw = generate_forecast( + model_name="auto_arima", + df=train_df, + h=h, + freq="h", + ).copy() + + point_candidates = [ + col + for col in raw.columns + if col not in {"unique_id", "ds"} and "-q-" not in col + ] + if "auto_arima" in point_candidates: + point_col = "auto_arima" + elif len(point_candidates) == 1: + point_col = point_candidates[0] + else: + raise ValueError( + "Could not detect AutoARIMA point forecast column. " + f"Candidates found: {point_candidates}" + ) + + fcst = ( + raw[["unique_id", "ds", point_col]] + .rename(columns={point_col: "y_hat"}) + .sort_values(["unique_id", "ds"]) + .reset_index(drop=True) + ) + return fcst + + +MODEL_REGISTRY = { + "naive_last_value": naive_last_value, + "seasonal_naive_24": seasonal_naive_24, + "auto_arima": auto_arima, +} diff --git a/src/pipeline.py b/src/pipeline.py new file mode 100644 index 0000000..28bc14b --- /dev/null +++ b/src/pipeline.py @@ -0,0 +1,32 @@ +from __future__ import annotations + +import argparse + +from src.dataset_registry import PIPELINE_DATASET_CHOICES, run_dataset_pipeline + + +def parse_args() -> argparse.Namespace: + parser = argparse.ArgumentParser() + parser.add_argument("--dataset", required=True, choices=PIPELINE_DATASET_CHOICES) + parser.add_argument("--model", required=True) + parser.add_argument("--cutoff", required=True) + parser.add_argument("--h", type=int, default=24) + parser.add_argument("--max-window-size", type=int, default=48) + parser.add_argument("--skip-aggregate", action="store_true") + return parser.parse_args() + + +def main() -> None: + args = parse_args() + run_dataset_pipeline( + dataset=args.dataset, + cutoff=args.cutoff, + model=args.model, + h=args.h, + max_window_size=args.max_window_size, + skip_aggregate=args.skip_aggregate, + ) + + +if __name__ == "__main__": + main() From 157587711987842e0b1c5923c6ec79cee4473c1f Mon Sep 17 00:00:00 2001 From: elmartinj Date: Mon, 20 Apr 2026 17:18:16 -0600 Subject: [PATCH 2/9] Reshape CENACE forecast and evaluation to follow GH structure --- src/cenace_pipeline.py | 69 ----------------- src/dataset_registry.py | 25 ------- src/evaluation/cenace/metrics.py | 47 ------------ src/forecast/cenace/models.py | 122 ------------------------------- src/pipeline.py | 32 -------- 5 files changed, 295 deletions(-) delete mode 100644 src/cenace_pipeline.py delete mode 100644 src/dataset_registry.py delete mode 100644 src/evaluation/cenace/metrics.py delete mode 100644 src/forecast/cenace/models.py delete mode 100644 src/pipeline.py diff --git a/src/cenace_pipeline.py b/src/cenace_pipeline.py deleted file mode 100644 index e437e7c..0000000 --- a/src/cenace_pipeline.py +++ /dev/null @@ -1,69 +0,0 @@ -from __future__ import annotations - -import argparse - -import pandas as pd - -from src.data.cenace.aggregate.core import build_hourly_partitions -from src.evaluation.cenace.core import run_evaluation -from src.forecast.cenace.core import run_forecast - - -def run_cenace_pipeline( - cutoff: str, - model: str, - h: int = 24, - max_window_size: int = 48, - skip_aggregate: bool = False, -) -> tuple[str, str]: - try: - cutoff_ts = pd.Timestamp(cutoff) - except Exception as exc: - raise ValueError(f"Invalid cutoff timestamp: {cutoff}") from exc - - if not skip_aggregate: - n_written = build_hourly_partitions() - print(f"Aggregated {n_written} partitions") - - forecast_path = run_forecast( - cutoff=cutoff_ts, - model=model, - h=h, - max_window_size=max_window_size, - ) - print(f"Forecasts saved to: {forecast_path}") - - eval_path = run_evaluation( - cutoff=cutoff_ts, - model=model, - h=h, - max_window_size=max_window_size, - ) - print(f"Metrics saved to: {eval_path}") - - return str(forecast_path), str(eval_path) - - -def parse_args() -> argparse.Namespace: - parser = argparse.ArgumentParser() - parser.add_argument("--cutoff", required=True) - parser.add_argument("--model", required=True) - parser.add_argument("--h", type=int, default=24) - parser.add_argument("--max-window-size", type=int, default=48) - parser.add_argument("--skip-aggregate", action="store_true") - return parser.parse_args() - - -def main() -> None: - args = parse_args() - run_cenace_pipeline( - cutoff=args.cutoff, - model=args.model, - h=args.h, - max_window_size=args.max_window_size, - skip_aggregate=args.skip_aggregate, - ) - - -if __name__ == "__main__": - main() diff --git a/src/dataset_registry.py b/src/dataset_registry.py deleted file mode 100644 index 87e674c..0000000 --- a/src/dataset_registry.py +++ /dev/null @@ -1,25 +0,0 @@ -from __future__ import annotations - -from collections.abc import Callable -from typing import Any - -from src.cenace_pipeline import run_cenace_pipeline - -DatasetRunner = Callable[..., tuple[str, str]] - -DATASET_REGISTRY: dict[str, DatasetRunner] = { - "cenace": run_cenace_pipeline, -} - -PIPELINE_DATASET_CHOICES = sorted(DATASET_REGISTRY) - - -def run_dataset_pipeline(dataset: str, **kwargs: Any) -> tuple[str, str]: - try: - runner = DATASET_REGISTRY[dataset] - except KeyError as exc: - raise ValueError( - f"Unsupported dataset: {dataset}. " f"Available: {PIPELINE_DATASET_CHOICES}" - ) from exc - - return runner(**kwargs) diff --git a/src/evaluation/cenace/metrics.py b/src/evaluation/cenace/metrics.py deleted file mode 100644 index 435fe70..0000000 --- a/src/evaluation/cenace/metrics.py +++ /dev/null @@ -1,47 +0,0 @@ -from __future__ import annotations - -import pandas as pd - - -def mae(y_true: pd.Series, y_pred: pd.Series) -> float: - return (y_true - y_pred).abs().mean() - - -def rmse(y_true: pd.Series, y_pred: pd.Series) -> float: - return ((y_true - y_pred) ** 2).mean() ** 0.5 - - -def smape(y_true: pd.Series, y_pred: pd.Series) -> float: - denom = (y_true.abs() + y_pred.abs()) / 2 - out = (y_true - y_pred).abs() / denom - out = out.where(denom != 0, 0.0) - return 100 * out.mean() - - -def evaluate_forecasts(merged: pd.DataFrame) -> pd.DataFrame: - per_uid = ( - merged.groupby("unique_id", as_index=False) - .apply( - lambda g: pd.Series( - { - "mae": mae(g["y"], g["y_hat"]), - "rmse": rmse(g["y"], g["y_hat"]), - "smape": smape(g["y"], g["y_hat"]), - } - ) - ) - .reset_index(drop=True) - ) - - overall = pd.DataFrame( - [ - { - "unique_id": "__overall__", - "mae": mae(merged["y"], merged["y_hat"]), - "rmse": rmse(merged["y"], merged["y_hat"]), - "smape": smape(merged["y"], merged["y_hat"]), - } - ] - ) - - return pd.concat([per_uid, overall], ignore_index=True) diff --git a/src/forecast/cenace/models.py b/src/forecast/cenace/models.py deleted file mode 100644 index 3a11d21..0000000 --- a/src/forecast/cenace/models.py +++ /dev/null @@ -1,122 +0,0 @@ -from __future__ import annotations - -import pandas as pd - -from src.forecast.forecast import generate_forecast - - -def naive_last_value(train_df: pd.DataFrame, cutoff: str, h: int = 24) -> pd.DataFrame: - cutoff = pd.Timestamp(cutoff) - - last_values = ( - train_df.sort_values(["unique_id", "ds"]) - .groupby("unique_id", as_index=False) - .tail(1)[["unique_id", "y"]] - .rename(columns={"y": "y_hat"}) - ) - - future_ds = pd.date_range( - cutoff + pd.Timedelta(hours=1), - periods=h, - freq="h", - ) - - out = [] - for ds in future_ds: - tmp = last_values.copy() - tmp["ds"] = ds - out.append(tmp) - - fcst = pd.concat(out, ignore_index=True) - return ( - fcst[["unique_id", "ds", "y_hat"]] - .sort_values(["unique_id", "ds"]) - .reset_index(drop=True) - ) - - -def seasonal_naive_24(train_df: pd.DataFrame, cutoff: str, h: int = 24) -> pd.DataFrame: - cutoff = pd.Timestamp(cutoff) - - expected_start = cutoff - pd.Timedelta(hours=23) - last_day = train_df.loc[ - (train_df["ds"] >= expected_start) & (train_df["ds"] <= cutoff), - ["unique_id", "ds", "y"], - ].copy() - - counts = last_day.groupby("unique_id")["ds"].count() - bad_ids = counts[counts != 24] - if not bad_ids.empty: - raise ValueError( - "seasonal_naive_24 needs exactly 24 hourly observations in the last day " - f"for every series. Bad series count: {len(bad_ids)}" - ) - - last_day["hour_ahead"] = ( - (last_day["ds"] - expected_start) / pd.Timedelta(hours=1) - ).astype(int) - profile = last_day[["unique_id", "hour_ahead", "y"]].rename(columns={"y": "y_hat"}) - - unique_ids = profile["unique_id"].drop_duplicates().sort_values().tolist() - future_ds = pd.date_range(cutoff + pd.Timedelta(hours=1), periods=h, freq="h") - - future_index = pd.DataFrame( - [(uid, ds, i % 24) for i, ds in enumerate(future_ds) for uid in unique_ids], - columns=["unique_id", "ds", "hour_ahead"], - ) - - fcst = ( - future_index.merge(profile, on=["unique_id", "hour_ahead"], how="left")[ - ["unique_id", "ds", "y_hat"] - ] - .sort_values(["unique_id", "ds"]) - .reset_index(drop=True) - ) - - if fcst["y_hat"].isna().any(): - raise ValueError( - "seasonal_naive_24" + " produced missing forecasts after profile merge." - ) - - return fcst - - -def auto_arima(train_df: pd.DataFrame, cutoff: str, h: int = 24) -> pd.DataFrame: - cutoff = pd.Timestamp(cutoff) - - raw = generate_forecast( - model_name="auto_arima", - df=train_df, - h=h, - freq="h", - ).copy() - - point_candidates = [ - col - for col in raw.columns - if col not in {"unique_id", "ds"} and "-q-" not in col - ] - if "auto_arima" in point_candidates: - point_col = "auto_arima" - elif len(point_candidates) == 1: - point_col = point_candidates[0] - else: - raise ValueError( - "Could not detect AutoARIMA point forecast column. " - f"Candidates found: {point_candidates}" - ) - - fcst = ( - raw[["unique_id", "ds", point_col]] - .rename(columns={point_col: "y_hat"}) - .sort_values(["unique_id", "ds"]) - .reset_index(drop=True) - ) - return fcst - - -MODEL_REGISTRY = { - "naive_last_value": naive_last_value, - "seasonal_naive_24": seasonal_naive_24, - "auto_arima": auto_arima, -} diff --git a/src/pipeline.py b/src/pipeline.py deleted file mode 100644 index 28bc14b..0000000 --- a/src/pipeline.py +++ /dev/null @@ -1,32 +0,0 @@ -from __future__ import annotations - -import argparse - -from src.dataset_registry import PIPELINE_DATASET_CHOICES, run_dataset_pipeline - - -def parse_args() -> argparse.Namespace: - parser = argparse.ArgumentParser() - parser.add_argument("--dataset", required=True, choices=PIPELINE_DATASET_CHOICES) - parser.add_argument("--model", required=True) - parser.add_argument("--cutoff", required=True) - parser.add_argument("--h", type=int, default=24) - parser.add_argument("--max-window-size", type=int, default=48) - parser.add_argument("--skip-aggregate", action="store_true") - return parser.parse_args() - - -def main() -> None: - args = parse_args() - run_dataset_pipeline( - dataset=args.dataset, - cutoff=args.cutoff, - model=args.model, - h=args.h, - max_window_size=args.max_window_size, - skip_aggregate=args.skip_aggregate, - ) - - -if __name__ == "__main__": - main() From cb21d1e5eedb2afd3fae03e79ce4c16b0418750b Mon Sep 17 00:00:00 2001 From: elmartinj Date: Mon, 20 Apr 2026 17:42:37 -0600 Subject: [PATCH 3/9] Align CENACE forecast and evaluation with GH structure --- src/data/cenace/extract/core.py | 107 ++++++++++++++++++++++++++++++++ src/evaluation/cenace/core.py | 14 +++-- src/forecast/cenace/core.py | 44 +++++-------- 3 files changed, 131 insertions(+), 34 deletions(-) create mode 100644 src/data/cenace/extract/core.py diff --git a/src/data/cenace/extract/core.py b/src/data/cenace/extract/core.py new file mode 100644 index 0000000..3388b7a --- /dev/null +++ b/src/data/cenace/extract/core.py @@ -0,0 +1,107 @@ +import argparse +from datetime import datetime, timedelta +from pathlib import Path +import requests +from bs4 import BeautifulSoup +import zipfile + +URL = "https://www.cenace.gob.mx/Paginas/SIM/Reportes/PreEnerServConMTR.aspx" + +session = requests.Session() + +HEADERS = { + "User-Agent": "Mozilla/5.0", + "Referer": URL, + "Origin": "https://www.cenace.gob.mx", + "Content-Type": "application/x-www-form-urlencoded", +} + +# repo root = imper/ +ROOT_DIR = Path(__file__).resolve().parents[4] +DEFAULT_BASE_DIR = ROOT_DIR / "data" / "cenace" + + +def get_form_state(): + r = session.get(URL, headers=HEADERS) + soup = BeautifulSoup(r.text, "html.parser") + + def get_value(name): + el = soup.find("input", {"name": name}) + return el.get("value") if el else "" + + return { + "__VIEWSTATE": get_value("__VIEWSTATE"), + "__VIEWSTATEGENERATOR": get_value("__VIEWSTATEGENERATOR"), + "__VIEWSTATEENCRYPTED": get_value("__VIEWSTATEENCRYPTED"), + "__EVENTVALIDATION": get_value("__EVENTVALIDATION"), + } + + +def download_and_extract(date, raw_dir, tmp_dir): + date_str = date.strftime("%d/%m/%Y") + period_str = f"{date_str} - {date_str}" + + state = get_form_state() + + payload = { + "ctl00$ContentPlaceHolder1$ddlReporte": "362,325", + "ctl00$ContentPlaceHolder1$ddlPeriodicidad": "D", + "ctl00$ContentPlaceHolder1$ddlSistema": "SIN", + "ctl00$ContentPlaceHolder1$txtPeriodo": period_str, + "ctl00$ContentPlaceHolder1$hdfStartDateSelected": date_str, + "ctl00$ContentPlaceHolder1$hdfEndDateSelected": date_str, + "ctl00$ContentPlaceHolder1$btnDescargarZIP": "Descargar ZIP", + "__VIEWSTATE": state["__VIEWSTATE"], + "__VIEWSTATEGENERATOR": state["__VIEWSTATEGENERATOR"], + "__VIEWSTATEENCRYPTED": state["__VIEWSTATEENCRYPTED"], + "__EVENTVALIDATION": state["__EVENTVALIDATION"], + "__EVENTTARGET": "", + "__EVENTARGUMENT": "", + } + + r = session.post(URL, data=payload, headers=HEADERS) + + size = len(r.content) + print(f"{date_str} | {size} bytes") + + if size < 10000: + print(f"Skipping {date_str}") + return + + raw_dir.mkdir(parents=True, exist_ok=True) + tmp_dir.mkdir(parents=True, exist_ok=True) + + zip_path = raw_dir / f"{date.strftime('%Y%m%d')}.zip" + + with open(zip_path, "wb") as f: + f.write(r.content) + + with zipfile.ZipFile(zip_path, "r") as z: + z.extractall(tmp_dir) + + +def run(start_date, end_date, base_dir): + raw_dir = base_dir / "raw" + tmp_dir = base_dir / "tmp" + + current = start_date + while current <= end_date: + try: + download_and_extract(current, raw_dir, tmp_dir) + except Exception as e: + print(f"Error on {current.strftime('%Y-%m-%d')}: {e}") + current += timedelta(days=1) + +if __name__ == "__main__": + parser = argparse.ArgumentParser() + parser.add_argument("--start-date", default="2023-01-01") + parser.add_argument("--end-date", required=True) + parser.add_argument("--out", default=str(DEFAULT_BASE_DIR)) + + args = parser.parse_args() + + run( + start_date=datetime.strptime(args.start_date, "%Y-%m-%d"), + end_date=datetime.strptime(args.end_date, "%Y-%m-%d"), + base_dir=Path(args.out).resolve(), + ) diff --git a/src/evaluation/cenace/core.py b/src/evaluation/cenace/core.py index b6207e9..1257ac1 100644 --- a/src/evaluation/cenace/core.py +++ b/src/evaluation/cenace/core.py @@ -11,7 +11,7 @@ PROCESSED_EVENTS_HOURLY_DIR, ) from src.data.cenace.utils.cenace_data import CENACEData -from src.evaluation.cenace.metrics import evaluate_forecasts +from src.evaluation.evaluate import evaluate_forecast def cutoff_partition(root: Path, cutoff: pd.Timestamp) -> Path: @@ -47,14 +47,16 @@ def run_evaluation( / "forecasts.parquet" ) + train = data.get_df(cutoff, max_window_size=max_window_size) actuals = data.get_actuals(cutoff, h=h) forecasts = pd.read_parquet(forecast_path) - merged = forecasts.merge(actuals, on=["unique_id", "ds"], how="inner") - if merged.empty: - raise ValueError("Merged forecasts/actuals is empty") - - metrics = evaluate_forecasts(merged) + metrics, _ = evaluate_forecast( + forecast_df=forecasts, + actuals_df=actuals, + train_df=train, + seasonality=24, + ) eval_root = EVALUATIONS_HOURLY_DIR / model out_dir = cutoff_partition(eval_root, cutoff) diff --git a/src/forecast/cenace/core.py b/src/forecast/cenace/core.py index 1bc8b72..2b31af1 100644 --- a/src/forecast/cenace/core.py +++ b/src/forecast/cenace/core.py @@ -7,7 +7,7 @@ from src.data.cenace.config import FORECASTS_HOURLY_DIR, PROCESSED_EVENTS_HOURLY_DIR from src.data.cenace.utils.cenace_data import CENACEData -from src.forecast.cenace.models import MODEL_REGISTRY +from src.forecast.forecast import generate_forecast def cutoff_partition(root: Path, cutoff: pd.Timestamp) -> Path: @@ -22,7 +22,7 @@ def cutoff_partition(root: Path, cutoff: pd.Timestamp) -> Path: def parse_args() -> argparse.Namespace: parser = argparse.ArgumentParser() parser.add_argument("--cutoff", required=True) - parser.add_argument("--model", required=True, choices=sorted(MODEL_REGISTRY)) + parser.add_argument("--model", required=True) parser.add_argument("--h", type=int, default=24) parser.add_argument("--max-window-size", type=int, default=48) return parser.parse_args() @@ -45,13 +45,17 @@ def run_forecast( train = data.get_df(cutoff, max_window_size=max_window_size) - if model not in MODEL_REGISTRY: - raise ValueError( - f"Unknown CENACE model: {model}. " f"Available: {sorted(MODEL_REGISTRY)}" - ) + model_name = "seasonal_naive" if model == "seasonal_naive_24" else model - model_fn = MODEL_REGISTRY[model] - forecasts = model_fn(train, cutoff=cutoff, h=h) + forecasts = generate_forecast( + model_name=model_name, + df=train, + h=h, + freq="h", + ) + + if "y_hat" in forecasts.columns: + forecasts = forecasts.rename(columns={"y_hat": model}) forecast_root = FORECASTS_HOURLY_DIR / model out_dir = cutoff_partition(forecast_root, cutoff) @@ -64,29 +68,13 @@ def run_forecast( def main() -> None: args = parse_args() - cutoff = pd.Timestamp(args.cutoff) - - data = CENACEData( - base_path=PROCESSED_EVENTS_HOURLY_DIR, - freq="hourly", + out_path = run_forecast( + cutoff=args.cutoff, + model=args.model, h=args.h, max_window_size=args.max_window_size, ) - - train = data.get_df(cutoff, max_window_size=args.max_window_size) - model_fn = MODEL_REGISTRY[args.model] - forecasts = model_fn(train, cutoff=cutoff, h=args.h) - - forecast_root = FORECASTS_HOURLY_DIR / args.model - out_dir = cutoff_partition(forecast_root, cutoff) - out_dir.mkdir(parents=True, exist_ok=True) - out_path = out_dir / "forecasts.parquet" - - forecasts.to_parquet(out_path, index=False) - - print(f"Saved forecasts: {out_path}") - print(forecasts.head()) - print(forecasts.shape) + print(f"Forecasts saved to: {out_path}") if __name__ == "__main__": From 3c2bf0eccba5cbb028af3e52c58da2dd0d22ded2 Mon Sep 17 00:00:00 2001 From: elmartinj Date: Mon, 20 Apr 2026 17:45:22 -0600 Subject: [PATCH 4/9] Add CENACE extractor with execution-date and backfill logic --- src/data/cenace/extract/core.py | 106 +++++++++++++++++++++++++------- 1 file changed, 84 insertions(+), 22 deletions(-) diff --git a/src/data/cenace/extract/core.py b/src/data/cenace/extract/core.py index 3388b7a..e9b633f 100644 --- a/src/data/cenace/extract/core.py +++ b/src/data/cenace/extract/core.py @@ -1,9 +1,12 @@ +from __future__ import annotations + import argparse from datetime import datetime, timedelta from pathlib import Path +import zipfile + import requests from bs4 import BeautifulSoup -import zipfile URL = "https://www.cenace.gob.mx/Paginas/SIM/Reportes/PreEnerServConMTR.aspx" @@ -16,16 +19,25 @@ "Content-Type": "application/x-www-form-urlencoded", } -# repo root = imper/ +# repo root = impermanent/ ROOT_DIR = Path(__file__).resolve().parents[4] DEFAULT_BASE_DIR = ROOT_DIR / "data" / "cenace" -def get_form_state(): +def target_date_for_execution(execution_date: datetime) -> datetime: + return execution_date + timedelta(days=1) + + +def raw_zip_path(date: datetime, raw_dir: Path) -> Path: + return raw_dir / f"{date.strftime('%Y%m%d')}.zip" + + +def get_form_state() -> dict[str, str]: r = session.get(URL, headers=HEADERS) + r.raise_for_status() soup = BeautifulSoup(r.text, "html.parser") - def get_value(name): + def get_value(name: str) -> str: el = soup.find("input", {"name": name}) return el.get("value") if el else "" @@ -37,7 +49,7 @@ def get_value(name): } -def download_and_extract(date, raw_dir, tmp_dir): +def download_and_extract(date: datetime, raw_dir: Path, tmp_dir: Path) -> bool: date_str = date.strftime("%d/%m/%Y") period_str = f"{date_str} - {date_str}" @@ -60,18 +72,19 @@ def download_and_extract(date, raw_dir, tmp_dir): } r = session.post(URL, data=payload, headers=HEADERS) + r.raise_for_status() size = len(r.content) print(f"{date_str} | {size} bytes") if size < 10000: - print(f"Skipping {date_str}") - return + print(f"Skipping {date_str}: file not published or response too small") + return False raw_dir.mkdir(parents=True, exist_ok=True) tmp_dir.mkdir(parents=True, exist_ok=True) - zip_path = raw_dir / f"{date.strftime('%Y%m%d')}.zip" + zip_path = raw_zip_path(date, raw_dir) with open(zip_path, "wb") as f: f.write(r.content) @@ -79,29 +92,78 @@ def download_and_extract(date, raw_dir, tmp_dir): with zipfile.ZipFile(zip_path, "r") as z: z.extractall(tmp_dir) + return True -def run(start_date, end_date, base_dir): + +def backfill_missing(start_date: datetime, end_date: datetime, base_dir: Path) -> None: raw_dir = base_dir / "raw" tmp_dir = base_dir / "tmp" current = start_date while current <= end_date: - try: - download_and_extract(current, raw_dir, tmp_dir) - except Exception as e: - print(f"Error on {current.strftime('%Y-%m-%d')}: {e}") + zip_path = raw_zip_path(current, raw_dir) + if zip_path.exists(): + print(f"Already have {current.strftime('%Y-%m-%d')}, skipping") + else: + try: + ok = download_and_extract(current, raw_dir, tmp_dir) + if not ok: + print(f"Stopping at {current.strftime('%Y-%m-%d')}") + break + except Exception as e: + print(f"Error on {current.strftime('%Y-%m-%d')}: {e}") + break current += timedelta(days=1) -if __name__ == "__main__": + +def run_execution_date(execution_date: datetime, base_dir: Path) -> bool: + raw_dir = base_dir / "raw" + tmp_dir = base_dir / "tmp" + target_date = target_date_for_execution(execution_date) + + zip_path = raw_zip_path(target_date, raw_dir) + if zip_path.exists(): + print(f"Already have {target_date.strftime('%Y-%m-%d')}, skipping") + return True + + return download_and_extract(target_date, raw_dir, tmp_dir) + + +def parse_args() -> argparse.Namespace: parser = argparse.ArgumentParser() - parser.add_argument("--start-date", default="2023-01-01") - parser.add_argument("--end-date", required=True) + parser.add_argument("--execution-date", default=None) + parser.add_argument("--start-date", default=None) + parser.add_argument("--end-date", default=None) parser.add_argument("--out", default=str(DEFAULT_BASE_DIR)) + return parser.parse_args() + + +def main() -> None: + args = parse_args() + base_dir = Path(args.out).resolve() + + if args.start_date: + start_date = datetime.strptime(args.start_date, "%Y-%m-%d") + if args.end_date: + end_date = datetime.strptime(args.end_date, "%Y-%m-%d") + elif args.execution_date: + end_date = target_date_for_execution( + datetime.strptime(args.execution_date, "%Y-%m-%d") + ) + else: + end_date = datetime.today() + backfill_missing(start_date=start_date, end_date=end_date, base_dir=base_dir) + return + + if args.execution_date: + execution_date = datetime.strptime(args.execution_date, "%Y-%m-%d") + ok = run_execution_date(execution_date=execution_date, base_dir=base_dir) + if not ok: + print("No new CENACE publication detected; stopping cleanly") + return - args = parser.parse_args() + raise ValueError("Provide either --start-date or --execution-date") - run( - start_date=datetime.strptime(args.start_date, "%Y-%m-%d"), - end_date=datetime.strptime(args.end_date, "%Y-%m-%d"), - base_dir=Path(args.out).resolve(), - ) + +if __name__ == "__main__": + main() From e60697d35ec9310e28d89db72e6a9cba9ed775cd Mon Sep 17 00:00:00 2001 From: elmartinj Date: Wed, 10 Jun 2026 16:11:17 -0600 Subject: [PATCH 5/9] Add CAISO data pipeline structure --- src/data/caiso/__init__.py | 0 src/data/caiso/aggregate/__init__.py | 0 src/data/caiso/aggregate/core.py | 0 src/data/caiso/config.py | 0 src/data/caiso/extract/__init__.py | 0 src/data/caiso/extract/core.py | 0 src/data/caiso/transform/__init__.py | 0 src/data/caiso/transform/core.py | 0 src/data/caiso/utils/__init__.py | 0 src/data/caiso/utils/caiso_data.py | 0 10 files changed, 0 insertions(+), 0 deletions(-) create mode 100644 src/data/caiso/__init__.py create mode 100644 src/data/caiso/aggregate/__init__.py create mode 100644 src/data/caiso/aggregate/core.py create mode 100644 src/data/caiso/config.py create mode 100644 src/data/caiso/extract/__init__.py create mode 100644 src/data/caiso/extract/core.py create mode 100644 src/data/caiso/transform/__init__.py create mode 100644 src/data/caiso/transform/core.py create mode 100644 src/data/caiso/utils/__init__.py create mode 100644 src/data/caiso/utils/caiso_data.py diff --git a/src/data/caiso/__init__.py b/src/data/caiso/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/src/data/caiso/aggregate/__init__.py b/src/data/caiso/aggregate/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/src/data/caiso/aggregate/core.py b/src/data/caiso/aggregate/core.py new file mode 100644 index 0000000..e69de29 diff --git a/src/data/caiso/config.py b/src/data/caiso/config.py new file mode 100644 index 0000000..e69de29 diff --git a/src/data/caiso/extract/__init__.py b/src/data/caiso/extract/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/src/data/caiso/extract/core.py b/src/data/caiso/extract/core.py new file mode 100644 index 0000000..e69de29 diff --git a/src/data/caiso/transform/__init__.py b/src/data/caiso/transform/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/src/data/caiso/transform/core.py b/src/data/caiso/transform/core.py new file mode 100644 index 0000000..e69de29 diff --git a/src/data/caiso/utils/__init__.py b/src/data/caiso/utils/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/src/data/caiso/utils/caiso_data.py b/src/data/caiso/utils/caiso_data.py new file mode 100644 index 0000000..e69de29 From 5fee22dddfd7a2926ae675251be2cf4ae80803d3 Mon Sep 17 00:00:00 2001 From: elmartinj Date: Wed, 10 Jun 2026 17:23:16 -0600 Subject: [PATCH 6/9] Add working CAISO forecasting pipeline, we have a go --- src/data/caiso/aggregate/core.py | 48 +++++++++++++++ src/data/caiso/config.py | 26 +++++++++ src/data/caiso/extract/core.py | 84 ++++++++++++++++++++++++++ src/data/caiso/transform/core.py | 74 +++++++++++++++++++++++ src/data/caiso/utils/caiso_data.py | 89 ++++++++++++++++++++++++++++ src/evaluation/caiso/__init__.py | 0 src/evaluation/caiso/core.py | 94 ++++++++++++++++++++++++++++++ src/forecast/caiso/__init__.py | 0 src/forecast/caiso/core.py | 81 +++++++++++++++++++++++++ 9 files changed, 496 insertions(+) create mode 100644 src/evaluation/caiso/__init__.py create mode 100644 src/evaluation/caiso/core.py create mode 100644 src/forecast/caiso/__init__.py create mode 100644 src/forecast/caiso/core.py diff --git a/src/data/caiso/aggregate/core.py b/src/data/caiso/aggregate/core.py index e69de29..44ff275 100644 --- a/src/data/caiso/aggregate/core.py +++ b/src/data/caiso/aggregate/core.py @@ -0,0 +1,48 @@ +from __future__ import annotations + +import pandas as pd + +from src.data.caiso.config import PROCESSED_CSV, PROCESSED_EVENTS_HOURLY_DIR + +INPUT_CSV = PROCESSED_CSV +OUTPUT_ROOT = PROCESSED_EVENTS_HOURLY_DIR + + +def build_hourly_partitions() -> int: + df = pd.read_csv(INPUT_CSV) + + df["ds"] = pd.to_datetime(df["ds"], errors="coerce") + df["y"] = pd.to_numeric(df["y"], errors="coerce") + + df = df.dropna(subset=["unique_id", "ds", "y"]).copy() + df = df.sort_values(["unique_id", "ds"]).drop_duplicates(["unique_id", "ds"]) + + df["year"] = df["ds"].dt.year + df["month"] = df["ds"].dt.month + df["day"] = df["ds"].dt.day + + OUTPUT_ROOT.mkdir(parents=True, exist_ok=True) + + n_written = 0 + for (year, month, day), part in df.groupby(["year", "month", "day"], sort=True): + part_dir = ( + OUTPUT_ROOT / f"year={year:04d}" / f"month={month:02d}" / f"day={day:02d}" + ) + part_dir.mkdir(parents=True, exist_ok=True) + + out_path = part_dir / "series.parquet" + part[["unique_id", "ds", "y"]].to_parquet(out_path, index=False) + + print(f"Saved: {out_path}") + n_written += 1 + + return n_written + + +def main() -> None: + n_written = build_hourly_partitions() + print(f"\nDone. Wrote {n_written} daily partitions.") + + +if __name__ == "__main__": + main() diff --git a/src/data/caiso/config.py b/src/data/caiso/config.py index e69de29..779d895 100644 --- a/src/data/caiso/config.py +++ b/src/data/caiso/config.py @@ -0,0 +1,26 @@ +from __future__ import annotations +from datetime import date +from dateutil.relativedelta import relativedelta + +from pathlib import Path + +ROOT = Path(__file__).resolve().parents[3] +DATA_ROOT = ROOT / "data" / "caiso" + +RAW_DIR = DATA_ROOT / "raw" +TMP_DIR = DATA_ROOT / "tmp" +PROCESSED_DIR = DATA_ROOT / "processed" +PROCESSED_CSV = PROCESSED_DIR / "caiso.csv" + +PROCESSED_EVENTS_HOURLY_DIR = DATA_ROOT / "processed-events" / "hourly" +FORECASTS_HOURLY_DIR = DATA_ROOT / "forecasts" / "hourly" +EVALUATIONS_HOURLY_DIR = DATA_ROOT / "evaluations" / "hourly" + +START_DATE = (date.today() - relativedelta(months=39)).isoformat() +MARKET_RUN_ID = "DAM" + +NODES = ( + "TH_NP15_GEN-APND", + "TH_SP15_GEN-APND", + "TH_ZP26_GEN-APND", +) diff --git a/src/data/caiso/extract/core.py b/src/data/caiso/extract/core.py index e69de29..dcb87ee 100644 --- a/src/data/caiso/extract/core.py +++ b/src/data/caiso/extract/core.py @@ -0,0 +1,84 @@ +from __future__ import annotations + +import argparse +import zipfile +from datetime import date, datetime, timedelta +from pathlib import Path +from zoneinfo import ZoneInfo + +import requests + +from src.data.caiso.config import MARKET_RUN_ID, NODES, RAW_DIR, START_DATE + +OASIS_URL = "https://oasis.caiso.com/oasisapi/SingleZip" +PACIFIC = ZoneInfo("America/Los_Angeles") +UTC = ZoneInfo("UTC") + + +def oasis_datetime(value: date) -> str: + local_midnight = datetime.combine(value, datetime.min.time(), PACIFIC) + return local_midnight.astimezone(UTC).strftime("%Y%m%dT%H:%M-0000") + + +def download_node( + node: str, + start: date, + end: date, + output_dir: Path = RAW_DIR, +) -> Path: + output_dir.mkdir(parents=True, exist_ok=True) + + output = output_dir / f"{node}_{start:%Y%m%d}_{end:%Y%m%d}.zip" + + if output.exists() and zipfile.is_zipfile(output): + print(f"Skipping existing: {output}") + return output + + params = { + "queryname": "PRC_LMP", + "startdatetime": oasis_datetime(start), + "enddatetime": oasis_datetime(end), + "version": "12", + "market_run_id": MARKET_RUN_ID, + # "node": node, + "node": ",".join(NODES), + "resultformat": "6", + } + + response = requests.get(OASIS_URL, params=params, timeout=120) + response.raise_for_status() + output.write_bytes(response.content) + + if not zipfile.is_zipfile(output): + output.unlink(missing_ok=True) + raise RuntimeError(f"CAISO returned a non-ZIP response for {node}") + + print(f"Downloaded: {output}") + return output + + +def extract_caiso(start: date, end: date) -> None: + current = start + + while current < end: + chunk_end = min(current + timedelta(days=31), end) + + download_node("CAISO_3_NODES", current, chunk_end) + + current = chunk_end + + +def parse_args() -> argparse.Namespace: + parser = argparse.ArgumentParser() + parser.add_argument("--start", default=START_DATE) + parser.add_argument("--end", required=True) + return parser.parse_args() + + +if __name__ == "__main__": + args = parse_args() + + extract_caiso( + start=date.fromisoformat(args.start), + end=date.fromisoformat(args.end), + ) diff --git a/src/data/caiso/transform/core.py b/src/data/caiso/transform/core.py index e69de29..7b6bd77 100644 --- a/src/data/caiso/transform/core.py +++ b/src/data/caiso/transform/core.py @@ -0,0 +1,74 @@ +from __future__ import annotations + +import argparse +import zipfile +from pathlib import Path + +import pandas as pd + +from src.data.caiso.config import NODES, PROCESSED_CSV, RAW_DIR + + +def transform_caiso( + raw_dir: Path = RAW_DIR, + output_path: Path = PROCESSED_CSV, +) -> Path: + frames: list[pd.DataFrame] = [] + + for zip_path in sorted(raw_dir.glob("CAISO_3_NODES_*.zip")): + with zipfile.ZipFile(zip_path) as archive: + for filename in archive.namelist(): + if not filename.endswith(".csv"): + continue + + with archive.open(filename) as file: + frame = pd.read_csv(file) + + frame = frame[ + (frame["LMP_TYPE"] == "LMP") + & (frame["XML_DATA_ITEM"] == "LMP_PRC") + & (frame["NODE"].isin(NODES)) + ] + + frames.append( + frame[["NODE", "INTERVALSTARTTIME_GMT", "MW"]] + ) + + if not frames: + raise RuntimeError(f"No CAISO ZIP files found in {raw_dir}") + + result = pd.concat(frames, ignore_index=True) + result = result.rename( + columns={ + "NODE": "unique_id", + "INTERVALSTARTTIME_GMT": "ds", + "MW": "y", + } + ) + + result["ds"] = pd.to_datetime(result["ds"], utc=True) + result["y"] = pd.to_numeric(result["y"], errors="coerce") + + result = ( + result.dropna(subset=["unique_id", "ds", "y"]) + .drop_duplicates(["unique_id", "ds"]) + .sort_values(["unique_id", "ds"]) + .reset_index(drop=True) + ) + + output_path.parent.mkdir(parents=True, exist_ok=True) + result.to_csv(output_path, index=False) + + print(result.groupby("unique_id").size()) + print(f"Saved {len(result)} rows to {output_path}") + + return output_path + + +if __name__ == "__main__": + parser = argparse.ArgumentParser() + parser.add_argument("--raw-dir", type=Path, default=RAW_DIR) + parser.add_argument("--output", type=Path, default=PROCESSED_CSV) + args = parser.parse_args() + + transform_caiso(args.raw_dir, args.output) diff --git a/src/data/caiso/utils/caiso_data.py b/src/data/caiso/utils/caiso_data.py index e69de29..7f1bbc6 100644 --- a/src/data/caiso/utils/caiso_data.py +++ b/src/data/caiso/utils/caiso_data.py @@ -0,0 +1,89 @@ +from __future__ import annotations + +from dataclasses import dataclass +from pathlib import Path + +import duckdb +import pandas as pd + + +@dataclass +class CAISOData: + base_path: Path + freq: str = "hourly" + h: int = 24 + max_window_size: int = 24 * 90 + + def __post_init__(self) -> None: + self.base_path = Path(self.base_path) + duckdb.sql("SET TimeZone='UTC'") + + def _date_to_partition(self, d: pd.Timestamp) -> Path: + return ( + self.base_path + / f"year={d.year:04d}" + / f"month={d.month:02d}" + / f"day={d.day:02d}" + / "series.parquet" + ) + + def _paths_for_range(self, start: pd.Timestamp, end: pd.Timestamp) -> list[str]: + days = pd.date_range(start.normalize(), end.normalize(), freq="D") + paths = [self._date_to_partition(d) for d in days] + existing = [str(p) for p in paths if p.exists()] + if not existing: + raise FileNotFoundError( + f"No parquet files found between {start} and \ + {end} under {self.base_path}" + ) + return existing + + def get_df( + self, + cutoff: str | pd.Timestamp, + max_window_size: int | None = None, + sort: bool = True, + ) -> pd.DataFrame: + cutoff = pd.Timestamp(cutoff) + window = max_window_size or self.max_window_size + start = cutoff - pd.Timedelta(hours=window - 1) + + paths = self._paths_for_range(start, cutoff) + + query = f""" + SELECT unique_id, ds, y + FROM read_parquet({paths}) + WHERE ds >= TIMESTAMP '{start}' + AND ds <= TIMESTAMP '{cutoff}' + """ + + df = duckdb.sql(query).df() + df["ds"] = pd.to_datetime(df["ds"]) + + if sort: + df = df.sort_values(["unique_id", "ds"]).reset_index(drop=True) + + return df + + def get_actuals( + self, cutoff: str | pd.Timestamp, h: int | None = None + ) -> pd.DataFrame: + cutoff = pd.Timestamp(cutoff) + horizon = h or self.h + + start = cutoff + pd.Timedelta(hours=1) + end = cutoff + pd.Timedelta(hours=horizon) + + paths = self._paths_for_range(start, end) + + query = f""" + SELECT unique_id, ds, y + FROM read_parquet({paths}) + WHERE ds >= TIMESTAMP '{start}' + AND ds <= TIMESTAMP '{end}' + """ + + df = duckdb.sql(query).df() + df["ds"] = pd.to_datetime(df["ds"]) + df = df.sort_values(["unique_id", "ds"]).reset_index(drop=True) + return df diff --git a/src/evaluation/caiso/__init__.py b/src/evaluation/caiso/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/src/evaluation/caiso/core.py b/src/evaluation/caiso/core.py new file mode 100644 index 0000000..d75843f --- /dev/null +++ b/src/evaluation/caiso/core.py @@ -0,0 +1,94 @@ +from __future__ import annotations + +import argparse +from pathlib import Path + +import pandas as pd + +from src.data.caiso.config import ( + EVALUATIONS_HOURLY_DIR, + FORECASTS_HOURLY_DIR, + PROCESSED_EVENTS_HOURLY_DIR, +) +from src.data.caiso.utils.caiso_data import CAISOData +from src.evaluation.evaluate import evaluate_forecast + + +def cutoff_partition(root: Path, cutoff: pd.Timestamp) -> Path: + return ( + root + / f"year={cutoff.year:04d}" + / f"month={cutoff.month:02d}" + / f"day={cutoff.day:02d}" + ) + + +def run_evaluation( + cutoff: str | pd.Timestamp, + model: str, + h: int = 24, + max_window_size: int = 48, +) -> Path: + cutoff = pd.Timestamp(cutoff) + + data = CAISOData( + base_path=PROCESSED_EVENTS_HOURLY_DIR, + freq="hourly", + h=h, + max_window_size=max_window_size, + ) + + forecast_path = ( + FORECASTS_HOURLY_DIR + / model + / f"year={cutoff.year:04d}" + / f"month={cutoff.month:02d}" + / f"day={cutoff.day:02d}" + / "forecasts.parquet" + ) + + train = data.get_df(cutoff, max_window_size=max_window_size) + actuals = data.get_actuals(cutoff, h=h) + forecasts = pd.read_parquet(forecast_path) + + metrics, _ = evaluate_forecast( + forecast_df=forecasts, + actuals_df=actuals, + train_df=train, + seasonality=24, + ) + + eval_root = EVALUATIONS_HOURLY_DIR / model + out_dir = cutoff_partition(eval_root, cutoff) + out_dir.mkdir(parents=True, exist_ok=True) + out_path = out_dir / "metrics.parquet" + + metrics.to_parquet(out_path, index=False) + return out_path + + +def parse_args() -> argparse.Namespace: + parser = argparse.ArgumentParser() + parser.add_argument("--cutoff", required=True) + parser.add_argument("--model", required=True) + parser.add_argument("--h", type=int, default=24) + parser.add_argument("--max-window-size", type=int, default=48) + return parser.parse_args() + + +def main() -> None: + args = parse_args() + out_path = run_evaluation( + cutoff=args.cutoff, + model=args.model, + h=args.h, + max_window_size=args.max_window_size, + ) + metrics = pd.read_parquet(out_path) + print(f"Saved metrics: {out_path}") + print(metrics.head()) + print(metrics.shape) + + +if __name__ == "__main__": + main() diff --git a/src/forecast/caiso/__init__.py b/src/forecast/caiso/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/src/forecast/caiso/core.py b/src/forecast/caiso/core.py new file mode 100644 index 0000000..73ba1fc --- /dev/null +++ b/src/forecast/caiso/core.py @@ -0,0 +1,81 @@ +from __future__ import annotations + +import argparse +from pathlib import Path + +import pandas as pd + +from src.data.caiso.config import FORECASTS_HOURLY_DIR, PROCESSED_EVENTS_HOURLY_DIR +from src.data.caiso.utils.caiso_data import CAISOData +from src.forecast.forecast import generate_forecast + + +def cutoff_partition(root: Path, cutoff: pd.Timestamp) -> Path: + return ( + root + / f"year={cutoff.year:04d}" + / f"month={cutoff.month:02d}" + / f"day={cutoff.day:02d}" + ) + + +def parse_args() -> argparse.Namespace: + parser = argparse.ArgumentParser() + parser.add_argument("--cutoff", required=True) + parser.add_argument("--model", required=True) + parser.add_argument("--h", type=int, default=24) + parser.add_argument("--max-window-size", type=int, default=48) + return parser.parse_args() + + +def run_forecast( + cutoff: str | pd.Timestamp, + model: str, + h: int = 24, + max_window_size: int = 48, +) -> Path: + cutoff = pd.Timestamp(cutoff) + + data = CAISOData( + base_path=PROCESSED_EVENTS_HOURLY_DIR, + freq="hourly", + h=h, + max_window_size=max_window_size, + ) + + train = data.get_df(cutoff, max_window_size=max_window_size) + + model_name = "seasonal_naive" if model == "seasonal_naive_24" else model + + forecasts = generate_forecast( + model_name=model_name, + df=train, + h=h, + freq="h", + ) + + if "y_hat" in forecasts.columns: + forecasts = forecasts.rename(columns={"y_hat": model}) + + forecast_root = FORECASTS_HOURLY_DIR / model + out_dir = cutoff_partition(forecast_root, cutoff) + out_dir.mkdir(parents=True, exist_ok=True) + out_path = out_dir / "forecasts.parquet" + + forecasts.to_parquet(out_path, index=False) + return out_path + + +def main() -> None: + args = parse_args() + out_path = run_forecast( + cutoff=args.cutoff, + model=args.model, + h=args.h, + max_window_size=args.max_window_size, + ) + print(f"Forecasts saved to: {out_path}") + + +if __name__ == "__main__": + main() From 320c42afe28febcd747397488704198e63750410 Mon Sep 17 00:00:00 2001 From: elmartinj Date: Fri, 12 Jun 2026 09:51:40 -0600 Subject: [PATCH 7/9] Add dateutil typing support to pre-commit and format correction process using pre-commits --- .pre-commit-config.yaml | 4 ++-- pyproject.toml | 1 + uv.lock | 26 +++++++++++++++++++------- 3 files changed, 22 insertions(+), 9 deletions(-) diff --git a/.pre-commit-config.yaml b/.pre-commit-config.yaml index 604256a..f417adb 100644 --- a/.pre-commit-config.yaml +++ b/.pre-commit-config.yaml @@ -11,7 +11,7 @@ repos: hooks: - id: mypy args: [--ignore-missing-imports, --check-untyped-defs] - additional_dependencies: [types-PyYAML, types-requests] + additional_dependencies: [types-PyYAML, types-requests,types-python-dateutil] - repo: https://github.com/pappasam/toml-sort rev: v0.24.2 hooks: @@ -20,4 +20,4 @@ repos: stages: [pre-commit] - id: toml-sort args: ["--all", "--trailing-comma-inline-array", "--in-place", "--check"] - stages: [manual] \ No newline at end of file + stages: [manual] diff --git a/pyproject.toml b/pyproject.toml index 79b05ca..94f0ebc 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -7,6 +7,7 @@ dev = [ "pre-commit", "pytest", "pytest-cov", + "types-python-dateutil>=2.9.0.20260518", ] [project] diff --git a/uv.lock b/uv.lock index af3fa47..b47ecde 100644 --- a/uv.lock +++ b/uv.lock @@ -1,5 +1,5 @@ version = 1 -revision = 2 +revision = 3 requires-python = ">=3.10" resolution-markers = [ "python_full_version >= '3.13' and sys_platform == 'linux'", @@ -2041,7 +2041,6 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/fe/65/5b235b40581ad75ab97dcd8b4218022ae8e3ab77c13c919f1a1dfe9171fd/greenlet-3.3.1-cp310-cp310-macosx_11_0_universal2.whl", hash = "sha256:04bee4775f40ecefcdaa9d115ab44736cd4b9c5fba733575bfe9379419582e13", size = 273723, upload-time = "2026-01-23T15:30:37.521Z" }, { url = "https://files.pythonhosted.org/packages/ce/ad/eb4729b85cba2d29499e0a04ca6fbdd8f540afd7be142fd571eea43d712f/greenlet-3.3.1-cp310-cp310-manylinux_2_24_aarch64.manylinux_2_28_aarch64.whl", hash = "sha256:50e1457f4fed12a50e427988a07f0f9df53cf0ee8da23fab16e6732c2ec909d4", size = 574874, upload-time = "2026-01-23T16:00:54.551Z" }, { url = "https://files.pythonhosted.org/packages/87/32/57cad7fe4c8b82fdaa098c89498ef85ad92dfbb09d5eb713adedfc2ae1f5/greenlet-3.3.1-cp310-cp310-manylinux_2_24_ppc64le.manylinux_2_28_ppc64le.whl", hash = "sha256:070472cd156f0656f86f92e954591644e158fd65aa415ffbe2d44ca77656a8f5", size = 586309, upload-time = "2026-01-23T16:05:25.18Z" }, - { url = "https://files.pythonhosted.org/packages/66/66/f041005cb87055e62b0d68680e88ec1a57f4688523d5e2fb305841bc8307/greenlet-3.3.1-cp310-cp310-manylinux_2_24_s390x.manylinux_2_28_s390x.whl", hash = "sha256:1108b61b06b5224656121c3c8ee8876161c491cbe74e5c519e0634c837cf93d5", size = 597461, upload-time = "2026-01-23T16:15:51.943Z" }, { url = "https://files.pythonhosted.org/packages/87/eb/8a1ec2da4d55824f160594a75a9d8354a5fe0a300fb1c48e7944265217e1/greenlet-3.3.1-cp310-cp310-manylinux_2_24_x86_64.manylinux_2_28_x86_64.whl", hash = "sha256:3a300354f27dd86bae5fbf7002e6dd2b3255cd372e9242c933faf5e859b703fe", size = 586985, upload-time = "2026-01-23T15:32:47.968Z" }, { url = "https://files.pythonhosted.org/packages/15/1c/0621dd4321dd8c351372ee8f9308136acb628600658a49be1b7504208738/greenlet-3.3.1-cp310-cp310-musllinux_1_2_aarch64.whl", hash = "sha256:e84b51cbebf9ae573b5fbd15df88887815e3253fc000a7d0ff95170e8f7e9729", size = 1547271, upload-time = "2026-01-23T16:04:18.977Z" }, { url = "https://files.pythonhosted.org/packages/9d/53/24047f8924c83bea7a59c8678d9571209c6bfe5f4c17c94a78c06024e9f2/greenlet-3.3.1-cp310-cp310-musllinux_1_2_x86_64.whl", hash = "sha256:e0093bd1a06d899892427217f0ff2a3c8f306182b8c754336d32e2d587c131b4", size = 1613427, upload-time = "2026-01-23T15:33:44.428Z" }, @@ -2049,7 +2048,6 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/ec/e8/2e1462c8fdbe0f210feb5ac7ad2d9029af8be3bf45bd9fa39765f821642f/greenlet-3.3.1-cp311-cp311-macosx_11_0_universal2.whl", hash = "sha256:5fd23b9bc6d37b563211c6abbb1b3cab27db385a4449af5c32e932f93017080c", size = 274974, upload-time = "2026-01-23T15:31:02.891Z" }, { url = "https://files.pythonhosted.org/packages/7e/a8/530a401419a6b302af59f67aaf0b9ba1015855ea7e56c036b5928793c5bd/greenlet-3.3.1-cp311-cp311-manylinux_2_24_aarch64.manylinux_2_28_aarch64.whl", hash = "sha256:09f51496a0bfbaa9d74d36a52d2580d1ef5ed4fdfcff0a73730abfbbbe1403dd", size = 577175, upload-time = "2026-01-23T16:00:56.213Z" }, { url = "https://files.pythonhosted.org/packages/8e/89/7e812bb9c05e1aaef9b597ac1d0962b9021d2c6269354966451e885c4e6b/greenlet-3.3.1-cp311-cp311-manylinux_2_24_ppc64le.manylinux_2_28_ppc64le.whl", hash = "sha256:cb0feb07fe6e6a74615ee62a880007d976cf739b6669cce95daa7373d4fc69c5", size = 590401, upload-time = "2026-01-23T16:05:26.365Z" }, - { url = "https://files.pythonhosted.org/packages/70/ae/e2d5f0e59b94a2269b68a629173263fa40b63da32f5c231307c349315871/greenlet-3.3.1-cp311-cp311-manylinux_2_24_s390x.manylinux_2_28_s390x.whl", hash = "sha256:67ea3fc73c8cd92f42467a72b75e8f05ed51a0e9b1d15398c913416f2dafd49f", size = 601161, upload-time = "2026-01-23T16:15:53.456Z" }, { url = "https://files.pythonhosted.org/packages/5c/ae/8d472e1f5ac5efe55c563f3eabb38c98a44b832602e12910750a7c025802/greenlet-3.3.1-cp311-cp311-manylinux_2_24_x86_64.manylinux_2_28_x86_64.whl", hash = "sha256:39eda9ba259cc9801da05351eaa8576e9aa83eb9411e8f0c299e05d712a210f2", size = 590272, upload-time = "2026-01-23T15:32:49.411Z" }, { url = "https://files.pythonhosted.org/packages/a8/51/0fde34bebfcadc833550717eade64e35ec8738e6b097d5d248274a01258b/greenlet-3.3.1-cp311-cp311-musllinux_1_2_aarch64.whl", hash = "sha256:e2e7e882f83149f0a71ac822ebf156d902e7a5d22c9045e3e0d1daf59cee2cc9", size = 1550729, upload-time = "2026-01-23T16:04:20.867Z" }, { url = "https://files.pythonhosted.org/packages/16/c9/2fb47bee83b25b119d5a35d580807bb8b92480a54b68fef009a02945629f/greenlet-3.3.1-cp311-cp311-musllinux_1_2_x86_64.whl", hash = "sha256:80aa4d79eb5564f2e0a6144fcc744b5a37c56c4a92d60920720e99210d88db0f", size = 1615552, upload-time = "2026-01-23T15:33:45.743Z" }, @@ -2058,7 +2056,6 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/f9/c8/9d76a66421d1ae24340dfae7e79c313957f6e3195c144d2c73333b5bfe34/greenlet-3.3.1-cp312-cp312-macosx_11_0_universal2.whl", hash = "sha256:7e806ca53acf6d15a888405880766ec84721aa4181261cd11a457dfe9a7a4975", size = 276443, upload-time = "2026-01-23T15:30:10.066Z" }, { url = "https://files.pythonhosted.org/packages/81/99/401ff34bb3c032d1f10477d199724f5e5f6fbfb59816ad1455c79c1eb8e7/greenlet-3.3.1-cp312-cp312-manylinux_2_24_aarch64.manylinux_2_28_aarch64.whl", hash = "sha256:d842c94b9155f1c9b3058036c24ffb8ff78b428414a19792b2380be9cecf4f36", size = 597359, upload-time = "2026-01-23T16:00:57.394Z" }, { url = "https://files.pythonhosted.org/packages/2b/bc/4dcc0871ed557792d304f50be0f7487a14e017952ec689effe2180a6ff35/greenlet-3.3.1-cp312-cp312-manylinux_2_24_ppc64le.manylinux_2_28_ppc64le.whl", hash = "sha256:20fedaadd422fa02695f82093f9a98bad3dab5fcda793c658b945fcde2ab27ba", size = 607805, upload-time = "2026-01-23T16:05:28.068Z" }, - { url = "https://files.pythonhosted.org/packages/3b/cd/7a7ca57588dac3389e97f7c9521cb6641fd8b6602faf1eaa4188384757df/greenlet-3.3.1-cp312-cp312-manylinux_2_24_s390x.manylinux_2_28_s390x.whl", hash = "sha256:c620051669fd04ac6b60ebc70478210119c56e2d5d5df848baec4312e260e4ca", size = 622363, upload-time = "2026-01-23T16:15:54.754Z" }, { url = "https://files.pythonhosted.org/packages/cf/05/821587cf19e2ce1f2b24945d890b164401e5085f9d09cbd969b0c193cd20/greenlet-3.3.1-cp312-cp312-manylinux_2_24_x86_64.manylinux_2_28_x86_64.whl", hash = "sha256:14194f5f4305800ff329cbf02c5fcc88f01886cadd29941b807668a45f0d2336", size = 609947, upload-time = "2026-01-23T15:32:51.004Z" }, { url = "https://files.pythonhosted.org/packages/a4/52/ee8c46ed9f8babaa93a19e577f26e3d28a519feac6350ed6f25f1afee7e9/greenlet-3.3.1-cp312-cp312-musllinux_1_2_aarch64.whl", hash = "sha256:7b2fe4150a0cf59f847a67db8c155ac36aed89080a6a639e9f16df5d6c6096f1", size = 1567487, upload-time = "2026-01-23T16:04:22.125Z" }, { url = "https://files.pythonhosted.org/packages/8f/7c/456a74f07029597626f3a6db71b273a3632aecb9afafeeca452cfa633197/greenlet-3.3.1-cp312-cp312-musllinux_1_2_x86_64.whl", hash = "sha256:49f4ad195d45f4a66a0eb9c1ba4832bb380570d361912fa3554746830d332149", size = 1636087, upload-time = "2026-01-23T15:33:47.486Z" }, @@ -2067,7 +2064,6 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/ec/ab/d26750f2b7242c2b90ea2ad71de70cfcd73a948a49513188a0fc0d6fc15a/greenlet-3.3.1-cp313-cp313-macosx_11_0_universal2.whl", hash = "sha256:7ab327905cabb0622adca5971e488064e35115430cec2c35a50fd36e72a315b3", size = 275205, upload-time = "2026-01-23T15:30:24.556Z" }, { url = "https://files.pythonhosted.org/packages/10/d3/be7d19e8fad7c5a78eeefb2d896a08cd4643e1e90c605c4be3b46264998f/greenlet-3.3.1-cp313-cp313-manylinux_2_24_aarch64.manylinux_2_28_aarch64.whl", hash = "sha256:65be2f026ca6a176f88fb935ee23c18333ccea97048076aef4db1ef5bc0713ac", size = 599284, upload-time = "2026-01-23T16:00:58.584Z" }, { url = "https://files.pythonhosted.org/packages/ae/21/fe703aaa056fdb0f17e5afd4b5c80195bbdab701208918938bd15b00d39b/greenlet-3.3.1-cp313-cp313-manylinux_2_24_ppc64le.manylinux_2_28_ppc64le.whl", hash = "sha256:7a3ae05b3d225b4155bda56b072ceb09d05e974bc74be6c3fc15463cf69f33fd", size = 610274, upload-time = "2026-01-23T16:05:29.312Z" }, - { url = "https://files.pythonhosted.org/packages/06/00/95df0b6a935103c0452dad2203f5be8377e551b8466a29650c4c5a5af6cc/greenlet-3.3.1-cp313-cp313-manylinux_2_24_s390x.manylinux_2_28_s390x.whl", hash = "sha256:12184c61e5d64268a160226fb4818af4df02cfead8379d7f8b99a56c3a54ff3e", size = 624375, upload-time = "2026-01-23T16:15:55.915Z" }, { url = "https://files.pythonhosted.org/packages/cb/86/5c6ab23bb3c28c21ed6bebad006515cfe08b04613eb105ca0041fecca852/greenlet-3.3.1-cp313-cp313-manylinux_2_24_x86_64.manylinux_2_28_x86_64.whl", hash = "sha256:6423481193bbbe871313de5fd06a082f2649e7ce6e08015d2a76c1e9186ca5b3", size = 612904, upload-time = "2026-01-23T15:32:52.317Z" }, { url = "https://files.pythonhosted.org/packages/c2/f3/7949994264e22639e40718c2daf6f6df5169bf48fb038c008a489ec53a50/greenlet-3.3.1-cp313-cp313-musllinux_1_2_aarch64.whl", hash = "sha256:33a956fe78bbbda82bfc95e128d61129b32d66bcf0a20a1f0c08aa4839ffa951", size = 1567316, upload-time = "2026-01-23T16:04:23.316Z" }, { url = "https://files.pythonhosted.org/packages/8d/6e/d73c94d13b6465e9f7cd6231c68abde838bb22408596c05d9059830b7872/greenlet-3.3.1-cp313-cp313-musllinux_1_2_x86_64.whl", hash = "sha256:4b065d3284be43728dd280f6f9a13990b56470b81be20375a207cdc814a983f2", size = 1636549, upload-time = "2026-01-23T15:33:48.643Z" }, @@ -2076,7 +2072,6 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/ae/fb/011c7c717213182caf78084a9bea51c8590b0afda98001f69d9f853a495b/greenlet-3.3.1-cp314-cp314-macosx_11_0_universal2.whl", hash = "sha256:bd59acd8529b372775cd0fcbc5f420ae20681c5b045ce25bd453ed8455ab99b5", size = 275737, upload-time = "2026-01-23T15:32:16.889Z" }, { url = "https://files.pythonhosted.org/packages/41/2e/a3a417d620363fdbb08a48b1dd582956a46a61bf8fd27ee8164f9dfe87c2/greenlet-3.3.1-cp314-cp314-manylinux_2_24_aarch64.manylinux_2_28_aarch64.whl", hash = "sha256:b31c05dd84ef6871dd47120386aed35323c944d86c3d91a17c4b8d23df62f15b", size = 646422, upload-time = "2026-01-23T16:01:00.354Z" }, { url = "https://files.pythonhosted.org/packages/b4/09/c6c4a0db47defafd2d6bab8ddfe47ad19963b4e30f5bed84d75328059f8c/greenlet-3.3.1-cp314-cp314-manylinux_2_24_ppc64le.manylinux_2_28_ppc64le.whl", hash = "sha256:02925a0bfffc41e542c70aa14c7eda3593e4d7e274bfcccca1827e6c0875902e", size = 658219, upload-time = "2026-01-23T16:05:30.956Z" }, - { url = "https://files.pythonhosted.org/packages/e2/89/b95f2ddcc5f3c2bc09c8ee8d77be312df7f9e7175703ab780f2014a0e781/greenlet-3.3.1-cp314-cp314-manylinux_2_24_s390x.manylinux_2_28_s390x.whl", hash = "sha256:3e0f3878ca3a3ff63ab4ea478585942b53df66ddde327b59ecb191b19dbbd62d", size = 671455, upload-time = "2026-01-23T16:15:57.232Z" }, { url = "https://files.pythonhosted.org/packages/80/38/9d42d60dffb04b45f03dbab9430898352dba277758640751dc5cc316c521/greenlet-3.3.1-cp314-cp314-manylinux_2_24_x86_64.manylinux_2_28_x86_64.whl", hash = "sha256:34a729e2e4e4ffe9ae2408d5ecaf12f944853f40ad724929b7585bca808a9d6f", size = 660237, upload-time = "2026-01-23T15:32:53.967Z" }, { url = "https://files.pythonhosted.org/packages/96/61/373c30b7197f9e756e4c81ae90a8d55dc3598c17673f91f4d31c3c689c3f/greenlet-3.3.1-cp314-cp314-musllinux_1_2_aarch64.whl", hash = "sha256:aec9ab04e82918e623415947921dea15851b152b822661cce3f8e4393c3df683", size = 1615261, upload-time = "2026-01-23T16:04:25.066Z" }, { url = "https://files.pythonhosted.org/packages/fd/d3/ca534310343f5945316f9451e953dcd89b36fe7a19de652a1dc5a0eeef3f/greenlet-3.3.1-cp314-cp314-musllinux_1_2_x86_64.whl", hash = "sha256:71c767cf281a80d02b6c1bdc41c9468e1f5a494fb11bc8688c360524e273d7b1", size = 1683719, upload-time = "2026-01-23T15:33:50.61Z" }, @@ -2085,7 +2080,6 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/28/24/cbbec49bacdcc9ec652a81d3efef7b59f326697e7edf6ed775a5e08e54c2/greenlet-3.3.1-cp314-cp314t-macosx_11_0_universal2.whl", hash = "sha256:3e63252943c921b90abb035ebe9de832c436401d9c45f262d80e2d06cc659242", size = 282706, upload-time = "2026-01-23T15:33:05.525Z" }, { url = "https://files.pythonhosted.org/packages/86/2e/4f2b9323c144c4fe8842a4e0d92121465485c3c2c5b9e9b30a52e80f523f/greenlet-3.3.1-cp314-cp314t-manylinux_2_24_aarch64.manylinux_2_28_aarch64.whl", hash = "sha256:76e39058e68eb125de10c92524573924e827927df5d3891fbc97bd55764a8774", size = 651209, upload-time = "2026-01-23T16:01:01.517Z" }, { url = "https://files.pythonhosted.org/packages/d9/87/50ca60e515f5bb55a2fbc5f0c9b5b156de7d2fc51a0a69abc9d23914a237/greenlet-3.3.1-cp314-cp314t-manylinux_2_24_ppc64le.manylinux_2_28_ppc64le.whl", hash = "sha256:c9f9d5e7a9310b7a2f416dd13d2e3fd8b42d803968ea580b7c0f322ccb389b97", size = 654300, upload-time = "2026-01-23T16:05:32.199Z" }, - { url = "https://files.pythonhosted.org/packages/7c/25/c51a63f3f463171e09cb586eb64db0861eb06667ab01a7968371a24c4f3b/greenlet-3.3.1-cp314-cp314t-manylinux_2_24_s390x.manylinux_2_28_s390x.whl", hash = "sha256:4b9721549a95db96689458a1e0ae32412ca18776ed004463df3a9299c1b257ab", size = 662574, upload-time = "2026-01-23T16:15:58.364Z" }, { url = "https://files.pythonhosted.org/packages/1d/94/74310866dfa2b73dd08659a3d18762f83985ad3281901ba0ee9a815194fb/greenlet-3.3.1-cp314-cp314t-manylinux_2_24_x86_64.manylinux_2_28_x86_64.whl", hash = "sha256:92497c78adf3ac703b57f1e3813c2d874f27f71a178f9ea5887855da413cd6d2", size = 653842, upload-time = "2026-01-23T15:32:55.671Z" }, { url = "https://files.pythonhosted.org/packages/97/43/8bf0ffa3d498eeee4c58c212a3905dd6146c01c8dc0b0a046481ca29b18c/greenlet-3.3.1-cp314-cp314t-musllinux_1_2_aarch64.whl", hash = "sha256:ed6b402bc74d6557a705e197d47f9063733091ed6357b3de33619d8a8d93ac53", size = 1614917, upload-time = "2026-01-23T16:04:26.276Z" }, { url = "https://files.pythonhosted.org/packages/89/90/a3be7a5f378fc6e84abe4dcfb2ba32b07786861172e502388b4c90000d1b/greenlet-3.3.1-cp314-cp314t-musllinux_1_2_x86_64.whl", hash = "sha256:59913f1e5ada20fde795ba906916aea25d442abcc0593fba7e26c92b7ad76249", size = 1676092, upload-time = "2026-01-23T15:33:52.176Z" }, @@ -2404,6 +2398,7 @@ dev = [ { name = "pre-commit" }, { name = "pytest" }, { name = "pytest-cov" }, + { name = "types-python-dateutil" }, ] [package.metadata] @@ -2425,6 +2420,7 @@ dev = [ { name = "pre-commit" }, { name = "pytest" }, { name = "pytest-cov" }, + { name = "types-python-dateutil", specifier = ">=2.9.0.20260518" }, ] [[package]] @@ -7898,6 +7894,13 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/0f/8b/4b61d6e13f7108f36910df9ab4b58fd389cc2520d54d81b88660804aad99/torch-2.10.0-2-cp311-none-macosx_11_0_arm64.whl", hash = "sha256:418997cb02d0a0f1497cf6a09f63166f9f5df9f3e16c8a716ab76a72127c714f", size = 79423467, upload-time = "2026-02-10T21:44:48.711Z" }, { url = "https://files.pythonhosted.org/packages/d3/54/a2ba279afcca44bbd320d4e73675b282fcee3d81400ea1b53934efca6462/torch-2.10.0-2-cp312-none-macosx_11_0_arm64.whl", hash = "sha256:13ec4add8c3faaed8d13e0574f5cd4a323c11655546f91fbe6afa77b57423574", size = 79498202, upload-time = "2026-02-10T21:44:52.603Z" }, { url = "https://files.pythonhosted.org/packages/ec/23/2c9fe0c9c27f7f6cb865abcea8a4568f29f00acaeadfc6a37f6801f84cb4/torch-2.10.0-2-cp313-none-macosx_11_0_arm64.whl", hash = "sha256:e521c9f030a3774ed770a9c011751fb47c4d12029a3d6522116e48431f2ff89e", size = 79498254, upload-time = "2026-02-10T21:44:44.095Z" }, + { url = "https://files.pythonhosted.org/packages/16/ee/efbd56687be60ef9af0c9c0ebe106964c07400eade5b0af8902a1d8cd58c/torch-2.10.0-3-cp310-cp310-manylinux_2_28_x86_64.whl", hash = "sha256:a1ff626b884f8c4e897c4c33782bdacdff842a165fee79817b1dd549fdda1321", size = 915510070, upload-time = "2026-03-11T14:16:39.386Z" }, + { url = "https://files.pythonhosted.org/packages/36/ab/7b562f1808d3f65414cd80a4f7d4bb00979d9355616c034c171249e1a303/torch-2.10.0-3-cp311-cp311-manylinux_2_28_x86_64.whl", hash = "sha256:ac5bdcbb074384c66fa160c15b1ead77839e3fe7ed117d667249afce0acabfac", size = 915518691, upload-time = "2026-03-11T14:15:43.147Z" }, + { url = "https://files.pythonhosted.org/packages/b3/7a/abada41517ce0011775f0f4eacc79659bc9bc6c361e6bfe6f7052a6b9363/torch-2.10.0-3-cp312-cp312-manylinux_2_28_x86_64.whl", hash = "sha256:98c01b8bb5e3240426dcde1446eed6f40c778091c8544767ef1168fc663a05a6", size = 915622781, upload-time = "2026-03-11T14:17:11.354Z" }, + { url = "https://files.pythonhosted.org/packages/ab/c6/4dfe238342ffdcec5aef1c96c457548762d33c40b45a1ab7033bb26d2ff2/torch-2.10.0-3-cp313-cp313-manylinux_2_28_x86_64.whl", hash = "sha256:80b1b5bfe38eb0e9f5ff09f206dcac0a87aadd084230d4a36eea5ec5232c115b", size = 915627275, upload-time = "2026-03-11T14:16:11.325Z" }, + { url = "https://files.pythonhosted.org/packages/d8/f0/72bf18847f58f877a6a8acf60614b14935e2f156d942483af1ffc081aea0/torch-2.10.0-3-cp313-cp313t-manylinux_2_28_x86_64.whl", hash = "sha256:46b3574d93a2a8134b3f5475cfb98e2eb46771794c57015f6ad1fb795ec25e49", size = 915523474, upload-time = "2026-03-11T14:17:44.422Z" }, + { url = "https://files.pythonhosted.org/packages/f4/39/590742415c3030551944edc2ddc273ea1fdfe8ffb2780992e824f1ebee98/torch-2.10.0-3-cp314-cp314-manylinux_2_28_x86_64.whl", hash = "sha256:b1d5e2aba4eb7f8e87fbe04f86442887f9167a35f092afe4c237dfcaaef6e328", size = 915632474, upload-time = "2026-03-11T14:15:13.666Z" }, + { url = "https://files.pythonhosted.org/packages/b6/8e/34949484f764dde5b222b7fe3fede43e4a6f0da9d7f8c370bb617d629ee2/torch-2.10.0-3-cp314-cp314t-manylinux_2_28_x86_64.whl", hash = "sha256:0228d20b06701c05a8f978357f657817a4a63984b0c90745def81c18aedfa591", size = 915523882, upload-time = "2026-03-11T14:14:46.311Z" }, { url = "https://files.pythonhosted.org/packages/0c/1a/c61f36cfd446170ec27b3a4984f072fd06dab6b5d7ce27e11adb35d6c838/torch-2.10.0-cp310-cp310-manylinux_2_28_aarch64.whl", hash = "sha256:5276fa790a666ee8becaffff8acb711922252521b28fbce5db7db5cf9cb2026d", size = 145992962, upload-time = "2026-01-21T16:24:14.04Z" }, { url = "https://files.pythonhosted.org/packages/b5/60/6662535354191e2d1555296045b63e4279e5a9dbad49acf55a5d38655a39/torch-2.10.0-cp310-cp310-manylinux_2_28_x86_64.whl", hash = "sha256:aaf663927bcd490ae971469a624c322202a2a1e68936eb952535ca4cd3b90444", size = 915599237, upload-time = "2026-01-21T16:23:25.497Z" }, { url = "https://files.pythonhosted.org/packages/40/b8/66bbe96f0d79be2b5c697b2e0b187ed792a15c6c4b8904613454651db848/torch-2.10.0-cp310-cp310-win_amd64.whl", hash = "sha256:a4be6a2a190b32ff5c8002a0977a25ea60e64f7ba46b1be37093c141d9c49aeb", size = 113720931, upload-time = "2026-01-21T16:24:23.743Z" }, @@ -8264,6 +8267,15 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/aa/43/58e75bac4219cbafee83179505ff44cae3153ec279be0e30583a73b8f108/types_protobuf-6.32.1.20251210-py3-none-any.whl", hash = "sha256:2641f78f3696822a048cfb8d0ff42ccd85c25f12f871fbebe86da63793692140", size = 77921, upload-time = "2025-12-10T03:14:24.477Z" }, ] +[[package]] +name = "types-python-dateutil" +version = "2.9.0.20260518" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/8d/e8/c01bdf0d7c3659428c091fbd693177093639565bcbc86bc20098e6d37cc6/types_python_dateutil-2.9.0.20260518.tar.gz", hash = "sha256:51f02dc03b61c7f6a07df45797d4dfe8a1aa47f0b7db9ad89f6fd3a1a70e1b51", size = 17082, upload-time = "2026-05-18T06:05:24.508Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/36/22/169273273ca34e9ab0ae2f387ba72ed7e09faaaf834da01d6b89c2bea71a/types_python_dateutil-2.9.0.20260518-py3-none-any.whl", hash = "sha256:d6a9c5bd0de61460c8fdef8ab2b400f956a1a1075cce08d4e2b4434e478c50b8", size = 18431, upload-time = "2026-05-18T06:05:23.641Z" }, +] + [[package]] name = "types-pyyaml" version = "6.0.12.20250915" From b48d30f41678201d4cec28000f30a2ea2128fb7c Mon Sep 17 00:00:00 2001 From: elmartinj Date: Fri, 12 Jun 2026 09:52:27 -0600 Subject: [PATCH 8/9] Add CAISO pipeline and pre-commit typing support --- pyproject.toml | 1 - src/data/caiso/config.py | 5 +++-- src/data/caiso/transform/core.py | 4 +--- src/data/cenace/extract/core.py | 2 +- uv.lock | 11 ----------- 5 files changed, 5 insertions(+), 18 deletions(-) diff --git a/pyproject.toml b/pyproject.toml index 94f0ebc..79b05ca 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -7,7 +7,6 @@ dev = [ "pre-commit", "pytest", "pytest-cov", - "types-python-dateutil>=2.9.0.20260518", ] [project] diff --git a/src/data/caiso/config.py b/src/data/caiso/config.py index 779d895..45fcdbe 100644 --- a/src/data/caiso/config.py +++ b/src/data/caiso/config.py @@ -1,9 +1,10 @@ from __future__ import annotations -from datetime import date -from dateutil.relativedelta import relativedelta +from datetime import date from pathlib import Path +from dateutil.relativedelta import relativedelta + ROOT = Path(__file__).resolve().parents[3] DATA_ROOT = ROOT / "data" / "caiso" diff --git a/src/data/caiso/transform/core.py b/src/data/caiso/transform/core.py index 7b6bd77..537df35 100644 --- a/src/data/caiso/transform/core.py +++ b/src/data/caiso/transform/core.py @@ -30,9 +30,7 @@ def transform_caiso( & (frame["NODE"].isin(NODES)) ] - frames.append( - frame[["NODE", "INTERVALSTARTTIME_GMT", "MW"]] - ) + frames.append(frame[["NODE", "INTERVALSTARTTIME_GMT", "MW"]]) if not frames: raise RuntimeError(f"No CAISO ZIP files found in {raw_dir}") diff --git a/src/data/cenace/extract/core.py b/src/data/cenace/extract/core.py index e9b633f..273dcdf 100644 --- a/src/data/cenace/extract/core.py +++ b/src/data/cenace/extract/core.py @@ -1,9 +1,9 @@ from __future__ import annotations import argparse +import zipfile from datetime import datetime, timedelta from pathlib import Path -import zipfile import requests from bs4 import BeautifulSoup diff --git a/uv.lock b/uv.lock index b47ecde..7156d5e 100644 --- a/uv.lock +++ b/uv.lock @@ -2398,7 +2398,6 @@ dev = [ { name = "pre-commit" }, { name = "pytest" }, { name = "pytest-cov" }, - { name = "types-python-dateutil" }, ] [package.metadata] @@ -2420,7 +2419,6 @@ dev = [ { name = "pre-commit" }, { name = "pytest" }, { name = "pytest-cov" }, - { name = "types-python-dateutil", specifier = ">=2.9.0.20260518" }, ] [[package]] @@ -8267,15 +8265,6 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/aa/43/58e75bac4219cbafee83179505ff44cae3153ec279be0e30583a73b8f108/types_protobuf-6.32.1.20251210-py3-none-any.whl", hash = "sha256:2641f78f3696822a048cfb8d0ff42ccd85c25f12f871fbebe86da63793692140", size = 77921, upload-time = "2025-12-10T03:14:24.477Z" }, ] -[[package]] -name = "types-python-dateutil" -version = "2.9.0.20260518" -source = { registry = "https://pypi.org/simple" } -sdist = { url = "https://files.pythonhosted.org/packages/8d/e8/c01bdf0d7c3659428c091fbd693177093639565bcbc86bc20098e6d37cc6/types_python_dateutil-2.9.0.20260518.tar.gz", hash = "sha256:51f02dc03b61c7f6a07df45797d4dfe8a1aa47f0b7db9ad89f6fd3a1a70e1b51", size = 17082, upload-time = "2026-05-18T06:05:24.508Z" } -wheels = [ - { url = "https://files.pythonhosted.org/packages/36/22/169273273ca34e9ab0ae2f387ba72ed7e09faaaf834da01d6b89c2bea71a/types_python_dateutil-2.9.0.20260518-py3-none-any.whl", hash = "sha256:d6a9c5bd0de61460c8fdef8ab2b400f956a1a1075cce08d4e2b4434e478c50b8", size = 18431, upload-time = "2026-05-18T06:05:23.641Z" }, -] - [[package]] name = "types-pyyaml" version = "6.0.12.20250915" From bdfa5aad13e15ae080c7a42239c64e251491c8a9 Mon Sep 17 00:00:00 2001 From: elmartinj Date: Fri, 26 Jun 2026 11:16:56 -0600 Subject: [PATCH 9/9] Add S3 and Modal support for CAISO benchmark --- Makefile | 16 +++++ src/data/caiso/aggregate/core.py | 24 +++++-- src/data/caiso/config.py | 3 +- src/data/caiso/extract/core.py | 18 ++++- src/data/caiso/modal_app.py | 53 +++++++++++++++ src/data/caiso/pipeline.py | 50 ++++++++++++++ src/evaluation/caiso/core.py | 6 +- src/forecast/caiso/core.py | 6 +- src/forecast/caiso/modal_app.py | 113 +++++++++++++++++++++++++++++++ 9 files changed, 278 insertions(+), 11 deletions(-) create mode 100644 src/data/caiso/modal_app.py create mode 100644 src/data/caiso/pipeline.py create mode 100644 src/forecast/caiso/modal_app.py diff --git a/Makefile b/Makefile index 0685a97..ec387d1 100644 --- a/Makefile +++ b/Makefile @@ -51,3 +51,19 @@ $(addprefix validate-evaluate-,$(EV_FREQUENCIES)): validate-evaluate-%: .PHONY: leaderboard leaderboard: # Build leaderboard parquet from all evaluation parquets $(MODAL) src.evaluation.gh_archive.modal_app::build_leaderboard + +## CAISO Data + +.PHONY: update-caiso-data +update-caiso-data: + $(MODAL) src.data.caiso.modal_app --start $(START) --end $(END) + +## CAISO Forecast/Evaluation + +.PHONY: update-caiso-forecast +update-caiso-forecast: + $(MODAL) src.forecast.caiso.modal_app::forecast --cutoff $(CUTOFF) + +.PHONY: update-caiso-evaluate +update-caiso-evaluate: + $(MODAL) src.forecast.caiso.modal_app::evaluate --cutoff $(CUTOFF) diff --git a/src/data/caiso/aggregate/core.py b/src/data/caiso/aggregate/core.py index 44ff275..187f75c 100644 --- a/src/data/caiso/aggregate/core.py +++ b/src/data/caiso/aggregate/core.py @@ -1,5 +1,9 @@ from __future__ import annotations +import shutil +import tempfile +from pathlib import Path + import pandas as pd from src.data.caiso.config import PROCESSED_CSV, PROCESSED_EVENTS_HOURLY_DIR @@ -8,9 +12,7 @@ OUTPUT_ROOT = PROCESSED_EVENTS_HOURLY_DIR -def build_hourly_partitions() -> int: - df = pd.read_csv(INPUT_CSV) - +def write_hourly_partitions(df: pd.DataFrame, output_root: Path = OUTPUT_ROOT) -> int: df["ds"] = pd.to_datetime(df["ds"], errors="coerce") df["y"] = pd.to_numeric(df["y"], errors="coerce") @@ -21,17 +23,19 @@ def build_hourly_partitions() -> int: df["month"] = df["ds"].dt.month df["day"] = df["ds"].dt.day - OUTPUT_ROOT.mkdir(parents=True, exist_ok=True) + output_root.mkdir(parents=True, exist_ok=True) n_written = 0 for (year, month, day), part in df.groupby(["year", "month", "day"], sort=True): part_dir = ( - OUTPUT_ROOT / f"year={year:04d}" / f"month={month:02d}" / f"day={day:02d}" + output_root / f"year={year:04d}" / f"month={month:02d}" / f"day={day:02d}" ) part_dir.mkdir(parents=True, exist_ok=True) out_path = part_dir / "series.parquet" - part[["unique_id", "ds", "y"]].to_parquet(out_path, index=False) + with tempfile.NamedTemporaryFile(suffix=".parquet") as tmp: + part[["unique_id", "ds", "y"]].to_parquet(tmp.name, index=False) + shutil.copyfile(tmp.name, out_path) print(f"Saved: {out_path}") n_written += 1 @@ -39,6 +43,14 @@ def build_hourly_partitions() -> int: return n_written +def build_hourly_partitions( + input_csv: Path = INPUT_CSV, + output_root: Path = OUTPUT_ROOT, +) -> int: + df = pd.read_csv(input_csv) + return write_hourly_partitions(df, output_root) + + def main() -> None: n_written = build_hourly_partitions() print(f"\nDone. Wrote {n_written} daily partitions.") diff --git a/src/data/caiso/config.py b/src/data/caiso/config.py index 45fcdbe..a83200a 100644 --- a/src/data/caiso/config.py +++ b/src/data/caiso/config.py @@ -1,12 +1,13 @@ from __future__ import annotations +import os from datetime import date from pathlib import Path from dateutil.relativedelta import relativedelta ROOT = Path(__file__).resolve().parents[3] -DATA_ROOT = ROOT / "data" / "caiso" +DATA_ROOT = Path(os.environ.get("CAISO_DATA_ROOT", ROOT / "data" / "caiso")) RAW_DIR = DATA_ROOT / "raw" TMP_DIR = DATA_ROOT / "tmp" diff --git a/src/data/caiso/extract/core.py b/src/data/caiso/extract/core.py index dcb87ee..180665e 100644 --- a/src/data/caiso/extract/core.py +++ b/src/data/caiso/extract/core.py @@ -53,17 +53,31 @@ def download_node( output.unlink(missing_ok=True) raise RuntimeError(f"CAISO returned a non-ZIP response for {node}") + with zipfile.ZipFile(output) as archive: + names = archive.namelist() + if names and names[0].endswith(".xml"): + with archive.open(names[0]) as file: + text = file.read().decode("utf-8", errors="replace") + if "" in text or "" in text: + output.unlink(missing_ok=True) + message = f"CAISO returned an error response for {node}: {text}" + raise RuntimeError(message) + print(f"Downloaded: {output}") return output -def extract_caiso(start: date, end: date) -> None: +def extract_caiso( + start: date, + end: date, + output_dir: Path = RAW_DIR, +) -> None: current = start while current < end: chunk_end = min(current + timedelta(days=31), end) - download_node("CAISO_3_NODES", current, chunk_end) + download_node("CAISO_3_NODES", current, chunk_end, output_dir=output_dir) current = chunk_end diff --git a/src/data/caiso/modal_app.py b/src/data/caiso/modal_app.py new file mode 100644 index 0000000..05bccca --- /dev/null +++ b/src/data/caiso/modal_app.py @@ -0,0 +1,53 @@ +from __future__ import annotations + +import modal + +CAISO_DATA_ROOT = "/s3-bucket/v0.1.0/caiso" + +app = modal.App(name="timecopilot-caiso-data") +image = ( + modal.Image.debian_slim(python_version="3.11") + .pip_install("uv") + .add_local_file("pyproject.toml", "/root/pyproject.toml", copy=True) + .add_local_file(".python-version", "/root/.python-version", copy=True) + .add_local_file("uv.lock", "/root/uv.lock", copy=True) + .workdir("/root") + .run_commands("uv pip install . --system --compile-bytecode") +) + +secret = modal.Secret.from_name( + "aws-secret", + required_keys=["AWS_ACCESS_KEY_ID", "AWS_SECRET_ACCESS_KEY"], +) + +volume = { + "/s3-bucket": modal.CloudBucketMount( + bucket_name="impermanent-benchmark", + secret=secret, + ) +} + + +@app.function( + image=image, + volumes=volume, + timeout=60 * 30, +) +def update_caiso_range(start: str, end: str) -> int: + import os + from datetime import date + + os.environ["CAISO_DATA_ROOT"] = CAISO_DATA_ROOT + + from src.data.caiso.pipeline import update_date_range + + return update_date_range( + start=date.fromisoformat(start), + end=date.fromisoformat(end), + ) + + +@app.local_entrypoint() +def update(start: str, end: str): + n_written = update_caiso_range.remote(start, end) + print(f"Done. Wrote {n_written} CAISO daily partitions.") diff --git a/src/data/caiso/pipeline.py b/src/data/caiso/pipeline.py new file mode 100644 index 0000000..4dd5188 --- /dev/null +++ b/src/data/caiso/pipeline.py @@ -0,0 +1,50 @@ +from __future__ import annotations + +import shutil +import tempfile +from datetime import date +from pathlib import Path + +from src.data.caiso.aggregate.core import build_hourly_partitions +from src.data.caiso.config import DATA_ROOT, RAW_DIR +from src.data.caiso.extract.core import extract_caiso +from src.data.caiso.transform.core import transform_caiso + + +def update_date_range( + start: date, + end: date, + data_root: Path = DATA_ROOT, + raw_dir: Path = RAW_DIR, +) -> int: + with tempfile.TemporaryDirectory() as tmp: + tmp_root = Path(tmp) + tmp_raw = tmp_root / "raw" + tmp_processed = tmp_root / "processed" + tmp_csv = tmp_processed / "caiso.csv" + + extract_caiso(start=start, end=end, output_dir=tmp_raw) + + raw_dir.mkdir(parents=True, exist_ok=True) + for source_zip in sorted(tmp_raw.glob("*.zip")): + shutil.copyfile(source_zip, raw_dir / source_zip.name) + + transform_caiso(raw_dir=tmp_raw, output_path=tmp_csv) + + output_root = data_root / "processed-events" / "hourly" + return build_hourly_partitions(input_csv=tmp_csv, output_root=output_root) + + +if __name__ == "__main__": + import argparse + + parser = argparse.ArgumentParser() + parser.add_argument("--start", required=True) + parser.add_argument("--end", required=True) + args = parser.parse_args() + + n_written = update_date_range( + start=date.fromisoformat(args.start), + end=date.fromisoformat(args.end), + ) + print(f"Done. Wrote {n_written} CAISO daily partitions.") diff --git a/src/evaluation/caiso/core.py b/src/evaluation/caiso/core.py index d75843f..a9eef57 100644 --- a/src/evaluation/caiso/core.py +++ b/src/evaluation/caiso/core.py @@ -1,6 +1,8 @@ from __future__ import annotations import argparse +import shutil +import tempfile from pathlib import Path import pandas as pd @@ -63,7 +65,9 @@ def run_evaluation( out_dir.mkdir(parents=True, exist_ok=True) out_path = out_dir / "metrics.parquet" - metrics.to_parquet(out_path, index=False) + with tempfile.NamedTemporaryFile(suffix=".parquet") as tmp: + metrics.to_parquet(tmp.name, index=False) + shutil.copyfile(tmp.name, out_path) return out_path diff --git a/src/forecast/caiso/core.py b/src/forecast/caiso/core.py index 73ba1fc..0fe94a0 100644 --- a/src/forecast/caiso/core.py +++ b/src/forecast/caiso/core.py @@ -1,6 +1,8 @@ from __future__ import annotations import argparse +import shutil +import tempfile from pathlib import Path import pandas as pd @@ -62,7 +64,9 @@ def run_forecast( out_dir.mkdir(parents=True, exist_ok=True) out_path = out_dir / "forecasts.parquet" - forecasts.to_parquet(out_path, index=False) + with tempfile.NamedTemporaryFile(suffix=".parquet") as tmp: + forecasts.to_parquet(tmp.name, index=False) + shutil.copyfile(tmp.name, out_path) return out_path diff --git a/src/forecast/caiso/modal_app.py b/src/forecast/caiso/modal_app.py new file mode 100644 index 0000000..eb8207f --- /dev/null +++ b/src/forecast/caiso/modal_app.py @@ -0,0 +1,113 @@ +from __future__ import annotations + +import modal + +CAISO_DATA_ROOT = "/s3-bucket/v0.1.0/caiso" + +CPU_MODELS = ( + "seasonal_naive", + "historic_average", + "auto_ets", + "auto_ces", + "dynamic_optimized_theta", +) + +app = modal.App(name="timecopilot-caiso-forecast") +image = ( + modal.Image.debian_slim(python_version="3.11") + .pip_install("uv") + .add_local_file("pyproject.toml", "/root/pyproject.toml", copy=True) + .add_local_file(".python-version", "/root/.python-version", copy=True) + .add_local_file("uv.lock", "/root/uv.lock", copy=True) + .workdir("/root") + .run_commands("uv pip install . --system --compile-bytecode") +) + +secret = modal.Secret.from_name( + "aws-secret", + required_keys=["AWS_ACCESS_KEY_ID", "AWS_SECRET_ACCESS_KEY"], +) + +volume = { + "/s3-bucket": modal.CloudBucketMount( + bucket_name="impermanent-benchmark", + secret=secret, + ) +} + + +@app.function( + image=image, + volumes=volume, + timeout=60 * 30, +) +def run_forecast_model(cutoff: str, model: str) -> tuple[str, str, str | None]: + import os + + os.environ["CAISO_DATA_ROOT"] = CAISO_DATA_ROOT + + from src.forecast.caiso.core import run_forecast + + try: + forecast_path = run_forecast( + cutoff=cutoff, + model=model, + h=24, + max_window_size=48, + ) + return model, str(forecast_path), None + except Exception as exc: + return model, "", repr(exc) + + +@app.function( + image=image, + volumes=volume, + timeout=60 * 30, +) +def run_evaluation_model(cutoff: str, model: str) -> tuple[str, str, str | None]: + import os + + os.environ["CAISO_DATA_ROOT"] = CAISO_DATA_ROOT + + from src.evaluation.caiso.core import run_evaluation + + try: + metrics_path = run_evaluation( + cutoff=cutoff, + model=model, + h=24, + max_window_size=48, + ) + return model, str(metrics_path), None + except Exception as exc: + return model, "", repr(exc) + + +def _print_results(results: list[tuple[str, str, str | None]]) -> None: + failures = [result for result in results if result[2] is not None] + + for model, path, error in results: + if error: + print(f"FAILED {model}: {error}") + else: + print(f"OK {model}: {path}") + + if failures: + raise RuntimeError(f"{len(failures)} CAISO model runs failed") + + +@app.local_entrypoint() +def forecast(cutoff: str): + results = list( + run_forecast_model.starmap([(cutoff, model) for model in CPU_MODELS]) + ) + _print_results(results) + + +@app.local_entrypoint() +def evaluate(cutoff: str): + results = list( + run_evaluation_model.starmap([(cutoff, model) for model in CPU_MODELS]) + ) + _print_results(results)