From 6f0eaad6a8b70341cf10efdf6f8add12c06faf20 Mon Sep 17 00:00:00 2001
From: AmandaloveYang <2569718930@qq.com>
Date: Sat, 21 Feb 2026 10:37:37 +0800
Subject: [PATCH] feat: Add a multi-source weather data collection module
supporting OpenWeatherMap, Visual Crossing, and METAR.
---
bot_listener.py | 116 +++++++++++++++++++++++++
src/data_collection/weather_sources.py | 88 +++++++++++++++++++
2 files changed, 204 insertions(+)
diff --git a/bot_listener.py b/bot_listener.py
index c4175260..fadd9d6b 100644
--- a/bot_listener.py
+++ b/bot_listener.py
@@ -57,6 +57,73 @@ def analyze_weather_trend(weather_data, temp_symbol):
local_date_str = datetime.now().strftime("%Y-%m-%d")
local_hour = datetime.now().hour
+ # === 模型共识评分 ===
+ labeled_forecasts = []
+ om_today = daily.get("temperature_2m_max", [None])[0]
+ if om_today is not None:
+ labeled_forecasts.append(("OM", om_today))
+ 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"]))
+ # 集合预报中位数 (如果有)
+ ensemble = weather_data.get("ensemble", {})
+ ens_median = ensemble.get("median")
+ if ens_median is not None:
+ labeled_forecasts.append(("ENS", ens_median))
+
+ consensus_level = "unknown"
+ consensus_spread = None
+ if len(labeled_forecasts) >= 2:
+ f_values = [v for _, v in labeled_forecasts]
+ f_max = max(f_values)
+ f_min = min(f_values)
+ consensus_spread = f_max - f_min
+ f_avg = sum(f_values) / len(f_values)
+
+ # 动态阈值:华氏度场景用更大的容差
+ is_f = (temp_symbol == "°F")
+ tight_threshold = 1.5 if is_f else 0.8 # 高共识
+ mid_threshold = 3.0 if is_f else 1.5 # 中共识
+
+ parts = " | ".join([f"{name} {val}{temp_symbol}" for name, val in labeled_forecasts])
+
+ if consensus_spread <= tight_threshold:
+ consensus_level = "high"
+ insights.append(
+ f"🎯 模型共识:高 ({len(labeled_forecasts)}/{len(labeled_forecasts)}) — "
+ f"{parts},极差仅 {consensus_spread:.1f}°,预报高度一致。"
+ )
+ elif consensus_spread <= mid_threshold:
+ consensus_level = "medium"
+ insights.append(
+ f"⚖️ 模型共识:中 ({len(labeled_forecasts)}源) — "
+ f"{parts},极差 {consensus_spread:.1f}°,有轻微分歧。"
+ )
+ else:
+ consensus_level = "low"
+ # 找出最高和最低的源
+ highest = max(labeled_forecasts, key=lambda x: x[1])
+ lowest = min(labeled_forecasts, key=lambda x: x[1])
+ insights.append(
+ f"⚠️ 模型共识:低 ({len(labeled_forecasts)}源) — "
+ f"{parts},极差 {consensus_spread:.1f}°!"
+ f"{highest[0]} 最高 ({highest[1]}{temp_symbol}) vs {lowest[0]} 最低 ({lowest[1]}{temp_symbol}),不确定性大。"
+ )
+
+ # 集合预报区间 (如果有)
+ ens_p10 = ensemble.get("p10")
+ ens_p90 = ensemble.get("p90")
+ if ens_p10 is not None and ens_p90 is not None and ens_median is not None:
+ ens_range = ens_p90 - ens_p10
+ insights.append(
+ f"📊 集合预报:中位数 {ens_median}{temp_symbol},"
+ f"90% 区间 [{ens_p10}{temp_symbol} - {ens_p90}{temp_symbol}],"
+ f"波动幅度 {ens_range:.1f}°。"
+ )
+
# === 核心判断:实测是否已超预报 ===
is_breakthrough = False
if max_so_far is not None and forecast_high is not None:
@@ -341,6 +408,55 @@ def analyze_weather_trend(weather_data, temp_symbol):
except (ValueError, IndexError):
pass
+ # 11. 入场时机信号
+ hours_to_peak = first_peak_h - local_hour if local_hour < first_peak_h else 0
+
+ # 综合评分:距离峰值越近 + 共识越高 + 实测越接近预报 → 越适合入场
+ timing_score = 0
+ timing_factors = []
+
+ if is_peak_passed:
+ timing_score += 3
+ timing_factors.append("最热已过")
+ elif hours_to_peak <= 2:
+ timing_score += 2
+ timing_factors.append(f"距峰值{hours_to_peak}h")
+ elif hours_to_peak <= 4:
+ timing_score += 1
+ timing_factors.append(f"距峰值{hours_to_peak}h")
+ else:
+ timing_factors.append(f"距峰值{hours_to_peak}h")
+
+ if consensus_level == "high":
+ timing_score += 2
+ timing_factors.append("模型一致")
+ elif consensus_level == "medium":
+ timing_score += 1
+ timing_factors.append("模型小分歧")
+ else:
+ timing_factors.append("模型分歧大")
+
+ if max_so_far is not None and forecast_high is not None:
+ gap = abs(max_so_far - forecast_high)
+ if gap <= 0.5:
+ timing_score += 2
+ timing_factors.append("实测≈预报")
+ elif gap <= 1.5:
+ timing_score += 1
+ timing_factors.append(f"差{gap:.1f}°")
+ else:
+ timing_factors.append(f"差{gap:.1f}°")
+
+ factors_str = ",".join(timing_factors)
+ if timing_score >= 5:
+ insights.append(f"⏰ 入场时机:理想 — {factors_str}。不确定性低,适合下注。")
+ elif timing_score >= 3:
+ insights.append(f"⏰ 入场时机:较好 — {factors_str}。可以考虑小仓位入场。")
+ elif timing_score >= 2:
+ insights.append(f"⏰ 入场时机:谨慎 — {factors_str}。建议继续观察。")
+ else:
+ insights.append(f"⏰ 入场时机:不建议 — {factors_str}。不确定性大,等更多数据。")
+
if not insights:
return ""
diff --git a/src/data_collection/weather_sources.py b/src/data_collection/weather_sources.py
index fe872f69..19e35099 100644
--- a/src/data_collection/weather_sources.py
+++ b/src/data_collection/weather_sources.py
@@ -535,6 +535,89 @@ class WeatherDataCollector:
logger.error(f"Open-Meteo forecast failed: {e}")
return None
+ def fetch_ensemble(
+ self,
+ lat: float,
+ lon: float,
+ use_fahrenheit: bool = False,
+ ) -> Optional[Dict]:
+ """
+ 从 Open-Meteo Ensemble API 获取 51 成员集合预报
+ 用于计算预报不确定性范围(散度)
+ """
+ try:
+ url = "https://ensemble-api.open-meteo.com/v1/ensemble"
+ params = {
+ "latitude": lat,
+ "longitude": lon,
+ "daily": "temperature_2m_max",
+ "timezone": "auto",
+ "forecast_days": 3,
+ "_t": int(time.time()),
+ }
+ if use_fahrenheit:
+ params["temperature_unit"] = "fahrenheit"
+ else:
+ params["temperature_unit"] = "celsius"
+
+ response = self.session.get(
+ url,
+ params=params,
+ headers={"Cache-Control": "no-cache"},
+ timeout=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)} 个成员")
+ 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}"
+ )
+ return result
+ except Exception as e:
+ logger.warning(f"Ensemble API 请求失败: {e}")
+ return None
+
def fetch_from_meteoblue(
self,
lat: float,
@@ -807,6 +890,11 @@ class WeatherDataCollector:
nws_data = self.fetch_nws(lat, lon)
if nws_data:
results["nws"] = nws_data
+
+ # 集合预报 (所有城市通用,用于不确定性分析)
+ ens_data = self.fetch_ensemble(lat, lon, use_fahrenheit=use_fahrenheit)
+ if ens_data:
+ results["ensemble"] = ens_data
else:
# Open-Meteo 失败时,仍然尝试获取 METAR 和 NWS
metar_data = self.fetch_metar(city, use_fahrenheit=use_fahrenheit)