Retire dual storage mode and default runtime state to sqlite

This commit is contained in:
2569718930@qq.com
2026-04-16 00:52:55 +08:00
parent 4bfd2534bb
commit 53d2ae4c10
11 changed files with 437 additions and 1661 deletions
+5 -8
View File
@@ -7,7 +7,6 @@ from src.analysis.settlement_rounding import apply_city_settlement
from loguru import logger
from src.database.runtime_state import (
DailyRecordRepository,
STATE_STORAGE_DUAL,
STATE_STORAGE_SQLITE,
TrainingFeatureRecordRepository,
TruthRecordRepository,
@@ -79,7 +78,7 @@ def load_history(filepath):
logger.error(f"Error loading daily records from sqlite, fallback to file: {e}")
if not os.path.exists(filepath):
if mode == STATE_STORAGE_DUAL:
if mode == STATE_STORAGE_SQLITE:
try:
data = _daily_record_repo.load_all()
_history_cache = data
@@ -113,13 +112,12 @@ def save_history(filepath, data):
_history_cache = data
mode = get_state_storage_mode()
if mode in {STATE_STORAGE_DUAL, STATE_STORAGE_SQLITE}:
if mode == STATE_STORAGE_SQLITE:
try:
_daily_record_repo.replace_all(data)
except Exception as e:
logger.error(f"Error saving daily records to sqlite: {e}")
if mode == STATE_STORAGE_SQLITE:
return
return
if mode == STATE_STORAGE_SQLITE:
return
@@ -891,15 +889,14 @@ def update_daily_record(
for d in old_dates:
del data[city][d]
if mode in {STATE_STORAGE_DUAL, STATE_STORAGE_SQLITE}:
if mode == STATE_STORAGE_SQLITE:
try:
_daily_record_repo.upsert_record(city_name, date_str, existing)
cutoff = (datetime.now() - timedelta(days=14)).strftime("%Y-%m-%d")
_daily_record_repo.delete_older_than(cutoff)
except Exception as e:
logger.error(f"Error upserting daily record to sqlite city={city_name} date={date_str}: {e}")
if mode == STATE_STORAGE_SQLITE:
raise
raise
if mode != STATE_STORAGE_SQLITE:
save_history(history_file, data)
+2 -3
View File
@@ -7,7 +7,6 @@ from typing import Any, Dict, List, Optional
from src.database.runtime_state import (
ProbabilitySnapshotRepository,
STATE_STORAGE_DUAL,
STATE_STORAGE_SQLITE,
TrainingFeatureRecordRepository,
get_state_storage_mode,
@@ -105,7 +104,7 @@ def load_snapshot_rows_for_day(
root_dir = os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
path = archive_path or os.path.join(root_dir, "data", "probability_training_snapshots.jsonl")
if not os.path.exists(path):
if mode == STATE_STORAGE_DUAL:
if mode == STATE_STORAGE_SQLITE:
return _snapshot_repo.load_rows_by_city_date(city_key, date_key)
return []
@@ -260,7 +259,7 @@ def append_probability_snapshot(
return
mode = get_state_storage_mode()
if mode in {STATE_STORAGE_DUAL, STATE_STORAGE_SQLITE}:
if mode == STATE_STORAGE_SQLITE:
_snapshot_repo.append_snapshot(payload)
_training_feature_repo.upsert_record(
city_key,
+2 -3
View File
@@ -7,7 +7,6 @@ import time
from loguru import logger
from src.database.runtime_state import (
OpenMeteoCacheRepository,
STATE_STORAGE_DUAL,
STATE_STORAGE_SQLITE,
get_state_storage_mode,
)
@@ -77,7 +76,7 @@ class OpenMeteoCacheMixin:
self._disk_cache_last_mtime = current_mtime
if loaded:
logger.info(f"✅ 从磁盘加载 Open-Meteo 缓存 {loaded} 条 ({self._disk_cache_path})")
if mode == STATE_STORAGE_DUAL:
if mode == STATE_STORAGE_SQLITE:
try:
_open_meteo_cache_repo.replace_payload(saved, self._disk_cache_max_age_sec)
except Exception as exc:
@@ -124,7 +123,7 @@ class OpenMeteoCacheMixin:
"multi_model": multi_model_snapshot,
"saved_at": time.time(),
}
if mode in {STATE_STORAGE_DUAL, STATE_STORAGE_SQLITE}:
if mode == STATE_STORAGE_SQLITE:
_open_meteo_cache_repo.replace_payload(payload, self._disk_cache_max_age_sec)
self._disk_cache_last_mtime = _open_meteo_cache_repo.latest_updated_at()
if mode != STATE_STORAGE_SQLITE:
+1 -2
View File
@@ -12,7 +12,6 @@ from loguru import logger
from src.database.runtime_state import (
OfficialIntradayObservationRepository,
STATE_STORAGE_DUAL,
STATE_STORAGE_SQLITE,
get_state_storage_mode,
)
@@ -183,7 +182,7 @@ class SettlementSourceMixin:
date_str = local_dt.strftime("%Y-%m-%d")
time_str = local_dt.strftime("%H:%M")
mode = get_state_storage_mode()
if mode not in {STATE_STORAGE_DUAL, STATE_STORAGE_SQLITE}:
if mode != STATE_STORAGE_SQLITE:
return [{"time": time_str, "temp": round(float(current_temp), 1)}]
lock = self._get_settlement_series_lock()
+9 -4
View File
@@ -16,9 +16,9 @@ from src.database.db_manager import DBManager
STATE_STORAGE_FILE = "file"
STATE_STORAGE_DUAL = "dual"
STATE_STORAGE_SQLITE = "sqlite"
DEFAULT_STATE_STORAGE_MODE = STATE_STORAGE_SQLITE
VALID_STATE_STORAGE_MODES = {
STATE_STORAGE_FILE,
STATE_STORAGE_DUAL,
STATE_STORAGE_SQLITE,
}
@@ -26,12 +26,17 @@ _LOGGED_MODES: set[str] = set()
def get_state_storage_mode() -> str:
raw = str(os.getenv("POLYWEATHER_STATE_STORAGE_MODE") or STATE_STORAGE_DUAL).strip().lower()
raw = str(os.getenv("POLYWEATHER_STATE_STORAGE_MODE") or DEFAULT_STATE_STORAGE_MODE).strip().lower()
if raw == STATE_STORAGE_DUAL:
logger.warning(
f"POLYWEATHER_STATE_STORAGE_MODE={STATE_STORAGE_DUAL!r} is deprecated, normalize to {STATE_STORAGE_SQLITE}"
)
raw = STATE_STORAGE_SQLITE
if raw not in VALID_STATE_STORAGE_MODES:
logger.warning(
f"invalid POLYWEATHER_STATE_STORAGE_MODE={raw!r}, fallback to {STATE_STORAGE_DUAL}"
f"invalid POLYWEATHER_STATE_STORAGE_MODE={raw!r}, fallback to {DEFAULT_STATE_STORAGE_MODE}"
)
raw = STATE_STORAGE_DUAL
raw = DEFAULT_STATE_STORAGE_MODE
if raw not in _LOGGED_MODES:
logger.info(f"runtime state storage mode={raw}")
_LOGGED_MODES.add(raw)
+3 -2
View File
@@ -289,13 +289,14 @@ def build_training_samples(
history_data: Optional[Dict[str, Any]] = None,
snapshot_index: Optional[Dict[Tuple[str, str], Dict[str, Any]]] = None,
) -> List[Dict[str, Any]]:
mode = get_state_storage_mode()
if isinstance(history_data, dict):
runtime_history = history_data
elif get_state_storage_mode() == STATE_STORAGE_SQLITE:
elif mode == STATE_STORAGE_SQLITE:
runtime_history = DailyRecordRepository().load_all()
else:
runtime_history = load_history(_history_file_path())
if get_state_storage_mode() == STATE_STORAGE_SQLITE:
if mode == STATE_STORAGE_SQLITE:
truth_history = TruthRecordRepository().load_all()
training_feature_history = TrainingFeatureRecordRepository().load_all()
else:
+1 -7
View File
@@ -10,7 +10,6 @@ from typing import Any, Dict, List, Optional, Tuple
from loguru import logger
from src.database.runtime_state import (
STATE_STORAGE_DUAL,
STATE_STORAGE_SQLITE,
TelegramAlertStateRepository,
get_state_storage_mode,
@@ -208,11 +207,6 @@ def _load_state(path: str) -> Dict[str, Any]:
except Exception as exc:
logger.error(f"failed to load telegram push state from sqlite: {exc}")
if not os.path.exists(path):
if mode == STATE_STORAGE_DUAL:
try:
return _telegram_state_repo.load_state()
except Exception:
return {"last_by_city": {}, "by_signature": {}}
return {"last_by_city": {}, "by_signature": {}}
try:
with open(path, "r", encoding="utf-8") as fh:
@@ -228,7 +222,7 @@ def _load_state(path: str) -> Dict[str, Any]:
def _save_state(path: str, state: Dict[str, Any]) -> None:
mode = get_state_storage_mode()
if mode in {STATE_STORAGE_DUAL, STATE_STORAGE_SQLITE}:
if mode == STATE_STORAGE_SQLITE:
_telegram_state_repo.save_state(state)
if mode == STATE_STORAGE_SQLITE:
return