| """Preprocess raw review + meta data: clean, join, stratified split.""" |
| import html |
| import logging |
| import re |
| from typing import Optional |
|
|
| import numpy as np |
| import pandas as pd |
| from sklearn.model_selection import train_test_split |
|
|
| from . import config as cfg |
|
|
|
|
| logger = logging.getLogger(__name__) |
|
|
|
|
| _URL_RE = re.compile(r"https?://\S+|www\.\S+") |
| _HTML_RE = re.compile(r"<[^>]+>") |
| _WHITESPACE_RE = re.compile(r"\s+") |
|
|
|
|
| def clean_text(text: str) -> str: |
| """Light cleaning suitable for BERT input. |
| |
| NOTE: do NOT remove stopwords or lemmatize - BERT needs the original context. |
| """ |
| if not isinstance(text, str): |
| return "" |
| text = html.unescape(text) |
| text = _URL_RE.sub(" ", text) |
| text = _HTML_RE.sub(" ", text) |
| text = _WHITESPACE_RE.sub(" ", text).strip() |
| return text |
|
|
|
|
| def _features_to_text(features) -> str: |
| """Flatten meta `features` field (which can be list or string) to one blob.""" |
| if features is None: |
| return "" |
| if isinstance(features, (list, tuple, np.ndarray)): |
| return " ".join(str(x) for x in features if x) |
| return str(features) |
|
|
|
|
| def _categories_to_text(categories) -> str: |
| if categories is None: |
| return "" |
| if isinstance(categories, (list, tuple, np.ndarray)): |
| return " > ".join(str(x) for x in categories if x) |
| return str(categories) |
|
|
|
|
| def join_and_clean(reviews_df: pd.DataFrame, meta_df: pd.DataFrame) -> pd.DataFrame: |
| """Inner join review + meta on parent_asin, clean text, build derived columns.""" |
| logger.info("Pre-join: %d reviews, %d meta items.", len(reviews_df), len(meta_df)) |
|
|
| |
| meta_cols = ["parent_asin", "title", "features", "categories", "main_category", "price", "average_rating", "rating_number"] |
| meta_cols = [c for c in meta_cols if c in meta_df.columns] |
| meta_sub = meta_df[meta_cols].copy() |
| meta_sub = meta_sub.rename(columns={"title": "product_title"}) |
|
|
| |
| if "features" in meta_sub.columns: |
| meta_sub["features_text"] = meta_sub["features"].apply(_features_to_text) |
| else: |
| meta_sub["features_text"] = "" |
| if "categories" in meta_sub.columns: |
| meta_sub["categories_text"] = meta_sub["categories"].apply(_categories_to_text) |
| meta_sub["leaf_category"] = meta_sub["categories"].apply( |
| lambda c: (c[-1] if isinstance(c, (list, tuple, np.ndarray)) and len(c) > 0 else "") |
| ) |
| else: |
| meta_sub["categories_text"] = "" |
| meta_sub["leaf_category"] = "" |
|
|
| |
| merged = reviews_df.merge(meta_sub, on="parent_asin", how="inner", suffixes=("", "_meta")) |
| logger.info("After inner join: %d rows", len(merged)) |
|
|
| |
| if "title" in merged.columns: |
| merged["title"] = merged["title"].fillna("").astype(str).apply(clean_text) |
| else: |
| merged["title"] = "" |
| merged["text"] = merged["text"].fillna("").astype(str).apply(clean_text) |
|
|
| |
| merged["full_text"] = (merged["title"] + " " + merged["text"]).str.strip() |
|
|
| |
| merged = merged[merged["full_text"].str.len() >= cfg.MIN_REVIEW_LENGTH].copy() |
|
|
| |
| merged["rating"] = pd.to_numeric(merged["rating"], errors="coerce") |
| merged = merged.dropna(subset=["rating"]).copy() |
| merged["rating"] = merged["rating"].astype(int) |
|
|
| |
| |
| |
| def rating_to_overall(r): |
| if r <= 2: return 0 |
| if r == 3: return 1 |
| return 2 |
| merged["overall_label"] = merged["rating"].apply(rating_to_overall) |
|
|
| logger.info("Final rows: %d", len(merged)) |
| logger.info("Rating distribution:\n%s", merged["rating"].value_counts().sort_index().to_string()) |
|
|
| return merged.reset_index(drop=True) |
|
|
|
|
| def stratified_split(df: pd.DataFrame, label_col: str = "rating"): |
| """Stratified split by rating to keep class balance across train/val/test.""" |
| train_val, test = train_test_split( |
| df, test_size=cfg.TEST_RATIO, stratify=df[label_col], random_state=cfg.RANDOM_SEED |
| ) |
| val_size_relative = cfg.VAL_RATIO / (cfg.TRAIN_RATIO + cfg.VAL_RATIO) |
| train, val = train_test_split( |
| train_val, test_size=val_size_relative, |
| stratify=train_val[label_col], random_state=cfg.RANDOM_SEED |
| ) |
| logger.info("Split sizes: train=%d val=%d test=%d", len(train), len(val), len(test)) |
| return train.reset_index(drop=True), val.reset_index(drop=True), test.reset_index(drop=True) |
|
|
|
|
| def preprocess_pipeline(): |
| """Run the full preprocessing pipeline from raw files.""" |
| logger.info("Loading raw data...") |
| reviews_df = pd.read_parquet(cfg.RAW_REVIEWS_PATH) |
| meta_df = pd.read_parquet(cfg.RAW_META_PATH) |
|
|
| processed = join_and_clean(reviews_df, meta_df) |
| processed.to_parquet(cfg.PROCESSED_PATH, index=False) |
| logger.info("Saved processed data to %s", cfg.PROCESSED_PATH) |
| return processed |
|
|