/
NeonBite
/
big-data-4
Обзор
Документация
Войти
/
NeonBite
/
big-data-4
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
model-cmp/main.py
754 строки
30 KB
NeonBite
fix: added modelcmp service and set logs and reports paths properly
03 дек 2025, 06:11
03 дек 2025, 06:11
44613d3
Код
Авторство
О чём код?
import datetime import io import logging import os import shutil from typing import Dict, List, Tuple import kagglehub import matplotlib.pyplot as plt import mlflow import mlflow.sklearn import numpy as np import pandas as pd from mlflow.models import infer_signature from sklearn.ensemble import GradientBoostingClassifier, RandomForestClassifier from sklearn.impute import SimpleImputer from sklearn.linear_model import LogisticRegression, SGDClassifier from sklearn.metrics import f1_score, roc_auc_score, log_loss from sklearn.model_selection import train_test_split, learning_curve from sklearn.pipeline import Pipeline from sklearn.preprocessing import LabelEncoder, StandardScaler from catboost import CatBoostClassifier from sklearn.cluster import MiniBatchKMeans DATA_DIR = os.environ.get("DATA_DIR", os.path.join(os.getcwd(), "data")) DATASET_NAME = "ratings_Electronics (1).csv" # Папка отчётов: для каждого запуска создаём отдельную подпапку с timestamp BASE_REPORT_DIR = os.environ.get("REPORT_DIR", os.path.join(os.getcwd(), "reports")) os.makedirs(BASE_REPORT_DIR, exist_ok=True) RUN_TS = datetime.datetime.now().strftime("%Y%m%d_%H%M%S") REPORT_DIR = os.path.join(BASE_REPORT_DIR, RUN_TS) os.makedirs(REPORT_DIR, exist_ok=True) REPORT_PATH = os.path.join(REPORT_DIR, "eda_report.md") # Логи LOG_DIR = os.environ.get("LOG_DIR", os.path.join(os.getcwd(), ".logs")) os.makedirs(LOG_DIR, exist_ok=True) LOG_PATH = os.path.join(LOG_DIR, f"pipeline_{RUN_TS}.log") class KMeansClassifier: """Простой обёрточный классификатор на MiniBatchKMeans для бинарного случая.""" def __init__(self, n_clusters: int = 2, random_state: int = 42, batch_size: int = 1024): self.n_clusters = n_clusters self.kmeans = MiniBatchKMeans( n_clusters=n_clusters, random_state=random_state, batch_size=batch_size ) self.cluster_to_label: Dict[int, int] = {} def fit(self, X: np.ndarray, y: np.ndarray) -> "KMeansClassifier": self.kmeans.fit(X) labels = self.kmeans.labels_ for c in range(self.n_clusters): mask = labels == c if mask.sum() == 0: self.cluster_to_label[c] = 0 else: # большинство классов в кластере counts = np.bincount(y[mask].astype(int), minlength=2) self.cluster_to_label[c] = int(counts.argmax()) return self def predict(self, X: np.ndarray) -> np.ndarray: dists = self.kmeans.transform(X) weights = 1.0 / np.maximum(dists, 1e-6) weights = weights / weights.sum(axis=1, keepdims=True) cluster_pred = weights.argmax(axis=1) return np.array([self.cluster_to_label.get(c, 0) for c in cluster_pred]) def predict_proba(self, X: np.ndarray) -> np.ndarray: # используем веса обратных расстояний как суррогат вероятностей кластеров dists = self.kmeans.transform(X) weights = 1.0 / np.maximum(dists, 1e-6) weights = weights / weights.sum(axis=1, keepdims=True) proba1 = [] for row_weights in weights: p1 = 0.0 for c, w in enumerate(row_weights): label = self.cluster_to_label.get(c, 0) p1 += w * label proba1.append(p1) proba1 = np.clip(np.array(proba1), 0.0, 1.0) proba0 = 1.0 - proba1 return np.vstack([proba0, proba1]).T def setup_logger() -> logging.Logger: logger = logging.getLogger("recsys_pipeline") logger.setLevel(logging.INFO) # сброс старых хендлеров при повторном запуске в одном процессе logger.handlers = [] formatter = logging.Formatter( "%(asctime)s [%(levelname)s] %(message)s", datefmt="%Y-%m-%d %H:%M:%S" ) file_handler = logging.FileHandler(LOG_PATH, encoding="utf-8") file_handler.setFormatter(formatter) logger.addHandler(file_handler) console_handler = logging.StreamHandler() console_handler.setFormatter(formatter) logger.addHandler(console_handler) return logger def append_section(lines: List[str], text: str) -> None: lines.append(text) def download_dataset_if_needed(logger: logging.Logger) -> None: if not os.path.exists(DATA_DIR) or len(os.listdir(DATA_DIR)) == 0: logger.info("DATA_DIR is empty, downloading dataset via kagglehub") os.makedirs(DATA_DIR, exist_ok=True) try: path = kagglehub.dataset_download("saurav9786/amazon-product-reviews") for name in os.listdir(path): src = os.path.join(path, name) dst = os.path.join(DATA_DIR, name) shutil.move(src, dst) logger.info("Dataset moved to %s", DATA_DIR) except Exception: logger.exception("Failed to download or move dataset") raise else: logger.info("Dataset already exists in %s, skipping download", DATA_DIR) def load_dataset(logger: logging.Logger) -> pd.DataFrame: csv_path = os.path.join(DATA_DIR, DATASET_NAME) try: df = pd.read_csv( csv_path, header=None, names=["user_id", "item_id", "rating", "timestamp"], ) logger.info("Dataset loaded from %s, shape=%s", csv_path, df.shape) return df except Exception: logger.exception("Failed to read dataset from %s", csv_path) raise def add_initial_eda(df: pd.DataFrame, report_lines: List[str], logger: logging.Logger) -> None: append_section(report_lines, "# EDA отчёт по датасету\n\n") append_section(report_lines, "## Структура датасета\n\n") # Пример строк append_section(report_lines, "### Пример строк\n\n```text\n") append_section(report_lines, df.head().to_string() + "\n") append_section(report_lines, "```\n\n") # info() buffer = io.StringIO() df.info(buf=buffer) append_section(report_lines, "### Информация о DataFrame\n\n```text\n") append_section(report_lines, buffer.getvalue()) append_section(report_lines, "```\n\n") # describe() append_section(report_lines, "### Описательная статистика по числовым колонкам\n\n```text\n") append_section(report_lines, df.describe().to_string() + "\n") append_section(report_lines, "```\n\n") logger.info("Initial EDA section prepared") def add_target_and_correlation( df: pd.DataFrame, report_lines: List[str], logger: logging.Logger ) -> None: df["label"] = (df["rating"] >= 4).astype(int) append_section(report_lines, "## Целевая переменная и рейтинги\n\n") # Распределение label label_prop = df["label"].value_counts(normalize=True).rename("label_proportion") label_counts = df["label"].value_counts() append_section(report_lines, "### Распределение label\n\n```text\n") append_section(report_lines, label_prop.to_string() + "\n\n") append_section(report_lines, label_counts.to_string() + "\n") append_section(report_lines, "```\n\n") # Распределение рейтингов append_section(report_lines, "### Распределение рейтингов\n\n```text\n") append_section(report_lines, df["rating"].value_counts().sort_index().to_string() + "\n\n") append_section(report_lines, df["rating"].describe().to_string() + "\n") append_section(report_lines, "```\n\n") logger.info("Target distribution section prepared") df["timestamp"] = pd.to_datetime(df["timestamp"], unit="s") df["year"] = df["timestamp"].dt.year corr = df[["rating", "label", "year"]].corr() append_section(report_lines, "## Корреляции\n\n") append_section(report_lines, "### Корреляционная матрица (rating, label, year)\n\n```text\n") append_section(report_lines, corr.to_string() + "\n") append_section(report_lines, "```\n\n") logger.info("Correlation matrix computed") # Построение графика корреляций try: corr_plot_path = os.path.join(REPORT_DIR, "corr_matrix.png") fig, ax = plt.subplots(figsize=(4, 3)) im = ax.imshow(corr.values, cmap="coolwarm", vmin=-1, vmax=1) ax.set_xticks(range(len(corr.columns))) ax.set_xticklabels(corr.columns) ax.set_yticks(range(len(corr.index))) ax.set_yticklabels(corr.index) plt.colorbar(im, ax=ax) plt.tight_layout() fig.savefig(corr_plot_path) plt.close(fig) append_section( report_lines, f"})\n\n", ) logger.info("Correlation heatmap saved to %s", corr_plot_path) except Exception: logger.exception("Failed to generate correlation heatmap") def add_data_issues(df: pd.DataFrame, report_lines: List[str], logger: logging.Logger) -> None: missing_per_col = df.isna().sum() dup_all_share = df.duplicated().mean() dup_user_item_share = df.duplicated(subset=["user_id", "item_id"]).mean() label_prop = df["label"].value_counts(normalize=True).rename("label_proportion") label_counts = df["label"].value_counts() user_activity = df["user_id"].value_counts() item_popularity = df["item_id"].value_counts() append_section(report_lines, "## Проблемы с данными\n\n") append_section(report_lines, "### Пропуски по столбцам\n\n```text\n") append_section(report_lines, missing_per_col.to_string() + "\n") append_section(report_lines, "```\n\n") append_section(report_lines, "### Дубликаты\n\n```text\n") append_section( report_lines, f"Доля полностью дублирующих строк: {dup_all_share:.4f}\n" f"Доля дубликатов по (user_id, item_id): {dup_user_item_share:.4f}\n", ) append_section(report_lines, "```\n\n") append_section(report_lines, "### Дисбаланс целевой переменной (label)\n\n```text\n") append_section(report_lines, label_prop.to_string() + "\n\n") append_section(report_lines, label_counts.to_string() + "\n") append_section(report_lines, "```\n\n") append_section(report_lines, "### Активность пользователей\n\n```text\n") append_section(report_lines, user_activity.describe().to_string() + "\n") append_section(report_lines, "```\n\n") append_section(report_lines, "### Популярность товаров\n\n```text\n") append_section(report_lines, item_popularity.describe().to_string() + "\n") append_section(report_lines, "```\n\n") logger.info("Data issues section prepared") def train_basic_models( df: pd.DataFrame, report_lines: List[str], logger: logging.Logger ) -> Tuple[pd.DataFrame, Dict[str, Pipeline], np.ndarray, np.ndarray, np.ndarray, np.ndarray]: sample_size = int(os.getenv("SAMPLE_SIZE", "500000")) if len(df) > sample_size: df_sample = df.sample(n=sample_size, random_state=42) else: df_sample = df user_encoder = LabelEncoder() item_encoder = LabelEncoder() user_ids_enc = user_encoder.fit_transform(df_sample["user_id"]) item_ids_enc = item_encoder.fit_transform(df_sample["item_id"]) years = df_sample["year"].astype("int32") X = np.vstack([user_ids_enc, item_ids_enc, years]).T y = df_sample["label"].values X_train, X_test, y_train, y_test = train_test_split( X, y, test_size=0.2, random_state=42, stratify=y ) logger.info("Basic features prepared for modeling, train size=%d, test size=%d", len(X_train), len(X_test)) models: Dict[str, Pipeline] = { "logistic_regression": Pipeline( [ ("imputer", SimpleImputer(strategy="median")), ("scaler", StandardScaler()), ("classifier", LogisticRegression(max_iter=200, n_jobs=-1)), ] ), "sgd_classifier": Pipeline( [ ("imputer", SimpleImputer(strategy="median")), ("scaler", StandardScaler()), ( "classifier", SGDClassifier( loss="log_loss", learning_rate="optimal", max_iter=5, tol=1e-3, random_state=42, ), ), ] ), "catboost": Pipeline( [ ("imputer", SimpleImputer(strategy="median")), ("scaler", StandardScaler()), ( "classifier", CatBoostClassifier( depth=6, learning_rate=0.1, iterations=200, loss_function="Logloss", verbose=False, random_state=42, ), ), ] ), "minibatch_kmeans": Pipeline( [ ("imputer", SimpleImputer(strategy="median")), ("scaler", StandardScaler()), ("classifier", KMeansClassifier(n_clusters=2, batch_size=2048)), ] ), "random_forest": Pipeline( [ ("imputer", SimpleImputer(strategy="median")), ("scaler", StandardScaler()), ( "classifier", RandomForestClassifier(n_estimators=100, n_jobs=-1, random_state=42), ), ] ), "gradient_boosting": Pipeline( [ ("imputer", SimpleImputer(strategy="median")), ("scaler", StandardScaler()), ("classifier", GradientBoostingClassifier(random_state=42)), ] ), } logger.info("Basic models initialized: %s", list(models.keys())) results = [] trained_models: Dict[str, Pipeline] = {} for name, model in models.items(): logger.info("Training basic model: %s", name) model.fit(X_train, y_train) trained_models[name] = model y_proba = model.predict_proba(X_test)[:, 1] y_pred = (y_proba >= 0.5).astype(int) auc = roc_auc_score(y_test, y_proba) f1 = f1_score(y_test, y_pred) ll = log_loss(y_test, y_proba) results.append({"model": name, "auc": auc, "f1": f1, "log_loss": ll}) logger.info("Basic model %s: AUC=%.4f, F1=%.4f, log_loss=%.4f", name, auc, f1, ll) results_df = pd.DataFrame(results).sort_values(by="auc", ascending=False) best_model_name = results_df.iloc[0]["model"] append_section(report_lines, "## Сравнение моделей\n\n```text\n") append_section(report_lines, results_df.to_string(index=False) + "\n") append_section(report_lines, "```\n\n") append_section(report_lines, f"Лучшая модель по AUC: {best_model_name}\n\n") logger.info("Basic model comparison section prepared, best=%s", best_model_name) return results_df, trained_models, X_train, X_test, y_train, y_test def train_enriched_models( df: pd.DataFrame, report_lines: List[str], logger: logging.Logger ) -> Tuple[pd.DataFrame, Dict[str, Pipeline], np.ndarray, np.ndarray, np.ndarray, np.ndarray]: # агрегаты по пользователю: средний рейтинг, количество оценок, доля высоких оценок user_stats = ( df.assign(high_label=(df["label"] == 1).astype(float)) .groupby("user_id") .agg( user_mean_rating=("rating", "mean"), user_rating_count=("rating", "count"), user_high_share=("high_label", "mean"), ) .reset_index() ) # агрегаты по товару: средний рейтинг, количество оценок, доля высоких оценок item_stats = ( df.assign(high_label=(df["label"] == 1).astype(float)) .groupby("item_id") .agg( item_mean_rating=("rating", "mean"), item_rating_count=("rating", "count"), item_high_share=("high_label", "mean"), ) .reset_index() ) df_enriched = df.merge(user_stats, on="user_id", how="left").merge( item_stats, on="item_id", how="left" ) # временная фича: давность взаимодействия (в днях от максимального timestamp) try: max_ts = df_enriched["timestamp"].max() df_enriched["recency_days"] = ( max_ts - df_enriched["timestamp"] ).dt.total_seconds() / 86400.0 except Exception: logger.exception("Failed to compute recency_days, filling with 0") df_enriched["recency_days"] = 0.0 logger.info("Enriched features with user/item aggregates and recency, shape=%s", df_enriched.shape) sample_size_enriched = int(os.getenv("SAMPLE_SIZE_ENRICHED", "500000")) if len(df_enriched) > sample_size_enriched: df_sample_enriched = df_enriched.sample(n=sample_size_enriched, random_state=42) else: df_sample_enriched = df_enriched user_encoder_en = LabelEncoder() item_encoder_en = LabelEncoder() user_ids_enc_en = user_encoder_en.fit_transform(df_sample_enriched["user_id"]) item_ids_enc_en = item_encoder_en.fit_transform(df_sample_enriched["item_id"]) years_en = df_sample_enriched["year"].astype("int32") recency_days = df_sample_enriched["recency_days"].astype("float32").values user_mean = df_sample_enriched["user_mean_rating"].values item_mean = df_sample_enriched["item_mean_rating"].values user_cnt = df_sample_enriched["user_rating_count"].values item_cnt = df_sample_enriched["item_rating_count"].values user_high_share = df_sample_enriched["user_high_share"].values item_high_share = df_sample_enriched["item_high_share"].values X_en = np.vstack( [ user_ids_enc_en, item_ids_enc_en, years_en, recency_days, user_mean, item_mean, user_cnt, item_cnt, user_high_share, item_high_share, ] ).T y_en = df_sample_enriched["label"].values X_train_en, X_test_en, y_train_en, y_test_en = train_test_split( X_en, y_en, test_size=0.2, random_state=42, stratify=y_en ) logger.info( "Enriched features prepared for modeling, train size=%d, test size=%d", len(X_train_en), len(X_test_en), ) models_en: Dict[str, Pipeline] = { "logistic_regression_en": Pipeline( [ ("imputer", SimpleImputer(strategy="median")), ("scaler", StandardScaler()), ("classifier", LogisticRegression(max_iter=200, n_jobs=-1)), ] ), "sgd_classifier_en": Pipeline( [ ("imputer", SimpleImputer(strategy="median")), ("scaler", StandardScaler()), ( "classifier", SGDClassifier( loss="log_loss", learning_rate="optimal", max_iter=5, tol=1e-3, random_state=42, ), ), ] ), "catboost_en": Pipeline( [ ("imputer", SimpleImputer(strategy="median")), ("scaler", StandardScaler()), ( "classifier", CatBoostClassifier( depth=6, learning_rate=0.1, iterations=200, loss_function="Logloss", verbose=False, random_state=42, ), ), ] ), "minibatch_kmeans_en": Pipeline( [ ("imputer", SimpleImputer(strategy="median")), ("scaler", StandardScaler()), ("classifier", KMeansClassifier(n_clusters=2, batch_size=2048)), ] ), "random_forest_en": Pipeline( [ ("imputer", SimpleImputer(strategy="median")), ("scaler", StandardScaler()), ( "classifier", RandomForestClassifier(n_estimators=100, n_jobs=-1, random_state=42), ), ] ), "gradient_boosting_en": Pipeline( [ ("imputer", SimpleImputer(strategy="median")), ("scaler", StandardScaler()), ("classifier", GradientBoostingClassifier(random_state=42)), ] ), } logger.info("Enriched models initialized: %s", list(models_en.keys())) results_en = [] for name, model in models_en.items(): logger.info("Training enriched model: %s", name) model.fit(X_train_en, y_train_en) y_proba_en = model.predict_proba(X_test_en)[:, 1] y_pred_en = (y_proba_en >= 0.5).astype(int) auc_en = roc_auc_score(y_test_en, y_proba_en) f1_en = f1_score(y_test_en, y_pred_en) ll_en = log_loss(y_test_en, y_proba_en) results_en.append({"model": name, "auc": auc_en, "f1": f1_en, "log_loss": ll_en}) logger.info( "Enriched model %s: AUC=%.4f, F1=%.4f, log_loss=%.4f", name, auc_en, f1_en, ll_en ) results_en_df = pd.DataFrame(results_en).sort_values(by="auc", ascending=False) append_section(report_lines, "## Сравнение моделей с новыми фичами\n\n```text\n") append_section(report_lines, results_en_df.to_string(index=False) + "\n") append_section(report_lines, "```\n\n") append_section(report_lines, "### Описание новых фичей\n\n") append_section( report_lines, "В качестве дополнительных признаков добавлены агрегаты и временные фичи:\n\n" "- user_mean_rating — средний рейтинг пользователя;\n" "- user_rating_count — количество оценок пользователя;\n" "- user_high_share — доля высоких оценок пользователя (label=1);\n" "- item_mean_rating — средний рейтинг товара;\n" "- item_rating_count — количество оценок товара;\n" "- item_high_share — доля высоких оценок товара (label=1);\n" "- recency_days — давность взаимодействия в днях относительно последнего события.\n\n" "Эти признаки дополнены закодированными идентификаторами user_id и item_id, " "а также признаком года взаимодействия (year).\n\n", ) append_section(report_lines, "### Вывод по расширенным фичам\n\n") append_section( report_lines, "По сравнению с моделями на сырых признаках (user_id, item_id, year) " "добавление пользовательских и товарных агрегатов даёт существенный прирост качества: " "AUC и F1 заметно растут, а log_loss по вероятностям снижается. " "Это означает, что информация о средней оценке и активности пользователя/товара " "является важным источником сигнала для модели рекомендаций.\n\n", ) logger.info("Enriched model comparison section prepared") # Кривая обучения для лучшей расширенной модели try: best_model_name = results_en_df.iloc[0]["model"] best_model = models_en[best_model_name] logger.info("Computing learning curve for enriched model: %s", best_model_name) train_sizes, train_scores, val_scores = learning_curve( best_model, X_train_en, y_train_en, cv=3, scoring="roc_auc", n_jobs=-1, train_sizes=np.linspace(0.2, 1.0, 5), shuffle=True, random_state=42, ) train_mean = train_scores.mean(axis=1) val_mean = val_scores.mean(axis=1) lc_path = os.path.join(REPORT_DIR, f"learning_curve_enriched_{best_model_name}.png") fig, ax = plt.subplots(figsize=(5, 4)) ax.plot(train_sizes, train_mean, label="train AUC") ax.plot(train_sizes, val_mean, label="cv AUC") ax.set_xlabel("train size") ax.set_ylabel("AUC") ax.set_title(f"Кривая обучения (enriched, {best_model_name})") ax.legend() plt.tight_layout() fig.savefig(lc_path) plt.close(fig) append_section( report_lines, f"### Кривая обучения для лучшей расширенной модели ({best_model_name})\n\n", ) append_section( report_lines, f"})\n\n", ) logger.info("Learning curve saved to %s", lc_path) except Exception: logger.exception("Failed to compute or save learning curve for enriched model") return results_en_df, models_en, X_train_en, X_test_en, y_train_en, y_test_en def log_enriched_models_to_mlflow( results_en_df: pd.DataFrame, models_en: Dict[str, Pipeline], X_train_en: np.ndarray, df_sample_enriched: pd.DataFrame, logger: logging.Logger, ) -> None: """Логирует в MLflow все расширенные модели из results_en_df.""" try: tracking_uri = os.getenv("MLFLOW_TRACKING_URI", "http://localhost:5000") mlflow.set_tracking_uri(tracking_uri) mlflow.set_experiment("recsys_baseline") logger.info("MLflow tracking URI set to %s", tracking_uri) for row in results_en_df.itertuples(index=False): model_name = row.model auc = float(row.auc) f1 = float(row.f1) ll = float(row.log_loss) model_pipeline = models_en.get(model_name) if model_pipeline is None: logger.warning("Model %s not found in models_en, skip logging", model_name) continue # Пример входа и сигнатура input_example = X_train_en[:5] try: y_example = model_pipeline.predict_proba(X_train_en[:5])[:, 1] signature = infer_signature(X_train_en, y_example) except Exception: logger.warning("Failed to infer model signature automatically for %s", model_name) signature = None try: with mlflow.start_run(run_name=f"{model_name}_enriched"): mlflow.log_param("model_name", model_name) mlflow.log_param("sample_size", int(len(df_sample_enriched))) mlflow.log_metric("auc", auc) mlflow.log_metric("f1", f1) mlflow.log_metric("log_loss", ll) mlflow.sklearn.log_model( model_pipeline, "model", input_example=input_example, signature=signature, ) logger.info("Model %s logged to MLflow (AUC=%.4f, F1=%.4f, log_loss=%.4f)", model_name, auc, f1, ll) except Exception: logger.exception("Failed to log model %s to MLflow", model_name) except Exception: logger.exception("MLflow logging failed; pipeline results are still available locally") def write_report(report_lines: List[str], logger: logging.Logger) -> None: try: with open(REPORT_PATH, "w", encoding="utf-8") as f: f.writelines(report_lines) logger.info("Markdown report written to %s", REPORT_PATH) except Exception: logger.exception("Failed to write markdown report to %s", REPORT_PATH) def main() -> None: logger = setup_logger() logger.info("Recommender pipeline started") report_lines: List[str] = [] try: download_dataset_if_needed(logger) df = load_dataset(logger) add_initial_eda(df, report_lines, logger) add_target_and_correlation(df, report_lines, logger) add_data_issues(df, report_lines, logger) basic_results_df, basic_models, X_train, X_test, y_train, y_test = train_basic_models( df, report_lines, logger ) ( enriched_results_df, models_en, X_train_en, X_test_en, y_train_en, y_test_en, ) = train_enriched_models(df, report_lines, logger) # Для MLflow логируем все расширенные модели log_enriched_models_to_mlflow( enriched_results_df, models_en, X_train_en, df, logger # df используется как размер выборки ) logger.info("Recommender pipeline finished successfully") except Exception: logger.exception("Recommender pipeline failed with an unhandled error") finally: write_report(report_lines, logger) if __name__ == "__main__": main()