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 ?