03f9ea2129
Tier 2 enhancements: - SpatialWeatherClient: multi-station Open-Meteo fetcher for all HK locations Extracts urban heat island delta, coastal-inland gradients, wind convergence, precipitation spatial heterogeneity, composite instability index - TyphoonModel: data-driven signal probability for T1/T3/T8/T10 Climatological base rates + conditional transition probabilities Currently active T1 signal → 25% T3/24h, 10% T8/72h, 22% T8/120h ENSO modulation, active storm proximity boost, month-specific seasonality - ERA5 download/process pipeline via CDS API Downloads hourly reanalysis for HK region, processes to daily training format Output schema matches Open-Meteo for seamless feature compatibility - PortfolioKelly: correlation-aware simultaneous Kelly sizing Covariance matrix from historical outcome correlations Prevents over-betting on correlated rain/temp/wind markets Σ⁻¹ μ vector formulation, regularized inversion, independent fallback - MLPredictor updated: integrates spatial + typhoon + portfolio Kelly record_outcome feeds both calibration AND portfolio correlation matrix
317 lines
10 KiB
Python
317 lines
10 KiB
Python
"""ERA5 historical data download pipeline for ML training.
|
|
|
|
Downloads ERA5 reanalysis data via the CDS API (Climate Data Store),
|
|
processes it into training data for the HK weather prediction models.
|
|
|
|
ERA5 provides:
|
|
- Global 0.25° hourly reanalysis from 1940-present
|
|
- All atmospheric variables needed for our NWP features
|
|
- Ground truth for training (what actually happened)
|
|
|
|
Requirements:
|
|
- CDS API key (free registration at https://cds.climate.copernicus.eu)
|
|
- Set CDSAPI_KEY and CDSAPI_URL in .env
|
|
|
|
Usage:
|
|
python ml/era5.py --download # Download raw data
|
|
python ml/era5.py --process # Process into training format
|
|
python ml/era5.py --download --process # Full pipeline
|
|
"""
|
|
|
|
import argparse
|
|
import json
|
|
import os
|
|
import sys
|
|
from datetime import datetime, timedelta
|
|
from pathlib import Path
|
|
from typing import List, Optional, Dict, Tuple
|
|
|
|
import numpy as np
|
|
import pandas as pd
|
|
|
|
sys.path.insert(0, str(Path(__file__).parent.parent))
|
|
|
|
from config import DATA_DIR, HK_BBOX, HK_COORDS, CDSAPI_KEY, CDSAPI_URL
|
|
|
|
ERA5_DIR = Path(DATA_DIR) / "era5"
|
|
ERA5_DIR.mkdir(parents=True, exist_ok=True)
|
|
|
|
# Variables to download (daily + sub-daily aggregates)
|
|
ERA5_VARIABLES = {
|
|
# 2m surface variables (daily)
|
|
"2m_temperature": {"name": "t2m", "agg": "mean_max_min"},
|
|
"2m_dewpoint_temperature": {"name": "d2m", "agg": "mean"},
|
|
"surface_pressure": {"name": "sp", "agg": "mean"},
|
|
"mean_sea_level_pressure": {"name": "msl", "agg": "mean"},
|
|
"10m_u_component_of_wind": {"name": "u10", "agg": "max"},
|
|
"10m_v_component_of_wind": {"name": "v10", "agg": "max"},
|
|
"100m_u_component_of_wind": {"name": "u100", "agg": "max"},
|
|
"100m_v_component_of_wind": {"name": "v100", "agg": "max"},
|
|
"2m_relative_humidity": {"name": "rh", "agg": "mean_min"},
|
|
"total_precipitation": {"name": "tp", "agg": "sum"},
|
|
"total_cloud_cover": {"name": "tcc", "agg": "mean"},
|
|
"surface_solar_radiation_downwards": {"name": "ssrd", "agg": "mean"},
|
|
}
|
|
|
|
|
|
def check_cds_credentials() -> bool:
|
|
"""Verify CDS API credentials are configured."""
|
|
key = CDSAPI_KEY or os.getenv("CDSAPI_KEY")
|
|
url = CDSAPI_URL or os.getenv("CDSAPI_URL") or "https://cds.climate.copernicus.eu/api"
|
|
if not key:
|
|
print("CDSAPI_KEY not set. Register at https://cds.climate.copernicus.eu")
|
|
print("Then set in .env: CDSAPI_KEY=your-key")
|
|
return False
|
|
# Write credentials file for cdsapi
|
|
creds_path = os.path.expanduser("~/.cdsapirc")
|
|
with open(creds_path, "w") as f:
|
|
f.write(f"url: {url}\nkey: {key}\n")
|
|
return True
|
|
|
|
|
|
def download_era5(
|
|
start_year: int = 2015,
|
|
end_year: int = 2025,
|
|
months: Optional[List[int]] = None,
|
|
area: Optional[List[float]] = None,
|
|
):
|
|
"""
|
|
Download ERA5 hourly data for Hong Kong region.
|
|
|
|
Parameters
|
|
----------
|
|
start_year, end_year : int
|
|
Year range for download
|
|
months : list[int], optional
|
|
Months to download (default: all)
|
|
area : list[float], optional
|
|
[N, W, S, E] bounding box (default: HK region)
|
|
"""
|
|
if not check_cds_credentials():
|
|
return False
|
|
|
|
try:
|
|
import cdsapi
|
|
except ImportError:
|
|
os.system(f"{sys.executable} -m pip install cdsapi")
|
|
import cdsapi
|
|
|
|
if months is None:
|
|
months = list(range(1, 13))
|
|
if area is None:
|
|
# HK + buffer: [N, W, S, E]
|
|
area = [
|
|
HK_BBOX["lat_max"] + 1,
|
|
HK_BBOX["lon_min"] - 1,
|
|
HK_BBOX["lat_min"] - 1,
|
|
HK_BBOX["lon_max"] + 1,
|
|
]
|
|
|
|
client = cdsapi.Client()
|
|
|
|
variable_names = list(ERA5_VARIABLES.keys())
|
|
|
|
for year in range(start_year, end_year + 1):
|
|
for month in months:
|
|
output_path = ERA5_DIR / f"era5_hk_{year}_{month:02d}.nc"
|
|
|
|
if output_path.exists():
|
|
print(f" Skipping {year}-{month:02d} (already exists)")
|
|
continue
|
|
|
|
print(f" Downloading {year}-{month:02d}...")
|
|
|
|
try:
|
|
client.retrieve(
|
|
"reanalysis-era5-single-levels",
|
|
{
|
|
"product_type": "reanalysis",
|
|
"format": "netcdf",
|
|
"variable": variable_names,
|
|
"year": str(year),
|
|
"month": f"{month:02d}",
|
|
"day": [f"{d:02d}" for d in range(1, 32)],
|
|
"time": [f"{h:02d}:00" for h in range(0, 24, 6)],
|
|
"area": area,
|
|
},
|
|
str(output_path),
|
|
)
|
|
print(f" Saved: {output_path}")
|
|
|
|
except Exception as e:
|
|
print(f" Failed: {e}")
|
|
|
|
return True
|
|
|
|
|
|
def process_era5_to_training(
|
|
years: Optional[List[int]] = None,
|
|
output_path: Optional[str] = None,
|
|
) -> Optional[pd.DataFrame]:
|
|
"""
|
|
Process downloaded ERA5 NetCDF files into training DataFrames.
|
|
|
|
Extracts daily statistics for the HK region and formats them
|
|
to match the Open-Meteo output structure for seamless feature
|
|
engineering compatibility.
|
|
"""
|
|
nc_files = sorted(ERA5_DIR.glob("era5_hk_*.nc"))
|
|
|
|
if not nc_files:
|
|
print(f"No ERA5 files found in {ERA5_DIR}")
|
|
print("Run: python ml/era5.py --download")
|
|
return None
|
|
|
|
print(f"Processing {len(nc_files)} ERA5 files...")
|
|
|
|
all_days = []
|
|
total_processed = 0
|
|
|
|
for nc_file in nc_files:
|
|
try:
|
|
import xarray as xr
|
|
ds = xr.open_dataset(nc_file, engine="h5netcdf")
|
|
|
|
# Extract HK region (single grid cell at 0.25°)
|
|
lat_center = HK_COORDS["hko_headquarters"][0]
|
|
lon_center = HK_COORDS["hko_headquarters"][1]
|
|
|
|
# Find nearest grid point
|
|
lat_idx = np.abs(ds.latitude.values - lat_center).argmin()
|
|
lon_idx = np.abs(ds.longitude.values - lon_center).argmin()
|
|
|
|
# Resample to daily
|
|
ds_hk = ds.isel(latitude=lat_idx, longitude=lon_idx)
|
|
|
|
# Convert to daily statistics
|
|
daily_ds = ds_hk.resample(time="1D").agg({
|
|
"t2m": ["max", "min", "mean"],
|
|
"d2m": "mean",
|
|
"sp": "mean",
|
|
"msl": "mean",
|
|
"tp": "sum",
|
|
"tcc": "mean",
|
|
"u10": "max",
|
|
"v10": "max",
|
|
"ssrd": "mean",
|
|
"r": "mean",
|
|
})
|
|
|
|
# Build DataFrame
|
|
df = daily_ds.to_dataframe()
|
|
df = df.reset_index()
|
|
df.columns = ['_'.join(col).strip('_') for col in df.columns]
|
|
|
|
# Rename to match Open-Meteo convention
|
|
rename_map = {
|
|
"t2m_max": "temperature_2m_max",
|
|
"t2m_min": "temperature_2m_min",
|
|
"t2m_mean": "temperature_2m_mean",
|
|
"d2m_mean": "dew_point_2m",
|
|
"sp_mean": "surface_pressure",
|
|
"msl_mean": "msl_pressure",
|
|
"tp_sum": "precipitation_sum",
|
|
"tcc_mean": "total_cloud_cover",
|
|
"u10_max": "wind_u_max",
|
|
"v10_max": "wind_v_max",
|
|
"ssrd_mean": "shortwave_radiation_sum",
|
|
"r_mean": "relative_humidity_2m",
|
|
}
|
|
|
|
df = df.rename(columns={k: v for k, v in rename_map.items() if k in df.columns})
|
|
|
|
# Derived columns
|
|
if "wind_u_max" in df.columns and "wind_v_max" in df.columns:
|
|
df["wind_speed_10m_max"] = np.sqrt(df["wind_u_max"]**2 + df["wind_v_max"]**2)
|
|
|
|
if "temperature_2m_max" in df.columns:
|
|
df["temperature_2m_max"] -= 273.15
|
|
if "temperature_2m_min" in df.columns:
|
|
df["temperature_2m_min"] -= 273.15
|
|
if "temperature_2m_mean" in df.columns:
|
|
df["temperature_2m_mean"] -= 273.15
|
|
if "dew_point_2m" in df.columns:
|
|
df["dew_point_2m"] -= 273.15
|
|
if "surface_pressure" in df.columns:
|
|
df["surface_pressure"] /= 100.0
|
|
|
|
if "precipitation_sum" in df.columns:
|
|
df["precipitation_sum"] *= 1000.0
|
|
|
|
# Precipitation probability (any rain?)
|
|
if "precipitation_sum" in df.columns:
|
|
df["precipitation_probability_max"] = (
|
|
df["precipitation_sum"] > 0.1
|
|
).astype(int) * 100.0
|
|
|
|
# Wind gusts (approximate: 1.4x mean max)
|
|
if "wind_speed_10m_max" in df.columns:
|
|
df["wind_gusts_10m_max"] = df["wind_speed_10m_max"] * 1.4
|
|
|
|
all_days.append(df)
|
|
total_processed += len(df)
|
|
|
|
except Exception as e:
|
|
print(f" Error processing {nc_file.name}: {e}")
|
|
|
|
if not all_days:
|
|
return None
|
|
|
|
combined = pd.concat(all_days, ignore_index=True)
|
|
|
|
if "time" in combined.columns:
|
|
combined["time"] = pd.to_datetime(combined["time"])
|
|
combined = combined.set_index("time").sort_index()
|
|
|
|
# Save
|
|
if output_path is None:
|
|
output_path = ERA5_DIR / "hk_era5_training.parquet"
|
|
|
|
combined.to_parquet(output_path)
|
|
print(f"\nSaved {total_processed} daily records to {output_path}")
|
|
print(f"Date range: {combined.index.min()} to {combined.index.max()}")
|
|
|
|
return combined
|
|
|
|
|
|
def load_training_data(path: Optional[str] = None) -> Optional[pd.DataFrame]:
|
|
"""Load processed ERA5 training data."""
|
|
if path is None:
|
|
# Find existing parquet files
|
|
pq_files = list(ERA5_DIR.glob("hk_era5_training*.parquet"))
|
|
if not pq_files:
|
|
print("No training data found. Run: python ml/era5.py --download --process")
|
|
return None
|
|
path = str(pq_files[0])
|
|
|
|
df = pd.read_parquet(path)
|
|
print(f"Loaded {len(df)} days from {path}")
|
|
return df
|
|
|
|
|
|
def main():
|
|
parser = argparse.ArgumentParser(description="ERA5 data pipeline")
|
|
parser.add_argument("--download", action="store_true", help="Download ERA5 data")
|
|
parser.add_argument("--process", action="store_true", help="Process into training format")
|
|
parser.add_argument("--start-year", type=int, default=2015)
|
|
parser.add_argument("--end-year", type=int, default=2024)
|
|
args = parser.parse_args()
|
|
|
|
if args.download:
|
|
print(f"Downloading ERA5 data ({args.start_year}-{args.end_year})...")
|
|
download_era5(start_year=args.start_year, end_year=args.end_year)
|
|
|
|
if args.process:
|
|
print("Processing ERA5 into training data...")
|
|
process_era5_to_training()
|
|
|
|
if not args.download and not args.process:
|
|
# Default: try to load existing data
|
|
df = load_training_data()
|
|
if df is not None:
|
|
print(df.describe())
|
|
|
|
|
|
if __name__ == "__main__":
|
|
main()
|