Add ML prediction pipeline — LightGBM, calibration fix, ensemble disagreement
Tier 1 ML enhancements: - Feature engineering (37 features across 5 groups: thermal, dynamic, moisture, temporal, interaction) from NWP model output - 7 LightGBM probability models for rain/temp/wind thresholds - Temperature-scaled probabilities to prevent overconfidence on bootstrap data - MLPredictor: unified inference pipeline replacing heuristic sigmoids - Ensemble disagreement signals (composite spread → edge amplification) - Fixed calibration loop: update_calibration() now functional (EMA of errors) - record_outcome() wired for post-resolution feedback - Nautilus strategy updated: ML predictions take priority, heuristics as fallback - Historical backtest engine with Sharpe/ROI/max-DD simulation - Bootstrap training data generator from HK climate normals Run: python ml/train.py && python ml/backtest.py --edge 50
This commit is contained in:
+385
@@ -0,0 +1,385 @@
|
||||
"""ML-powered signal generator for HK weather prediction markets.
|
||||
|
||||
Replaces heuristic sigmoids with LightGBM probability models.
|
||||
Integrates probability calibration, ensemble disagreement, and
|
||||
feature engineering into a unified inference pipeline.
|
||||
|
||||
Usage:
|
||||
predictor = MLPredictor()
|
||||
probs = predictor.predict("tomorrow") # All targets for tomorrow
|
||||
signal = predictor.generate_signal("temp_gt_30c_24h", market_price=0.45)
|
||||
"""
|
||||
|
||||
import sys
|
||||
from pathlib import Path
|
||||
from datetime import datetime, timedelta
|
||||
from typing import Dict, Optional, Tuple, List
|
||||
|
||||
import numpy as np
|
||||
import pandas as pd
|
||||
|
||||
sys.path.insert(0, str(Path(__file__).parent.parent))
|
||||
|
||||
from ml.features import FeatureEngine
|
||||
from ml.model import ModelEnsemble, TARGET_DEFINITIONS
|
||||
from weather.openmeteo_client import OpenMeteoClient
|
||||
from weather.hko_client import HKOClient
|
||||
from strategy.calibrator import ProbabilityCalibrator
|
||||
from strategy.kelly import KellyCriterion
|
||||
from config import HK_COORDS, MIN_EDGE_BPS, KELLY_FRACTION, MAX_POSITION_USDC
|
||||
|
||||
|
||||
class MLPredictor:
|
||||
"""
|
||||
ML-based weather probability predictor for Polymarket trading.
|
||||
|
||||
Combines:
|
||||
1. Feature engineering from NWP model output
|
||||
2. Trained LightGBM probability models
|
||||
3. Platt scaling calibration on historical outcomes
|
||||
4. Ensemble disagreement as edge amplifier
|
||||
5. Kelly criterion position sizing
|
||||
"""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
bankroll_usdc: float = 1000.0,
|
||||
min_edge_bps: float = MIN_EDGE_BPS,
|
||||
kelly_fraction: float = KELLY_FRACTION,
|
||||
):
|
||||
self.engine = FeatureEngine()
|
||||
self.ensemble = ModelEnsemble()
|
||||
self.calibrator = ProbabilityCalibrator()
|
||||
self.kelly = KellyCriterion(bankroll_usdc=bankroll_usdc, fraction=kelly_fraction)
|
||||
self.openmeteo = OpenMeteoClient()
|
||||
self.hko = HKOClient()
|
||||
|
||||
self.min_edge_bps = min_edge_bps
|
||||
|
||||
# Ensemble disagreement tracking
|
||||
self._last_forecast: Optional[pd.DataFrame] = None
|
||||
self._last_hourly: Optional[pd.DataFrame] = None
|
||||
self._last_features: Optional[np.ndarray] = None
|
||||
self._last_predictions: Optional[Dict[str, float]] = None
|
||||
self._ensemble_spread: Optional[Dict[str, float]] = None
|
||||
|
||||
# Load trained models
|
||||
self.models_loaded = self._load_models()
|
||||
|
||||
def _load_models(self) -> bool:
|
||||
"""Load trained models if available."""
|
||||
try:
|
||||
self.ensemble.load_all()
|
||||
return len(self.ensemble.models) > 0
|
||||
except Exception as e:
|
||||
print(f"ML models not loaded (train first): {e}")
|
||||
return False
|
||||
|
||||
def fetch_and_predict(self, target_date: Optional[str] = None) -> Dict[str, float]:
|
||||
"""
|
||||
Fetch latest forecast and predict all targets.
|
||||
|
||||
Returns dict of {target_name: calibrated_probability_0_100}
|
||||
"""
|
||||
# Fetch data
|
||||
daily = self.openmeteo.get_forecast(lead_days=7)
|
||||
if daily is None:
|
||||
print("MLPredictor: No forecast data available")
|
||||
return self._fallback_predictions()
|
||||
|
||||
# Get hourly data from the client's internal cache
|
||||
hourly = getattr(self.openmeteo, '_last_hourly', None)
|
||||
|
||||
self._last_forecast = daily
|
||||
self._last_hourly = hourly
|
||||
|
||||
# Feature engineering
|
||||
X = self.engine.transform(daily, hourly)
|
||||
self._last_features = X
|
||||
|
||||
# If models loaded, use ML predictions
|
||||
if self.models_loaded and len(self.ensemble.models) > 0:
|
||||
predictions = {}
|
||||
for day_idx in range(min(len(daily), 7)):
|
||||
date_str = daily.index[day_idx].strftime("%Y-%m-%d") if hasattr(daily.index[day_idx], 'strftime') else str(daily.index[day_idx])
|
||||
if target_date and date_str != target_date and day_idx > 1:
|
||||
continue
|
||||
|
||||
# Predict all loaded targets for this day
|
||||
X_day = X[day_idx].reshape(1, -1)
|
||||
day_probs = self.ensemble.predict_all(X_day)
|
||||
|
||||
# Calibrate
|
||||
calibrated = {}
|
||||
for target, raw_prob in day_probs.items():
|
||||
calibrated[target] = self.calibrator.calibrate(target, raw_prob)
|
||||
|
||||
# Compute ensemble disagreement (multi-level)
|
||||
spread = self._compute_multimodel_spread(daily, hourly, day_idx)
|
||||
self._ensemble_spread = spread
|
||||
|
||||
# Amplify edge based on spread
|
||||
for target in calibrated:
|
||||
adjusted = self._adjust_with_spread(
|
||||
calibrated[target], target, spread
|
||||
)
|
||||
calibrated[target] = adjusted
|
||||
|
||||
# Store date-str tagged predictions
|
||||
if day_idx <= 2: # Keep near-term predictions
|
||||
for target, prob in calibrated.items():
|
||||
predictions[f"{target}_{date_str}"] = prob
|
||||
|
||||
# Also store as raw target key (overwrites with latest)
|
||||
if day_idx == 1: # Tomorrow
|
||||
for target, prob in calibrated.items():
|
||||
predictions[target] = prob
|
||||
|
||||
self._last_predictions = predictions
|
||||
return predictions
|
||||
|
||||
# Fallback: use heuristic predictions
|
||||
return self._fallback_predictions()
|
||||
|
||||
def _fallback_predictions(self) -> Dict[str, float]:
|
||||
"""Fallback heuristic predictions when no ML models loaded."""
|
||||
if self._last_forecast is None or len(self._last_forecast) == 0:
|
||||
return {}
|
||||
|
||||
d1 = self._last_forecast.iloc[min(1, len(self._last_forecast) - 1)]
|
||||
predictions = {}
|
||||
|
||||
for target, tdef in TARGET_DEFINITIONS.items():
|
||||
var = tdef["variable"]
|
||||
threshold = tdef["threshold"]
|
||||
if var in self._last_forecast.columns:
|
||||
val = float(d1.get(var, 0))
|
||||
prob = self._heuristic_prob(val, threshold, target)
|
||||
predictions[target] = prob
|
||||
|
||||
return predictions
|
||||
|
||||
def _heuristic_prob(self, value: float, threshold: float, target: str) -> float:
|
||||
"""Fallback heuristic: sigmoid-based probability."""
|
||||
if "temp" in target:
|
||||
# Temperature: wider sigmoid, calibrated to HK summer
|
||||
excess = value - threshold
|
||||
return float(np.clip(50 + excess * 15, 3, 97))
|
||||
elif "rain" in target:
|
||||
# Rain probabilities from Open-Meteo directly
|
||||
if threshold == 0:
|
||||
return float(np.clip(value, 0.5, 99.5))
|
||||
else:
|
||||
return float(np.clip(value * 0.8 if threshold < 10 else value * 0.5, 1, 95))
|
||||
elif "wind" in target:
|
||||
excess = value - threshold
|
||||
return float(np.clip(50 + excess * 5, 3, 97))
|
||||
return 50.0
|
||||
|
||||
def _compute_multimodel_spread(
|
||||
self, daily: pd.DataFrame, hourly: pd.DataFrame, day_idx: int
|
||||
) -> Dict[str, float]:
|
||||
"""Compute ensemble disagreement metrics across model outputs.
|
||||
|
||||
When multiple model outputs are available (GFS, ECMWF, WeatherNext),
|
||||
disagreement signifies uncertainty that the market may misprice.
|
||||
"""
|
||||
spread = {}
|
||||
|
||||
# 1. Inter-day variability (persistence disagreement)
|
||||
if day_idx > 0 and len(daily) > day_idx:
|
||||
d0 = daily.iloc[day_idx - 1]
|
||||
d1 = daily.iloc[day_idx]
|
||||
|
||||
spread["t2m_max_day_change"] = abs(
|
||||
float(d1.get("temperature_2m_max", 0)) -
|
||||
float(d0.get("temperature_2m_max", 0))
|
||||
)
|
||||
spread["precip_prob_day_change"] = abs(
|
||||
float(d1.get("precipitation_probability_max", 0)) -
|
||||
float(d0.get("precipitation_probability_max", 0))
|
||||
)
|
||||
spread["pressure_day_change"] = abs(
|
||||
float(hourly["surface_pressure"].mean() if "surface_pressure" in hourly else 1013) -
|
||||
float(hourly["surface_pressure"].iloc[max(0, day_idx * 24 - 24)] if "surface_pressure" in hourly else 1013)
|
||||
) if hourly is not None and len(hourly) > 0 else 0.0
|
||||
|
||||
# 2. Wind direction variability (storm potential indicator)
|
||||
if hourly is not None and len(hourly) > 0 and "wind_direction_10m" in hourly:
|
||||
day_hourly = hourly.iloc[day_idx * 24:(day_idx + 1) * 24] if len(hourly) > (day_idx + 1) * 24 else hourly
|
||||
if len(day_hourly) > 0:
|
||||
wd = day_hourly["wind_direction_10m"].values
|
||||
spread["wind_dir_variance"] = float(np.var(wd)) if len(wd) > 1 else 0.0
|
||||
|
||||
# 3. Cloud structure complexity (convection proxy)
|
||||
if hourly is not None and len(hourly) > 0:
|
||||
for level in ["cloud_cover_low", "cloud_cover_mid", "cloud_cover_high"]:
|
||||
if level in hourly.columns:
|
||||
day_hourly = hourly.iloc[day_idx * 24:(day_idx + 1) * 24] if len(hourly) > (day_idx + 1) * 24 else hourly
|
||||
if len(day_hourly) > 0:
|
||||
spread[f"{level}_std"] = float(day_hourly[level].std())
|
||||
|
||||
# 4. Compute composite spread score (0-1)
|
||||
indicators = []
|
||||
for k, v in spread.items():
|
||||
if "t2m" in k:
|
||||
indicators.append(np.clip(v / 5.0, 0, 1)) # 5°C change = full signal
|
||||
elif "precip" in k:
|
||||
indicators.append(np.clip(v / 50.0, 0, 1)) # 50% change = full signal
|
||||
elif "pressure" in k:
|
||||
indicators.append(np.clip(v / 10.0, 0, 1)) # 10 hPa = full signal
|
||||
elif "variance" in k:
|
||||
indicators.append(np.clip(v / 5000.0, 0, 1))
|
||||
elif "_std" in k:
|
||||
indicators.append(np.clip(v / 30.0, 0, 1))
|
||||
|
||||
spread["composite_spread"] = float(np.mean(indicators)) if indicators else 0.0
|
||||
return spread
|
||||
|
||||
def _adjust_with_spread(
|
||||
self, probability: float, target: str, spread: Dict[str, float]
|
||||
) -> float:
|
||||
"""
|
||||
Adjust probability based on ensemble disagreement.
|
||||
|
||||
When models disagree → higher uncertainty → wider confidence interval.
|
||||
In prediction markets, this often means the market price is LESS accurate
|
||||
(traders anchor on the wrong model or over-weight consensus).
|
||||
|
||||
We amplify our edge when spread is high: push our probability
|
||||
further from 50% to reflect our confidence in the direction.
|
||||
"""
|
||||
composite = spread.get("composite_spread", 0.0)
|
||||
|
||||
if composite < 0.1:
|
||||
return probability
|
||||
|
||||
# Direction: is our prediction above or below 50%?
|
||||
direction = 1 if probability > 50 else -1
|
||||
|
||||
# Amplification: move probability up to spread * 20 bps further from 50
|
||||
# High spread = more uncertainty = wider market spread = more edge
|
||||
amplification = min(composite * 20, 20) # Cap at 20 percentage points
|
||||
|
||||
adjusted = probability + direction * amplification
|
||||
return float(np.clip(adjusted, 0.5, 99.5))
|
||||
|
||||
def generate_signal(
|
||||
self,
|
||||
target: str,
|
||||
market_probability: float,
|
||||
outcome: str = "YES",
|
||||
) -> Dict:
|
||||
"""
|
||||
Generate a trading signal for a specific market.
|
||||
|
||||
Parameters
|
||||
----------
|
||||
target : str
|
||||
Target name (e.g., 'temp_gt_30c_24h')
|
||||
market_probability : float
|
||||
Market-implied probability of the outcome (0-100)
|
||||
outcome : str
|
||||
Which outcome to bet on ('YES' or 'NO')
|
||||
|
||||
Returns
|
||||
-------
|
||||
Dict with model_prob, market_prob, edge_bps, kelly_size, side
|
||||
"""
|
||||
if not self._last_predictions:
|
||||
self.fetch_and_predict()
|
||||
|
||||
model_prob = (self._last_predictions or {}).get(target, 50.0)
|
||||
cal_prob = self.calibrator.calibrate(target, model_prob)
|
||||
|
||||
edge_bps = (cal_prob - market_probability)
|
||||
|
||||
if abs(edge_bps) < self.min_edge_bps:
|
||||
return {
|
||||
"signal": "pass",
|
||||
"model_prob": cal_prob,
|
||||
"market_prob": market_probability,
|
||||
"edge_bps": edge_bps,
|
||||
"size_usdc": 0.0,
|
||||
}
|
||||
|
||||
side = "buy_yes" if edge_bps > 0 else "buy_no"
|
||||
kelly_result = self.kelly.size_bet(
|
||||
our_probability=cal_prob,
|
||||
market_probability=market_probability,
|
||||
side=side,
|
||||
)
|
||||
|
||||
return {
|
||||
"signal": side,
|
||||
"model_prob": cal_prob,
|
||||
"market_prob": market_probability,
|
||||
"edge_bps": edge_bps,
|
||||
"size_usdc": kelly_result.size_usdc if kelly_result.kelly_active else 0.0,
|
||||
"kelly_fraction": kelly_result.fractional_kelly,
|
||||
"ensemble_spread": (self._ensemble_spread or {}).get("composite_spread", 0.0),
|
||||
}
|
||||
|
||||
def record_outcome(self, target: str, predicted_prob: float, actual: bool):
|
||||
"""Record resolved market outcome for calibration."""
|
||||
self.calibrator.record_outcome(
|
||||
date=datetime.now().strftime("%Y-%m-%d"),
|
||||
variable=target,
|
||||
predicted_probability=predicted_prob,
|
||||
actual_outcome=actual,
|
||||
)
|
||||
|
||||
def get_top_signals(
|
||||
self, markets: List[Dict], default_market_prob: float = 50.0
|
||||
) -> List[Dict]:
|
||||
"""Scan a list of market definitions and generate ranked signals."""
|
||||
predictions = self._last_predictions or self.fetch_and_predict()
|
||||
|
||||
signals = []
|
||||
for market in markets:
|
||||
target = market.get("target", "")
|
||||
if target not in predictions:
|
||||
continue
|
||||
|
||||
market_prob = market.get("market_probability", default_market_prob)
|
||||
signal = self.generate_signal(target, market_prob)
|
||||
|
||||
if signal["signal"] != "pass":
|
||||
signals.append({
|
||||
**signal,
|
||||
"target": target,
|
||||
"description": TARGET_DEFINITIONS.get(target, {}).get("description", ""),
|
||||
"question": market.get("question", ""),
|
||||
"condition_id": market.get("condition_id", ""),
|
||||
})
|
||||
|
||||
signals.sort(key=lambda s: abs(s["edge_bps"]), reverse=True)
|
||||
return signals
|
||||
|
||||
def summary(self) -> str:
|
||||
"""Human-readable summary of current predictions."""
|
||||
predictions = self._last_predictions or {}
|
||||
|
||||
lines = []
|
||||
lines.append(f"\n=== ML Weather Predictions ({datetime.now():%Y-%m-%d %H:%M}) ===")
|
||||
lines.append(f" Models loaded: {len(self.ensemble.models)}")
|
||||
lines.append(f" Calibration records: {self.calibrator.get_calibration_stats(list(predictions.keys())[0] if predictions else 'temp_gt_30c_24h').get('n_observations', 0)}")
|
||||
|
||||
if not predictions:
|
||||
lines.append(" No predictions available.")
|
||||
return "\n".join(lines)
|
||||
|
||||
lines.append(f"\n Tomorrow's targets:")
|
||||
for target in TARGET_DEFINITIONS:
|
||||
if target in predictions:
|
||||
tdef = TARGET_DEFINITIONS[target]
|
||||
prob = predictions[target]
|
||||
lines.append(f" {tdef['description']}: {prob:.1f}%")
|
||||
|
||||
spread = self._ensemble_spread or {}
|
||||
if spread.get("composite_spread", 0) > 0.1:
|
||||
lines.append(f"\n Ensemble disagreement: {spread['composite_spread']:.2f} (amplified edge)")
|
||||
if spread.get("t2m_max_day_change", 0) > 0:
|
||||
lines.append(f" ΔTmax: {spread.get('t2m_max_day_change', 0):.1f}°C")
|
||||
|
||||
return "\n".join(lines)
|
||||
Reference in New Issue
Block a user