Add Moscow Russia station web nearby provider with stale filtering

This commit is contained in:
2569718930@qq.com
2026-04-10 18:49:50 +08:00
parent 373e7cfec6
commit 95e1fd50d2
4 changed files with 450 additions and 1 deletions
+54
View File
@@ -41,6 +41,8 @@ def _provider_code_for_city(city: str) -> str:
return "turkey_mgm"
if normalized in {"busan", "seoul"}:
return "korea_kma"
if normalized == "moscow":
return "russia_station_web"
if settlement_source == "hko":
return "hongkong_hko"
if settlement_source == "cwa":
@@ -235,6 +237,34 @@ def _kma_rows(raw: Dict[str, Any], city: str) -> List[Dict[str, Any]]:
return out
def _ru_rows(raw: Dict[str, Any], city: str) -> List[Dict[str, Any]]:
rows = raw.get("ru_official_nearby") or []
out: List[Dict[str, Any]] = []
for row in rows:
if not isinstance(row, dict):
continue
out.append(
_normalize_station_row(
station_code=row.get("station_code") or row.get("icao") or row.get("istNo"),
station_label=row.get("station_label") or row.get("name"),
temp=row.get("temp"),
lat=row.get("lat"),
lon=row.get("lon"),
obs_time=row.get("obs_time"),
source_code="ru_station_web",
source_label="Russia station web",
is_official=True,
is_airport_station=_bool(row.get("is_airport_station")),
is_settlement_anchor=False,
extra={
"distance_km": _safe_float(row.get("distance_km")),
"page_url": row.get("page_url"),
},
)
)
return out
def _mgm_rows(raw: Dict[str, Any], city: str) -> List[Dict[str, Any]]:
meta = _city_meta(city)
rows = raw.get("mgm_nearby") or []
@@ -497,6 +527,28 @@ class KoreaKmaNetworkProvider(CountryNetworkProvider):
}
class RussiaStationWebNetworkProvider(CountryNetworkProvider):
def __init__(self) -> None:
super().__init__("russia_station_web", "Russia station web")
def official_nearby_current(self, city: str, raw: Dict[str, Any]) -> List[Dict[str, Any]]:
rows = _ru_rows(raw, city)
if rows:
return rows
return _metar_cluster_rows(raw)
def official_network_status(self, city: str, raw: Dict[str, Any]) -> Dict[str, Any]:
rows = self.official_nearby_current(city, raw)
has_ru = bool(_ru_rows(raw, city))
return {
"provider_code": self.provider_code,
"provider_label": self.provider_label,
"available": has_ru,
"mode": "official_web_crawl" if has_ru else ("fallback_metar_cluster" if rows else "reference_only"),
"row_count": len(rows),
}
class HongKongHkoNetworkProvider(CountryNetworkProvider):
def __init__(self) -> None:
super().__init__("hongkong_hko", "HKO")
@@ -521,6 +573,8 @@ def get_country_network_provider(city: str) -> CountryNetworkProvider:
return TurkeyMgmNetworkProvider()
if provider_code == "korea_kma":
return KoreaKmaNetworkProvider()
if provider_code == "russia_station_web":
return RussiaStationWebNetworkProvider()
if provider_code == "japan_jma":
return JapanJmaNetworkProvider()
if provider_code == "china_cma":
@@ -0,0 +1,331 @@
from __future__ import annotations
import html
import math
import re
import time
from datetime import datetime, timezone
from typing import Any, Dict, List, Optional
from loguru import logger
from src.utils.metrics import record_source_call
RUSSIA_MOSCOW_STATIONS: Dict[str, Dict[str, Any]] = {
"27524": {
"station_code": "27524",
"station_label": "Vnukovo",
"lat": 55.5870,
"lon": 37.2500,
},
"27518": {
"station_code": "27518",
"station_label": "Podmoskovnaya",
"lat": 55.7084,
"lon": 37.1823,
},
"27515": {
"station_code": "27515",
"station_label": "Nemchinovka",
"lat": 55.7065,
"lon": 37.3719,
},
"27504": {
"station_code": "27504",
"station_label": "Moscow (Butovo)",
"lat": 55.5780,
"lon": 37.5541,
},
"27416": {
"station_code": "27416",
"station_label": "Moscow (Strogino)",
"lat": 55.7976,
"lon": 37.3982,
},
"27605": {
"station_code": "27605",
"station_label": "Moscow (Balchug)",
"lat": 55.7455,
"lon": 37.6300,
},
}
class RussiaStationSourceMixin:
def _ru_http_get_text(self, url: str) -> str:
getter = getattr(self, "_http_get", None)
if callable(getter):
response = getter(url)
else:
response = self.session.get(url, timeout=self.timeout)
response.raise_for_status()
return response.text
@staticmethod
def _ru_safe_float(value: Any) -> Optional[float]:
try:
if value in (None, "", "-", ""):
return None
text = str(value).strip().replace(",", ".")
return float(text)
except Exception:
return None
@staticmethod
def _ru_distance_km(
lat1: Optional[float],
lon1: Optional[float],
lat2: Optional[float],
lon2: Optional[float],
) -> Optional[float]:
if None in (lat1, lon1, lat2, lon2):
return None
try:
r = 6371.0
d_lat = math.radians(float(lat2) - float(lat1))
d_lon = math.radians(float(lon2) - float(lon1))
a = (
math.sin(d_lat / 2) ** 2
+ math.cos(math.radians(float(lat1)))
* math.cos(math.radians(float(lat2)))
* math.sin(d_lon / 2) ** 2
)
c = 2 * math.atan2(math.sqrt(a), math.sqrt(1 - a))
return round(r * c, 2)
except Exception:
return None
@staticmethod
def _ru_clean_cell(cell_html: str) -> str:
text = re.sub(r"<[^>]+>", " ", str(cell_html or ""))
text = html.unescape(text)
text = text.replace("\xa0", " ")
text = re.sub(r"\s+", " ", text)
return text.strip()
@classmethod
def _ru_parse_table_rows(cls, table_html: str) -> List[List[str]]:
rows: List[List[str]] = []
for row_html in re.findall(r"<tr[^>]*>(.*?)</tr>", table_html, flags=re.S | re.I):
cells = re.findall(r"<t[dh][^>]*>(.*?)</t[dh]>", row_html, flags=re.S | re.I)
cleaned = [cls._ru_clean_cell(cell) for cell in cells]
if cleaned:
rows.append(cleaned)
return rows
@staticmethod
def _ru_build_obs_time(hour_text: str, day_month_text: str) -> Optional[str]:
hour_match = re.search(r"(\d{1,2})", str(hour_text or ""))
day_match = re.search(r"(\d{1,2})\.(\d{1,2})", str(day_month_text or ""))
if not hour_match or not day_match:
return None
hour = int(hour_match.group(1))
day = int(day_match.group(1))
month = int(day_match.group(2))
now_utc = datetime.now(timezone.utc)
year = now_utc.year
try:
candidate = datetime(year, month, day, hour, 0, tzinfo=timezone.utc)
except ValueError:
return None
if candidate > now_utc and (candidate - now_utc).days > 40:
candidate = datetime(year - 1, month, day, hour, 0, tzinfo=timezone.utc)
return candidate.isoformat()
def _ru_parse_station_current_from_weather_html(self, html_text: str) -> Optional[Dict[str, Any]]:
tables = re.findall(r"<table[^>]*>(.*?)</table>", html_text, flags=re.S | re.I)
if len(tables) < 2:
return None
time_rows = self._ru_parse_table_rows(tables[0])
data_rows = self._ru_parse_table_rows(tables[1])
if len(time_rows) < 2 or len(data_rows) < 2:
return None
pair_count = min(len(time_rows), len(data_rows)) - 1
for idx in range(pair_count):
time_row = time_rows[idx + 1]
data_row = data_rows[idx + 1]
if len(time_row) < 2 or len(data_row) < 6:
continue
temp_c = self._ru_safe_float(data_row[5])
if temp_c is None:
continue
obs_time = self._ru_build_obs_time(time_row[0], time_row[1])
return {
"temp_c": round(temp_c, 1),
"obs_time": obs_time,
"raw_hour": time_row[0],
"raw_day_month": time_row[1],
}
return None
def _ru_cached_station_current(
self,
station_code: str,
station_meta: Dict[str, Any],
use_fahrenheit: bool = False,
) -> Optional[Dict[str, Any]]:
cache_key = f"{station_code}:{use_fahrenheit}"
now_ts = time.time()
with self._ru_station_cache_lock:
cached = self._ru_station_cache.get(cache_key)
if cached and now_ts - cached["t"] < self.ru_station_cache_ttl_sec:
return cached["d"]
started = time.perf_counter()
try:
url = f"https://www.pogodaiklimat.ru/weather.php?id={station_code}"
html_text = self._ru_http_get_text(url)
parsed = self._ru_parse_station_current_from_weather_html(html_text)
if not parsed:
record_source_call(
"ru_station_web",
"current",
"empty",
(time.perf_counter() - started) * 1000.0,
)
return None
obs_time = parsed.get("obs_time")
max_stale_sec = max(
0,
int(getattr(self, "ru_station_max_stale_sec", 72 * 3600)),
)
if obs_time and max_stale_sec > 0:
try:
obs_dt = datetime.fromisoformat(str(obs_time))
obs_age_sec = (datetime.now(timezone.utc) - obs_dt).total_seconds()
if obs_age_sec > max_stale_sec:
logger.info(
"Russia station web row is stale station={} obs_time={} age_hours={:.1f}",
station_code,
obs_time,
obs_age_sec / 3600.0,
)
record_source_call(
"ru_station_web",
"current",
"stale_row",
(time.perf_counter() - started) * 1000.0,
)
return None
except Exception:
pass
temp_c = parsed.get("temp_c")
if temp_c is None:
record_source_call(
"ru_station_web",
"current",
"no_temperature",
(time.perf_counter() - started) * 1000.0,
)
return None
temp = round(temp_c * 9 / 5 + 32, 1) if use_fahrenheit else round(temp_c, 1)
result = {
"station_code": station_code,
"station_label": station_meta.get("station_label") or f"RU {station_code}",
"name": station_meta.get("station_label") or f"RU {station_code}",
"lat": station_meta.get("lat"),
"lon": station_meta.get("lon"),
"temp": temp,
"obs_time": parsed.get("obs_time"),
"source": "ru_station_web",
"source_label": "Russia station web",
"source_code": "ru_station_web",
"is_official": True,
"is_airport_station": station_code == "27524",
"is_settlement_anchor": False,
"page_url": f"https://www.pogodaiklimat.ru/weather.php?id={station_code}",
}
with self._ru_station_cache_lock:
self._ru_station_cache[cache_key] = {"d": result, "t": now_ts}
record_source_call(
"ru_station_web",
"current",
"success",
(time.perf_counter() - started) * 1000.0,
)
return result
except Exception as exc:
logger.warning("Russia station web fetch failed station={} error={}", station_code, exc)
with self._ru_station_cache_lock:
stale = self._ru_station_cache.get(cache_key)
if stale:
record_source_call(
"ru_station_web",
"current",
"stale_cache",
(time.perf_counter() - started) * 1000.0,
)
return stale["d"]
record_source_call(
"ru_station_web",
"current",
"error",
(time.perf_counter() - started) * 1000.0,
)
return None
def fetch_russia_moscow_official_nearby(
self,
city: str,
use_fahrenheit: bool = False,
) -> List[Dict[str, Any]]:
started = time.perf_counter()
city_key = str(city or "").strip().lower()
if city_key != "moscow":
record_source_call(
"ru_station_web",
"nearby",
"unsupported_city",
(time.perf_counter() - started) * 1000.0,
)
return []
city_meta = self.CITY_REGISTRY.get(city_key) or {}
anchor_lat = self._ru_safe_float(city_meta.get("lat"))
anchor_lon = self._ru_safe_float(city_meta.get("lon"))
rows: List[Dict[str, Any]] = []
try:
for station_code, station_meta in RUSSIA_MOSCOW_STATIONS.items():
current = self._ru_cached_station_current(
station_code,
station_meta,
use_fahrenheit=use_fahrenheit,
)
if not current:
continue
row = dict(current)
row["distance_km"] = self._ru_distance_km(
anchor_lat,
anchor_lon,
self._ru_safe_float(row.get("lat")),
self._ru_safe_float(row.get("lon")),
)
row["icao"] = station_code
row["istNo"] = station_code
rows.append(row)
rows.sort(
key=lambda item: (
item.get("distance_km") is None,
item.get("distance_km") if item.get("distance_km") is not None else 9999,
item.get("station_label") or "",
)
)
trimmed = rows[:6]
record_source_call(
"ru_station_web",
"nearby",
"success" if trimmed else "empty",
(time.perf_counter() - started) * 1000.0,
)
return trimmed
except Exception as exc:
logger.warning("Russia station nearby fetch failed city={} error={}", city_key, exc)
record_source_call(
"ru_station_web",
"nearby",
"error",
(time.perf_counter() - started) * 1000.0,
)
return []
+27 -1
View File
@@ -12,11 +12,12 @@ from src.data_collection.metar_sources import MetarSourceMixin
from src.data_collection.mgm_sources import MgmSourceMixin
from src.data_collection.kma_station_sources import KmaStationSourceMixin
from src.data_collection.jma_amedas_sources import JmaAmedasSourceMixin
from src.data_collection.russia_station_sources import RussiaStationSourceMixin
from src.data_collection.nmc_sources import NmcSourceMixin
from src.data_collection.nws_open_meteo_sources import NwsOpenMeteoSourceMixin
class WeatherDataCollector(OpenMeteoCacheMixin, SettlementSourceMixin, MetarSourceMixin, MgmSourceMixin, KmaStationSourceMixin, JmaAmedasSourceMixin, NmcSourceMixin, NwsOpenMeteoSourceMixin):
class WeatherDataCollector(OpenMeteoCacheMixin, SettlementSourceMixin, MetarSourceMixin, MgmSourceMixin, KmaStationSourceMixin, JmaAmedasSourceMixin, RussiaStationSourceMixin, NmcSourceMixin, NwsOpenMeteoSourceMixin):
"""
Multi-source weather data collector
@@ -179,6 +180,14 @@ class WeatherDataCollector(OpenMeteoCacheMixin, SettlementSourceMixin, MetarSour
)
self._kma_cache: Dict[str, Dict] = {}
self._kma_cache_lock = threading.Lock()
self.ru_station_cache_ttl_sec = int(
os.getenv("RU_STATION_CACHE_TTL_SEC", "300")
)
self.ru_station_max_stale_sec = int(
os.getenv("RU_STATION_MAX_STALE_SEC", str(72 * 3600))
)
self._ru_station_cache: Dict[str, Dict] = {}
self._ru_station_cache_lock = threading.Lock()
self.settlement_cache_ttl_sec = int(
os.getenv("SETTLEMENT_SOURCE_CACHE_TTL_SEC", "120")
)
@@ -798,6 +807,21 @@ class WeatherDataCollector(OpenMeteoCacheMixin, SettlementSourceMixin, MetarSour
results["mgm_nearby"] = official_rows
results["nearby_source"] = "kma"
def _attach_russia_official_nearby(
self, results: Dict, city_lower: str, use_fahrenheit: bool
) -> None:
if city_lower != "moscow":
return
official_rows = self.fetch_russia_moscow_official_nearby(
city_lower, use_fahrenheit=use_fahrenheit
)
if not official_rows:
return
results["ru_official_nearby"] = official_rows
if "mgm_nearby" not in results:
results["mgm_nearby"] = official_rows
results["nearby_source"] = "ru_station_web"
def _attach_warsaw_official_nearby(
self, results: Dict, use_fahrenheit: bool
) -> None:
@@ -895,6 +919,7 @@ class WeatherDataCollector(OpenMeteoCacheMixin, SettlementSourceMixin, MetarSour
self._attach_china_official_nearby(results, city_lower, use_fahrenheit)
self._attach_japan_official_nearby(results, city_lower, use_fahrenheit)
self._attach_korea_official_nearby(results, city_lower, use_fahrenheit)
self._attach_russia_official_nearby(results, city_lower, use_fahrenheit)
if city_lower == "warsaw":
self._attach_warsaw_official_nearby(results, use_fahrenheit)
self._attach_global_nearby_cluster(
@@ -924,6 +949,7 @@ class WeatherDataCollector(OpenMeteoCacheMixin, SettlementSourceMixin, MetarSour
self._attach_china_official_nearby(results, city_lower, use_fahrenheit)
self._attach_japan_official_nearby(results, city_lower, use_fahrenheit)
self._attach_korea_official_nearby(results, city_lower, use_fahrenheit)
self._attach_russia_official_nearby(results, city_lower, use_fahrenheit)
if city_lower == "warsaw":
self._attach_warsaw_official_nearby(results, use_fahrenheit)
self._attach_global_nearby_cluster(
+38
View File
@@ -108,6 +108,44 @@ def test_hko_provider_marks_explicit_official_station_as_anchor():
assert snapshot["official_nearby"][0]["station_code"] == "LFS"
def test_russia_provider_prefers_official_web_rows_when_available():
raw = {
"metar": {
"observation_time": "2026-04-06T10:00:00.000Z",
"current": {"temp": 11.0},
},
"ru_official_nearby": [
{
"station_code": "27524",
"station_label": "Vnukovo",
"lat": 55.5870,
"lon": 37.2500,
"temp": 12.3,
"obs_time": "2026-04-06T09:00:00+00:00",
"is_airport_station": True,
"page_url": "https://www.pogodaiklimat.ru/weather.php?id=27524",
}
],
"mgm_nearby": [
{
"name": "Sheremetyevo",
"icao": "UUEE",
"lat": 55.97,
"lon": 37.41,
"temp": 12.0,
}
],
}
snapshot = build_country_network_snapshot("moscow", raw)
assert snapshot["provider_code"] == "russia_station_web"
assert snapshot["official_network_status"]["available"] is True
assert snapshot["official_network_status"]["mode"] == "official_web_crawl"
assert snapshot["official_nearby"][0]["source_code"] == "ru_station_web"
assert snapshot["official_nearby"][0]["is_official"] is True
def test_city_detail_payload_exposes_airport_and_official_network_layers():
payload = _build_city_detail_payload(
{