Files
ramseshk 03f9ea2129 Add spatial features, typhoon model, ERA5 pipeline, portfolio Kelly
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
2026-08-10 17:56:02 +08:00

269 lines
9.4 KiB
Python

"""Correlation-aware portfolio Kelly sizing.
When trading multiple weather markets simultaneously, outcomes are correlated:
- Rain ↔ Temperature (rain suppresses temperature)
- Wind ↔ Typhoon (wind is a prerequisite for T3/T8)
- Rain tomorrow ↔ Rain day-after-tomorrow (persistence)
Simple Kelly treats each market independently, over-betting when outcomes
are correlated. This module implements simultaneous Kelly using a covariance
matrix of historical outcomes.
f* = Σ⁻¹ μ (vector form, where Σ is the covariance of returns, μ is edge vector)
Reference: MacLean, Thorp, Ziemba (2011) "The Kelly Capital Growth Investment Criterion"
"""
import json
from pathlib import Path
from typing import Dict, List, Optional, Tuple
import numpy as np
from config import DATA_DIR
COV_DIR = Path(DATA_DIR) / "covariance"
COV_DIR.mkdir(parents=True, exist_ok=True)
class PortfolioKelly:
"""
Simultaneous Kelly sizing for correlated prediction market bets.
Uses historical outcome correlation matrix to adjust individual
Kelly fractions, preventing over-betting on correlated markets.
Usage:
pk = PortfolioKelly()
pk.update_correlation({"rain_yes": True, "temp_above_30": False})
sizes = pk.simultaneous_kelly(edges, sigmas, bankroll=1000)
"""
def __init__(self, fraction: float = 0.25, max_position_per_market: float = 500.0):
self.fraction = fraction
self.max_position_per_market = max_position_per_market
# Outcome tracking
self.outcome_history: Dict[str, List[int]] = {}
self.target_names: List[str] = []
self.correlation_matrix: Optional[np.ndarray] = None
self.n_observations: int = 0
self._load_correlation()
def record_outcomes(self, outcomes: Dict[str, bool]):
"""
Record resolved market outcomes.
Parameters
----------
outcomes : dict
{target_name: bool} e.g., {"temp_gt_30c_24h": True, "rain_gt_0mm_24h": False}
"""
for name, outcome in outcomes.items():
if name not in self.outcome_history:
self.outcome_history[name] = []
self.outcome_history[name].append(1 if outcome else 0)
# Align all vectors to same length
lengths = [len(v) for v in self.outcome_history.values()]
if lengths:
min_len = min(lengths)
if min_len > self.n_observations:
self.n_observations = min_len
self._recompute_correlation()
def _recompute_correlation(self):
"""Compute correlation matrix from outcome history."""
self.target_names = list(self.outcome_history.keys())
if len(self.target_names) < 2 or self.n_observations < 10:
self.correlation_matrix = None
return
# Build outcome matrix
n = min(self.n_observations, min(len(v) for v in self.outcome_history.values()))
O = np.zeros((n, len(self.target_names)))
for j, name in enumerate(self.target_names):
O[:, j] = self.outcome_history[name][:n]
# Pearson correlation of binary outcomes
self.correlation_matrix = np.corrcoef(O.T)
# Save
self._save_correlation()
def simultaneous_kelly(
self,
edges: Dict[str, float],
sigmas: Optional[Dict[str, float]] = None,
bankroll: float = 1000.0,
) -> Dict[str, float]:
"""
Compute simultaneous Kelly bet sizes.
Parameters
----------
edges : dict
{target_name: edge_in_decimal} e.g., {"temp_gt_30c": 0.15}
sigmas : dict, optional
{target_name: outcome_std} — if None, uses sqrt(p * (1-p))
bankroll : float
Current bankroll
Returns
-------
dict
{target_name: bet_size_in_usdc}
"""
targets = list(edges.keys())
n = len(targets)
if n == 0:
return {}
# Edge vector
mu = np.array([edges[t] for t in targets])
# Variance vector (binary outcome variance)
if sigmas:
sigma_diag = np.array([sigmas[t] for t in targets])
else:
# Binary variance: p(1-p) for outcomes at price p
sigma_diag = np.array([0.25] * n) # Max variance at p=0.5
# Build covariance matrix
if self.correlation_matrix is not None and n >= 2:
idxs = []
for t in targets:
if t in self.target_names:
idxs.append(self.target_names.index(t))
else:
idxs.append(None)
Sigma = np.zeros((n, n))
for i in range(n):
for j in range(n):
if i == j:
Sigma[i, j] = sigma_diag[i]
elif idxs[i] is not None and idxs[j] is not None:
rho = self.correlation_matrix[idxs[i], idxs[j]]
Sigma[i, j] = rho * np.sqrt(sigma_diag[i] * sigma_diag[j])
else:
Sigma[i, j] = 0.0 # Unknown correlation → assume 0
else:
# No correlation data → diagonal
Sigma = np.diag(sigma_diag)
# Regularize to ensure invertibility
Sigma += np.eye(n) * 1e-6
try:
# Simultaneous Kelly: f* = Σ⁻¹ μ
f_star = np.linalg.solve(Sigma, mu)
# Apply fractional Kelly
f_fractional = f_star * self.fraction
# Convert to USD sizes
sizes = {}
for i, name in enumerate(targets):
size = f_fractional[i] * bankroll
# Cap at per-market max, floor at 0
size = max(0, min(size, self.max_position_per_market))
sizes[name] = float(size)
return sizes
except np.linalg.LinAlgError:
# Singular matrix fallback: independent Kelly
sizes = {}
for name in targets:
size = (edges[name] * bankroll * self.fraction / max(sigma_diag[list(edges.keys()).index(name)], 0.01))
sizes[name] = float(max(0, min(size, self.max_position_per_market)))
return sizes
def get_correlation(self, target_a: str, target_b: str) -> float:
"""Get correlation between two target variables."""
if self.correlation_matrix is None:
return 0.0
try:
i = self.target_names.index(target_a)
j = self.target_names.index(target_b)
return float(self.correlation_matrix[i, j])
except (ValueError, IndexError):
return 0.0
def correlation_summary(self) -> str:
"""Human-readable correlation summary."""
if self.correlation_matrix is None or len(self.target_names) < 2:
return "No correlation data (need 10+ observations on 2+ targets)"
lines = ["Correlation Matrix:", " " + " ".join(f"{n:>12s}" for n in self.target_names)]
for i, name in enumerate(self.target_names):
row = f" {name:<12s}"
for j in range(len(self.target_names)):
row += f"{self.correlation_matrix[i,j]:>12.3f}"
lines.append(row)
return "\n".join(lines)
def _save_correlation(self):
"""Persist correlation data."""
path = COV_DIR / "correlation.json"
data = {
"target_names": self.target_names,
"n_observations": self.n_observations,
"correlation_matrix": self.correlation_matrix.tolist() if self.correlation_matrix is not None else None,
"outcome_history": self.outcome_history,
}
with open(path, "w") as f:
json.dump(data, f)
def _load_correlation(self):
"""Load persisted correlation data."""
path = COV_DIR / "correlation.json"
if not path.exists():
return
try:
with open(path) as f:
data = json.load(f)
self.target_names = data.get("target_names", [])
self.n_observations = data.get("n_observations", 0)
self.outcome_history = data.get("outcome_history", {})
corr = data.get("correlation_matrix")
if corr is not None:
self.correlation_matrix = np.array(corr)
except Exception:
pass
def independent_kelly(
self,
edges: Dict[str, float],
bankroll: float = 1000.0,
) -> Dict[str, float]:
"""Independent (simple) Kelly sizing — for comparison."""
sizes = {}
for name, edge in edges.items():
# Simple Kelly: f* = edge / sigma^2 for binary outcomes
sigma_sq = 0.25
f_star = edge / max(sigma_sq, 0.01)
size = f_star * self.fraction * bankroll
sizes[name] = float(max(0, min(size, self.max_position_per_market)))
return sizes
def compare(self, edges: Dict[str, float], bankroll: float = 1000.0) -> Dict:
"""Compare independent vs simultaneous Kelly allocations."""
ind = self.independent_kelly(edges, bankroll)
sim = self.simultaneous_kelly(edges, bankroll=bankroll)
comparison = {}
for name in edges:
comparison[name] = {
"independent": ind.get(name, 0),
"simultaneous": sim.get(name, 0),
"reduction": ind.get(name, 0) - sim.get(name, 0),
}
return comparison