/
NeonBite
/
big-data-6
Обзор
Документация
Войти
/
NeonBite
/
big-data-6
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
monitoring/data_drift.py
107 строк
3 KB
NeonBite
feat: report done
20 дек 2025, 19:28
20 дек 2025, 19:28
d12a6ed
Код
Авторство
О чём код?
import datetime as dt import logging import math import os import urllib.request logger = logging.getLogger(__name__) def setup_logging(): log_dir = os.environ.get("LOG_DIR", ".logs") os.makedirs(log_dir, exist_ok=True) script_name = os.path.splitext(os.path.basename(__file__))[0] timestamp = dt.datetime.utcnow().strftime("%Y%m%d-%H%M%S") log_path = os.path.join(log_dir, f"{script_name}-{timestamp}.log") logging.basicConfig( level=os.environ.get("LOG_LEVEL", "INFO"), format="%(asctime)s %(levelname)s %(name)s - %(message)s", handlers=[ logging.FileHandler(log_path), logging.StreamHandler(), ], ) logger.info("Logging to %s", log_path) def load_ratings(path: str, sample_size: int): ratings = [] with open(path, "r", encoding="utf-8") as f: for i, line in enumerate(f): if sample_size and i >= sample_size: break parts = line.strip().split(",") if len(parts) != 4: continue try: rating = float(parts[2]) except ValueError: continue ratings.append(rating) return ratings def psi(ref, cur, bins): eps = 1e-6 ref_counts = [0] * (len(bins) - 1) cur_counts = [0] * (len(bins) - 1) for v in ref: for i in range(len(bins) - 1): if bins[i] <= v < bins[i + 1]: ref_counts[i] += 1 break for v in cur: for i in range(len(bins) - 1): if bins[i] <= v < bins[i + 1]: cur_counts[i] += 1 break ref_total = sum(ref_counts) or 1 cur_total = sum(cur_counts) or 1 score = 0.0 for r, c in zip(ref_counts, cur_counts): r_pct = max(r / ref_total, eps) c_pct = max(c / cur_total, eps) score += (c_pct - r_pct) * math.log(c_pct / r_pct) return score def push_to_gateway(url: str, metric_name: str, value: float): body = f"{metric_name} {value}\n" req = urllib.request.Request(url, data=body.encode("utf-8"), method="POST") with urllib.request.urlopen(req, timeout=5) as resp: logger.info("Pushed metric to %s status=%s", url, resp.status) def main(): setup_logging() ref_path = os.environ.get("REF_DATASET_PATH", "data/ratings_Electronics.csv") cur_path = os.environ.get("CUR_DATASET_PATH", "data/ratings_Electronics.csv") sample_size = int(os.environ.get("DRIFT_SAMPLE_SIZE", "50000")) pushgateway_url = os.environ.get( "PUSHGATEWAY_URL", "http://pushgateway:9091/metrics/job/data_drift", ) ref_ratings = load_ratings(ref_path, sample_size) cur_ratings = load_ratings(cur_path, sample_size) if not ref_ratings or not cur_ratings: logger.warning("Empty ratings for drift check") return bins = [0.0, 1.5, 2.5, 3.5, 4.5, 5.5] drift_score = psi(ref_ratings, cur_ratings, bins) logger.info("Data drift score (PSI)=%.6f", drift_score) try: push_to_gateway(pushgateway_url, "data_drift_score", drift_score) except Exception: logger.exception("Failed to push drift score to pushgateway") if __name__ == "__main__": main()