砍掉机场快照,改为 per-city 独立告警;新增安卡拉 MGM 17128 高频监控

快照定时推送无实际价值(没变化也报),改为仅温度急变时触发告警。
各城市独立检测、独立冷却,不再捆绑推送。
同时接入安卡拉 Esenboğa 机场 MGM 站点 17128 的实时温度数据。

Constraint: 快照已运行验证格式无误,砍掉不影响现有急变告警功能
Scope-risk: 中低,仅影响高频通道逻辑,主循环不变
Tested: ruff check 通过
@
This commit is contained in:
2569718930@qq.com
2026-05-12 17:51:40 +08:00
parent 5a6a487a97
commit 12b0c76caf
4 changed files with 46 additions and 156 deletions
+29 -55
View File
@@ -4,11 +4,12 @@
### 现有数据
| 城市 | 机场 | ICAO | 数据源 | 刷新频率 | 缓存 TTL | 数据字段 |
|------|------|------|--------|---------|----------|---------|
| 城市 | 机场 | ICAO/站点 | 数据源 | 刷新频率 | 缓存 TTL | 数据字段 |
|------|------|-----------|--------|---------|----------|---------|
| 首尔 | 仁川国际 | RKSI | AMOS (`global.amo.go.kr`) | 1 分钟 | 60s | 温度/露点/风速/气压/能见度/RVR/跑道对温度/METAR/TAF |
| 釜山 | 金海国际 | RKPK | AMOS (`global.amo.go.kr`) | 1 分钟 | 60s | 同上 |
| 东京 | 羽田 | RJTT | JMA AMeDAS (`jma.go.jp`) | 10 分钟 | 120s | 10 分钟温度观测(当前仅提取 temp |
| 东京 | 羽田 | RJTT | JMA AMeDAS (`jma.go.jp`) | 10 分钟 | 120s | 10 分钟温度观测 |
| 安卡拉 | Esenboğa | 17128 | MGM (`servis.mgm.gov.tr`) | ~10 分钟 | 无独立缓存 | 温度/风速/气压/湿度/降水 |
### 现有 Telegram 推送系统
@@ -20,9 +21,7 @@
### 问题
高频机场数据已就绪,但现有推送系统 30 分钟一轮对所有城市一视同仁。首尔/釜山 1 分钟级 AMOS 和东京 10 分钟级 JMA 接近交易高峰期时,温度变化可能比 30 分钟窗口更快,需要更灵敏的监控。
同时,DEB 每日最高温预测是结算的核心参照,将机场实时温度与 DEB 预测并排对比,可以快速发现实际温度偏离预测的程度。
四座机场城市的实时数据已就绪,但现有推送系统 30 分钟一轮对所有城市一视同仁。1-10 分钟级高频数据在接近交易高峰期时,温度变化可能比 30 分钟窗口更快,需要更灵敏的监控。
---
@@ -30,15 +29,15 @@
### 核心思路
**在现有 30 分钟主循环之上叠加高频通道**,对三大机场城市用更短间隔推送"实时温度 + DEB 预测最高温"快报。不做市场分析、不输出 AI 建议,只报温度数据本身
在现有 30 分钟主循环之上叠加高频通道,对四座机场城市用 10 分钟间隔独立检测温度急变。温度波动达到阈值时推送告警,包含当前温度 + DEB 预测最高温。不做市场分析、不输出 AI 建议、不约定时快照
### 1. 高频机场城市快速通道
在现有 30 分钟主循环之外,为 `{seoul, busan, tokyo}` 单独跑一个 10 分钟间隔的子循环,仅检查温度动量突变。
在现有 30 分钟主循环之外,为 `{seoul, busan, tokyo, ankara}` 单独跑一个 10 分钟间隔的子循环,每个城市独立检测温度急变。
**配置(写死在代码中)**:
```python
HIGH_FREQ_AIRPORT_CITIES = {"seoul", "busan", "tokyo"}
HIGH_FREQ_AIRPORT_CITIES = {"seoul", "busan", "tokyo", "ankara"}
HIGH_FREQ_PUSH_INTERVAL_SEC = 600 # 10 分钟
HIGH_FREQ_MOMENTUM_THRESHOLD_C = 0.5 # 比默认 0.8°C 更灵敏
HIGH_FREQ_COOLDOWN_SEC = 7200 # 同一城市冷却 2 小时
@@ -46,15 +45,13 @@ HIGH_FREQ_COOLDOWN_SEC = 7200 # 同一城市冷却 2 小时
**逻辑**:
- 主循环 30 分钟照常跑全部城市(不变)
- 每 10 分钟额外跑一轮高频城市子集
- 高频轮次仅检查温度动量突变一条规则
- 高频告警独立冷却期 2 小时
- **最高温已锁定则跳过**:当日最高温已出现且之后持续下降,不再推送
- 每 10 分钟对四座机场城市各检查一次温度急变
- 高频轮次仅检查 `airport_rapid_temp_change` 一条规则
- 各城市独立冷却,触发后 2 小时内同一城市不再重复推送
- **最高温已锁定则跳过**:当日最高已过且持续下降,不再推送
### 2. 机场观测积累与趋势检测
**问题**: 当前 AMOS/JMA 每次只返回最新一条观测,无法计算温度斜率做动量检测。
**新增数据库表**: `airport_obs_log`
```sql
CREATE TABLE IF NOT EXISTS airport_obs_log (
@@ -71,52 +68,29 @@ CREATE INDEX IF NOT EXISTS idx_airport_obs_log_icao_time
ON airport_obs_log(icao, created_at DESC);
```
**写入**: 在 AMOS/JMA 成功获取数据后自动调用 `append_airport_obs()` 写入。自动清理 2 小时前的旧数据。
**写入**: 在 AMOS/JMA/MGM 成功获取数据后自动调用 `append_airport_obs()` 写入。自动清理 2 小时前的旧数据。
**读取**: `get_airport_obs_recent(icao, minutes=30)` 返回最近 N 分钟观测列表,用于计算温度变化斜率。
### 3. 温度突变即时告警
基于积累的观测日志,新增告警规则 `airport_rapid_temp_change`
基于积累的观测日志,新增告警规则 `airport_rapid_temp_change`每条告警 per-city 独立推送。
| 参数 | 值 | 说明 |
|------|-----|------|
| 滑动窗口 | 20 分钟 | 取最近 20 分钟内的观测 |
| 最少样本 | 3 条 | 确保有足够数据点 |
| 触发阈值 | > 0.5°C/10min | 高频通道专用,比默认 0.8°C 更灵敏 |
| 触发阈值 | > 0.5°C/10min | 比默认 0.8°C/30min 更灵敏 |
| 冷却期 | 2 小时 | 同城市两次推送最小间隔 |
| 锁定跳过 | 最高温已锁定 | 当日最高已过且持续下降,不推送 |
**告警消息示例**:
```
🚨 首尔/仁川 温度急变
🚨 首尔/仁川 温度急变 🔥
跑道中位数 14.9°C → 16.5°C (+1.6°C / 15min)
220° 14kt 暖平流增强
DEB 预测今日最高 18.2°C | 目前距预测差 1.7°C
```
### 4. 机场实况快照摘要(含 DEB 预测)
每隔 30 分钟,向市场监控频道推送三座机场的"当前温度 + DEB 预测最高温"摘要。
**推送时段**: 仅各城市当地 08:00-20:00(避免半夜噪音)
**锁定跳过**: 如该城市当日最高温已锁定,跳过该城市不推送(三城中仍有未锁定的则正常推送其余)
**消息格式**:
```
🛫 机场实况 14:30 CST
🇰🇷 首尔/仁川 RKSI
跑道中位数 14.9°C · DEB 预测最高 18.2°C
风 220° 14kt | QNH 1015.2 hPa | 能见度 ≥10km
🇰🇷 釜山/金海 RKPK
跑道中位数 13.8°C · DEB 预测最高 16.5°C
风 180° 8kt
🇯🇵 东京/羽田 RJTT
当前 12.4°C (14:20 JST) · DEB 预测最高 15.1°C
跑道温度 14.9°C → 16.5°C (+1.6°C / 15min)
14kt
DEB 预测今日最高 18.2°C | 差距 +1.7°C
```
---
@@ -126,26 +100,26 @@ DEB 预测今日最高 18.2°C | 目前距预测差 1.7°C
| 优先级 | 文件 | 改动 |
|--------|------|------|
| 1 | `src/database/db_manager.py` | 新增 `airport_obs_log` 表、`append_airport_obs()``get_airport_obs_recent()` |
| 2 | `src/data_collection/amos_station_sources.py` | 成功获取后调用 `append_airport_obs()` 写日志 |
| 3 | `src/data_collection/jma_amedas_sources.py` | 成功获取后调用 `append_airport_obs()` 写日志,可选扩展提取更多字段(风速/气压) |
| 4 | `src/analysis/market_alert_engine.py` | 新增 `airport_rapid_temp_change` 规则 |
| 5 | `src/utils/telegram_push.py` | 高频快速通道子循环、机场快照推送(含 DEB 预测)、温度急变告警集成 |
| 2 | `src/data_collection/weather_sources.py` | AMOS/JMA/MGM 成功后调用 `append_airport_obs()` 写日志 |
| 3 | `src/analysis/market_alert_engine.py` | 新增 `airport_rapid_temp_change` 规则 |
| 4 | `src/utils/telegram_push.py` | 10 分钟高频子循环、温度急变告警推送、最高温锁定跳过 |
| 5 | `src/bot/runtime_coordinator.py` | 注册机场高频推送循环 |
---
## 实施顺序
1. **Phase 1 — DB 层**: `airport_obs_log` 表 + 读写方法
2. **Phase 2 — 采集层**: AMOS/JMA 成功后自动写日志,部署观察 1-2 天确认数据积累正常
2. **Phase 2 — 采集层**: AMOS/JMA/MGM 成功后自动写日志,部署观察 1-2 天确认数据积累正常
3. **Phase 3 — 告警引擎**: `airport_rapid_temp_change` 规则 + 单元测试
4. **Phase 4 — 推送层**: 高频快速通道 + 机场快照 + DEB 预测,直接推送市场监控频道
4. **Phase 4 — 推送层**: 高频快速通道,直接推送市场监控频道
5. **Phase 5 — 调参**: 观察 3-7 天调整阈值
---
## 风险与注意事项
- **AMOS/JMA 站点可用性**: `global.amo.go.kr``jma.go.jp` 可能偶发性不可用,需容错处理
- **AMOS/JMA/MGM 站点可用性**: 各数据源可能偶发性不可用,需容错处理
- **告警频率控制**: 高频循环可能产生过多告警,需要严格的冷却期和去重机制
- **数据库体积**: `airport_obs_log` 每 1-10 分钟写入 3 条记录,2 小时约 36-360 条,自动清理后体积可控
- **东京 JMA 数据完整性**: 当前只提取 temp,风速/气压需要额外解析(JMA JSON 中有 `wind``pressure` 数组但未使用)
- **数据库体积**: `airport_obs_log` 每 1-10 分钟写入 4 条记录,2 小时约 48-480 条,自动清理后体积可控
- **安卡拉 MGM 刷新频率**: `servis.mgm.gov.tr` 实际更新间隔待确认,暂按 ~10 分钟预估
+1 -1
View File
@@ -813,7 +813,7 @@ def _build_advice_cn(
return "".join(parts) + ""
_AIRPORT_ICAO_MAP = {"seoul": "RKSI", "busan": "RKPK", "tokyo": "RJTT"}
_AIRPORT_ICAO_MAP = {"seoul": "RKSI", "busan": "RKPK", "tokyo": "RJTT", "ankara": "17128"}
def _calc_airport_rapid_temp_change(
+13
View File
@@ -805,6 +805,19 @@ class WeatherDataCollector(OpenMeteoCacheMixin, SettlementSourceMixin, MetarSour
if not mgm_data:
return
results["mgm"] = mgm_data
# Log airport obs for high-freq monitoring (Ankara 17128)
try:
mgm_current = mgm_data.get("current") or {}
DBManager().append_airport_obs(
icao=str(istno),
city=city_lower,
temp_c=mgm_current.get("temp"),
wind_kt=mgm_current.get("wind_speed_kt"),
pressure_hpa=mgm_current.get("pressure"),
obs_time=mgm_data.get("obs_time") or datetime.now().isoformat(),
)
except Exception:
logger.exception("airport_obs_log append failed for mgm city={}", city_lower)
if include_nearby:
results["nearby_source"] = "mgm"
nearby = self.fetch_mgm_nearby_stations(province, root_ist_no=istno)
+3 -100
View File
@@ -690,12 +690,11 @@ def start_trade_alert_push_loop(bot: Any, config: Dict[str, Any]) -> Optional[th
# ── high-freq airport push loop ──
HIGH_FREQ_AIRPORT_CITIES = {"seoul", "busan", "tokyo"}
HIGH_FREQ_AIRPORT_ICAO = {"seoul": "RKSI", "busan": "RKPK", "tokyo": "RJTT"}
HIGH_FREQ_AIRPORT_CITIES = {"seoul", "busan", "tokyo", "ankara"}
HIGH_FREQ_AIRPORT_ICAO = {"seoul": "RKSI", "busan": "RKPK", "tokyo": "RJTT", "ankara": "17128"}
HIGH_FREQ_PUSH_INTERVAL_SEC = 600
HIGH_FREQ_MOMENTUM_THRESHOLD_C = 0.5
HIGH_FREQ_COOLDOWN_SEC = 7200
AIRPORT_SNAPSHOT_INTERVAL_SEC = 1800
_AIRPORT_PUSH_STATE_PATH = os.path.join(
os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))),
@@ -758,7 +757,7 @@ def _build_airport_rapid_change_message(
deb_pred: Optional[float],
) -> str:
city_display = city.title()
airport_label = {"seoul": "首尔/仁川", "busan": "釜山/金海", "tokyo": "东京/羽田"}.get(city, city_display)
airport_label = {"seoul": "首尔/仁川", "busan": "釜山/金海", "tokyo": "东京/羽田", "ankara": "安卡拉/Esenboğa"}.get(city, city_display)
emoji = "🔥" if rule.get("direction") == "up" else "❄️"
first = rule.get("first_temp", 0)
last = rule.get("last_temp", 0)
@@ -781,42 +780,6 @@ def _build_airport_rapid_change_message(
return "\n".join(lines)
def _build_airport_snapshot_message(snapshots: list[dict[str, Any]], local_time: str = "") -> str:
flag = {"seoul": "🇰🇷", "busan": "🇰🇷", "tokyo": "🇯🇵"}
name = {"seoul": "首尔/仁川 RKSI", "busan": "釜山/金海 RKPK", "tokyo": "东京/羽田 RJTT"}
lines = [f"🛫 机场实况 {local_time}"]
for s in snapshots:
city = s.get("city", "")
temp = s.get("temp_c")
deb = s.get("deb_prediction")
runway_pairs = s.get("runway_pairs") or []
runway_temps = s.get("runway_temps") or []
lines.append("")
header = f"{flag.get(city, '')} {name.get(city, city)}"
# Show individual runway pair temps for Seoul/Busan
if runway_pairs and runway_temps and len(runway_pairs) == len(runway_temps):
runway_lines = []
for (r1, r2), (t, d) in zip(runway_pairs, runway_temps):
dew_str = f" / 露点 {d:.1f}°C" if d is not None else ""
runway_lines.append(f"{r1}/{r2}: {t:.1f}°C{dew_str}")
line = header + "\n" + "\n".join(runway_lines)
elif temp is not None:
line = header + f"\n当前 {temp:.1f}°C"
else:
line = header
if deb is not None:
line += f"\nDEB 预测最高 {deb:.1f}°C"
if s.get("wind_kt") is not None:
line += f" | 风 {s['wind_kt']:.0f}kt"
if s.get("pressure_hpa") is not None:
line += f" | QNH {s['pressure_hpa']:.1f} hPa"
lines.append(line)
return "\n".join(lines)
def _run_high_freq_airport_cycle(
bot: Any,
config: Dict[str, Any],
@@ -827,9 +790,6 @@ def _run_high_freq_airport_cycle(
now_ts = int(time.time())
last_by_city = state.setdefault("last_by_city", {})
snapshots: list[dict[str, Any]] = []
snapshot_due = now_ts - int(state.get("last_snapshot_ts") or 0) >= AIRPORT_SNAPSHOT_INTERVAL_SEC
for city in sorted(HIGH_FREQ_AIRPORT_CITIES):
try:
icao = HIGH_FREQ_AIRPORT_ICAO.get(city, "")
@@ -844,7 +804,6 @@ def _run_high_freq_airport_cycle(
except Exception:
latest_obs = {}
# Fetch city_weather once and extract DEB
city_weather: Dict[str, Any] = {}
deb_pred: Optional[float] = None
try:
@@ -856,51 +815,13 @@ def _run_high_freq_airport_cycle(
except Exception:
pass
# Extract airport-level data from city_weather for snapshot
amos = city_weather.get("amos") or {}
runway_data = (amos.get("runway_obs") or {}) if amos else {}
runway_pairs = runway_data.get("runway_pairs") or []
runway_temps = runway_data.get("temperatures") or []
# Tokyo JMA data lives in mgm_nearby (set by _attach_japan_official_nearby)
mgm_nearby = city_weather.get("mgm_nearby") or []
nearby_row = mgm_nearby[0] if mgm_nearby else {}
current_temp = (
amos.get("temp_c")
or nearby_row.get("temp")
or latest_obs.get("temp_c")
or (city_weather.get("current") or {}).get("temp")
)
wind_kt = amos.get("wind_kt") or latest_obs.get("wind_kt")
pressure_hpa = amos.get("pressure_hpa") or latest_obs.get("pressure_hpa")
# Check high-locked
if _check_high_locked(city):
if snapshot_due:
snapshots.append({
"city": city,
"temp_c": current_temp,
"wind_kt": wind_kt,
"pressure_hpa": pressure_hpa,
"deb_prediction": deb_pred,
"runway_pairs": runway_pairs,
"runway_temps": runway_temps,
})
continue
from src.analysis.market_alert_engine import _calc_airport_rapid_temp_change
rule = _calc_airport_rapid_temp_change(city_weather, "°C")
if snapshot_due:
snapshots.append({
"city": city,
"temp_c": current_temp,
"wind_kt": wind_kt,
"pressure_hpa": pressure_hpa,
"deb_prediction": deb_pred,
"runway_pairs": runway_pairs,
"runway_temps": runway_temps,
})
if not rule.get("triggered"):
continue
@@ -926,24 +847,6 @@ def _run_high_freq_airport_cycle(
except Exception:
logger.exception("high freq airport cycle failed for city={}", city)
if snapshot_due and snapshots:
from datetime import timedelta
kst_hour = (datetime.utcnow() + timedelta(hours=9)).hour
if 8 <= kst_hour < 20:
kst_now = datetime.utcnow() + timedelta(hours=9)
local_time = kst_now.strftime("%H:%M KST")
snap_message = _build_airport_snapshot_message(snapshots, local_time)
for chat_id in chat_ids:
try:
bot.send_message(chat_id, snap_message)
except Exception as exc:
logger.warning("airport snapshot push failed chat_id={}: {}", chat_id, exc)
logger.info("airport snapshot pushed cities={}", len(snapshots))
else:
logger.info("airport snapshot skipped: outside 08:00-20:00 window (KST hour={})", kst_hour)
state["last_snapshot_ts"] = now_ts
state_dirty = True
return state_dirty