.. _chap-tpFouilleFluxPy: ############################################## Travaux pratiques - Fouille de flux de données ############################################## .. only:: html .. container:: notebook .. image:: _static/jupyter_logo.png :class: svg-inline `Cahier Jupyter en Python `_ Références externes utiles : * `Documentation Spark Structured Streaming `_ * `Documentation API Spark en Python `_ * `Tutoriel Spark avec exemples `_ * `Wikimedia EventStreams `_ **L'objectif** de cette séance de TP est d'introduire l'utilisation de *Spark Structured Streaming* pour le traitement de données en flux réel. Nous utiliserons le flux public de modifications de Wikipédia (Wikimedia EventStreams) comme source de données, et appliquerons des traitements d'agrégation sur les tranches successives du flux. Contrairement à l'ancienne API *DStream*, *Structured Streaming* repose sur l'API *DataFrame*, ce qui permet d'employer la même syntaxe que pour les traitements par lots. Architecture du TP ================== La figure ci-dessous illustre l'architecture employée dans ce TP. .. code-block:: none ┌─────────────────────────────┐ ┌──────────────────────────────────┐ │ Thread producteur Python │ │ Spark Structured Streaming │ │ │ socket │ │ │ lit Wikimedia EventStreams ├───────►│ readStream (format "socket") │ │ (SSE sur HTTPS) │ TCP │ → parse JSON │ │ → pousse JSON ligne/ligne │ │ → agrégations, fenêtres │ └─────────────────────────────┘ │ → writeStream (console) │ └──────────────────────────────────┘ Un thread Python lit en continu le flux SSE de Wikimedia et envoie chaque événement JSON sur un *socket* TCP local. Spark Structured Streaming se connecte à ce *socket* et traite les données reçues tranche par tranche. Installation de la bibliothèque cliente SSE ============================================ La bibliothèque ``sseclient-py`` permet de lire facilement un flux SSE (*Server-Sent Events*). Sur JupyterHub, installez-la dans le notebook : .. code-block:: python !pip install sseclient-py --quiet Initialisation de la session Spark ==================================== .. code-block:: python import sys, os os.environ['PYSPARK_PYTHON'] = sys.executable os.environ['PYSPARK_DRIVER_PYTHON'] = sys.executable from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("WikimediaStreaming") \ .config("spark.sql.shuffle.partitions", "4") \ .getOrCreate() spark.sparkContext.setLogLevel("ERROR") La configuration ``spark.sql.shuffle.partitions`` est réduite à 4 (au lieu de 200 par défaut) car nous travaillons sur une machine locale avec un flux de débit modéré. Démarrage du producteur de flux ================================ Le code suivant démarre un serveur *socket* TCP sur le port 9999 en arrière-plan. Quand Spark se connecte, le thread producteur commence à lire le flux Wikimedia EventStreams et à envoyer chaque événement JSON sur le socket, une ligne par événement. .. code-block:: python import socket import threading import json import sseclient import urllib.request WIKIMEDIA_URL = "https://stream.wikimedia.org/v2/stream/recentchange" PORT = 9999 def produire_flux(conn): """Lit le flux SSE Wikimedia et envoie chaque événement JSON sur le socket.""" try: with urllib.request.urlopen(WIKIMEDIA_URL) as response: client = sseclient.SSEClient(response) for event in client.events(): if event.data and event.data.strip(): try: # Vérification que c'est bien du JSON valide json.loads(event.data) ligne = event.data.replace('\n', ' ') + '\n' conn.sendall(ligne.encode('utf-8')) except (json.JSONDecodeError, BrokenPipeError): break except Exception as e: print(f"Producteur arrêté : {e}") finally: conn.close() def demarrer_serveur(port): """Ouvre un socket serveur et attend la connexion de Spark.""" srv = socket.socket(socket.AF_INET, socket.SOCK_STREAM) srv.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) srv.bind(('localhost', port)) srv.listen(1) print(f"Serveur en attente sur le port {port}...") conn, addr = srv.accept() print(f"Spark connecté depuis {addr}") t = threading.Thread(target=produire_flux, args=(conn,), daemon=True) t.start() thread_serveur = threading.Thread(target=demarrer_serveur, args=(PORT,), daemon=True) thread_serveur.start() .. admonition:: Note Le thread est démarré avec ``daemon=True`` : il s'arrêtera automatiquement quand le processus principal (le noyau Jupyter) s'arrête, sans laisser de ressources pendantes. Lecture du flux dans Spark Structured Streaming ================================================= Spark Structured Streaming peut lire depuis un *socket* TCP avec le format ``socket``. Chaque ligne reçue sera une ligne d'un *DataFrame* de flux. .. code-block:: python from pyspark.sql.types import StructType, StringType, LongType, BooleanType, IntegerType from pyspark.sql.functions import from_json, col, window, to_timestamp # Schéma des événements Wikimedia RecentChange (champs principaux) schema = StructType() \ .add("wiki", StringType()) \ .add("type", StringType()) \ .add("title", StringType()) \ .add("user", StringType()) \ .add("bot", BooleanType()) \ .add("timestamp", LongType()) \ .add("length", StructType() .add("old", IntegerType()) .add("new", IntegerType())) # Lecture depuis le socket TCP lignes_brutes = spark.readStream \ .format("socket") \ .option("host", "localhost") \ .option("port", 9999) \ .load() # Parsing JSON evenements = lignes_brutes \ .select(from_json(col("value"), schema).alias("e")) \ .select("e.*") \ .filter(col("wiki").isNotNull()) \ .withColumn("heure", to_timestamp(col("timestamp"))) Les champs principaux de chaque événement sont les suivants : - ``wiki`` : identifiant du wiki concerné (``frwiki``, ``enwiki``, ``wikidata``, etc.) - ``type`` : type de modification (``edit``, ``new``, ``log``, etc.) - ``title`` : titre de la page modifiée - ``user`` : nom d'utilisateur ou adresse IP - ``bot`` : booléen indiquant si la modification a été faite par un bot - ``timestamp`` : horodatage Unix de la modification - ``length.old`` / ``length.new`` : taille en octets de la page avant et après la modification Requête 1 : affichage des titres des pages modifiées ====================================================== Commençons par une requête simple qui affiche, pour chaque tranche de 10 secondes, les titres des pages récemment modifiées (hors bots, hors wikidata qui génère un très grand nombre de modifications automatiques). .. code-block:: python titres = evenements \ .filter((col("bot") == False) & (col("wiki") != "wikidatawiki")) \ .filter(col("type") == "edit") \ .select("wiki", "title", "user") requete1 = titres.writeStream \ .outputMode("append") \ .format("console") \ .option("truncate", True) \ .option("numRows", 10) \ .trigger(processingTime="10 seconds") \ .start() Après quelques secondes, des résultats commencent à s'afficher dans la console. Vous verrez passer des modifications de pages Wikipedia dans de nombreuses langues différentes. .. admonition:: Question Observez les résultats pendant 30 à 60 secondes. Quels wikis sont les plus actifs ? Y a-t-il des pages modifiées plusieurs fois de suite sur une courte période ? .. admonition:: Note La source ``socket`` de Spark Structured Streaming ne garantit pas la tolérance aux pannes (*fault tolerance*) en cas de défaillance du driver. Elle est appropriée pour des démonstrations et des travaux pratiques, mais pas pour un environnement de production. Arrêt de la requête précédente ================================ Avant de lancer une nouvelle requête, arrêtez la précédente : .. code-block:: python requete1.stop() Requête 2 : comptage des modifications par wiki (fenêtre glissante) ===================================================================== Comptons maintenant le nombre de modifications par wiki sur une fenêtre glissante de 60 secondes (avançant toutes les 10 secondes). Cette requête illustre l'utilisation des fenêtres temporelles dans *Structured Streaming*. .. code-block:: python from pyspark.sql.functions import count modifications_par_wiki = evenements \ .filter(col("type") == "edit") \ .groupBy( window(col("heure"), "60 seconds", "10 seconds"), col("wiki") ) \ .agg(count("*").alias("nb_modifications")) \ .orderBy("window", col("nb_modifications").desc()) requete2 = modifications_par_wiki.writeStream \ .outputMode("complete") \ .format("console") \ .option("truncate", False) \ .option("numRows", 15) \ .trigger(processingTime="10 seconds") \ .start() .. admonition:: Question Le mode de sortie ``complete`` réaffiche la table entière à chaque tranche. Pourquoi est-il nécessaire ici, alors que ``append`` suffisait dans la requête précédente ? (Indice : réfléchissez à ce que signifie « ajouter une ligne » pour une agrégation sur une fenêtre glissante.) .. ifconfig:: tpscalapython in ('public') .. admonition:: Correction Le mode ``append`` n'émet une ligne de résultat que lorsqu'elle ne peut plus être modifiée (c'est-à-dire quand la fenêtre correspondante est définitivement fermée). Avec une agrégation sur une fenêtre glissante, chaque nouvelle tranche peut modifier le compte d'une fenêtre encore ouverte. Le mode ``complete`` réécrit l'intégralité de la table de résultats à chaque tranche, ce qui permet de voir les agrégations en cours de construction. Le mode ``update``, qui n'émet que les lignes modifiées, serait également possible ici. Arrêt de la requête : .. code-block:: python requete2.stop() Requête 3 : détection des grandes modifications ================================================ Calculons maintenant la taille des modifications (différence entre la taille de la page après et avant la modification). Une modification très grande en valeur absolue peut indiquer un ajout ou une suppression importante de contenu. .. code-block:: python from pyspark.sql.functions import abs as spark_abs, when grandes_modifs = evenements \ .filter((col("type") == "edit") & col("length.old").isNotNull() & col("length.new").isNotNull()) \ .withColumn("delta", col("length.new") - col("length.old")) \ .withColumn("sens", when(col("delta") > 0, "ajout").otherwise("suppression")) \ .filter(spark_abs(col("delta")) > 5000) \ .select("wiki", "title", "user", "bot", "delta", "sens") requete3 = grandes_modifs.writeStream \ .outputMode("append") \ .format("console") \ .option("truncate", True) \ .option("numRows", 10) \ .trigger(processingTime="10 seconds") \ .start() .. admonition:: Question Observez les résultats. Les grandes modifications sont-elles plutôt des ajouts ou des suppressions ? Sont-elles plutôt faites par des bots ou par des humains ? Modifiez le seuil (par ex. 500 octets au lieu de 5 000) et observez l'impact sur le volume de résultats. .. ifconfig:: tpscalapython in ('public') .. admonition:: Correction Les grandes suppressions sont souvent des actes de vandalisme (réversion de pages) ou des restructurations de contenu, parfois effectuées par des bots de maintenance. Les grands ajouts correspondent à des créations de sections ou à des imports de contenu. En abaissant le seuil à 500 octets, le nombre de modifications visibles augmente nettement. Arrêt de la requête : .. code-block:: python requete3.stop() Requête 4 : fraction des modifications faites par des bots =========================================================== Utilisons une fenêtre tumblante (*tumbling window*) de 30 secondes pour estimer, sur chaque tranche, la proportion de modifications effectuées par des bots versus des humains. .. code-block:: python from pyspark.sql.functions import sum as spark_sum, round as spark_round bots_vs_humains = evenements \ .filter(col("type") == "edit") \ .groupBy(window(col("heure"), "30 seconds")) \ .agg( count("*").alias("total"), spark_sum(col("bot").cast("int")).alias("nb_bots") ) \ .withColumn("pct_bots", spark_round(col("nb_bots") * 100.0 / col("total"), 1)) \ .select("window", "total", "nb_bots", "pct_bots") requete4 = bots_vs_humains.writeStream \ .outputMode("update") \ .format("console") \ .option("truncate", False) \ .trigger(processingTime="10 seconds") \ .start() .. admonition:: Question Observez la proportion de modifications effectuées par des bots. Que remarquez-vous ? Cette proportion vous semble-t-elle stable ou variable dans le temps ? .. ifconfig:: tpscalapython in ('public') .. admonition:: Correction La proportion de modifications effectuées par des bots sur l'ensemble des wikis Wikimedia est généralement élevée (souvent 30 à 50 % des modifications visibles dans ce flux, avec de fortes variations selon l'heure). Elle correspond notamment aux bots de maintenance de Wikidata et aux bots anti-vandalisme. Arrêt de la requête et du producteur ====================================== .. code-block:: python requete4.stop() # Le thread producteur est daemon=True et s'arrêtera avec le noyau Jupyter. # Pour arrêter explicitement toutes les requêtes en cours : for q in spark.streams.active: q.stop() print("Toutes les requêtes arrêtées.") .. admonition:: Note Il n'est pas nécessaire d'arrêter la session Spark elle-même entre les requêtes. Si vous souhaitez relancer le producteur depuis le début (par exemple après une interruption réseau), relancez les cellules de démarrage du serveur socket et des requêtes. .. admonition:: Question (bilan) En résumé, quelle est la différence entre le modèle de programmation de *Spark Structured Streaming* et celui des traitements par lots (*batch*) sur des *DataFrames* que vous avez utilisés dans les séances précédentes ? Qu'est-ce qui change, et qu'est-ce qui reste identique ? .. ifconfig:: tpscalapython in ('public') .. admonition:: Correction La syntaxe des transformations (``filter``, ``groupBy``, ``agg``, ``withColumn``, etc.) est identique à celle des traitements par lots. Ce qui change est la *source* des données (``readStream`` au lieu de ``read``) et la *destination* des résultats (``writeStream`` avec un mode de sortie et un déclencheur, au lieu d'une action comme ``show`` ou ``write``). Le moteur Spark gère de façon transparente la découpe du flux en tranches et la mise à jour incrémentale des résultats.