{"metadata":{"kernelspec":{"display_name":"Python 3","language":"python","name":"python3"},"language_info":{"name":"python","version":"3.12.12","mimetype":"text/x-python","codemirror_mode":{"name":"ipython","version":3},"pygments_lexer":"ipython3","nbconvert_exporter":"python","file_extension":".py"},"kaggle":{"accelerator":"none","dataSources":[{"sourceType":"competition","sourceId":31254,"databundleVersionId":3103714}],"dockerImageVersionId":31328,"isInternetEnabled":true,"language":"python","sourceType":"notebook","isGpuEnabled":false}},"nbformat_minor":4,"nbformat":4,"cells":[{"cell_type":"markdown","source":"# Fashion Intelligence Engine\n## Spark-Based Trend Analytics & Recommendation System\n**Dataset:** H&M Personalized Fashion Recommendations  \n**Tools:** PySpark · RDD · DataFrame · Spark SQL · MLlib · Plotly","metadata":{}},{"cell_type":"markdown","source":"---\n## ⚙️ Section 1: Environment Setup","metadata":{}},{"cell_type":"code","source":"# Install PySpark on Kaggle\n!pip install pyspark --quiet\n\nimport os\nimport time\nimport warnings\nwarnings.filterwarnings('ignore')\n\nimport pandas as pd\nimport numpy as np\nimport plotly.express as px\nimport plotly.graph_objects as go\nfrom plotly.subplots import make_subplots\n\nfrom pyspark.sql import SparkSession\nfrom pyspark.sql import functions as F\nfrom pyspark.sql.types import *\nfrom pyspark.sql.window import Window\nfrom pyspark import StorageLevel","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:13:53.581061Z","iopub.execute_input":"2026-06-12T12:13:53.581403Z","iopub.status.idle":"2026-06-12T12:14:00.173877Z","shell.execute_reply.started":"2026-06-12T12:13:53.581378Z","shell.execute_reply":"2026-06-12T12:14:00.173064Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# Initialize Spark Session\nspark = SparkSession.builder \\\n    .appName('FashionIntelligenceEngine') \\\n    .config('spark.driver.memory', '4g') \\\n    .config('spark.executor.memory', '4g') \\\n    .config('spark.sql.shuffle.partitions', '50') \\\n    .config('spark.sql.adaptive.enabled', 'true') \\\n    .config('spark.sql.adaptive.coalescePartitions.enabled', 'true') \\\n    .getOrCreate()\n\nspark.sparkContext.setLogLevel('ERROR')\nprint(f'Spark version: {spark.version}')\nprint(f'SparkContext: {spark.sparkContext.appName}')","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:14:00.175281Z","iopub.execute_input":"2026-06-12T12:14:00.175781Z","iopub.status.idle":"2026-06-12T12:14:09.122915Z","shell.execute_reply.started":"2026-06-12T12:14:00.175744Z","shell.execute_reply":"2026-06-12T12:14:09.121947Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"---\n## 📂 Section 2: Data Loading & Schema Definition","metadata":{}},{"cell_type":"code","source":"import os\n\nBASE_PATH = '/kaggle/input/competitions/h-and-m-personalized-fashion-recommendations'\n\nCUSTOMERS_PATH = f'{BASE_PATH}/customers.csv'\nARTICLES_PATH = f'{BASE_PATH}/articles.csv'\nTRANSACTIONS_PATH = f'{BASE_PATH}/transactions_train.csv'\n\nprint('Paths configured:')\nfor p in [CUSTOMERS_PATH, ARTICLES_PATH, TRANSACTIONS_PATH]:\n    print(f'{\"Exsist: \" if os.path.exists(p) else \"Not exsist: \"} {p}')","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:14:09.124893Z","iopub.execute_input":"2026-06-12T12:14:09.125501Z","iopub.status.idle":"2026-06-12T12:14:09.132817Z","shell.execute_reply.started":"2026-06-12T12:14:09.125472Z","shell.execute_reply":"2026-06-12T12:14:09.132217Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# Define schemas explicitly for performance\n\ncustomer_schema = StructType([\n    StructField('customer_id',         StringType(),  True),\n    StructField('FN',                  FloatType(),   True),\n    StructField('Active',              FloatType(),   True),\n    StructField('club_member_status',  StringType(),  True),\n    StructField('fashion_news_frequency', StringType(), True),\n    StructField('age',                 FloatType(),   True),\n    StructField('postal_code',         StringType(),  True)\n])\n\narticle_schema = StructType([\n    StructField('article_id',              IntegerType(), True),\n    StructField('product_code',            IntegerType(), True),\n    StructField('prod_name',               StringType(),  True),\n    StructField('product_type_no',         IntegerType(), True),\n    StructField('product_type_name',       StringType(),  True),\n    StructField('product_group_name',      StringType(),  True),\n    StructField('graphical_appearance_no', IntegerType(), True),\n    StructField('graphical_appearance_name', StringType(), True),\n    StructField('colour_group_code',       IntegerType(), True),\n    StructField('colour_group_name',       StringType(),  True),\n    StructField('perceived_colour_value_id', IntegerType(), True),\n    StructField('perceived_colour_value_name', StringType(), True),\n    StructField('perceived_colour_master_id', IntegerType(), True),\n    StructField('perceived_colour_master_name', StringType(), True),\n    StructField('department_no',           IntegerType(), True),\n    StructField('department_name',         StringType(),  True),\n    StructField('index_code',              StringType(),  True),\n    StructField('index_name',              StringType(),  True),\n    StructField('index_group_no',          IntegerType(), True),\n    StructField('index_group_name',        StringType(),  True),\n    StructField('section_no',              IntegerType(), True),\n    StructField('section_name',            StringType(),  True),\n    StructField('garment_group_no',        IntegerType(), True),\n    StructField('garment_group_name',      StringType(),  True),\n    StructField('detail_desc',             StringType(),  True)\n])\n\ntransaction_schema = StructType([\n    StructField('t_dat',        StringType(),  True),\n    StructField('customer_id',  StringType(),  True),\n    StructField('article_id',   IntegerType(), True),\n    StructField('price',        DoubleType(),  True),\n    StructField('sales_channel_id', IntegerType(), True)\n])","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:14:15.713935Z","iopub.execute_input":"2026-06-12T12:14:15.714671Z","iopub.status.idle":"2026-06-12T12:14:15.722778Z","shell.execute_reply.started":"2026-06-12T12:14:15.714642Z","shell.execute_reply":"2026-06-12T12:14:15.722164Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# Load DataFrames\nt0 = time.time()\n\ncustomers_df = spark.read.csv(CUSTOMERS_PATH, header=True, schema=customer_schema)\narticles_df  = spark.read.csv(ARTICLES_PATH,  header=True, schema=article_schema)\ntrans_df     = spark.read.csv(TRANSACTIONS_PATH, header=True, schema=transaction_schema)\n\n# Convert date column\ntrans_df = trans_df.withColumn('t_dat', F.to_date('t_dat', 'yyyy-MM-dd'))\n\nprint(f'   Data loaded in {time.time()-t0:.2f}s')\nprint(f'   Customers:    {customers_df.count():>12,}')\nprint(f'   Articles:     {articles_df.count():>12,}')\nprint(f'   Transactions: {trans_df.count():>12,}')","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:14:28.660098Z","iopub.execute_input":"2026-06-12T12:14:28.660665Z","iopub.status.idle":"2026-06-12T12:14:48.564825Z","shell.execute_reply.started":"2026-06-12T12:14:28.660634Z","shell.execute_reply":"2026-06-12T12:14:48.562413Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"---\n## 🧹 Section 3: Data Cleaning","metadata":{}},{"cell_type":"code","source":"# Check nulls\nprint('=== NULL COUNTS ===')\nfor col_name in customers_df.columns:\n    n = customers_df.filter(F.col(col_name).isNull()).count()\n    if n > 0:\n        print(f'  customers.{col_name}: {n:,}')\n\nfor col_name in trans_df.columns:\n    n = trans_df.filter(F.col(col_name).isNull()).count()\n    if n > 0:\n        print(f'  transactions.{col_name}: {n:,}')","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:14:48.566323Z","iopub.execute_input":"2026-06-12T12:14:48.566657Z","iopub.status.idle":"2026-06-12T12:16:33.325969Z","shell.execute_reply.started":"2026-06-12T12:14:48.566625Z","shell.execute_reply":"2026-06-12T12:16:33.324977Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# Clean customers: fill missing age with median, fill club_member_status\nmedian_age = customers_df.approxQuantile('age', [0.5], 0.01)[0]\n\ncustomers_clean = customers_df \\\n    .fillna({'age': median_age,\n             'club_member_status': 'UNKNOWN',\n             'fashion_news_frequency': 'NONE'}) \\\n    .withColumn('age_group',\n        F.when(F.col('age') < 25, 'Gen-Z (< 25)')\n         .when(F.col('age') < 35, 'Millennial (25-34)')\n         .when(F.col('age') < 50, 'Gen-X (35-49)')\n         .when(F.col('age') < 65, 'Boomer (50-64)')\n         .otherwise('Senior (65+)'))\n\n# Clean transactions: drop rows with null price or customer\ntrans_clean = trans_df \\\n    .dropna(subset=['customer_id', 'article_id', 'price']) \\\n    .filter(F.col('price') > 0) \\\n    .withColumn('year',    F.year('t_dat')) \\\n    .withColumn('month',   F.month('t_dat')) \\\n    .withColumn('quarter', F.quarter('t_dat')) \\\n    .withColumn('season',\n        F.when(F.col('month').isin([12,1,2]),  'Winter')\n         .when(F.col('month').isin([3,4,5]),   'Spring')\n         .when(F.col('month').isin([6,7,8]),   'Summer')\n         .otherwise('Autumn'))\n\nprint(f'Cleaned customers: {customers_clean.count():,}')\nprint(f'Cleaned transactions: {trans_clean.count():,}')","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:16:33.327638Z","iopub.execute_input":"2026-06-12T12:16:33.328722Z","iopub.status.idle":"2026-06-12T12:17:00.283777Z","shell.execute_reply.started":"2026-06-12T12:16:33.328684Z","shell.execute_reply":"2026-06-12T12:17:00.283144Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# Cache frequently used DataFrames\ncustomers_clean.cache()\narticles_df.cache()\ntrans_clean.cache()\n\n# Register as SQL temp views\ncustomers_clean.createOrReplaceTempView('customers')\narticles_df.createOrReplaceTempView('articles')\ntrans_clean.createOrReplaceTempView('transactions')","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:17:00.284598Z","iopub.execute_input":"2026-06-12T12:17:00.284873Z","iopub.status.idle":"2026-06-12T12:17:00.482036Z","shell.execute_reply.started":"2026-06-12T12:17:00.284841Z","shell.execute_reply":"2026-06-12T12:17:00.481229Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"---\n## 🔵 Section 4: RDD Phase Implementation","metadata":{}},{"cell_type":"code","source":"# ── RDD Task 1: map() and filter() to remove nulls ──\nprint('=== RDD TASK 1: map() + filter() ===')\nt0 = time.time()\n\ntrans_rdd = spark.sparkContext.textFile(TRANSACTIONS_PATH)\n\n# Skip header\nheader = trans_rdd.first()\ntrans_rdd_clean = trans_rdd \\\n    .filter(lambda row: row != header) \\\n    .map(lambda row: row.split(',')) \\\n    .filter(lambda fields: len(fields) == 5) \\\n    .filter(lambda fields: all(f.strip() != '' for f in fields)) \\\n    .map(lambda fields: {\n        'date':        fields[0].strip(),\n        'customer_id': fields[1].strip(),\n        'article_id':  fields[2].strip(),\n        'price':       float(fields[3].strip()),\n        'channel':     fields[4].strip()\n    })\n\nrdd_count = trans_rdd_clean.count()\nrdd_time  = time.time() - t0\nprint(f'  Records after cleaning: {rdd_count:,}')\nprint(f'  Time: {rdd_time:.2f}s')","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:17:00.483688Z","iopub.execute_input":"2026-06-12T12:17:00.483957Z","iopub.status.idle":"2026-06-12T12:18:01.216533Z","shell.execute_reply.started":"2026-06-12T12:17:00.483928Z","shell.execute_reply":"2026-06-12T12:18:01.215697Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ── RDD Task 2: Pair RDDs + manual join ──\nprint('=== RDD TASK 2: Pair RDDs + Join ===')\nt0 = time.time()\n\n# article_id -> price pair RDD\narticle_price_rdd = trans_rdd_clean \\\n    .map(lambda r: (r['article_id'], r['price']))\n\n# article_id -> count pair RDD\narticle_count_rdd = trans_rdd_clean \\\n    .map(lambda r: (r['article_id'], 1))\n\n# Join: (article_id, (total_price, count))\njoined_rdd = article_price_rdd.join(article_count_rdd) \\\n    .map(lambda x: (x[0], x[1][0] + x[1][1]))  # combine price + count as demo\n\nsample = joined_rdd.take(5)\nrdd_join_time = time.time() - t0\nprint(f'  Sample joined records: {sample}')\nprint(f'  Time: {rdd_join_time:.2f}s')","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:18:01.217700Z","iopub.execute_input":"2026-06-12T12:18:01.218131Z","iopub.status.idle":"2026-06-12T12:21:22.780982Z","shell.execute_reply.started":"2026-06-12T12:18:01.218097Z","shell.execute_reply":"2026-06-12T12:21:22.780144Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ── RDD Task 3: reduceByKey() for category revenue ──\nprint('=== RDD TASK 3: reduceByKey() ===')\nt0 = time.time()\n\n# Load articles to a dict for broadcast\narticles_pd = articles_df.select('article_id', 'product_group_name').toPandas()\narticle_to_category = dict(zip(\n    articles_pd['article_id'].astype(str),\n    articles_pd['product_group_name']\n))\nbroadcast_categories = spark.sparkContext.broadcast(article_to_category)\n\n# Map: (category, price) → reduceByKey: sum\ncategory_revenue_rdd = trans_rdd_clean \\\n    .map(lambda r: (\n        broadcast_categories.value.get(r['article_id'], 'Unknown'),\n        r['price']\n    )) \\\n    .reduceByKey(lambda a, b: a + b) \\\n    .sortBy(lambda x: -x[1])\n\ntop_categories_rdd = category_revenue_rdd.take(10)\nrdd_reduce_time = time.time() - t0\n\nprint(f'  Top categories by revenue (RDD):')\nfor cat, rev in top_categories_rdd:\n    print(f'    {cat:<35} ${rev:,.2f}')\nprint(f'  Time: {rdd_reduce_time:.2f}s')","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:21:22.782404Z","iopub.execute_input":"2026-06-12T12:21:22.782704Z","iopub.status.idle":"2026-06-12T12:22:58.885280Z","shell.execute_reply.started":"2026-06-12T12:21:22.782681Z","shell.execute_reply":"2026-06-12T12:22:58.884624Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ── RDD : flatMap() for customer-article pairs ──\nprint('=== RDD: flatMap() ===')\n\ncustomer_items_rdd = trans_rdd_clean \\\n    .map(lambda r: (r['customer_id'], r['article_id'])) \\\n    .groupByKey() \\\n    .flatMap(lambda kv: [(kv[0], item) for item in list(kv[1])[:5]])\n\nprint(f'  Customer-item pairs sample: {customer_items_rdd.take(5)}')","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:22:58.886109Z","iopub.execute_input":"2026-06-12T12:22:58.886393Z","iopub.status.idle":"2026-06-12T12:24:28.179989Z","shell.execute_reply.started":"2026-06-12T12:22:58.886354Z","shell.execute_reply":"2026-06-12T12:24:28.179170Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"---\n## 🟢 Section 5: DataFrame Implementation","metadata":{}},{"cell_type":"code","source":"# Enrich transactions with article & customer info\nenriched_df = trans_clean \\\n    .join(F.broadcast(articles_df.select(\n        'article_id', 'prod_name', 'product_group_name',\n        'section_name', 'index_group_name', 'colour_group_name'\n    )), on='article_id', how='left') \\\n    .join(F.broadcast(customers_clean.select(\n        'customer_id', 'age', 'age_group', 'club_member_status'\n    )), on='customer_id', how='left')\n\nenriched_df.cache()\nenriched_df.createOrReplaceTempView('enriched')\n\nprint(' Enriched DataFrame created & cached')\nprint(f'   Columns: {enriched_df.columns}')","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:24:28.181375Z","iopub.execute_input":"2026-06-12T12:24:28.181743Z","iopub.status.idle":"2026-06-12T12:24:28.342089Z","shell.execute_reply.started":"2026-06-12T12:24:28.181718Z","shell.execute_reply":"2026-06-12T12:24:28.341073Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"---\n## 🔴 Section 6: Ten Required Queries","metadata":{}},{"cell_type":"code","source":"# ──────────────────────────────────────────────────\n# Q1: Filter customers purchasing items above price threshold\n# ──────────────────────────────────────────────────\nprint('=== Q1: High-Price Filter ===')\n\nPRICE_THRESHOLD = 0.05  # H&M prices are normalized\n\n# DataFrame API\nq1_df = trans_clean \\\n    .filter(F.col('price') > PRICE_THRESHOLD) \\\n    .select('customer_id', 'article_id', 'price', 't_dat') \\\n    .orderBy(F.desc('price'))\n\nq1_count = q1_df.count()\nprint(f'  [DataFrame] Transactions above {PRICE_THRESHOLD}: {q1_count:,}')\n\n# SQL\nq1_sql = spark.sql(f\"\"\"\n    SELECT customer_id, article_id, price, t_dat\n    FROM transactions\n    WHERE price > {PRICE_THRESHOLD}\n    ORDER BY price DESC\n\"\"\")\nprint(f'  [SQL]       Transactions above {PRICE_THRESHOLD}: {q1_sql.count():,}')\nq1_sql.show(5)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:26:02.858234Z","iopub.execute_input":"2026-06-12T12:26:02.858674Z","iopub.status.idle":"2026-06-12T12:26:05.637638Z","shell.execute_reply.started":"2026-06-12T12:26:02.858628Z","shell.execute_reply":"2026-06-12T12:26:05.636785Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ──────────────────────────────────────────────────\n# Q2: Total revenue by category\n# ──────────────────────────────────────────────────\nprint('=== Q2: Revenue by Category ===')\nt0 = time.time()\n\n# DataFrame API\nq2_df = enriched_df \\\n    .groupBy('product_group_name') \\\n    .agg(\n        F.sum('price').alias('total_revenue'),\n        F.count('*').alias('num_transactions'),\n        F.avg('price').alias('avg_price')\n    ) \\\n    .orderBy(F.desc('total_revenue'))\n\ndf_time = time.time() - t0\nprint(f'  [DataFrame] Time: {df_time:.2f}s')\n\n# SQL\nt0 = time.time()\nq2_sql = spark.sql(\"\"\"\n    SELECT a.product_group_name,\n           SUM(t.price)   AS total_revenue,\n           COUNT(*)       AS num_transactions,\n           AVG(t.price)   AS avg_price\n    FROM transactions t\n    JOIN articles a ON t.article_id = a.article_id\n    GROUP BY a.product_group_name\n    ORDER BY total_revenue DESC\n\"\"\")\nsql_time = time.time() - t0\nprint(f'  [SQL]       Time: {sql_time:.2f}s')\nq2_sql.show(10)\n\n# Visualize\nq2_pd = q2_sql.limit(10).toPandas()\nfig = px.bar(q2_pd, x='product_group_name', y='total_revenue',\n             title='Top 10 Product Groups by Total Revenue',\n             color='total_revenue', color_continuous_scale='Viridis',\n             labels={'product_group_name': 'Category', 'total_revenue': 'Revenue'})\nfig.update_layout(xaxis_tickangle=-30)\nfig.show()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:26:14.237132Z","iopub.execute_input":"2026-06-12T12:26:14.237869Z","iopub.status.idle":"2026-06-12T12:26:24.765402Z","shell.execute_reply.started":"2026-06-12T12:26:14.237836Z","shell.execute_reply":"2026-06-12T12:26:24.764619Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ──────────────────────────────────────────────────\n# Q3: Average purchase amount by age group\n# ──────────────────────────────────────────────────\nprint('=== Q3: Avg Purchase by Age Group ===')\n\nq3_df = enriched_df \\\n    .groupBy('age_group') \\\n    .agg(\n        F.avg('price').alias('avg_purchase'),\n        F.sum('price').alias('total_spent'),\n        F.count('*').alias('num_purchases'),\n        F.countDistinct('customer_id').alias('unique_customers')\n    ) \\\n    .orderBy('age_group')\n\nq3_df.show()\n\nq3_pd = q3_df.toPandas()\nfig = px.bar(q3_pd, x='age_group', y='avg_purchase',\n             title='Average Purchase Amount by Age Group',\n             color='avg_purchase', color_continuous_scale='RdYlGn',\n             text='avg_purchase')\nfig.update_traces(texttemplate='$%{text:.4f}', textposition='outside')\nfig.show()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:26:24.766632Z","iopub.execute_input":"2026-06-12T12:26:24.766895Z","iopub.status.idle":"2026-06-12T12:28:41.543468Z","shell.execute_reply.started":"2026-06-12T12:26:24.766873Z","shell.execute_reply":"2026-06-12T12:28:41.542661Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ──────────────────────────────────────────────────\n# Q4: Group by category and season\n# ──────────────────────────────────────────────────\nprint('=== Q4: Category × Season Analysis ===')\n\nq4_df = spark.sql(\"\"\"\n    SELECT\n        a.product_group_name AS category,\n        t.season,\n        COUNT(*)             AS transactions,\n        SUM(t.price)         AS revenue,\n        AVG(t.price)         AS avg_price\n    FROM enriched t\n    JOIN articles a ON t.article_id = a.article_id\n    GROUP BY a.product_group_name, t.season\n    ORDER BY revenue DESC\n\"\"\")\nq4_df.show(15)\n\n# Pivot for heatmap\nq4_pd = q4_df.toPandas()\nq4_pivot = q4_pd.pivot_table(index='category', columns='season', values='revenue', aggfunc='sum').fillna(0)\nfig = px.imshow(q4_pivot, title='Revenue Heatmap: Category × Season',\n                color_continuous_scale='Blues', aspect='auto')\nfig.show()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:28:41.544744Z","iopub.execute_input":"2026-06-12T12:28:41.545097Z","iopub.status.idle":"2026-06-12T12:28:52.854268Z","shell.execute_reply.started":"2026-06-12T12:28:41.545072Z","shell.execute_reply":"2026-06-12T12:28:52.853281Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ──────────────────────────────────────────────────\n# Q5: Sort products by popularity\n# ──────────────────────────────────────────────────\nprint('=== Q5: Products by Popularity ===')\n\nq5_df = spark.sql(\"\"\"\n    SELECT\n        t.article_id,\n        a.prod_name,\n        a.product_group_name,\n        COUNT(*)                          AS purchase_count,\n        COUNT(DISTINCT t.customer_id)     AS unique_buyers,\n        SUM(t.price)                      AS total_revenue\n    FROM transactions t\n    JOIN articles a ON t.article_id = a.article_id\n    GROUP BY t.article_id, a.prod_name, a.product_group_name\n    ORDER BY purchase_count DESC\n    LIMIT 20\n\"\"\")\nq5_df.show(20)\n\nq5_pd = q5_df.toPandas()\nfig = px.scatter(q5_pd, x='unique_buyers', y='total_revenue',\n                 size='purchase_count', color='product_group_name',\n                 hover_name='prod_name',\n                 title='Top 20 Products: Buyers vs Revenue')\nfig.show()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:28:52.855585Z","iopub.execute_input":"2026-06-12T12:28:52.856037Z","iopub.status.idle":"2026-06-12T12:30:15.017426Z","shell.execute_reply.started":"2026-06-12T12:28:52.855971Z","shell.execute_reply":"2026-06-12T12:30:15.016662Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ──────────────────────────────────────────────────\n# Q6: Rank top-selling products using Window Functions\n# ──────────────────────────────────────────────────\nprint('=== Q6: Window Function Ranking ===')\n\nfrom pyspark.sql.window import Window\n\n# Aggregate sales per product\ncategory_sales = enriched_df \\\n    .groupBy('product_group_name', 'article_id', 'prod_name') \\\n    .agg(\n        F.count('*').alias('purchase_count'),\n        F.sum('price').alias('revenue')\n    )\n\n# Define window\nw_purchase = Window.partitionBy('product_group_name') \\\n                   .orderBy(F.desc('purchase_count'))\n\nw_revenue = Window.partitionBy('product_group_name') \\\n                  .orderBy(F.desc('revenue'))\n\n# Apply ranking\nq6_df = category_sales \\\n    .withColumn(\n        'rank_in_category',\n        F.rank().over(w_purchase)\n    ) \\\n    .withColumn(\n        'revenue_rank',\n        F.dense_rank().over(w_revenue)\n    ) \\\n    .withColumn(\n        'pct_rank',\n        F.percent_rank().over(w_purchase)\n    ) \\\n    .filter(F.col('rank_in_category') <= 3) \\\n    .orderBy('product_group_name', 'rank_in_category')\n\nq6_df.show(30, truncate=False)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:30:15.018959Z","iopub.execute_input":"2026-06-12T12:30:15.019327Z","iopub.status.idle":"2026-06-12T12:30:27.233096Z","shell.execute_reply.started":"2026-06-12T12:30:15.019298Z","shell.execute_reply":"2026-06-12T12:30:27.231787Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ──────────────────────────────────────────────────\n# Q7: Moving Average Sales (7-day window)\n# ──────────────────────────────────────────────────\nprint('=== Q7: Moving Average Sales ===')\n\ndaily_sales = trans_clean \\\n    .groupBy('t_dat') \\\n    .agg(\n        F.sum('price').alias('daily_revenue'),\n        F.count('*').alias('daily_transactions')\n    ) \\\n    .orderBy('t_dat')\n\n# Window: 7-day moving average (FIXED for Spark 4)\nw7 = Window.orderBy(\n    F.unix_date('t_dat')\n).rowsBetween(-6, 0)\n\nq7_df = daily_sales \\\n    .withColumn(\n        'ma7_revenue',\n        F.avg('daily_revenue').over(w7)\n    ) \\\n    .withColumn(\n        'ma7_transactions',\n        F.avg('daily_transactions').over(w7)\n    )\n\nq7_df.show(10)\n\nq7_pd = q7_df.toPandas()\n\nfig = go.Figure()\n\nfig.add_trace(\n    go.Scatter(\n        x=q7_pd['t_dat'],\n        y=q7_pd['daily_revenue'],\n        name='Daily Revenue',\n        opacity=0.4\n    )\n)\n\nfig.add_trace(\n    go.Scatter(\n        x=q7_pd['t_dat'],\n        y=q7_pd['ma7_revenue'],\n        name='7-Day MA'\n    )\n)\n\nfig.update_layout(\n    title='Daily Revenue with 7-Day Moving Average',\n    xaxis_title='Date',\n    yaxis_title='Revenue'\n)\n\nfig.show()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:30:27.234865Z","iopub.execute_input":"2026-06-12T12:30:27.235214Z","iopub.status.idle":"2026-06-12T12:30:31.832874Z","shell.execute_reply.started":"2026-06-12T12:30:27.235177Z","shell.execute_reply":"2026-06-12T12:30:31.832110Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ──────────────────────────────────────────────────\n# Q8: Nested Query — High-Value Customers\n# ──────────────────────────────────────────────────\nprint('=== Q8: High-Value Customers (Nested Query) ===')\n\nq8_df = spark.sql(\"\"\"\n    SELECT\n        hv.customer_id,\n        hv.total_spent,\n        hv.num_purchases,\n        hv.avg_order_value,\n        c.age,\n        c.age_group,\n        c.club_member_status\n    FROM (\n        SELECT\n            customer_id,\n            SUM(price)   AS total_spent,\n            COUNT(*)     AS num_purchases,\n            AVG(price)   AS avg_order_value\n        FROM transactions\n        GROUP BY customer_id\n        HAVING SUM(price) > (\n            SELECT AVG(customer_total) * 2\n            FROM (\n                SELECT customer_id, SUM(price) AS customer_total\n                FROM transactions\n                GROUP BY customer_id\n            ) sub\n        )\n    ) hv\n    JOIN customers c ON hv.customer_id = c.customer_id\n    ORDER BY hv.total_spent DESC\n\"\"\")\nq8_count = q8_df.count()\nprint(f'  High-value customers: {q8_count:,}')\nq8_df.show(10)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:30:31.833807Z","iopub.execute_input":"2026-06-12T12:30:31.834145Z","iopub.status.idle":"2026-06-12T12:31:18.287975Z","shell.execute_reply.started":"2026-06-12T12:30:31.834110Z","shell.execute_reply":"2026-06-12T12:31:18.287074Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ──────────────────────────────────────────────────\n# Q9: Broadcast Join — Customer × Transactions\n# ──────────────────────────────────────────────────\nprint('=== Q9: Broadcast Join ===')\nt0 = time.time()\n\nq9_broadcast = trans_clean \\\n    .join(F.broadcast(customers_clean.select(\n        'customer_id', 'age_group', 'club_member_status'\n    )), on='customer_id', how='inner') \\\n    .groupBy('age_group', 'club_member_status') \\\n    .agg(\n        F.count('*').alias('transactions'),\n        F.sum('price').alias('total_revenue'),\n        F.avg('price').alias('avg_price')\n    ) \\\n    .orderBy(F.desc('total_revenue'))\n\nbroadcast_time = time.time() - t0\nprint(f'  Broadcast join time: {broadcast_time:.2f}s')\nq9_broadcast.show()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:31:18.289206Z","iopub.execute_input":"2026-06-12T12:31:18.289566Z","iopub.status.idle":"2026-06-12T12:31:28.130382Z","shell.execute_reply.started":"2026-06-12T12:31:18.289530Z","shell.execute_reply":"2026-06-12T12:31:28.129266Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# ──────────────────────────────────────────────────\n# Q10: Broadcast Join vs Sort-Merge Join Performance\n# ──────────────────────────────────────────────────\nprint('=== Q10: Join Performance Comparison ===')\n\n# Sort-Merge Join (disable broadcast)\nspark.conf.set('spark.sql.autoBroadcastJoinThreshold', '-1')\nt0 = time.time()\n\nq10_smj = trans_clean \\\n    .join(customers_clean.select('customer_id', 'age_group'), on='customer_id', how='inner') \\\n    .agg(F.count('*').alias('total'))\n\nsmj_count = q10_smj.collect()[0]['total']\nsmj_time  = time.time() - t0\n\n# Re-enable broadcast\nspark.conf.set('spark.sql.autoBroadcastJoinThreshold', '10485760')\nt0 = time.time()\n\nq10_bcast = trans_clean \\\n    .join(F.broadcast(customers_clean.select('customer_id', 'age_group')), on='customer_id', how='inner') \\\n    .agg(F.count('*').alias('total'))\n\nbcast_count = q10_bcast.collect()[0]['total']\nbcast_time  = time.time() - t0\n\nprint(f'  Sort-Merge Join: {smj_count:,} records | {smj_time:.2f}s')\nprint(f'  Broadcast Join:  {bcast_count:,} records | {bcast_time:.2f}s')\nprint(f'  Speedup: {smj_time/bcast_time:.2f}x faster with broadcast')\n\n# Bar chart\nfig = px.bar(\n    x=['Sort-Merge Join', 'Broadcast Join'],\n    y=[smj_time, bcast_time],\n    title='Join Strategy Performance Comparison',\n    labels={'x': 'Join Type', 'y': 'Execution Time (s)'},\n    color=['Sort-Merge Join', 'Broadcast Join'],\n    color_discrete_sequence=['#e74c3c', '#2ecc71']\n)\nfig.show()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:31:28.131442Z","iopub.execute_input":"2026-06-12T12:31:28.132644Z","iopub.status.idle":"2026-06-12T12:31:56.764168Z","shell.execute_reply.started":"2026-06-12T12:31:28.132618Z","shell.execute_reply":"2026-06-12T12:31:56.763434Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"---\n## 📈 Section 7: Scalability Experiments","metadata":{}},{"cell_type":"code","source":"# Experiment 1: Different partition sizes\nprint('=== EXPERIMENT 1: Partition Size Effect ===')\n\npartition_results = []\nfor n_partitions in [10, 50, 100, 200]:\n    t0 = time.time()\n    df_part = trans_clean.repartition(n_partitions)\n    rev = df_part.agg(F.sum('price')).collect()[0][0]\n    elapsed = time.time() - t0\n    partition_results.append({'partitions': n_partitions, 'time_sec': elapsed})\n    print(f'  {n_partitions:>4} partitions → {elapsed:.2f}s')\n\nfig = px.line(pd.DataFrame(partition_results), x='partitions', y='time_sec',\n              title='Execution Time vs Number of Partitions',\n              markers=True)\nfig.show()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:31:56.765489Z","iopub.execute_input":"2026-06-12T12:31:56.765919Z","iopub.status.idle":"2026-06-12T12:32:47.518175Z","shell.execute_reply.started":"2026-06-12T12:31:56.765895Z","shell.execute_reply":"2026-06-12T12:32:47.517395Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# Experiment 2: CSV vs Parquet\nprint('=== EXPERIMENT 2: CSV vs Parquet ===')\n\nPARQUET_PATH = '/kaggle/working/transactions.parquet'\n\n# Write Parquet\nt0 = time.time()\ntrans_clean.write.mode('overwrite').parquet(PARQUET_PATH)\nwrite_time = time.time() - t0\nprint(f'  Parquet write time: {write_time:.2f}s')\n\n# Read Parquet\nt0 = time.time()\ntrans_parquet = spark.read.parquet(PARQUET_PATH)\npq_rev = trans_parquet.agg(F.sum('price')).collect()[0][0]\npq_time = time.time() - t0\n\n# Read CSV\nt0 = time.time()\ntrans_csv2 = spark.read.csv(TRANSACTIONS_PATH, header=True, schema=transaction_schema)\ncsv_rev = trans_csv2.agg(F.sum('price')).collect()[0][0]\ncsv_time = time.time() - t0\n\nprint(f'  CSV  read + aggregate: {csv_time:.2f}s')\nprint(f'  Parquet read + aggregate: {pq_time:.2f}s')\nprint(f'  Parquet speedup: {csv_time/pq_time:.2f}x')\n\nfig = px.bar(x=['CSV', 'Parquet'], y=[csv_time, pq_time],\n             title='Storage Format: CSV vs Parquet Read Performance',\n             color=['CSV', 'Parquet'],\n             color_discrete_sequence=['#e74c3c', '#3498db'])\nfig.show()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:32:47.520658Z","iopub.execute_input":"2026-06-12T12:32:47.520951Z","iopub.status.idle":"2026-06-12T12:33:33.615630Z","shell.execute_reply.started":"2026-06-12T12:32:47.520920Z","shell.execute_reply":"2026-06-12T12:33:33.615082Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# Experiment 3: Caching ON vs OFF\nprint('=== EXPERIMENT 3: Caching Effect ===')\n\ndef build_expensive_df():\n    return (\n        trans_clean\n        .join(articles_df, 'article_id')\n        .join(customers_clean, 'customer_id')\n        .groupBy('customer_id', 'product_group_name', 'season')\n        .agg(\n            F.sum('price').alias('total_spent'),\n            F.count('*').alias('purchase_count')\n        )\n    )\n\n# NO CACHE — recomputes the full join+agg every time\nno_cache_df = build_expensive_df()\n\nt0 = time.time()\nno_cache_df.agg(F.sum('total_spent')).collect()\nno_cache_time = time.time() - t0\n\nt0 = time.time()\nno_cache_df.agg(F.avg('purchase_count')).collect()\nno_cache_time2 = time.time() - t0\n\n# WITH CACHE \ncached_df = build_expensive_df()\ncached_df.cache()\ncached_df.count() \n\nt0 = time.time()\ncached_df.agg(F.sum('total_spent')).collect()\ncache_time = time.time() - t0\n\nt0 = time.time()\ncached_df.agg(F.avg('purchase_count')).collect()\ncache_time2 = time.time() - t0\n\ncached_df.unpersist()\n\n# RESULTS\nprint(f'  No Cache - Query 1: {no_cache_time:.2f}s  | Query 2: {no_cache_time2:.2f}s  | Total: {no_cache_time + no_cache_time2:.2f}s')\nprint(f'  Cached   - Query 1: {cache_time:.2f}s  | Query 2: {cache_time2:.2f}s  | Total: {cache_time + cache_time2:.2f}s')\nprint(f'  Cache benefit on Query 1: {no_cache_time  - cache_time:+.2f}s')\nprint(f'  Cache benefit on Query 2: {no_cache_time2 - cache_time2:+.2f}s')\nprint(f'  Total savings           : {(no_cache_time + no_cache_time2) - (cache_time + cache_time2):+.2f}s')","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:33:47.430971Z","iopub.execute_input":"2026-06-12T12:33:47.431978Z","iopub.status.idle":"2026-06-12T12:35:47.415227Z","shell.execute_reply.started":"2026-06-12T12:33:47.431946Z","shell.execute_reply":"2026-06-12T12:35:47.414251Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"---\n## 🤖 Section 8: Recommendation System (ALS)\n","metadata":{}},{"cell_type":"code","source":"from pyspark.ml.recommendation import ALS\nfrom pyspark.ml.evaluation import RegressionEvaluator\nfrom pyspark.sql.functions import col\n\n# Prepare data: need integer user/item IDs\nfrom pyspark.ml.feature import StringIndexer\n\nuser_indexer = StringIndexer(inputCol='customer_id', outputCol='user_idx')\nuser_model   = user_indexer.fit(trans_clean)\ntrans_indexed = user_model.transform(trans_clean) \\\n    .withColumn('user_idx',    F.col('user_idx').cast(IntegerType())) \\\n    .withColumn('article_idx', F.col('article_id').cast(IntegerType())) \\\n    .withColumn('rating',      F.lit(1.0))\n\n# Use a sample for speed on Kaggle\nals_data = trans_indexed.select('user_idx', 'article_idx', 'rating').sample(0.1, seed=42)\n\ntrain, test = als_data.randomSplit([0.8, 0.2], seed=42)\nprint(f'  Train: {train.count():,} | Test: {test.count():,}')","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:35:47.416663Z","iopub.execute_input":"2026-06-12T12:35:47.417092Z","iopub.status.idle":"2026-06-12T12:37:12.101508Z","shell.execute_reply.started":"2026-06-12T12:35:47.417057Z","shell.execute_reply":"2026-06-12T12:37:12.100223Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"from pyspark.ml.recommendation import ALS\nfrom pyspark.ml.evaluation import RegressionEvaluator\nfrom pyspark.sql import functions as F\n\n# ── Step 1: Clear all cached data from previous attempts ──────────────────\nspark.catalog.clearCache()\n\n# ── Step 2: Sample early, project early, write to parquet to break lineage ─\n(\n    trans_clean\n    .sample(0.05, seed=42)                          # 5% — even smaller to be safe\n    .select(\n        (F.hash(F.col(\"customer_id\")) % 1_000_000).alias(\"user_idx\"),\n        F.col(\"article_id\").cast(\"int\").alias(\"article_idx\"),\n        F.log1p(F.col(\"price\")).alias(\"rating\")\n    )\n    .filter(F.col(\"rating\") > 0)\n    .write.mode(\"overwrite\").parquet(\"/tmp/als_input\")\n)\n\n# ── Step 3: Read back from parquet (clean slate, no lineage) ───────────────\nals_input = spark.read.parquet(\"/tmp/als_input\")\n\n# ── Step 4: Split ──────────────────────────────────────────────────────────\ntrain_raw, test_raw = als_input.randomSplit([0.8, 0.2], seed=42)\n\n# ── Step 5: Filter sparse users/items ─────────────────────────────────────\nuser_counts = train_raw.groupBy(\"user_idx\").count()\nitem_counts = train_raw.groupBy(\"article_idx\").count()\n\ntrain_final = (\n    train_raw\n    .join(user_counts.filter(\"count >= 2\"), \"user_idx\")\n    .join(item_counts.filter(\"count >= 2\"), \"article_idx\")\n    .select(\"user_idx\", \"article_idx\", \"rating\")\n)\n\ntest_final = test_raw.select(\"user_idx\", \"article_idx\", \"rating\")\n\n# ── Step 6: Cache ──────────────────────────────────────────────────────────\ntrain_final.cache()\ntest_final.cache()\nprint(\"train:\", train_final.count(), \"  test:\", test_final.count())\n\n# ── Step 7: ALS ────────────────────────────────────────────────────────────\nals = ALS(\n    maxIter=3, rank=6, regParam=0.15,\n    userCol=\"user_idx\", itemCol=\"article_idx\", ratingCol=\"rating\",\n    coldStartStrategy=\"drop\", implicitPrefs=True\n)\nmodel = als.fit(train_final)\n\n# ── Step 8: Evaluate ───────────────────────────────────────────────────────\npred = model.transform(test_final.sample(0.05, seed=42))\nevaluator = RegressionEvaluator(\n    labelCol=\"rating\", predictionCol=\"prediction\", metricName=\"rmse\"\n)\nprint(\"RMSE:\", evaluator.evaluate(pred))","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:37:12.102793Z","iopub.execute_input":"2026-06-12T12:37:12.103366Z","iopub.status.idle":"2026-06-12T12:37:59.897182Z","shell.execute_reply.started":"2026-06-12T12:37:12.103317Z","shell.execute_reply":"2026-06-12T12:37:59.896180Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# Run this to inspect what you have available\nprint([x for x in dir() if not x.startswith('_')])","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:37:59.899422Z","iopub.execute_input":"2026-06-12T12:37:59.899711Z","iopub.status.idle":"2026-06-12T12:37:59.905161Z","shell.execute_reply.started":"2026-06-12T12:37:59.899685Z","shell.execute_reply":"2026-06-12T12:37:59.903800Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# The model trained above is called `model`, not `als_model`\nuser_recs = model.recommendForAllUsers(5)\nprint('=== Top-5 Recommendations per User ===')\nuser_recs.show(10, truncate=False)","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:37:59.906133Z","iopub.execute_input":"2026-06-12T12:37:59.906405Z","iopub.status.idle":"2026-06-12T12:42:16.353061Z","shell.execute_reply.started":"2026-06-12T12:37:59.906381Z","shell.execute_reply":"2026-06-12T12:42:16.352092Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"---\n## 📊 Section 9: Trend Prediction & Customer Segmentation","metadata":{}},{"cell_type":"code","source":"from pyspark.ml.clustering import KMeans\nfrom pyspark.ml.feature import VectorAssembler, StandardScaler\n\n# Build customer features\ncustomer_features = trans_clean \\\n    .groupBy('customer_id') \\\n    .agg(\n        F.count('*').alias('num_purchases'),\n        F.sum('price').alias('total_spent'),\n        F.avg('price').alias('avg_price'),\n        F.countDistinct('article_id').alias('unique_articles'),\n        F.max('t_dat').alias('last_purchase')\n    ) \\\n    .dropna()\n\nassembler = VectorAssembler(\n    inputCols=['num_purchases', 'total_spent', 'avg_price', 'unique_articles'],\n    outputCol='features'\n)\nscaler = StandardScaler(inputCol='features', outputCol='scaled_features')\n\nfeat_df      = assembler.transform(customer_features)\nscaler_model = scaler.fit(feat_df)\nfeat_scaled  = scaler_model.transform(feat_df)\n\n# KMeans clustering\nkmeans   = KMeans(k=4, seed=42, featuresCol='scaled_features', maxIter=20)\nkm_model = kmeans.fit(feat_scaled)\nclustered = km_model.transform(feat_scaled)\n\nseg_summary = clustered \\\n    .groupBy('prediction') \\\n    .agg(\n        F.count('*').alias('customers'),\n        F.avg('num_purchases').alias('avg_purchases'),\n        F.avg('total_spent').alias('avg_spent'),\n        F.avg('avg_price').alias('avg_price')\n    ) \\\n    .orderBy('prediction')\n\nseg_summary.show()\n\nseg_pd = seg_summary.toPandas()\nfig = px.scatter(seg_pd, x='avg_purchases', y='avg_spent',\n                 size='customers', color='prediction',\n                 title='Customer Segments (KMeans k=4)',\n                 labels={'prediction': 'Segment'})\nfig.show()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:42:16.354279Z","iopub.execute_input":"2026-06-12T12:42:16.354609Z","iopub.status.idle":"2026-06-12T12:52:57.724505Z","shell.execute_reply.started":"2026-06-12T12:42:16.354572Z","shell.execute_reply":"2026-06-12T12:52:57.723861Z"}},"outputs":[],"execution_count":null},{"cell_type":"code","source":"# Fashion trend analysis: monthly top categories\nmonthly_trends = spark.sql(\"\"\"\n    SELECT\n        YEAR(t.t_dat)          AS year,\n        MONTH(t.t_dat)         AS month,\n        a.index_group_name     AS category,\n        COUNT(*)               AS purchases,\n        SUM(t.price)           AS revenue\n    FROM transactions t\n    JOIN articles a ON t.article_id = a.article_id\n    GROUP BY YEAR(t.t_dat), MONTH(t.t_dat), a.index_group_name\n    ORDER BY year, month, revenue DESC\n\"\"\")\n\nmonthly_pd = monthly_trends.toPandas()\nmonthly_pd['year_month'] = monthly_pd['year'].astype(str) + '-' + monthly_pd['month'].astype(str).str.zfill(2)\n\nfig = px.line(monthly_pd, x='year_month', y='purchases',\n              color='category',\n              title='Monthly Purchase Trends by Category',\n              labels={'year_month': 'Month', 'purchases': 'Purchases'})\nfig.update_xaxes(tickangle=-45)\nfig.show()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:52:57.725546Z","iopub.execute_input":"2026-06-12T12:52:57.725787Z","iopub.status.idle":"2026-06-12T12:53:41.297378Z","shell.execute_reply.started":"2026-06-12T12:52:57.725766Z","shell.execute_reply":"2026-06-12T12:53:41.296403Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"---\n## 📋 Section 10: Performance Summary Dashboard","metadata":{}},{"cell_type":"code","source":"# Collect all timing results\nperf_data = {\n    'Operation': [\n        'RDD filter+map',\n        'RDD reduceByKey',\n        'DataFrame Q2 (groupBy)',\n        'SQL Q2 (GROUP BY)',\n        'Sort-Merge Join',\n        'Broadcast Join',\n        'Parquet read',\n        'CSV read'\n    ],\n    'Time (s)': [\n        rdd_time,\n        rdd_reduce_time,\n        df_time,\n        sql_time,\n        smj_time,\n        bcast_time,\n        pq_time,\n        csv_time\n    ],\n    'API': ['RDD','RDD','DataFrame','SQL','DataFrame','DataFrame','DataFrame','DataFrame']\n}\n\nperf_df = pd.DataFrame(perf_data)\nprint('=== PERFORMANCE SUMMARY ===')\nprint(perf_df.to_string(index=False))\n\nfig = px.bar(perf_df, x='Operation', y='Time (s)', color='API',\n             title='Performance Comparison: All Operations',\n             color_discrete_map={'RDD':'#e74c3c', 'DataFrame':'#3498db', 'SQL':'#2ecc71'})\nfig.update_layout(xaxis_tickangle=-30)\nfig.show()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:53:41.298620Z","iopub.execute_input":"2026-06-12T12:53:41.298924Z","iopub.status.idle":"2026-06-12T12:53:41.358856Z","shell.execute_reply.started":"2026-06-12T12:53:41.298895Z","shell.execute_reply":"2026-06-12T12:53:41.358320Z"}},"outputs":[],"execution_count":null},{"cell_type":"markdown","source":"---\n## ✅ Section 11: Cleanup","metadata":{}},{"cell_type":"code","source":"# Unpersist cached DataFrames\ncustomers_clean.unpersist()\narticles_df.unpersist()\ntrans_clean.unpersist()\nenriched_df.unpersist()\n\nspark.stop()","metadata":{"trusted":true,"execution":{"iopub.status.busy":"2026-06-12T12:53:41.359782Z","iopub.execute_input":"2026-06-12T12:53:41.360874Z","iopub.status.idle":"2026-06-12T12:53:43.359930Z","shell.execute_reply.started":"2026-06-12T12:53:41.360835Z","shell.execute_reply":"2026-06-12T12:53:43.359081Z"}},"outputs":[],"execution_count":null}]}