From 6cc16c938d19e57d060c4a0a6d7a8b41deea7e19 Mon Sep 17 00:00:00 2001
From: AmandaloveYang <2569718930@qq.com>
Date: Sun, 22 Feb 2026 09:54:47 +0800
Subject: [PATCH] feat: introduce `WeatherDataCollector` for multi-source
weather data retrieval from OpenWeatherMap, Visual Crossing, and METAR.
---
bot_listener.py | 102 +++++++++++------
src/data_collection/weather_sources.py | 149 +++++++++++++++++++++----
2 files changed, 193 insertions(+), 58 deletions(-)
diff --git a/bot_listener.py b/bot_listener.py
index 182a2484..f65fb352 100644
--- a/bot_listener.py
+++ b/bot_listener.py
@@ -37,8 +37,10 @@ def analyze_weather_trend(weather_data, temp_symbol):
forecast_highs.append(mb["today_high"])
if nws.get("today_high") is not None:
forecast_highs.append(nws["today_high"])
- if mgm.get("today_high") is not None:
- forecast_highs.append(mgm["today_high"])
+ # 加入多模型预报 (ECMWF, GFS, ICON, GEM, JMA)
+ for mv in weather_data.get("multi_model", {}).get("forecasts", {}).values():
+ if mv is not None:
+ forecast_highs.append(mv)
forecast_highs = [h for h in forecast_highs if h is not None]
# 取预报中的最高值作为风险防御基准
@@ -58,17 +60,25 @@ def analyze_weather_trend(weather_data, temp_symbol):
local_hour = datetime.now().hour
# === 模型共识评分 ===
+ # 主要来源: 多模型预报 (ECMWF, GFS, ICON, GEM, JMA)
+ multi_model = weather_data.get("multi_model", {})
+ mm_forecasts = multi_model.get("forecasts", {})
+
labeled_forecasts = []
- om_today = daily.get("temperature_2m_max", [None])[0]
- if om_today is not None:
- labeled_forecasts.append(("OM", om_today))
+ for model_name, model_val in mm_forecasts.items():
+ if model_val is not None:
+ labeled_forecasts.append((model_name, model_val))
+
+ # 额外独立源 (如有)
if mb.get("today_high") is not None:
labeled_forecasts.append(("MB", mb["today_high"]))
if nws.get("today_high") is not None:
labeled_forecasts.append(("NWS", nws["today_high"]))
- if mgm.get("today_high") is not None:
- labeled_forecasts.append(("MGM", mgm["today_high"]))
- # 集合预报数据 (仅用于不确定性区间展示,不参与共识评分,避免与 OM 双重计数)
+
+ # Open-Meteo 确定性预报(用于后续偏差检测,不重复加入共识)
+ om_today = daily.get("temperature_2m_max", [None])[0]
+
+ # 集合预报数据 (仅用于不确定性区间展示)
ensemble = weather_data.get("ensemble", {})
ens_median = ensemble.get("median")
@@ -128,20 +138,36 @@ def analyze_weather_trend(weather_data, temp_symbol):
)
# 确定性预报 vs 集合分布偏差检测
if om_today is not None:
+ actual_reached = max_so_far is not None and max_so_far >= om_today - 0.5
if om_today > ens_p90:
- delta = om_today - ens_median
- insights.append(
- f"⚡ 预报偏高警告:确定性预报 {om_today}{temp_symbol} "
- f"超过了集合 90% 上限 ({ens_p90}{temp_symbol}),"
- f"比中位数高 {delta:.1f}°。实际高温更可能接近 {ens_median}{temp_symbol}。"
- )
+ if actual_reached:
+ # 实测已达到预报值 → 确定性预报是对的,集合偏保守
+ insights.append(
+ f"✅ 预报验证:确定性预报 {om_today}{temp_symbol} 已被实测验证 "
+ f"(实测最高 {max_so_far}{temp_symbol}),集合预报偏保守。"
+ )
+ else:
+ # 还没到最高温,存在偏高风险
+ delta = om_today - ens_median
+ insights.append(
+ f"⚡ 预报偏高警告:确定性预报 {om_today}{temp_symbol} "
+ f"超过了集合 90% 上限 ({ens_p90}{temp_symbol}),"
+ f"比中位数高 {delta:.1f}°。实际高温更可能接近 {ens_median}{temp_symbol}。"
+ )
elif om_today < ens_p10:
- delta = ens_median - om_today
- insights.append(
- f"⚡ 预报偏低警告:确定性预报 {om_today}{temp_symbol} "
- f"低于集合 90% 下限 ({ens_p10}{temp_symbol}),"
- f"比中位数低 {delta:.1f}°。实际高温更可能接近 {ens_median}{temp_symbol}。"
- )
+ if max_so_far is not None and max_so_far >= ens_median:
+ # 实测已超过中位数 → 确定性预报偏低,集合更准
+ insights.append(
+ f"✅ 预报验证:实测最高 {max_so_far}{temp_symbol} "
+ f"已超过确定性预报 {om_today}{temp_symbol},集合中位数 {ens_median}{temp_symbol} 更准确。"
+ )
+ else:
+ delta = ens_median - om_today
+ insights.append(
+ f"⚡ 预报偏低警告:确定性预报 {om_today}{temp_symbol} "
+ f"低于集合 90% 下限 ({ens_p10}{temp_symbol}),"
+ f"比中位数低 {delta:.1f}°。实际高温更可能接近 {ens_median}{temp_symbol}。"
+ )
# === 核心判断:实测是否已超预报 ===
is_breakthrough = False
@@ -224,13 +250,14 @@ def analyze_weather_trend(weather_data, temp_symbol):
# 正在峰值窗口内
if is_breakthrough:
insights.append(f"🔥 极端升温:正处于最热时段,温度已经超过所有预报,还在继续往上走!")
- elif diff_max <= 0.8:
+ elif max_so_far is not None and forecast_high - max_so_far <= 0.8:
insights.append(f"⚖️ 到顶了:正处于最热时段,温度基本到位,接下来会在这个水平上下浮动。")
else:
insights.append(f"⏳ 最热时段进行中:虽然在最热时段了,但离预报最高温还差一些,继续观察。")
elif local_hour < first_peak_h:
# 还没到峰值窗口
- if diff_max > 1.2:
+ gap_to_high = forecast_high - (max_so_far if max_so_far is not None else curr_temp)
+ if gap_to_high > 1.2:
insights.append(f"📈 还在升温:离最热时段还有 {first_peak_h - local_hour} 小时,温度还会继续往上走。")
else:
insights.append(f"🌅 快到最热了:马上就要进入最热时段,温度已经接近预报高位了。")
@@ -257,7 +284,7 @@ def analyze_weather_trend(weather_data, temp_symbol):
# 4. 云层遮挡分析 (仅在升温期/峰值期有意义)
clouds = metar.get("current", {}).get("clouds", [])
- if clouds and local_hour <= last_peak_h + 1:
+ if clouds and not is_peak_passed:
main_cloud = clouds[-1]
cover = main_cloud.get("cover", "")
if cover == "OVC":
@@ -265,8 +292,7 @@ def analyze_weather_trend(weather_data, temp_symbol):
elif cover == "BKN":
insights.append(f"🌥️ 云比较多:天空大部分被云挡住了,日照不足,升温会比较慢。")
elif cover in ["SKC", "CLR", "FEW"]:
- if not is_peak_passed:
- insights.append(f"☀️ 大晴天:阳光直射,没什么云,有利于温度继续往上冲。")
+ insights.append(f"☀️ 大晴天:阳光直射,没什么云,有利于温度继续往上冲。")
# 5. 特殊天气现象
wx_desc = metar.get("current", {}).get("wx_desc")
@@ -324,12 +350,14 @@ def analyze_weather_trend(weather_data, temp_symbol):
if 315 <= wd or wd <= 45:
insights.append(f"🌬️ 吹北风({wind_source} {wd:.0f}°):从北方来的冷空气,会压制升温。")
elif 135 <= wd <= 225:
- if diff_max > 0.5 or (is_breakthrough and curr_temp >= max_so_far):
- if is_peak_passed and not is_breakthrough:
- insights.append(f"🔥 吹南风({wind_source} {wd:.0f}°):南方的暖空气还在吹过来,但最热时段已过,后劲不足了。")
- else:
- status = "温度还有继续上涨的空间" if not is_breakthrough else "可能把温度推得更高"
- insights.append(f"🔥 吹南风({wind_source} {wd:.0f}°):南方的暖空气正在吹过来,{status}。")
+ gap_to_forecast = forecast_high - (max_so_far if max_so_far is not None else curr_temp)
+ if is_peak_passed and not is_breakthrough:
+ insights.append(f"🔥 吹南风({wind_source} {wd:.0f}°):南方的暖空气还在吹过来,但最热时段已过,后劲不足了。")
+ elif gap_to_forecast > 0.5 or is_breakthrough:
+ status = "温度还有继续上涨的空间" if not is_breakthrough else "可能把温度推得更高"
+ insights.append(f"🔥 吹南风({wind_source} {wd:.0f}°):南方的暖空气正在吹过来,{status}。")
+ else:
+ insights.append(f"🔥 吹南风({wind_source} {wd:.0f}°):南方的暖空气正在吹过来,但温度已接近预报峰值。")
elif 225 < wd < 315:
if wd <= 260:
insights.append(f"🌬️ 吹西南风({wind_source} {wd:.0f}°):带有一定暖湿气流,对升温有轻微帮助。")
@@ -351,20 +379,22 @@ def analyze_weather_trend(weather_data, temp_symbol):
except (TypeError, ValueError):
pass
- # 7. 模型准确度预警
+ # 7. 模型准确度预警(使用多模型数据)
if is_peak_passed and max_so_far is not None:
model_checks = []
- if om_high and om_high > max_so_far + 1.5:
- model_checks.append(f"Open-Meteo ({om_high}{temp_symbol})")
+ for m_name, m_val in mm_forecasts.items():
+ if m_val is not None and m_val > max_so_far + 1.5:
+ model_checks.append(f"{m_name} ({m_val}{temp_symbol})")
+ # 附加源也查一下
mb_h = mb.get("today_high")
if mb_h and mb_h > max_so_far + 1.5:
- model_checks.append(f"Meteoblue ({mb_h}{temp_symbol})")
+ model_checks.append(f"MB ({mb_h}{temp_symbol})")
nws_h = nws.get("today_high")
if nws_h and nws_h > max_so_far + 1.5:
model_checks.append(f"NWS ({nws_h}{temp_symbol})")
if model_checks:
- insights.append(f"⚠️ 预报偏高了:实测远低于 " + "、".join(model_checks) + ",这些预报今天报高了。")
+ insights.append(f"⚠️ 预报偏高了:实测远低于 " + "、".join(model_checks) + ",这些模型今天报高了。")
# 8. MGM 气压分析 (仅安卡拉)
mgm_pressure = mgm.get("current", {}).get("pressure")
diff --git a/src/data_collection/weather_sources.py b/src/data_collection/weather_sources.py
index 19e35099..4aa5e73b 100644
--- a/src/data_collection/weather_sources.py
+++ b/src/data_collection/weather_sources.py
@@ -369,14 +369,33 @@ class WeatherDataCollector:
"station_name": latest.get("istasyonAd") or latest.get("adi") or latest.get("merkezAd") or "Ankara Esenboğa"
}
- # 2. 每日预报
- daily_resp = self.session.get(f"{base_url}/tahminler/gunluk?istno={istno}", headers=headers, timeout=self.timeout)
- if daily_resp.status_code == 200:
- forecasts = daily_resp.json()
- if forecasts and isinstance(forecasts, list):
- today = forecasts[0]
- results["today_high"] = today.get("enYuksekGun1")
- results["today_low"] = today.get("enDusukGun1")
+ # 2. 每日预报(尝试两个可能的 API 路径)
+ forecast_urls = [
+ f"{base_url}/tahminler/gunluk?istno={istno}",
+ f"https://servis.mgm.gov.tr/api/tahminler/gunluk?istno={istno}",
+ ]
+ for forecast_url in forecast_urls:
+ try:
+ daily_resp = self.session.get(forecast_url, headers=headers, timeout=self.timeout)
+ if daily_resp.status_code == 200:
+ forecasts = daily_resp.json()
+ if forecasts and isinstance(forecasts, list):
+ today = forecasts[0]
+ high_val = today.get("enYuksekGun1")
+ low_val = today.get("enDusukGun1")
+ if high_val is not None:
+ results["today_high"] = high_val
+ results["today_low"] = low_val
+ logger.info(f"📋 MGM 每日预报: 最高 {high_val}°C, 最低 {low_val}°C (from {forecast_url})")
+ break
+ else:
+ # 记录所有可用字段,方便调试
+ available_keys = [k for k in today.keys() if "yuksek" in k.lower() or "sicaklik" in k.lower() or "gun" in k.lower()]
+ logger.warning(f"MGM 每日预报: enYuksekGun1 为空,可用字段: {available_keys}")
+ else:
+ logger.debug(f"MGM forecast URL {forecast_url} returned {daily_resp.status_code}")
+ except Exception as e:
+ logger.debug(f"MGM forecast URL {forecast_url} failed: {e}")
return results if "current" in results else None
except Exception as e:
@@ -618,6 +637,86 @@ class WeatherDataCollector:
logger.warning(f"Ensemble API 请求失败: {e}")
return None
+ def fetch_multi_model(
+ self,
+ lat: float,
+ lon: float,
+ use_fahrenheit: bool = False,
+ ) -> Optional[Dict]:
+ """
+ 从 Open-Meteo 获取多个独立 NWP 模型的预报
+ 用于真正的多模型共识评分
+
+ 模型列表:
+ - ECMWF IFS (欧洲中期天气预报中心)
+ - GFS (美国 NOAA)
+ - ICON (德国气象局 DWD)
+ - GEM (加拿大气象局)
+ - JMA (日本气象厅)
+ """
+ try:
+ url = "https://api.open-meteo.com/v1/forecast"
+ models = "ecmwf_ifs025,gfs_seamless,icon_seamless,gem_seamless,jma_seamless"
+ params = {
+ "latitude": lat,
+ "longitude": lon,
+ "daily": "temperature_2m_max",
+ "models": models,
+ "timezone": "auto",
+ "forecast_days": 1,
+ "_t": int(time.time()),
+ }
+ if use_fahrenheit:
+ params["temperature_unit"] = "fahrenheit"
+
+ response = self.session.get(
+ url,
+ params=params,
+ headers={"Cache-Control": "no-cache"},
+ timeout=self.timeout,
+ )
+ response.raise_for_status()
+ data = response.json()
+
+ # Open-Meteo 多模型返回格式:
+ # "daily": {
+ # "temperature_2m_max_ecmwf_ifs025": [12.3],
+ # "temperature_2m_max_gfs_seamless": [11.8],
+ # ...
+ # }
+ daily = data.get("daily", {})
+
+ model_labels = {
+ "ecmwf_ifs025": "ECMWF",
+ "gfs_seamless": "GFS",
+ "icon_seamless": "ICON",
+ "gem_seamless": "GEM",
+ "jma_seamless": "JMA",
+ }
+
+ forecasts = {}
+ for model_key, label in model_labels.items():
+ key = f"temperature_2m_max_{model_key}"
+ values = daily.get(key, [])
+ if values and values[0] is not None:
+ forecasts[label] = round(values[0], 1)
+
+ if not forecasts:
+ logger.warning("Multi-model: 无有效模型数据")
+ return None
+
+ labels_str = ", ".join([f"{k}={v}" for k, v in forecasts.items()])
+ logger.info(f"🔬 Multi-model ({len(forecasts)}个): {labels_str}")
+
+ return {
+ "source": "multi_model",
+ "forecasts": forecasts, # {"ECMWF": 12.3, "GFS": 11.8, ...}
+ "unit": "fahrenheit" if use_fahrenheit else "celsius",
+ }
+ except Exception as e:
+ logger.warning(f"Multi-model API 请求失败: {e}")
+ return None
+
def fetch_from_meteoblue(
self,
lat: float,
@@ -725,22 +824,23 @@ class WeatherDataCollector:
"""
使用 Open-Meteo Geocoding API 获取城市坐标 (免费, 无需 Key)
"""
- # 预设常用城市坐标,避免网络波动导致启动失败
+ # 坐标使用 METAR 机场位置(Polymarket 以机场数据结算)
static_coords = {
- "london": {"lat": 51.5074, "lon": -0.1278},
- "new york": {"lat": 40.7128, "lon": -74.0060},
+ "london": {"lat": 51.5053, "lon": 0.0553}, # EGLC London City
+ "paris": {"lat": 49.0097, "lon": 2.5478}, # LFPG Charles de Gaulle
+ "new york": {"lat": 40.7750, "lon": -73.8750}, # KLGA LaGuardia
"new york's central park": {"lat": 40.7812, "lon": -73.9665},
- "nyc": {"lat": 40.7128, "lon": -74.0060},
- "seattle": {"lat": 47.6062, "lon": -122.3321},
- "chicago": {"lat": 41.8781, "lon": -87.6298},
- "dallas": {"lat": 32.7767, "lon": -96.7970},
- "miami": {"lat": 25.7617, "lon": -80.1918},
- "atlanta": {"lat": 33.7490, "lon": -84.3880},
- "seoul": {"lat": 37.5665, "lon": 126.9780},
- "toronto": {"lat": 43.6532, "lon": -79.3832},
- "ankara": {"lat": 39.9334, "lon": 32.8597},
- "wellington": {"lat": -41.2865, "lon": 174.7762},
- "buenos aires": {"lat": -34.6037, "lon": -58.3816},
+ "nyc": {"lat": 40.7750, "lon": -73.8750}, # KLGA LaGuardia
+ "seattle": {"lat": 47.4499, "lon": -122.3118}, # KSEA Sea-Tac
+ "chicago": {"lat": 41.9769, "lon": -87.9081}, # KORD O'Hare
+ "dallas": {"lat": 32.8459, "lon": -96.8509}, # KDAL Love Field
+ "miami": {"lat": 25.7933, "lon": -80.2906}, # KMIA International
+ "atlanta": {"lat": 33.6367, "lon": -84.4281}, # KATL Hartsfield-Jackson
+ "seoul": {"lat": 37.4691, "lon": 126.4510}, # RKSI Incheon
+ "toronto": {"lat": 43.6759, "lon": -79.6294}, # CYYZ Pearson
+ "ankara": {"lat": 40.1281, "lon": 32.9950}, # LTAC Esenboğa
+ "wellington": {"lat": -41.3272, "lon": 174.8053}, # NZWN Wellington
+ "buenos aires": {"lat": -34.8222, "lon": -58.5358}, # SAEZ Ezeiza
}
normalized_city = city.lower().strip()
@@ -895,6 +995,11 @@ class WeatherDataCollector:
ens_data = self.fetch_ensemble(lat, lon, use_fahrenheit=use_fahrenheit)
if ens_data:
results["ensemble"] = ens_data
+
+ # 多模型预报 (所有城市通用,用于共识评分)
+ mm_data = self.fetch_multi_model(lat, lon, use_fahrenheit=use_fahrenheit)
+ if mm_data:
+ results["multi_model"] = mm_data
else:
# Open-Meteo 失败时,仍然尝试获取 METAR 和 NWS
metar_data = self.fetch_metar(city, use_fahrenheit=use_fahrenheit)