946 lines
40 KiB
Python
946 lines
40 KiB
Python
from __future__ import annotations
|
|
|
|
import time
|
|
from datetime import datetime, timedelta
|
|
from typing import Any, Dict, Optional
|
|
|
|
from loguru import logger
|
|
|
|
from src.utils.metrics import record_source_call
|
|
|
|
|
|
OPEN_METEO_MULTI_MODEL_SPECS: Dict[str, Dict[str, Any]] = {
|
|
"ecmwf_ifs025": {
|
|
"label": "ECMWF",
|
|
"provider": "ECMWF",
|
|
"model": "IFS",
|
|
"tier": "global",
|
|
"resolution_km": 25,
|
|
"horizon": "15d",
|
|
},
|
|
"ecmwf_aifs025_single": {
|
|
"label": "ECMWF AIFS",
|
|
"provider": "ECMWF",
|
|
"model": "AIFS single",
|
|
"tier": "aifs_global",
|
|
"resolution_km": 25,
|
|
"horizon": "15d",
|
|
},
|
|
"gfs_seamless": {
|
|
"label": "GFS",
|
|
"provider": "NOAA",
|
|
"model": "GFS Seamless",
|
|
"tier": "global",
|
|
"resolution_km": None,
|
|
"horizon": "16d",
|
|
},
|
|
"icon_seamless": {
|
|
"label": "ICON",
|
|
"provider": "DWD",
|
|
"model": "ICON Seamless",
|
|
"tier": "global",
|
|
"resolution_km": None,
|
|
"horizon": "7d",
|
|
},
|
|
"icon_eu": {
|
|
"label": "ICON-EU",
|
|
"provider": "DWD",
|
|
"model": "ICON-EU",
|
|
"tier": "regional_europe",
|
|
"resolution_km": 7,
|
|
"horizon": "5d",
|
|
},
|
|
"icon_d2": {
|
|
"label": "ICON-D2",
|
|
"provider": "DWD",
|
|
"model": "ICON-D2",
|
|
"tier": "short_range_europe",
|
|
"resolution_km": 2.2,
|
|
"horizon": "2d",
|
|
},
|
|
"gem_seamless": {
|
|
"label": "GEM",
|
|
"provider": "ECCC",
|
|
"model": "GEM Seamless",
|
|
"tier": "global",
|
|
"resolution_km": None,
|
|
"horizon": "10d",
|
|
},
|
|
"gem_global": {
|
|
"label": "GDPS",
|
|
"provider": "ECCC",
|
|
"model": "GEM Global / GDPS",
|
|
"tier": "global",
|
|
"resolution_km": 15,
|
|
"horizon": "10d",
|
|
},
|
|
"gem_regional": {
|
|
"label": "RDPS",
|
|
"provider": "ECCC",
|
|
"model": "GEM Regional / RDPS",
|
|
"tier": "regional_north_america",
|
|
"resolution_km": 10,
|
|
"horizon": "3d",
|
|
},
|
|
"gem_hrdps_continental": {
|
|
"label": "HRDPS",
|
|
"provider": "ECCC",
|
|
"model": "HRDPS Continental",
|
|
"tier": "short_range_north_america",
|
|
"resolution_km": 2.5,
|
|
"horizon": "2d",
|
|
},
|
|
"jma_seamless": {
|
|
"label": "JMA",
|
|
"provider": "JMA",
|
|
"model": "JMA Seamless",
|
|
"tier": "global",
|
|
"resolution_km": None,
|
|
"horizon": "11d",
|
|
},
|
|
"meteofrance_arome_france_hd": {
|
|
"label": "AROME HD",
|
|
"provider": "Meteo-France",
|
|
"model": "AROME France HD",
|
|
"tier": "short_range_europe",
|
|
"resolution_km": 1.3,
|
|
"horizon": "1.5d",
|
|
},
|
|
}
|
|
|
|
OPEN_METEO_MULTI_MODEL_ORDER = tuple(OPEN_METEO_MULTI_MODEL_SPECS.keys())
|
|
|
|
|
|
def _parse_open_meteo_multi_model_daily(
|
|
daily: Dict[str, Any],
|
|
) -> tuple[list[str], Dict[str, Dict[str, float]], Dict[str, Dict[str, Any]], Dict[str, str]]:
|
|
dates = daily.get("time", [])
|
|
if not isinstance(dates, list):
|
|
dates = []
|
|
|
|
model_metadata: Dict[str, Dict[str, Any]] = {}
|
|
model_keys: Dict[str, str] = {}
|
|
daily_forecasts: Dict[str, Dict[str, float]] = {}
|
|
|
|
for day_idx, date_str in enumerate(dates):
|
|
day_data: Dict[str, float] = {}
|
|
for model_key, spec in OPEN_METEO_MULTI_MODEL_SPECS.items():
|
|
label = str(spec["label"])
|
|
key = f"temperature_2m_max_{model_key}"
|
|
values = daily.get(key, [])
|
|
if not isinstance(values, list):
|
|
continue
|
|
if day_idx >= len(values) or values[day_idx] is None:
|
|
continue
|
|
try:
|
|
day_data[label] = round(float(values[day_idx]), 1)
|
|
except (TypeError, ValueError):
|
|
continue
|
|
model_metadata[label] = {
|
|
**spec,
|
|
"open_meteo_model": model_key,
|
|
}
|
|
model_keys[label] = model_key
|
|
if day_data:
|
|
daily_forecasts[str(date_str)] = day_data
|
|
|
|
return [str(d) for d in dates], daily_forecasts, model_metadata, model_keys
|
|
|
|
|
|
def _parse_open_meteo_multi_model_hourly(
|
|
hourly: Dict[str, Any],
|
|
) -> tuple[list[str], Dict[str, list[Optional[float]]]]:
|
|
times = hourly.get("time", [])
|
|
if not isinstance(times, list):
|
|
times = []
|
|
|
|
hourly_times = [str(t) for t in times]
|
|
hourly_forecasts: Dict[str, list[Optional[float]]] = {}
|
|
|
|
for model_key, spec in OPEN_METEO_MULTI_MODEL_SPECS.items():
|
|
label = str(spec["label"])
|
|
key = f"temperature_2m_{model_key}"
|
|
values = hourly.get(key, [])
|
|
if not isinstance(values, list) or not values:
|
|
continue
|
|
|
|
parsed_values = []
|
|
has_valid = False
|
|
for v in values:
|
|
if v is None:
|
|
parsed_values.append(None)
|
|
else:
|
|
try:
|
|
parsed_values.append(round(float(v), 1))
|
|
has_valid = True
|
|
except (TypeError, ValueError):
|
|
parsed_values.append(None)
|
|
|
|
if has_valid:
|
|
if len(parsed_values) < len(hourly_times):
|
|
parsed_values.extend([None] * (len(hourly_times) - len(parsed_values)))
|
|
elif len(parsed_values) > len(hourly_times):
|
|
parsed_values = parsed_values[:len(hourly_times)]
|
|
hourly_forecasts[label] = parsed_values
|
|
|
|
return hourly_times, hourly_forecasts
|
|
|
|
|
|
def _count_models(values: Any) -> int:
|
|
if not isinstance(values, dict):
|
|
return 0
|
|
count = 0
|
|
for value in values.values():
|
|
try:
|
|
float(value)
|
|
count += 1
|
|
except (TypeError, ValueError):
|
|
continue
|
|
return count
|
|
|
|
|
|
def _get_hourly_by_time(
|
|
times: list[str],
|
|
forecasts: Dict[str, list[Optional[float]]],
|
|
) -> Dict[str, Dict[str, float]]:
|
|
hourly_by_time: Dict[str, Dict[str, float]] = {}
|
|
for i, t_str in enumerate(times):
|
|
hour_data: Dict[str, float] = {}
|
|
for label, values in forecasts.items():
|
|
if i < len(values) and values[i] is not None:
|
|
hour_data[label] = values[i]
|
|
if hour_data:
|
|
hourly_by_time[t_str] = hour_data
|
|
return hourly_by_time
|
|
|
|
|
|
def _reconstruct_hourly_forecasts(
|
|
merged_times: list[str],
|
|
merged_hourly_by_time: Dict[str, Dict[str, float]],
|
|
) -> Dict[str, list[Optional[float]]]:
|
|
labels = set()
|
|
for hour_data in merged_hourly_by_time.values():
|
|
labels.update(hour_data.keys())
|
|
|
|
hourly_forecasts: Dict[str, list[Optional[float]]] = {}
|
|
for label in sorted(labels):
|
|
values = []
|
|
for t_str in merged_times:
|
|
values.append(merged_hourly_by_time[t_str].get(label))
|
|
hourly_forecasts[label] = values
|
|
return hourly_forecasts
|
|
|
|
|
|
def _merge_multi_model_result_with_cache(
|
|
cached: Dict[str, Any],
|
|
fresh: Dict[str, Any],
|
|
) -> Dict[str, Any]:
|
|
"""Avoid downgrading a rich multi-model cache when Open-Meteo returns a sparse subset."""
|
|
if not isinstance(cached, dict):
|
|
return fresh
|
|
cached_forecasts = cached.get("forecasts") if isinstance(cached.get("forecasts"), dict) else {}
|
|
fresh_forecasts = fresh.get("forecasts") if isinstance(fresh.get("forecasts"), dict) else {}
|
|
if _count_models(fresh_forecasts) >= _count_models(cached_forecasts):
|
|
return fresh
|
|
|
|
merged_daily: Dict[str, Dict[str, float]] = {}
|
|
cached_daily = cached.get("daily_forecasts") if isinstance(cached.get("daily_forecasts"), dict) else {}
|
|
fresh_daily = fresh.get("daily_forecasts") if isinstance(fresh.get("daily_forecasts"), dict) else {}
|
|
for date in sorted({*cached_daily.keys(), *fresh_daily.keys()}):
|
|
cached_day = cached_daily.get(date) if isinstance(cached_daily.get(date), dict) else {}
|
|
fresh_day = fresh_daily.get(date) if isinstance(fresh_daily.get(date), dict) else {}
|
|
merged_daily[str(date)] = (
|
|
dict(fresh_day)
|
|
if _count_models(fresh_day) >= _count_models(cached_day)
|
|
else {**dict(cached_day), **dict(fresh_day)}
|
|
)
|
|
|
|
# Merge hourly forecasts
|
|
cached_hourly_times = cached.get("hourly_times") or []
|
|
cached_hourly_forecasts = cached.get("hourly_forecasts") or {}
|
|
fresh_hourly_times = fresh.get("hourly_times") or []
|
|
fresh_hourly_forecasts = fresh.get("hourly_forecasts") or {}
|
|
|
|
cached_hourly_by_time = _get_hourly_by_time(cached_hourly_times, cached_hourly_forecasts)
|
|
fresh_hourly_by_time = _get_hourly_by_time(fresh_hourly_times, fresh_hourly_forecasts)
|
|
|
|
merged_hourly_by_time = {}
|
|
all_hourly_times = sorted({*cached_hourly_by_time.keys(), *fresh_hourly_by_time.keys()})
|
|
|
|
# Cutoff to keep cache size bounded: discard times before fresh starts
|
|
if fresh_hourly_times:
|
|
cutoff = fresh_hourly_times[0]
|
|
all_hourly_times = [t for t in all_hourly_times if t >= cutoff]
|
|
|
|
for t_str in all_hourly_times:
|
|
c_hour = cached_hourly_by_time.get(t_str) or {}
|
|
f_hour = fresh_hourly_by_time.get(t_str) or {}
|
|
merged_hourly_by_time[t_str] = (
|
|
dict(f_hour)
|
|
if _count_models(f_hour) >= _count_models(c_hour)
|
|
else {**dict(c_hour), **dict(f_hour)}
|
|
)
|
|
|
|
merged_hourly_forecasts = _reconstruct_hourly_forecasts(all_hourly_times, merged_hourly_by_time)
|
|
|
|
merged = {
|
|
**cached,
|
|
**fresh,
|
|
"forecasts": {**dict(cached_forecasts), **dict(fresh_forecasts)},
|
|
"daily_forecasts": merged_daily or cached_daily or fresh_daily,
|
|
"hourly_times": all_hourly_times,
|
|
"hourly_forecasts": merged_hourly_forecasts,
|
|
"model_metadata": {
|
|
**(cached.get("model_metadata") if isinstance(cached.get("model_metadata"), dict) else {}),
|
|
**(fresh.get("model_metadata") if isinstance(fresh.get("model_metadata"), dict) else {}),
|
|
},
|
|
"model_keys": {
|
|
**(cached.get("model_keys") if isinstance(cached.get("model_keys"), dict) else {}),
|
|
**(fresh.get("model_keys") if isinstance(fresh.get("model_keys"), dict) else {}),
|
|
},
|
|
"partial_refresh": True,
|
|
"partial_refresh_model_count": _count_models(fresh_forecasts),
|
|
"preserved_cached_model_count": _count_models(cached_forecasts),
|
|
}
|
|
merged_dates = list(dict.fromkeys([*(fresh.get("dates") or []), *(cached.get("dates") or [])]))
|
|
if merged_dates:
|
|
merged["dates"] = merged_dates
|
|
return merged
|
|
|
|
|
|
class NwsOpenMeteoSourceMixin:
|
|
def fetch_nws(self, lat: float, lon: float) -> Optional[Dict]:
|
|
"""
|
|
从 NWS (美国国家气象局) 获取高精度预报
|
|
仅适用于美国城市,全球 VPS 均可访问
|
|
"""
|
|
started = time.perf_counter()
|
|
try:
|
|
# 1. 获取网格点
|
|
points_url = f"https://api.weather.gov/points/{lat},{lon}"
|
|
headers = {"User-Agent": "PolyWeather/1.0 (weather-bot)"}
|
|
|
|
points_resp = self._http_get(
|
|
points_url, headers=headers, timeout=self.timeout
|
|
)
|
|
points_resp.raise_for_status()
|
|
points_data = points_resp.json()
|
|
|
|
properties = points_data.get("properties", {})
|
|
forecast_url = properties.get("forecast")
|
|
hourly_url = properties.get("forecastHourly")
|
|
if not forecast_url:
|
|
record_source_call("nws", "forecast", "empty", (time.perf_counter() - started) * 1000.0)
|
|
return None
|
|
|
|
# 2. 获取预报
|
|
forecast_resp = self._http_get(
|
|
forecast_url, headers=headers, timeout=self.timeout
|
|
)
|
|
forecast_resp.raise_for_status()
|
|
forecast_data = forecast_resp.json()
|
|
|
|
periods = forecast_data.get("properties", {}).get("periods", [])
|
|
if not periods:
|
|
record_source_call("nws", "forecast", "empty", (time.perf_counter() - started) * 1000.0)
|
|
return None
|
|
|
|
hourly_periods = []
|
|
if hourly_url:
|
|
hourly_resp = self._http_get(
|
|
hourly_url, headers=headers, timeout=self.timeout
|
|
)
|
|
hourly_resp.raise_for_status()
|
|
hourly_data = hourly_resp.json()
|
|
hourly_periods = hourly_data.get("properties", {}).get("periods", [])[:48]
|
|
|
|
active_alerts = []
|
|
try:
|
|
alerts_resp = self._http_get(
|
|
"https://api.weather.gov/alerts/active",
|
|
params={"point": f"{lat},{lon}"},
|
|
headers=headers,
|
|
timeout=self.timeout,
|
|
)
|
|
alerts_resp.raise_for_status()
|
|
alerts_data = alerts_resp.json()
|
|
for feature in alerts_data.get("features", [])[:8]:
|
|
ap = feature.get("properties", {})
|
|
active_alerts.append(
|
|
{
|
|
"event": ap.get("event"),
|
|
"headline": ap.get("headline"),
|
|
"severity": ap.get("severity"),
|
|
"certainty": ap.get("certainty"),
|
|
"urgency": ap.get("urgency"),
|
|
"effective": ap.get("effective"),
|
|
"ends": ap.get("ends"),
|
|
}
|
|
)
|
|
except Exception:
|
|
active_alerts = []
|
|
|
|
# 3. 提取今日最高温(找 isDaytime=True 的第一个)
|
|
today_high = None
|
|
for p in periods:
|
|
if p.get("isDaytime") and "High" in p.get("name", ""):
|
|
today_high = p.get("temperature")
|
|
break
|
|
# 如果没有明确的 High,取第一个 daytime 的温度
|
|
if today_high is None:
|
|
for p in periods:
|
|
if p.get("isDaytime"):
|
|
today_high = p.get("temperature")
|
|
break
|
|
|
|
result = {
|
|
"source": "nws",
|
|
"today_high": today_high,
|
|
"unit": "fahrenheit",
|
|
"forecast_periods": [
|
|
{
|
|
"name": p.get("name"),
|
|
"start_time": p.get("startTime"),
|
|
"end_time": p.get("endTime"),
|
|
"is_daytime": p.get("isDaytime"),
|
|
"temperature": p.get("temperature"),
|
|
"temperature_trend": p.get("temperatureTrend"),
|
|
"wind_speed": p.get("windSpeed"),
|
|
"wind_direction": p.get("windDirection"),
|
|
"short_forecast": p.get("shortForecast"),
|
|
"detailed_forecast": p.get("detailedForecast"),
|
|
"precipitation_probability": (p.get("probabilityOfPrecipitation") or {}).get("value"),
|
|
}
|
|
for p in periods[:14]
|
|
],
|
|
"hourly_periods": [
|
|
{
|
|
"start_time": p.get("startTime"),
|
|
"end_time": p.get("endTime"),
|
|
"temperature": p.get("temperature"),
|
|
"temperature_unit": p.get("temperatureUnit"),
|
|
"wind_speed": p.get("windSpeed"),
|
|
"wind_direction": p.get("windDirection"),
|
|
"short_forecast": p.get("shortForecast"),
|
|
"precipitation_probability": (p.get("probabilityOfPrecipitation") or {}).get("value"),
|
|
}
|
|
for p in hourly_periods
|
|
],
|
|
"active_alerts": active_alerts,
|
|
}
|
|
record_source_call("nws", "forecast", "success", (time.perf_counter() - started) * 1000.0)
|
|
return result
|
|
except Exception as e:
|
|
logger.warning(f"NWS 请求失败: {e}")
|
|
record_source_call("nws", "forecast", "error", (time.perf_counter() - started) * 1000.0)
|
|
return None
|
|
|
|
def fetch_from_open_meteo(
|
|
self,
|
|
lat: float,
|
|
lon: float,
|
|
forecast_days: int = 14,
|
|
use_fahrenheit: bool = False,
|
|
) -> Optional[Dict]:
|
|
"""
|
|
Fetch weather from Open-Meteo with forecast data
|
|
|
|
Args:
|
|
lat: Latitude
|
|
lon: Longitude
|
|
forecast_days: Number of forecast days to fetch (default 14 to cover all market dates)
|
|
use_fahrenheit: Whether to return temperatures in Fahrenheit (for US markets)
|
|
"""
|
|
started = time.perf_counter()
|
|
cache_key = (
|
|
f"{round(float(lat), 4)}:{round(float(lon), 4)}:"
|
|
f"{forecast_days}:{'f' if use_fahrenheit else 'c'}"
|
|
)
|
|
self._maybe_reload_open_meteo_disk_cache()
|
|
now_ts = time.time()
|
|
# ── 429 冷却期检查(所有 Open-Meteo 端点共享)─────────────────
|
|
with self._open_meteo_rl_lock:
|
|
cooldown_until = self._refresh_open_meteo_rate_limit_until()
|
|
if now_ts < cooldown_until:
|
|
remaining = int(cooldown_until - now_ts)
|
|
logger.debug(f"Open-Meteo 冷却期中,跳过请求,还需 {remaining}s")
|
|
with self._open_meteo_cache_lock:
|
|
stale = self._open_meteo_cache.get(cache_key)
|
|
if stale and isinstance(stale.get("data"), dict):
|
|
record_source_call("open_meteo", "forecast", "stale_cache", (time.perf_counter() - started) * 1000.0)
|
|
return dict(stale["data"])
|
|
# Memory miss: force-reload from disk and retry once
|
|
self._load_open_meteo_disk_cache()
|
|
with self._open_meteo_cache_lock:
|
|
stale2 = self._open_meteo_cache.get(cache_key)
|
|
if stale2 and isinstance(stale2.get("data"), dict):
|
|
record_source_call("open_meteo", "forecast", "disk_fallback", (time.perf_counter() - started) * 1000.0)
|
|
return dict(stale2["data"])
|
|
record_source_call("open_meteo", "forecast", "cooldown_skip", (time.perf_counter() - started) * 1000.0)
|
|
return None
|
|
with self._open_meteo_cache_lock:
|
|
cached = self._open_meteo_cache.get(cache_key)
|
|
if (
|
|
cached
|
|
and now_ts - float(cached.get("t", 0)) < self.open_meteo_cache_ttl_sec
|
|
):
|
|
cached_data = cached.get("data")
|
|
if isinstance(cached_data, dict):
|
|
record_source_call("open_meteo", "forecast", "cache_hit", (time.perf_counter() - started) * 1000.0)
|
|
return dict(cached_data)
|
|
try:
|
|
url = "https://api.open-meteo.com/v1/forecast"
|
|
params = {
|
|
"latitude": lat,
|
|
"longitude": lon,
|
|
"current_weather": "true",
|
|
"hourly": (
|
|
"temperature_2m,shortwave_radiation,dew_point_2m,pressure_msl,"
|
|
"wind_speed_10m,wind_direction_10m,wind_speed_180m,wind_direction_180m,"
|
|
"precipitation_probability,cloud_cover,cape,convective_inhibition,"
|
|
"lifted_index,boundary_layer_height"
|
|
),
|
|
"daily": "temperature_2m_max,apparent_temperature_max,sunrise,sunset,sunshine_duration",
|
|
"timezone": "auto",
|
|
"forecast_days": forecast_days,
|
|
"past_days": 1,
|
|
}
|
|
|
|
# 显式指定单位,防止 API 默认行为漂移
|
|
if use_fahrenheit:
|
|
params["temperature_unit"] = "fahrenheit"
|
|
else:
|
|
params["temperature_unit"] = "celsius"
|
|
|
|
self._wait_open_meteo_slot("forecast")
|
|
response = self._http_get(
|
|
url,
|
|
params=params,
|
|
timeout=getattr(self, "open_meteo_timeout_sec", self.timeout),
|
|
)
|
|
response.raise_for_status()
|
|
data = response.json()
|
|
|
|
current = data.get("current_weather", {})
|
|
utc_offset = data.get("utc_offset_seconds", 0)
|
|
timezone_name = data.get("timezone", "UTC")
|
|
|
|
# 处理多模型数据 (如果请求了 models 参数,返回结构会变化)
|
|
daily_data = data.get("daily", {})
|
|
if "temperature_2m_max_ecmwf_ifs04" in daily_data:
|
|
ecmwf_max = daily_data.get("temperature_2m_max_ecmwf_ifs04", [])
|
|
hrrr_max = daily_data.get("temperature_2m_max_ncep_hrrr_conus", [])
|
|
|
|
# 记录今日模型分歧
|
|
daily_data["model_split"] = {
|
|
"ecmwf": ecmwf_max[0] if ecmwf_max else None,
|
|
"hrrr": hrrr_max[0] if hrrr_max else None,
|
|
}
|
|
|
|
# 智能合并:HRRR 仅覆盖 48 小时,远期用 ECMWF 补全
|
|
merged_max = []
|
|
for i in range(len(ecmwf_max)):
|
|
hrrr_val = hrrr_max[i] if i < len(hrrr_max) else None
|
|
ecmwf_val = ecmwf_max[i] if i < len(ecmwf_max) else None
|
|
|
|
# 优先 HRRR,其次 ECMWF,都没有就跳过
|
|
if hrrr_val is not None:
|
|
merged_max.append(hrrr_val)
|
|
elif ecmwf_val is not None:
|
|
merged_max.append(ecmwf_val)
|
|
else:
|
|
# 两个都没有,用占位符 (理论上不应该发生)
|
|
merged_max.append(ecmwf_val) # None
|
|
daily_data["temperature_2m_max"] = merged_max
|
|
|
|
# 映射逐小时数据
|
|
hourly_data = data.get("hourly", {})
|
|
if "temperature_2m_ncep_hrrr_conus" in hourly_data:
|
|
hourly_data["temperature_2m"] = hourly_data[
|
|
"temperature_2m_ncep_hrrr_conus"
|
|
]
|
|
|
|
# 计算精确的当地时间
|
|
now_utc = datetime.utcnow()
|
|
local_now = now_utc + timedelta(seconds=utc_offset)
|
|
local_time_str = local_now.strftime("%Y-%m-%d %H:%M")
|
|
|
|
result = {
|
|
"source": "open-meteo",
|
|
"timestamp": now_utc.isoformat(),
|
|
"timezone": timezone_name,
|
|
"utc_offset": utc_offset,
|
|
"current": {
|
|
"temp": current.get("temperature"),
|
|
"local_time": local_time_str,
|
|
},
|
|
"hourly": hourly_data,
|
|
"daily": daily_data,
|
|
"unit": "fahrenheit" if use_fahrenheit else "celsius",
|
|
}
|
|
with self._open_meteo_cache_lock:
|
|
self._open_meteo_cache[cache_key] = {
|
|
"t": time.time(),
|
|
"data": dict(result),
|
|
}
|
|
self._flush_open_meteo_disk_cache()
|
|
record_source_call("open_meteo", "forecast", "success", (time.perf_counter() - started) * 1000.0)
|
|
return result
|
|
except Exception as e:
|
|
status_code = getattr(getattr(e, "response", None), "status_code", None)
|
|
if status_code == 429:
|
|
retry_after_str = getattr(e.response, "headers", {}).get("Retry-After")
|
|
cooldown_to_use = self._open_meteo_rl_cooldown
|
|
if retry_after_str:
|
|
try:
|
|
parsed = int(retry_after_str)
|
|
if parsed > 0:
|
|
cooldown_to_use = min(parsed + 60, 3600) # Add 60s buffer, max 1 hour
|
|
logger.info(f"Open-Meteo 响应包含 Retry-After: {retry_after_str}s")
|
|
except ValueError:
|
|
pass
|
|
logger.warning(
|
|
f"Open-Meteo rate limited (429), fallback to cache if available: lat={lat}, lon={lon}"
|
|
)
|
|
# 设置全局冷却期,避免短时内重复触发 429
|
|
with self._open_meteo_rl_lock:
|
|
self._set_open_meteo_rate_limit_until(time.time() + cooldown_to_use, reason="forecast_429")
|
|
logger.warning(f"Open-Meteo 触发限流,设置 {cooldown_to_use}s 冷却期")
|
|
else:
|
|
logger.error(f"Open-Meteo forecast failed: {e}")
|
|
with self._open_meteo_cache_lock:
|
|
stale = self._open_meteo_cache.get(cache_key)
|
|
if stale and isinstance(stale.get("data"), dict):
|
|
fallback = dict(stale["data"])
|
|
fallback["stale_cache"] = True
|
|
record_source_call("open_meteo", "forecast", "stale_cache", (time.perf_counter() - started) * 1000.0)
|
|
return fallback
|
|
record_source_call("open_meteo", "forecast", "error", (time.perf_counter() - started) * 1000.0)
|
|
return None
|
|
|
|
def fetch_ensemble(
|
|
self,
|
|
lat: float,
|
|
lon: float,
|
|
use_fahrenheit: bool = False,
|
|
) -> Optional[Dict]:
|
|
"""
|
|
从 Open-Meteo Ensemble API 获取 51 成员集合预报
|
|
用于计算预报不确定性范围(散度)
|
|
"""
|
|
started = time.perf_counter()
|
|
cache_key = (
|
|
f"{round(float(lat), 4)}:{round(float(lon), 4)}:"
|
|
f"{'f' if use_fahrenheit else 'c'}"
|
|
)
|
|
self._maybe_reload_open_meteo_disk_cache()
|
|
now_ts = time.time()
|
|
# ── 429 冷却期检查(所有 Open-Meteo 端点共享)─────────────────
|
|
with self._open_meteo_rl_lock:
|
|
cooldown_until = self._refresh_open_meteo_rate_limit_until()
|
|
if now_ts < cooldown_until:
|
|
remaining = int(cooldown_until - now_ts)
|
|
logger.debug(f"Open-Meteo Ensemble 冷却期中,跳过请求,还需 {remaining}s")
|
|
with self._ensemble_cache_lock:
|
|
stale = self._ensemble_cache.get(cache_key)
|
|
if stale and isinstance(stale.get("data"), dict):
|
|
record_source_call("open_meteo", "ensemble", "stale_cache", (time.perf_counter() - started) * 1000.0)
|
|
return dict(stale["data"])
|
|
self._load_open_meteo_disk_cache()
|
|
with self._ensemble_cache_lock:
|
|
stale2 = self._ensemble_cache.get(cache_key)
|
|
if stale2 and isinstance(stale2.get("data"), dict):
|
|
record_source_call("open_meteo", "ensemble", "disk_fallback", (time.perf_counter() - started) * 1000.0)
|
|
return dict(stale2["data"])
|
|
record_source_call("open_meteo", "ensemble", "cooldown_skip", (time.perf_counter() - started) * 1000.0)
|
|
return None
|
|
|
|
with self._ensemble_cache_lock:
|
|
cached = self._ensemble_cache.get(cache_key)
|
|
if (
|
|
cached
|
|
and now_ts - float(cached.get("t", 0))
|
|
< self.open_meteo_ensemble_cache_ttl_sec
|
|
):
|
|
cached_data = cached.get("data")
|
|
if isinstance(cached_data, dict):
|
|
record_source_call("open_meteo", "ensemble", "cache_hit", (time.perf_counter() - started) * 1000.0)
|
|
return dict(cached_data)
|
|
try:
|
|
url = "https://ensemble-api.open-meteo.com/v1/ensemble"
|
|
params = {
|
|
"latitude": lat,
|
|
"longitude": lon,
|
|
"daily": "temperature_2m_max",
|
|
"timezone": "auto",
|
|
"forecast_days": 3,
|
|
}
|
|
if use_fahrenheit:
|
|
params["temperature_unit"] = "fahrenheit"
|
|
else:
|
|
params["temperature_unit"] = "celsius"
|
|
|
|
self._wait_open_meteo_slot("ensemble")
|
|
response = self._http_get(
|
|
url,
|
|
params=params,
|
|
timeout=getattr(self, "open_meteo_timeout_sec", self.timeout),
|
|
)
|
|
response.raise_for_status()
|
|
data = response.json()
|
|
|
|
daily = data.get("daily", {})
|
|
# 每个成员都会返回一组 temperature_2m_max
|
|
# 格式: {"time": [...], "temperature_2m_max_member01": [...], ...}
|
|
today_highs = []
|
|
for key, values in daily.items():
|
|
if key.startswith("temperature_2m_max") and key != "temperature_2m_max":
|
|
if values and values[0] is not None:
|
|
today_highs.append(values[0])
|
|
|
|
# 也检查非成员键(有些返回格式不同)
|
|
if not today_highs:
|
|
raw_max = daily.get("temperature_2m_max", [])
|
|
if isinstance(raw_max, list) and raw_max:
|
|
if isinstance(raw_max[0], list):
|
|
# 嵌套列表格式: [[member1_day1, member1_day2], [member2_day1, ...]]
|
|
today_highs = [m[0] for m in raw_max if m and m[0] is not None]
|
|
elif raw_max[0] is not None:
|
|
today_highs = [raw_max[0]]
|
|
|
|
if len(today_highs) < 3:
|
|
logger.warning(f"Ensemble 数据不足: 仅获取 {len(today_highs)} 个成员")
|
|
record_source_call("open_meteo", "ensemble", "empty", (time.perf_counter() - started) * 1000.0)
|
|
return None
|
|
|
|
today_highs.sort()
|
|
n = len(today_highs)
|
|
median = today_highs[n // 2]
|
|
p10 = today_highs[max(0, int(n * 0.1))]
|
|
p90 = today_highs[min(n - 1, int(n * 0.9))]
|
|
|
|
result = {
|
|
"source": "ensemble",
|
|
"members": n,
|
|
"median": round(median, 1),
|
|
"p10": round(p10, 1),
|
|
"p90": round(p90, 1),
|
|
"min": round(today_highs[0], 1),
|
|
"max": round(today_highs[-1], 1),
|
|
"unit": "fahrenheit" if use_fahrenheit else "celsius",
|
|
}
|
|
|
|
logger.info(
|
|
f"📊 Ensemble ({n} members): median={median:.1f}, "
|
|
f"p10={p10:.1f}, p90={p90:.1f}"
|
|
)
|
|
with self._ensemble_cache_lock:
|
|
self._ensemble_cache[cache_key] = {
|
|
"t": time.time(),
|
|
"data": dict(result),
|
|
}
|
|
self._flush_open_meteo_disk_cache()
|
|
record_source_call("open_meteo", "ensemble", "success", (time.perf_counter() - started) * 1000.0)
|
|
return result
|
|
except Exception as e:
|
|
status_code = getattr(getattr(e, "response", None), "status_code", None)
|
|
if status_code == 429:
|
|
retry_after_str = getattr(e.response, "headers", {}).get("Retry-After")
|
|
cooldown_to_use = self._open_meteo_rl_cooldown
|
|
if retry_after_str:
|
|
try:
|
|
parsed = int(retry_after_str)
|
|
if parsed > 0:
|
|
cooldown_to_use = min(parsed + 60, 3600)
|
|
except ValueError:
|
|
pass
|
|
logger.warning(
|
|
f"Ensemble API rate limited (429), fallback to cache if available: lat={lat}, lon={lon}"
|
|
)
|
|
with self._open_meteo_rl_lock:
|
|
self._set_open_meteo_rate_limit_until(time.time() + cooldown_to_use, reason="ensemble_429")
|
|
else:
|
|
logger.warning(f"Ensemble API 请求失败: {e}")
|
|
with self._ensemble_cache_lock:
|
|
stale = self._ensemble_cache.get(cache_key)
|
|
if stale and isinstance(stale.get("data"), dict):
|
|
fallback = dict(stale["data"])
|
|
fallback["stale_cache"] = True
|
|
record_source_call("open_meteo", "ensemble", "stale_cache", (time.perf_counter() - started) * 1000.0)
|
|
return fallback
|
|
record_source_call("open_meteo", "ensemble", "error", (time.perf_counter() - started) * 1000.0)
|
|
return None
|
|
|
|
def fetch_multi_model(
|
|
self,
|
|
lat: float,
|
|
lon: float,
|
|
city: str = "",
|
|
use_fahrenheit: bool = False,
|
|
) -> Optional[Dict]:
|
|
"""
|
|
从 Open-Meteo 获取多个独立 NWP 模型的预报
|
|
用于真正的多模型共识评分
|
|
|
|
模型列表:
|
|
- ECMWF IFS / AIFS
|
|
- GFS (美国 NOAA)
|
|
- ICON Seamless / ICON-EU / ICON-D2 (德国气象局 DWD)
|
|
- GEM Seamless / GDPS / RDPS / HRDPS (加拿大 ECCC)
|
|
- JMA (日本气象厅)
|
|
|
|
返回 3 天的预报数据,支持今日+明日共识分析
|
|
"""
|
|
started = time.perf_counter()
|
|
cache_city = str(city or "").strip().lower()
|
|
cache_key = (
|
|
f"{round(float(lat), 4)}:{round(float(lon), 4)}:{cache_city}:"
|
|
f"{'f' if use_fahrenheit else 'c'}:{self.multi_model_cache_version}"
|
|
)
|
|
self._maybe_reload_open_meteo_disk_cache()
|
|
now_ts = time.time()
|
|
# ── 429 冷却期检查(所有 Open-Meteo 端点共享)─────────────────
|
|
with self._open_meteo_rl_lock:
|
|
cooldown_until = self._refresh_open_meteo_rate_limit_until()
|
|
if now_ts < cooldown_until:
|
|
remaining = int(cooldown_until - now_ts)
|
|
logger.debug(f"Open-Meteo Multi-model 冷却期中,跳过请求,还需 {remaining}s")
|
|
with self._multi_model_cache_lock:
|
|
stale = self._multi_model_cache.get(cache_key)
|
|
if stale and isinstance(stale.get("data"), dict):
|
|
record_source_call("open_meteo", "multi_model", "stale_cache", (time.perf_counter() - started) * 1000.0)
|
|
return dict(stale["data"])
|
|
self._load_open_meteo_disk_cache()
|
|
with self._multi_model_cache_lock:
|
|
stale2 = self._multi_model_cache.get(cache_key)
|
|
if stale2 and isinstance(stale2.get("data"), dict):
|
|
record_source_call("open_meteo", "multi_model", "disk_fallback", (time.perf_counter() - started) * 1000.0)
|
|
return dict(stale2["data"])
|
|
record_source_call("open_meteo", "multi_model", "cooldown_skip", (time.perf_counter() - started) * 1000.0)
|
|
return None
|
|
|
|
with self._multi_model_cache_lock:
|
|
cached = self._multi_model_cache.get(cache_key)
|
|
if (
|
|
cached
|
|
and now_ts - float(cached.get("t", 0))
|
|
< self.open_meteo_multi_model_cache_ttl_sec
|
|
):
|
|
cached_data = cached.get("data")
|
|
if isinstance(cached_data, dict):
|
|
record_source_call("open_meteo", "multi_model", "cache_hit", (time.perf_counter() - started) * 1000.0)
|
|
return dict(cached_data)
|
|
try:
|
|
url = "https://api.open-meteo.com/v1/forecast"
|
|
models = ",".join(OPEN_METEO_MULTI_MODEL_ORDER)
|
|
params = {
|
|
"latitude": lat,
|
|
"longitude": lon,
|
|
"daily": "temperature_2m_max",
|
|
"hourly": "temperature_2m",
|
|
"models": models,
|
|
"timezone": "auto",
|
|
"forecast_days": 3,
|
|
}
|
|
if use_fahrenheit:
|
|
params["temperature_unit"] = "fahrenheit"
|
|
|
|
self._wait_open_meteo_slot("multi-model")
|
|
response = self._http_get(
|
|
url,
|
|
params=params,
|
|
timeout=getattr(self, "open_meteo_timeout_sec", self.timeout),
|
|
)
|
|
response.raise_for_status()
|
|
data = response.json()
|
|
|
|
daily = data.get("daily", {})
|
|
if not isinstance(daily, dict):
|
|
daily = {}
|
|
dates, daily_forecasts, model_metadata, model_keys = (
|
|
_parse_open_meteo_multi_model_daily(daily)
|
|
)
|
|
|
|
hourly = data.get("hourly", {})
|
|
if not isinstance(hourly, dict):
|
|
hourly = {}
|
|
hourly_times, hourly_forecasts = (
|
|
_parse_open_meteo_multi_model_hourly(hourly)
|
|
)
|
|
|
|
if not daily_forecasts:
|
|
logger.warning("Multi-model: 无有效模型数据")
|
|
record_source_call("open_meteo", "multi_model", "empty", (time.perf_counter() - started) * 1000.0)
|
|
return None
|
|
|
|
# 今天的预报 (向后兼容)
|
|
today_date = dates[0] if dates else None
|
|
forecasts = daily_forecasts.get(today_date, {})
|
|
|
|
labels_str = ", ".join([f"{k}={v}" for k, v in forecasts.items()])
|
|
logger.info(
|
|
f"🔬 Multi-model ({len(forecasts)}个, {len(daily_forecasts)}天): {labels_str}"
|
|
)
|
|
|
|
result = {
|
|
"source": "multi_model",
|
|
"provider": "open-meteo",
|
|
"forecasts": forecasts, # 今天 {"ECMWF": 12.3, "GFS": 11.8, ...} (向后兼容)
|
|
"daily_forecasts": daily_forecasts, # 按天 {"2026-02-23": {...}, "2026-02-24": {...}}
|
|
"hourly_times": hourly_times,
|
|
"hourly_forecasts": hourly_forecasts,
|
|
"model_metadata": model_metadata,
|
|
"model_keys": model_keys,
|
|
"dates": dates,
|
|
"unit": "fahrenheit" if use_fahrenheit else "celsius",
|
|
"attribution": "Open-Meteo forecast model API; underlying models from ECMWF, DWD, ECCC, NOAA and JMA.",
|
|
}
|
|
with self._multi_model_cache_lock:
|
|
previous = self._multi_model_cache.get(cache_key)
|
|
previous_data = previous.get("data") if isinstance(previous, dict) else None
|
|
if isinstance(previous_data, dict):
|
|
result = _merge_multi_model_result_with_cache(previous_data, result)
|
|
if result.get("partial_refresh"):
|
|
logger.info(
|
|
"🔬 Multi-model sparse refresh preserved cache: fresh={} cached={}",
|
|
result.get("partial_refresh_model_count"),
|
|
result.get("preserved_cached_model_count"),
|
|
)
|
|
with self._multi_model_cache_lock:
|
|
self._multi_model_cache[cache_key] = {
|
|
"t": time.time(),
|
|
"data": dict(result),
|
|
}
|
|
self._flush_open_meteo_disk_cache()
|
|
record_source_call("open_meteo", "multi_model", "success", (time.perf_counter() - started) * 1000.0)
|
|
return result
|
|
except Exception as e:
|
|
status_code = getattr(getattr(e, "response", None), "status_code", None)
|
|
if status_code == 429:
|
|
retry_after_str = getattr(e.response, "headers", {}).get("Retry-After")
|
|
cooldown_to_use = self._open_meteo_rl_cooldown
|
|
if retry_after_str:
|
|
try:
|
|
parsed = int(retry_after_str)
|
|
if parsed > 0:
|
|
cooldown_to_use = min(parsed + 60, 3600)
|
|
except ValueError:
|
|
pass
|
|
logger.warning(
|
|
f"Multi-model API rate limited (429), fallback to cache if available: lat={lat}, lon={lon}"
|
|
)
|
|
with self._open_meteo_rl_lock:
|
|
self._set_open_meteo_rate_limit_until(time.time() + cooldown_to_use, reason="multi_model_429")
|
|
else:
|
|
logger.warning(f"Multi-model API 请求失败: {e}")
|
|
with self._multi_model_cache_lock:
|
|
stale = self._multi_model_cache.get(cache_key)
|
|
if stale and isinstance(stale.get("data"), dict):
|
|
fallback = dict(stale["data"])
|
|
fallback["stale_cache"] = True
|
|
record_source_call("open_meteo", "multi_model", "stale_cache", (time.perf_counter() - started) * 1000.0)
|
|
return fallback
|
|
record_source_call("open_meteo", "multi_model", "error", (time.perf_counter() - started) * 1000.0)
|
|
return None
|
|
|