Travaux pratiques - Fouille de flux de données¶
Références externes utiles :
L’objectif de cette séance de TP est d’introduire l’emploi de Spark Structured Streaming pour le traitement de données en flux. Nous utilisons le flux public de modifications de Wikipédia (Wikimedia EventStreams) comme source de données,* et appliquons des traitements d’agrégation sur les tranches successives du flux.
Contrairement à l’ancienne API DStream qui représentait un flux par une succession de RDD, Structured Streaming repose sur l’API DataFrame et nous permet d’employer la même syntaxe que pour les traitements batch (par lots).
Architecture du TP¶
La figure ci-dessous illustre l’architecture employée dans ce TP :
┌─────────────────────────────┐ ┌──────────────────────────────────┐
│ Thread producteur Python │ │ Spark Structured Streaming │
│ │ socket │ │
│ lit Wikimedia EventStreams ├───────►│ readStream (format "socket") │
│ (SSE sur HTTPS) │ TCP │ → parse JSON │
│ → envoie JSON ligne/ligne │ │ → agrégations, fenêtres │
└─────────────────────────────┘ │ → writeStream (console) │
└──────────────────────────────────┘
Un thread Python lit en continu le flux SSE (Server-Sent Events) 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 l’environnement du notebook :
%pip install sseclient-py
Initialisation de la session Spark¶
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 en arrière-plan un serveur socket TCP sur le premier port disponible affecté par le système (→ indiquer port 0 dans srv.bind). 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).
ATTENTION ! AVANT d’exécuter la cellule suivante remplacez dans HEADERS prenom.nom@lecnam.net par votre adresse de courriel @lecnam.net !
import socket
import threading
import json
import sseclient
import urllib.request
WIKIMEDIA_URL = "https://stream.wikimedia.org/v2/stream/recentchange"
HEADERS = {
"User-Agent": "RCP216-TP-FouilleFlux/1.1 (Cnam; contact: prenom.nom@lecnam.net)",
"Accept": "text/event-stream",
}
def produire_flux(conn):
"""Lit le flux SSE Wikimedia et envoie chaque événement JSON sur le socket."""
try:
req = urllib.request.Request(WIKIMEDIA_URL, headers=HEADERS)
with urllib.request.urlopen(req) as response:
client = sseclient.SSEClient(response)
for event in client.events():
if event.data and event.data.strip():
try:
# Vérifie que le format est 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 boucle_accept(srv):
"""Lance un producteur pour chaque connexion de Spark."""
while True:
conn, addr = srv.accept()
print(f"Spark connecté depuis {addr}")
threading.Thread(target=produire_flux, args=(conn,), daemon=True).start()
if not globals().get("serveur_demarre"):
srv = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
srv.bind(('127.0.0.1', 0))
srv.listen(5)
PORT = srv.getsockname()[1]
threading.Thread(target=boucle_accept, args=(srv,), daemon=True).start()
serveur_demarre = True
print(f"Serveur en écoute sur le port {PORT}")
Note
Le thread est démarré avec daemon=True, il s’arrête donc automatiquement quand le processus principal (le noyau Jupyter) s’arrête, sans laisser de ressources encore occupées.
Lecture du flux dans Spark Structured Streaming¶
Spark Structured Streaming peut lire depuis un socket TCP avec le format socket. Chaque ligne reçue est une ligne d’un DataFrame de flux.
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", "127.0.0.1") \
.option("port", PORT) \
.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. Nous excluons ici les modifications faites par des bots et surtout celles apportées à wikidata car elles sont en très grand nombre.
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 (il faut patienter) des résultats commencent à s’afficher. On peut voir passer des modifications de pages Wikipedia dans de nombreuses langues différentes.
Question
Observez les résultats pendant 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 ?
Note
La source socket de Spark Structured Streaming ne garantit pas la tolérance aux pannes en cas de défaillance du driver Spark. Elle est appropriée pour des démonstrations et des travaux pratiques, 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 :
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.
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(col("nb_modifications").desc())
requete2 = modifications_par_wiki.writeStream \
.outputMode("complete") \
.format("console") \
.option("truncate", False) \
.option("numRows", 15) \
.trigger(processingTime="10 seconds") \
.start()
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.)
Arrêt de la requête :
requete2.stop()
Requête 3 : détection des grandes modifications¶
Nous pouvons calculer 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 indique un ajout ou une suppression importante de contenu.
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()
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.
Arrêt de la requête :
requete3.stop()
Requête 4 : fraction des modifications faites par des bots¶
Utilisons des fenêtres successives disjointes (tumbling windows) de 30 secondes pour estimer, dans chaque fenêtre, la proportion de modifications effectuées par des bots.
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()
Question
Observez la proportion de modifications effectuées par des bots. Cette proportion vous semble stable ou variable dans le temps ?
Arrêt de la requête et du producteur¶
requete4.stop()
# Le thread producteur est daemon=True et s'arrête donc 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.")
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.
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 ?