{
 "cells": [
  {
   "cell_type": "markdown",
   "id": "bcb8f39d",
   "metadata": {},
   "source": [
    "\n",
    "<a id='chap-tppyspark'></a>"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "7fcf9dc9",
   "metadata": {},
   "source": [
    "# Travaux pratiques - Introduction à PySpark et aux DataFrames Spark\n",
    "\n",
    "Références externes utiles :\n",
    "\n",
    "> - [Documentation Spark](https://spark.apache.org/docs/latest/)  \n",
    "- [Démarrage avec les DataFrames Spark en Python](https://spark.apache.org/docs/latest/api/python/getting_started/quickstart_df.html)  \n",
    "- [Documentation API PySpark](https://spark.apache.org/docs/latest/api/python/index.html)  \n",
    "- [Tutoriel Spark avec exemples](https://sparkbyexamples.com/pyspark-tutorial/)  \n",
    "- [Learning Spark, 2e édition (O’Reilly)](https://pages.databricks.com/rs/094-YMS-629/images/LearningSpark2.0.pdf) (disponible gratuitement)  "
   ]
  },
  {
   "cell_type": "markdown",
   "id": "9663378b",
   "metadata": {},
   "source": [
    "# Objectif de cette séance\n",
    "\n",
    "Cette séance introduit PySpark, l’interface Python d’Apache Spark, avec pour objectifs :\n",
    "\n",
    "1. Comprendre comment démarrer une session Spark et s’y connecter depuis Jupyter.  \n",
    "1. Créer, lire, écrire et inspecter des `DataFrames` **Spark**.  \n",
    "1. Maîtriser les opérations fondamentales : sélection, filtrage, colonnes calculées,\n",
    "  agrégation, jointure.  \n",
    "1. Comprendre la distinction entre transformations et actions, ainsi que l’évaluation paresseuse (*lazy*).  \n",
    "1. Utiliser Spark SQL pour interroger des DataFrames.  \n",
    "1. Mettre en œuvre l”**échantillonnage** dans Spark.  \n",
    "1. Utiliser `cache()` et `persist()` à bon escient.  \n",
    "\n",
    "\n",
    "Nous travaillons sur des données peu volumineuses mais les mécanismes mis en œuvre sont les mêmes que\n",
    "sur un vrai *cluster* à grande échelle."
   ]
  },
  {
   "cell_type": "markdown",
   "id": "23c8e895",
   "metadata": {},
   "source": [
    "# Mise en place de l’environnement\n",
    "\n",
    "JupyterHub Cnam\n",
    "\n",
    "Accédez au serveur JupyterHub via [Moodle](https://lecnam.net) (Mes\n",
    "enseignements > RCP216 > Accès à JupyterHub). PySpark est déjà installé.\n",
    "Importez le cahier Jupyter de ce TP (bouton de téléversement).\n",
    "\n",
    "Machine personnelle\n",
    "\n",
    "Installez PySpark avec `pip install pyspark`. Une installation Java (JDK 21\n",
    "ou 25) est nécessaire. Suivez les\n",
    "[instructions d’installation](http://cedric.cnam.fr/vertigo/Cours/RCP216/installationSpark.html)\n",
    "en cas de difficulté.\n",
    "\n",
    "Dans Jupyter, initialisez la session Spark ainsi :"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "845e31f7",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "import os, sys\n",
    "os.environ['PYSPARK_PYTHON']        = sys.executable\n",
    "os.environ['PYSPARK_DRIVER_PYTHON'] = sys.executable\n",
    "\n",
    "from pyspark.sql import SparkSession\n",
    "spark = SparkSession.builder \\\n",
    "            .master(\"local[*]\") \\\n",
    "            .appName(\"RCP216-TP3\") \\\n",
    "            .getOrCreate()\n",
    "\n",
    "print(spark.version)"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "fc45baef",
   "metadata": {},
   "source": [
    "`local[*]` indique à Spark d’utiliser tous les cœurs disponibles sur la machine\n",
    "locale. Sur un vrai *cluster* on remplacerait cette URL par l’adresse du *cluster manager*."
   ]
  },
  {
   "cell_type": "markdown",
   "id": "5fb5d42c",
   "metadata": {},
   "source": [
    "## Premiers pas : texte et évaluation paresseuse"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "f96bde50",
   "metadata": {},
   "source": [
    "### Préparation des données\n",
    "\n",
    "Téléchargeons d’abord quelques fichiers de travail."
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "8dda2dee",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "import os\n",
    "os.system(\"mkdir -p tpPySpark/data\")\n",
    "os.system(\"wget -nc https://cedric.cnam.fr/vertigo/Cours/RCP216/docs/sophie.txt \"\n",
    "          \"-P tpPySpark/data/\")\n",
    "os.system(\"wget -nc http://cedric.cnam.fr/vertigo/Cours/RCP216/docs/geyser.csv \"\n",
    "          \"-P tpPySpark/data/\")\n",
    "os.system(\"wget -nc --no-check-certificate https://files.grouplens.org/datasets/movielens/ml-latest-small.zip \"\n",
    "          \"-P tpPySpark/\")\n",
    "os.system(\"unzip -n tpPySpark/ml-latest-small.zip -d tpPySpark/\")\n",
    "os.system(\"cp tpPySpark/ml-latest-small/*.csv tpPySpark/data/\")"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "6b2d86dc",
   "metadata": {},
   "source": [
    "### DataFrame de texte : actions et transformations\n",
    "\n",
    "Créons un premier *DataFrame* à partir du fichier texte `sophie.txt` (première partie\n",
    "des « Malheurs de Sophie » de la comtesse de Ségur) :"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "a2d37c37",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "texteSophie = spark.read.text(\"tpPySpark/data/sophie.txt\")\n",
    "\n",
    "# Inspecter la structure\n",
    "texteSophie.printSchema()   # une seule colonne \"value\" de type string\n",
    "texteSophie.show(5)"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "edd60be5",
   "metadata": {},
   "source": [
    "Spark crée un *DataFrame* dont chaque ligne correspond à une ligne du fichier. La\n",
    "lecture est **paresseuse** : le fichier n’est pas encore chargé en mémoire.\n",
    "\n",
    "**Actions** : déclenchent l’exécution (et la lecture) :"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "e1309b55",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "texteSophie.count()          # nombre de lignes (déclenche la lecture + comptage)"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "4e07caaf",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "texteSophie.first()          # première ligne"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "cee55c87",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "texteSophie.take(3)          # 3 premières lignes sous forme de liste Python"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "9cbb0e5f",
   "metadata": {},
   "source": [
    "### Attention à `collect()`\n",
    "\n",
    "`texteSophie.collect()` ramène **toutes** les données sur le nœud *driver*. Sur un\n",
    "vrai cluster avec des téraoctets de données, cette opération provoquerait un\n",
    "débordement mémoire. À éviter donc, sauf pour de très petits *DataFrames* de résultats finaux.\n",
    "\n",
    "**Transformations** : paresseuses, ne déclenchent pas de calcul immédiat :"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "72596d7f",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "from pyspark.sql.functions import length, col\n",
    "\n",
    "# Filtrer les lignes contenant \"poupée\"\n",
    "lignesPoupee = texteSophie.filter(texteSophie.value.contains(\"poupée\"))\n",
    "\n",
    "# Calculer la longueur de chaque ligne\n",
    "longueursLignes = texteSophie.select(length(texteSophie.value).alias(\"longueur\"))"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "a3b2c02b",
   "metadata": {},
   "source": [
    "À ce stade, aucun calcul n’a encore lieu. Ce n’est qu’à l’appel d’une **action** que\n",
    "Spark planifie et exécute l’ensemble du graphe de transformations :"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "d79840cd",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "longueursLignes.show(5)      # action : déclenche lecture + calcul des longueurs"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "7c1f982a",
   "metadata": {},
   "source": [
    "### Question\n",
    "\n",
    "1. Combien de lignes du fichier contiennent le mot « poupée » ?  \n",
    "1. Enchaînez les transformations et l’action en une seule expression pour calculer\n",
    "  la **longueur totale** du texte (somme des longueurs de toutes les lignes).\n",
    "  Utilisez `from pyspark.sql.functions import sum`.  "
   ]
  },
  {
   "cell_type": "markdown",
   "id": "7dc8e21b",
   "metadata": {},
   "source": [
    "### Évaluation paresseuse et optimisation\n",
    "\n",
    "La puissance de l’évaluation paresseuse est visible lorsqu’on enchaîne plusieurs\n",
    "transformations : Spark les optimise globalement avant exécution."
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "899e3f78",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "# Ces trois transformations ne font rien immédiatement\n",
    "etape1 = texteSophie.filter(length(texteSophie.value) > 20)\n",
    "etape2 = etape1.select(texteSophie.value, length(texteSophie.value).alias(\"longueur\"))\n",
    "etape3 = etape2.filter(col(\"longueur\") < 80)\n",
    "\n",
    "# Le plan (optimisé par Catalyst) est exécuté seulement ici\n",
    "etape3.show(5)\n",
    "\n",
    "# Afficher le plan d'exécution physique optimisé\n",
    "etape3.explain()"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "d931836d",
   "metadata": {},
   "source": [
    "En lisant la sortie de `explain()` on constate que Catalyst a fusionné les deux\n",
    "filtres en un seul passage sur les données (*filter pushdown*)."
   ]
  },
  {
   "cell_type": "markdown",
   "id": "ef1ff424",
   "metadata": {},
   "source": [
    "### Persistance : `cache()` et `persist()`\n",
    "\n",
    "Lorsqu’un DataFrame est utilisé plusieurs fois dans des calculs **différents**, Spark le\n",
    "**recalcule entièrement à chaque action** (en rejouant le *lineage*). Pour éviter ce coût,\n",
    "on peut demander à Spark de **conserver** le DataFrame en mémoire :"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "7790f9a2",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "# Sans cache : le fichier est relu deux fois\n",
    "nb_lignes   = texteSophie.count()\n",
    "nb_poupee   = texteSophie.filter(texteSophie.value.contains(\"poupée\")).count()\n",
    "\n",
    "# Avec cache : le fichier est lu une seule fois, le DataFrame est conservé en mémoire\n",
    "texteSophie = spark.read.text(\"tpPySpark/data/sophie.txt\").cache()\n",
    "nb_lignes   = texteSophie.count()      # premier appel : lecture + mise en cache\n",
    "nb_poupee   = texteSophie.filter(      # deuxième appel : lecture depuis le cache\n",
    "    texteSophie.value.contains(\"poupée\")).count()\n",
    "\n",
    "# Libérer le cache quand on n'en a plus besoin\n",
    "texteSophie.unpersist()"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "dd676131",
   "metadata": {},
   "source": [
    "`persist()` offre plus de contrôle sur le niveau de stockage :"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "1051be45",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "from pyspark import StorageLevel\n",
    "texteSophie.persist(StorageLevel.MEMORY_AND_DISK)\n",
    "# MEMORY_ONLY, MEMORY_AND_DISK, DISK_ONLY, MEMORY_ONLY_2 (réplication x2), ..."
   ]
  },
  {
   "cell_type": "markdown",
   "id": "aa3cd2ba",
   "metadata": {},
   "source": [
    "### Question\n",
    "\n",
    "Reprenez les calculs `nb_lignes` et `nb_poupee` en activant `cache()`.\n",
    "Observez dans l’interface Spark UI (accessible sur `http://localhost:4040`\n",
    "si Spark tourne en local sur votre ordinateur) la différence entre le premier\n",
    "appel (lecture depuis disque) et le second (lecture depuis le cache). Quel\n",
    "onglet de Spark UI montre les *DataFrames* mis en cache ?"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "4218e337",
   "metadata": {},
   "source": [
    "## DataFrames numériques : lecture, écriture, manipulation"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "ba9ee65d",
   "metadata": {},
   "source": [
    "### Lecture de fichiers CSV"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "deec70b1",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "# Lecture basique : colonnes nommées _c0, _c1, type string par défaut\n",
    "geyser = spark.read.csv(\"tpPySpark/data/geyser.csv\")\n",
    "geyser.printSchema()\n",
    "geyser.show(5)"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "3d3bf05d",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "# Lecture avec en-tête et inférence de types\n",
    "geyser = (spark.read\n",
    "          .option(\"header\", True)\n",
    "          .option(\"inferSchema\", True)\n",
    "          .csv(\"tpPySpark/data/geyser.csv\"))\n",
    "geyser.printSchema()    # colonnes duration et interval de type double\n",
    "geyser.show(5)\n",
    "geyser.describe().show()"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "cd4c7afe",
   "metadata": {},
   "source": [
    "### Note sur `inferSchema`\n",
    "\n",
    "L’inférence de types nécessite une **seconde passe** sur les données pour déterminer\n",
    "le type de chaque colonne. Sur un très grand fichier il est préférable de spécifier le\n",
    "schéma explicitement :"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "7c92d75f",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "from pyspark.sql.types import StructType, StructField, DoubleType\n",
    "\n",
    "schema = StructType([\n",
    "    StructField(\"duration\", DoubleType(), True),\n",
    "    StructField(\"interval\", DoubleType(), True),\n",
    "])\n",
    "geyser = spark.read.schema(schema).csv(\"tpPySpark/data/geyser.csv\",\n",
    "                                       header=True)\n",
    "geyser.printSchema()"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "0ce7a84b",
   "metadata": {},
   "source": [
    "### Lecture des fichiers CSV MovieLens\n",
    "\n",
    "Chargeons dans Spark les données MovieLens utilisées dans la séance précédente :"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "100ca466",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "ratings = (spark.read\n",
    "           .option(\"header\", True)\n",
    "           .option(\"inferSchema\", True)\n",
    "           .csv(\"tpPySpark/data/ratings.csv\"))\n",
    "\n",
    "movies = (spark.read\n",
    "          .option(\"header\", True)\n",
    "          .option(\"inferSchema\", True)\n",
    "          .csv(\"tpPySpark/data/movies.csv\"))\n",
    "\n",
    "ratings.printSchema()\n",
    "ratings.show(5)\n",
    "print(f\"Nombre de notes : {ratings.count()}\")\n",
    "print(f\"Nombre de films : {movies.count()}\")"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "ed8f5093",
   "metadata": {},
   "source": [
    "### Écriture de DataFrames"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "04af149e",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "# Écriture en CSV (produit un répertoire contenant un fichier par partition)\n",
    "geyser.write.mode(\"overwrite\").csv(\"tpPySpark/data/geyser_out\", header=True)\n",
    "\n",
    "# Pour obtenir un seul fichier CSV : forcer une seule partition\n",
    "geyser.repartition(1).write.mode(\"overwrite\").csv(\"tpPySpark/data/geyser_1fichier\",\n",
    "                                                   header=True)\n",
    "\n",
    "# Écriture en Parquet (format binaire compressé, recommandé pour la production)\n",
    "ratings.write.mode(\"overwrite\").parquet(\"tpPySpark/data/ratings_parquet\")\n",
    "\n",
    "# Relecture depuis Parquet\n",
    "ratings_parquet = spark.read.parquet(\"tpPySpark/data/ratings_parquet\")\n",
    "ratings_parquet.printSchema()"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "3ad82e94",
   "metadata": {},
   "source": [
    "### Note sur les formats\n",
    "\n",
    "En production, **Parquet** est préféré à CSV : il est binaire (plus compact),\n",
    "supporte la compression, permet le *predicate pushdown* (Spark ne lit que les\n",
    "colonnes et lignes nécessaires) et préserve les types. Delta Lake et Iceberg sont\n",
    "des évolutions de Parquet qui ajoutent, entre autres, le support des transactions\n",
    "(ACID : *atomicity, consistency, isolation, and durability*)."
   ]
  },
  {
   "cell_type": "markdown",
   "id": "c52c1500",
   "metadata": {},
   "source": [
    "## Sélection, filtrage et colonnes calculées"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "17d7db4f",
   "metadata": {},
   "source": [
    "### Sélection de colonnes"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "5e1cf099",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "from pyspark.sql.functions import col\n",
    "\n",
    "# Plusieurs syntaxes équivalentes pour sélectionner des colonnes\n",
    "ratings.select(\"userId\", \"rating\").show(5)\n",
    "ratings.select(col(\"userId\"), col(\"rating\")).show(5)\n",
    "ratings.select(ratings.userId, ratings.rating).show(5)\n",
    "\n",
    "# Sélection dynamique à partir d'une liste\n",
    "colonnes = [\"userId\", \"movieId\", \"rating\"]\n",
    "ratings.select(*colonnes).show(5)"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "50c20dc2",
   "metadata": {},
   "source": [
    "### Filtrage"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "46c86f5f",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "# Notes >= 4.5\n",
    "bonnes_notes = ratings.filter(col(\"rating\") >= 4.5)\n",
    "bonnes_notes.count()"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "67567b7e",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "# Conditions multiples (& pour ET, | pour OU, ~ pour NON)\n",
    "bons_recents = ratings.filter(\n",
    "    (col(\"rating\") >= 4.0) & (col(\"timestamp\") > 1_000_000_000)\n",
    ")\n",
    "bons_recents.show(5)"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "9da7fcee",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "# Filtrage sur une chaîne de caractères\n",
    "dramas = movies.filter(col(\"genres\").contains(\"Drama\"))\n",
    "dramas.count()"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "d23e2cc7",
   "metadata": {},
   "source": [
    "### Colonnes calculées avec `withColumn`"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "d61a02d6",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "from pyspark.sql.functions import from_unixtime, year, round as spark_round\n",
    "\n",
    "# Convertir le timestamp Unix en date lisible, extraire l'année\n",
    "ratings = (ratings\n",
    "           .withColumn(\"date\", from_unixtime(col(\"timestamp\")))\n",
    "           .withColumn(\"annee\", year(from_unixtime(col(\"timestamp\")))))\n",
    "\n",
    "ratings.select(\"userId\", \"movieId\", \"rating\", \"date\", \"annee\").show(5)\n",
    "\n",
    "# Arrondir la note à l'entier le plus proche\n",
    "ratings = ratings.withColumn(\"rating_arrondi\", spark_round(col(\"rating\")))\n",
    "\n",
    "# Renommer et supprimer une colonne\n",
    "ratings = ratings.withColumnRenamed(\"rating\", \"note\")\n",
    "ratings = ratings.drop(\"timestamp\")\n",
    "ratings.printSchema()"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "acdc7179",
   "metadata": {},
   "source": [
    "### Question\n",
    "\n",
    "1. Créez une nouvelle colonne `note_normalisee` contenant la note ramenée entre\n",
    "  0 et 1 (divisez par 5). Vérifiez avec `describe()`.  \n",
    "1. Combien de films du genre « Animation » ont reçu au moins une note dans\n",
    "  `ratings` ? (filtrez `movies` sur `genres`, ensuite faites une jointure\n",
    "  avec `.join()`)  "
   ]
  },
  {
   "cell_type": "markdown",
   "id": "d8dd6650",
   "metadata": {},
   "source": [
    "## Agrégation et groupBy"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "d2e61bd5",
   "metadata": {},
   "source": [
    "### Agrégations de base"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "a9bc2e03",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "from pyspark.sql.functions import avg, count, min as spark_min, max as spark_max, stddev\n",
    "\n",
    "# Note moyenne globale\n",
    "ratings.select(avg(\"note\")).show()\n",
    "\n",
    "# Plusieurs agrégations en une passe\n",
    "ratings.agg(\n",
    "    avg(\"note\").alias(\"note_moy\"),\n",
    "    count(\"*\").alias(\"nb_notes\"),\n",
    "    spark_min(\"note\").alias(\"note_min\"),\n",
    "    spark_max(\"note\").alias(\"note_max\"),\n",
    "    stddev(\"note\").alias(\"note_std\"),\n",
    ").show()"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "fc86ea6a",
   "metadata": {},
   "source": [
    "### Agrégation par groupe"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "557bf4f0",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "# Statistiques par film\n",
    "stats_par_film = (ratings\n",
    "                  .groupBy(\"movieId\")\n",
    "                  .agg(\n",
    "                      avg(\"note\").alias(\"note_moy\"),\n",
    "                      count(\"*\").alias(\"nb_notes\"),\n",
    "                      stddev(\"note\").alias(\"note_std\"),\n",
    "                  ))\n",
    "stats_par_film.show(10)\n",
    "stats_par_film.orderBy(col(\"nb_notes\").desc()).show(10)"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "8eed7047",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "# Statistiques par année\n",
    "stats_par_annee = (ratings\n",
    "                   .groupBy(\"annee\")\n",
    "                   .agg(\n",
    "                       avg(\"note\").alias(\"note_moy\"),\n",
    "                       count(\"*\").alias(\"nb_notes\"),\n",
    "                   )\n",
    "                   .orderBy(\"annee\"))\n",
    "stats_par_annee.show()"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "08271239",
   "metadata": {},
   "source": [
    "### Question\n",
    "\n",
    "1. Calculez la note moyenne **par utilisateur** et identifiez les 5 utilisateurs\n",
    "  les plus sévères (note moyenne la plus basse, parmi ceux ayant noté au moins\n",
    "  20 films).  \n",
    "1. Calculez, par **genre principal** (premier genre de la colonne `genres`), la\n",
    "  note moyenne et le nombre de films distincts. Triez par note moyenne décroissante.\n",
    "  Indice : utilisez `split(col(\"genres\"), \"\\\\|\").getItem(0)` pour extraire le\n",
    "  premier genre.  "
   ]
  },
  {
   "cell_type": "markdown",
   "id": "625e52de",
   "metadata": {},
   "source": [
    "## Jointures"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "90c1426b",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "# Jointure gauche : enrichir les notes avec les titres et genres\n",
    "notes_enrichies = ratings.join(movies, \"movieId\", \"left\")\n",
    "notes_enrichies.show(5)\n",
    "\n",
    "# Vérification : aucune ligne perdue ?\n",
    "print(f\"Avant jointure : {ratings.count()}\")\n",
    "print(f\"Après jointure : {notes_enrichies.count()}\")"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "d64af5fa",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "# Top 10 films les plus populaires avec leur note moyenne\n",
    "top_films = (notes_enrichies\n",
    "             .groupBy(\"movieId\", \"title\")\n",
    "             .agg(\n",
    "                 avg(\"note\").alias(\"note_moy\"),\n",
    "                 count(\"*\").alias(\"nb_notes\"),\n",
    "             )\n",
    "             .filter(col(\"nb_notes\") >= 50)\n",
    "             .orderBy(col(\"note_moy\").desc()))\n",
    "top_films.show(10, truncate=False)"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "744aad9a",
   "metadata": {},
   "source": [
    "## Note sur les jointures distribuées\n",
    "\n",
    "Dans Spark, deux stratégies principales existent pour les jointures :\n",
    "\n",
    "- *Sort-merge join* (par défaut pour les grands *DataFrames*) : les deux *DataFrames*\n",
    "  sont triés sur la clé de jointure puis fusionnés. Exige un *shuffle* (transfert\n",
    "  de données entre nœuds) coûteux.  \n",
    "- *Broadcast join* : si un des *DataFrames* est petit, Spark peut l’envoyer\n",
    "  (*broadcast*) en entier à chaque nœud *worker*, évitant le *shuffle* coûteux.\n",
    "  Spark le fait automatiquement en-dessous d’un seuil configurable\n",
    "  (`spark.sql.autoBroadcastJoinThreshold`, 10 Mo par défaut).  "
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "e2c7bfad",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "from pyspark.sql.functions import broadcast\n",
    "\n",
    "# Forcer un broadcast join sur le petit DataFrame movies\n",
    "notes_enrichies = ratings.join(broadcast(movies), \"movieId\", \"left\")\n",
    "notes_enrichies.explain()   # vérifier que \"BroadcastHashJoin\" apparaît"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "64a1bc4b",
   "metadata": {},
   "source": [
    "## Spark SQL\n",
    "\n",
    "Spark permet d’interroger des *DataFrames* via des requêtes SQL standard, il suffit pour cela\n",
    "d’enregistrer le *DataFrame* comme une vue temporaire :"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "64cc5d4a",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "# Enregistrement des DataFrames comme vues SQL temporaires\n",
    "ratings.createOrReplaceTempView(\"notes\")\n",
    "movies.createOrReplaceTempView(\"films\")\n",
    "\n",
    "# Requête SQL : top 10 films les mieux notés (≥ 50 notes)\n",
    "top_sql = spark.sql(\"\"\"\n",
    "    SELECT f.title,\n",
    "           AVG(n.note)  AS note_moy,\n",
    "           COUNT(*)     AS nb_notes\n",
    "    FROM   notes  n\n",
    "    JOIN   films  f ON n.movieId = f.movieId\n",
    "    GROUP  BY f.movieId, f.title\n",
    "    HAVING COUNT(*) >= 50\n",
    "    ORDER  BY note_moy DESC\n",
    "    LIMIT  10\n",
    "\"\"\")\n",
    "top_sql.show(truncate=False)"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "87f79539",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "# Mélanger SQL et API DataFrame : le résultat de spark.sql() est un DataFrame\n",
    "films_action = spark.sql(\n",
    "    \"SELECT movieId, title FROM films WHERE genres LIKE '%Action%'\"\n",
    ")\n",
    "films_action.count()\n",
    "\n",
    "# Continuer avec l'API DataFrame\n",
    "films_action.join(ratings, \"movieId\").groupBy(\"title\") \\\n",
    "            .agg(avg(\"note\").alias(\"note_moy\")).orderBy(col(\"note_moy\").desc()) \\\n",
    "            .show(10, truncate=False)"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "f0ec0ca5",
   "metadata": {},
   "source": [
    "## Question\n",
    "\n",
    "Écrivez en Spark SQL la requête suivante : pour chaque utilisateur ayant noté plus\n",
    "de 50 films, calculez la moyenne de ses notes. Retournez les 5 utilisateurs avec les\n",
    "notes moyennes les plus élevées. Comparez avec la version API *DataFrame* équivalente."
   ]
  },
  {
   "cell_type": "markdown",
   "id": "29e510b5",
   "metadata": {},
   "source": [
    "## Échantillonnage dans Spark\n",
    "\n",
    "Comme vu en cours, l’échantillonnage est une première approche de réduction du volume\n",
    "de données. Spark propose deux méthodes directement sur les *DataFrames*."
   ]
  },
  {
   "cell_type": "markdown",
   "id": "49fda171",
   "metadata": {},
   "source": [
    "### Échantillonnage simple"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "e22a1fbf",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "# Échantillon d'environ 10 % des notes, sans remise, reproductible (seed=42)\n",
    "echantillon = ratings.sample(withReplacement=False, fraction=0.1, seed=42)\n",
    "print(f\"Taille originale  : {ratings.count()}\")\n",
    "print(f\"Taille échantillon: {echantillon.count()}\")\n",
    "print(f\"Taux réel         : {echantillon.count() / ratings.count():.3f}\")\n",
    "\n",
    "# Vérification : la distribution des notes est-elle préservée ?\n",
    "ratings.groupBy(\"note\").count().orderBy(\"note\").show()\n",
    "echantillon.groupBy(\"note\").count().orderBy(\"note\").show()"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "562a2f04",
   "metadata": {},
   "source": [
    "### Question\n",
    "\n",
    "Comparez la note moyenne sur l’ensemble des données et sur l’échantillon à 10 %.\n",
    "Répétez avec des fractions de 1 %, 5 %, 10 % et 50 %. À partir de quelle fraction\n",
    "l’estimation de la note moyenne est-elle stable (écart < 0.01) ?"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "a73afcb4",
   "metadata": {},
   "source": [
    "### Échantillonnage stratifié\n",
    "\n",
    "L’échantillonnage stratifié garantit une représentation contrôlée de chaque strate.\n",
    "Ici nous l’appliquons sur les valeurs de note (strate = valeur de note) pour garantir\n",
    "que toutes les valeurs de note sont représentées dans l’échantillon dans des proportions\n",
    "choisies :"
   ]
  },
  {
   "cell_type": "code",
   "execution_count": null,
   "id": "1c52cf51",
   "metadata": {
    "hide-output": false
   },
   "outputs": [],
   "source": [
    "# Fractions par valeur de note : même taux pour toutes les strates\n",
    "valeurs_notes = [row[\"note\"] for row in\n",
    "                 ratings.select(\"note\").distinct().collect()]\n",
    "fractions = {note: 0.1 for note in valeurs_notes}\n",
    "\n",
    "echantillon_strat = ratings.stat.sampleBy(\"note\", fractions=fractions, seed=42)\n",
    "print(f\"Taille : {echantillon_strat.count()}\")\n",
    "\n",
    "# Comparaison des distributions\n",
    "print(\"Distribution originale :\")\n",
    "ratings.groupBy(\"note\").count() \\\n",
    "       .withColumn(\"pct\", col(\"count\") / ratings.count() * 100) \\\n",
    "       .orderBy(\"note\").show()\n",
    "\n",
    "print(\"Distribution échantillon stratifié :\")\n",
    "echantillon_strat.groupBy(\"note\").count() \\\n",
    "                 .withColumn(\"pct\", col(\"count\") / echantillon_strat.count() * 100) \\\n",
    "                 .orderBy(\"note\").show()"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "a518a844",
   "metadata": {},
   "source": [
    "### Question\n",
    "\n",
    "Créez un échantillon stratifié sur la valeur de note en **sur-représentant** les\n",
    "notes extrêmes (0.5 et 5.0 à 50 %) et en **sous-représentant** les notes médianes\n",
    "(2.5 et 3.0 à 5 %). Pour quelle application de fouille de données ce type\n",
    "d’échantillonnage serait-il utile ?"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "c7d0ff69",
   "metadata": {},
   "source": [
    "## Exercice de synthèse optionnel sur MovieLens"
   ]
  },
  {
   "cell_type": "markdown",
   "id": "a43f1c8d",
   "metadata": {},
   "source": [
    "## Question 1 : Exploration\n",
    "\n",
    "En utilisant l’API *DataFrame* (pas SQL) :\n",
    "\n",
    "1. Chargez `ratings.csv`, `movies.csv` et `tags.csv` avec les bons types.  \n",
    "1. Combien y a-t-il d’utilisateurs distincts, de films distincts, et de tags distincts ?  \n",
    "1. Quelle est la plage temporelle des notes (dates min et max) ?  "
   ]
  },
  {
   "cell_type": "markdown",
   "id": "997ed5d4",
   "metadata": {},
   "source": [
    "## Question 2 : Passage à l’échelle\n",
    "\n",
    "Comparez les performances de Spark et Pandas pour l’agrégation `groupBy(\"movieId\")`\n",
    "avec calcul de la note moyenne, sur le jeu de données MovieLens complet. Mesurez\n",
    "avec `%%time`. Quel outil est le plus rapide ici, et pourquoi le résultat\n",
    "peut-il surprendre par rapport à la comparaison avec Polars la séance précédente ?"
   ]
  }
 ],
 "metadata": {
  "date": 1789050001.569119,
  "filename": "tpPySpark.rst",
  "kernelspec": {
   "display_name": "Python 3",
   "language": "python",
   "name": "python"
  },
  "title": "Travaux pratiques - Introduction à PySpark et aux DataFrames Spark"
 },
 "nbformat": 4,
 "nbformat_minor": 5
}