File size: 5,166 Bytes
422d4ca
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
"""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))

    # Subset meta to columns we'll use
    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"})

    # Build flattened meta text fields
    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"] = ""

    # Inner join
    merged = reviews_df.merge(meta_sub, on="parent_asin", how="inner", suffixes=("", "_meta"))
    logger.info("After inner join: %d rows", len(merged))

    # Clean review text
    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)

    # Combine title + body for model input
    merged["full_text"] = (merged["title"] + " " + merged["text"]).str.strip()

    # Filter junk
    merged = merged[merged["full_text"].str.len() >= cfg.MIN_REVIEW_LENGTH].copy()

    # Ensure rating is numeric
    merged["rating"] = pd.to_numeric(merged["rating"], errors="coerce")
    merged = merged.dropna(subset=["rating"]).copy()
    merged["rating"] = merged["rating"].astype(int)

    # Overall sentiment label from rating (3-class, used for baseline comparison)
    # 1-2 = Negative (2), 3 = Neutral (1), 4-5 = Positive (0)
    # We'll align with: 0=Negative, 1=Neutral, 2=Positive for baseline
    def rating_to_overall(r):
        if r <= 2: return 0  # Negative
        if r == 3: return 1  # Neutral
        return 2  # Positive
    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