Travaux pratiques - Fouille de textes 1 : SparkNLP et GloVe

Références externes utiles :

L’objectif de cette séance est double. La première partie est une introduction à SparkNLP : construction d’un pipeline NLP (tokenisation, étiquetage morpho-syntaxique, lemmatisation, entités nommées) et observation des résultats sur un texte réel. La seconde partie met en œuvre la classification de textes avec des plongements lexicaux GloVe sur le jeu de données AG News de classification thématique de dépêches de presse en 4 catégories.

Note

SparkNLP télécharge des ressources (modèles pré-entraînés) depuis Internet lors de la première utilisation. Ces téléchargements peuvent prendre quelques minutes, surtout après le réveil des utilisateurs d’Amérique du nord. Il est possible de télécharger les ressources en amont pour une exécution hors connexion (documentation offline).

Mise en place

%%bash
mkdir -p tpFouilleTexte/data
# Téléchargement du jeu AG News (format CSV)
wget -nc https://raw.githubusercontent.com/mhjabreel/CharCnn_Keras/master/data/ag_news_csv/train.csv -O tpFouilleTexte/data/agnews_train.csv
wget -nc https://raw.githubusercontent.com/mhjabreel/CharCnn_Keras/master/data/ag_news_csv/test.csv  -O tpFouilleTexte/data/agnews_test.csv
pyspark --packages com.johnsnowlabs.nlp:spark-nlp_2.12:6.4.2

Partie 1 : Pipeline NLP de base

Nous travaillons dans cette partie sur un court texte en anglais écrit directement dans le code, afin d’observer les résultats de chaque étape du pipeline NLP.

from pyspark.sql import SparkSession
from pyspark.ml import Pipeline
from sparknlp.base import DocumentAssembler, Finisher
from sparknlp.annotator import (Tokenizer, LemmatizerModel,
                                 PerceptronModel, NerDLModel,
                                 NerConverter, SentenceDetector)

spark = SparkSession.builder.getOrCreate()

# Texte de démonstration
texte_demo = spark.createDataFrame([
    ["Apple is looking at buying U.K. startup for $1 billion. "
     "The company, based in Cupertino, was founded by Steve Jobs in 1976."]
], ["text"])

Construction du pipeline

SparkNLP s’organise autour d”annotateurs (Annotators) chaînés dans un pipeline. Chaque annotateur consomme une ou plusieurs colonnes et produit une nouvelle colonne de type annotation.

# 1. DocumentAssembler : convertit la colonne texte au format SparkNLP
documentAssembler = DocumentAssembler() \
    .setInputCol("text") \
    .setOutputCol("document") \
    .setCleanupMode("shrink")

# 2. SentenceDetector : découpe le texte en phrases
sentenceDetector = SentenceDetector() \
    .setInputCols(["document"]) \
    .setOutputCol("sentence")

# 3. Tokenizer : découpe chaque phrase en tokens
tokenizer = Tokenizer() \
    .setInputCols(["sentence"]) \
    .setOutputCol("token")

# 4. POS tagger pré-entraîné (Penn Treebank tagset)
posTagger = PerceptronModel.pretrained("pos_anc", "en") \
    .setInputCols(["sentence", "token"]) \
    .setOutputCol("pos")

# 5. Lemmatizer pré-entraîné
lemmatizer = LemmatizerModel.pretrained("lemma_antbnc", "en") \
    .setInputCols(["token"]) \
    .setOutputCol("lemma")

# 6. NER (reconnaissance d'entités nommées) pré-entraîné
ner = NerDLModel.pretrained("ner_dl", "en") \
    .setInputCols(["sentence", "token", "glove_embeddings"]) \
    .setOutputCol("ner")

Note

Le modèle NER (Named Entity Recognition)``ner_dl`` nécessite des embeddings GloVe en entrée. Nous construisons d’abord un pipeline simplifié sans NER, puis un pipeline complet avec GloVe dans la deuxième partie.

# Pipeline simplifié (sans NER)
from sparknlp.annotator import LemmatizerModel

pipeline_nlp = Pipeline(stages=[
    documentAssembler,
    sentenceDetector,
    tokenizer,
    posTagger,
    lemmatizer,
])

# Finisher : rend les résultats lisibles
finisher = Finisher() \
    .setInputCols(["token", "pos", "lemma"]) \
    .setOutputCols(["tokens", "pos_tags", "lemmas"])

pipeline_complet = Pipeline(stages=[pipeline_nlp, finisher])

Application et observation

modele_nlp = pipeline_complet.fit(texte_demo)
resultat = modele_nlp.transform(texte_demo)
resultat.select("tokens", "pos_tags", "lemmas").show(truncate=False)

Question :

Examinez les colonnes tokens, pos_tags et lemmas. Pour chaque token, retrouvez l’étiquette POS et le lemme correspondants. Les lemmes correspondent-ils toujours à la forme canonique attendue ? Identifiez un cas où la lemmatisation et la racinisation donneraient des résultats différents.

# Afficher les phrases détectées
resultat_brut = pipeline_nlp.fit(texte_demo).transform(texte_demo)
resultat_brut.select("sentence").show(truncate=False)

Question :

Combien de phrases le SentenceDetector a-t-il identifiées dans le texte de démonstration ? Y a-t-il des cas où la segmentation en phrases est difficile (cas d’ambiguïté sur le rôle du point) ?

Pipeline avec entités nommées

Nous ajoutons maintenant les embeddings GloVe nécessaires au modèle NER.

from sparknlp.annotator import WordEmbeddingsModel, NerConverter

# GloVe 100d (plus léger pour les TP)
glove = WordEmbeddingsModel.pretrained("glove_100d", "en") \
    .setInputCols(["sentence", "token"]) \
    .setOutputCol("glove_embeddings")

ner = NerDLModel.pretrained("ner_dl", "en") \
    .setInputCols(["sentence", "token", "glove_embeddings"]) \
    .setOutputCol("ner")

# NerConverter regroupe les tokens BIO en entités
nerConverter = NerConverter() \
    .setInputCols(["sentence", "token", "ner"]) \
    .setOutputCol("entities")

finisher_ner = Finisher() \
    .setInputCols(["ner", "entities"]) \
    .setOutputCols(["ner_tags", "named_entities"])

pipeline_ner = Pipeline(stages=[
    documentAssembler, sentenceDetector, tokenizer,
    glove, ner, nerConverter, finisher_ner
])

modele_ner = pipeline_ner.fit(texte_demo)
resultat_ner = modele_ner.transform(texte_demo)
resultat_ner.select("ner_tags", "named_entities").show(truncate=False)

Question :

Quelles entités nommées sont détectées dans le texte de démonstration ? Quels sont leurs types (PER, ORG, LOC, MISC) ? Y a-t-il des erreurs ou des omissions ?

Partie 2 : Classification thématique avec GloVe

Nous utilisons maintenant le jeu de données AG News pour une tâche de classification thématique d’articles en 4 catégories : World, Sports, Business, Sci/Tech.

Chargement et exploration des données

from pyspark.sql.functions import col, when, concat_ws
from pyspark.sql.types import DoubleType

# Le CSV AG News n'a pas d'en-tête : colonnes = class (1-4), title, description
schema = "class INT, title STRING, description STRING"
train_raw = spark.read.csv("tpFouilleTexte/data/agnews_train.csv", schema=schema)
test_raw  = spark.read.csv("tpFouilleTexte/data/agnews_test.csv",  schema=schema)

# Concaténer titre et description, convertir la classe en 0-3
def preparer(df):
    return df.withColumn("text",
            concat_ws(" ", col("title"), col("description"))) \
         .withColumn("label", (col("class") - 1).cast(DoubleType())) \
         .select("label", "text")

train_data = preparer(train_raw).cache()
test_data  = preparer(test_raw).cache()

print(f"Entraînement : {train_data.count()} | Test : {test_data.count()}")
train_data.groupBy("label").count().orderBy("label").show()
test_data.groupBy("label").count().orderBy("label").show()

Question :

Le jeu de données est-il équilibré entre les 4 catégories ? Qu’est-ce que cela implique pour le choix de la métrique d’évaluation ? Affichez quelques exemples d’articles pour vous familiariser avec le contenu de chaque catégorie.

Pour garder des temps de calcul raisonnables, nous travaillons sur un échantillon équilibré (obtenu ici par échantillonnage stratifié, bien que la stratification ne soit pas indispensable) :

# Échantillon stratifié : 5 000 articles par classe pour l'entraînement
fractions = {float(i): 5000/30000 for i in range(4)}
train_ech = train_data.stat.sampleBy("label", fractions, seed=42).cache()
print(f"Taille de l'échantillon : {train_ech.count()} articles")

Pipeline GloVe -> embeddings de phrases

Construisons le pipeline SparkNLP qui produit un embedding de phrase avec GloVe (comme la moyenne des vecteurs GloVe des tokens) pour chaque article.

from sparknlp.base import DocumentAssembler, EmbeddingsFinisher
from sparknlp.annotator import (Tokenizer, StopWordsCleaner,
                                WordEmbeddingsModel, SentenceEmbeddings)

# DocumentAssembler
documentAssembler = DocumentAssembler() \
    .setInputCol("text") \
    .setOutputCol("document") \
    .setCleanupMode("shrink")

# Tokenizer
tokenizer = Tokenizer() \
    .setInputCols(["document"]) \
    .setOutputCol("token")

# Suppression des stop words
stopWords = StopWordsCleaner.pretrained("stopwords_en", "en") \
    .setInputCols(["token"]) \
    .setOutputCol("token_clean")

# Embeddings GloVe 300d (attention : téléchargement ~450 Mo)
glove = WordEmbeddingsModel.pretrained("glove_6B_300", "xx") \
    .setInputCols(["document", "token_clean"]) \
    .setOutputCol("glove_embeddings")

# Embedding de phrase : moyenne des embeddings des tokens
sentenceEmbeddings = SentenceEmbeddings() \
    .setInputCols(["document", "glove_embeddings"]) \
    .setOutputCol("sentence_embeddings") \
    .setPoolingStrategy("AVERAGE")

# Finisher : convertit en vecteur Spark ML
embeddingsFinisher = EmbeddingsFinisher() \
    .setInputCols(["sentence_embeddings"]) \
    .setOutputCols(["features"]) \
    .setOutputAsVector(True) \
    .setCleanAnnotations(False)

pipeline_glove = Pipeline(stages=[
    documentAssembler, tokenizer, stopWords,
    glove, sentenceEmbeddings, embeddingsFinisher
])

Question :

Quel est l’effet de l’étape StopWordsCleaner sur les embeddings de phrases ? Dans quels cas la suppression des stop words peut-elle être nocive pour la classification ?

# Application du pipeline sur l'échantillon d'entraînement et les données de test
# (opération coûteuse, donc mise en cache recommandée)
import time
t0 = time.time()
modele_glove = pipeline_glove.fit(train_ech)
train_emb = modele_glove.transform(train_ech) \
                        .select("label", col("features")[0].alias("features")).cache()
test_emb  = modele_glove.transform(test_data) \
                        .select("label", col("features")[0].alias("features"), "text").cache()
# Forcer le calcul pour mesurer le temps
n_train = train_emb.count()
n_test  = test_emb.count()
print(f"Embeddings calculés en {time.time()-t0:.1f} s "
      f"({n_train} train, {n_test} test)")

Question :

Pourquoi utilise-t-on .cache() après le calcul des embeddings ? Que se passerait-il sinon lorsque le classifieur est appliqué plusieurs fois ?

Classification avec SVM linéaire

from pyspark.ml.classification import LinearSVC
from pyspark.ml.classification import OneVsRest
from pyspark.ml.tuning import CrossValidator, ParamGridBuilder
from pyspark.ml.evaluation import MulticlassClassificationEvaluator
import numpy as np

# SVM linéaire OvR (4 classes)
svm = LinearSVC(featuresCol="features", labelCol="label", maxIter=20)
ovr = OneVsRest(classifier=svm, labelCol="label",
                featuresCol="features", predictionCol="prediction")

pipeline_svm = Pipeline(stages=[ovr])

grid = ParamGridBuilder() \
           .addGrid(svm.regParam, [0.01, 0.1, 0.5]) \
           .build()

evaluateur = MulticlassClassificationEvaluator(
                 labelCol="label", predictionCol="prediction",
                 metricName="accuracy")

cv = CrossValidator(estimator=pipeline_svm,
                    estimatorParamMaps=grid,
                    evaluator=evaluateur,
                    numFolds=5, seed=42)

t0 = time.time()
modele_svm = cv.fit(train_emb)
print(f"Entraînement SVM : {time.time()-t0:.1f} s")

meilleur_regParam = modele_svm.getEstimatorParamMaps() \
                        [np.argmax(modele_svm.avgMetrics)][svm.regParam]
print(f"Meilleur regParam : {meilleur_regParam}")
print(f"Accuracy (validation croisée) : {max(modele_svm.avgMetrics):.4f}")

acc_test = evaluateur.evaluate(modele_svm.transform(test_emb))
print(f"Accuracy sur test : {acc_test:.4f}")

Analyse des erreurs

from pyspark.sql.functions import col as F_col

predictions = modele_svm.transform(test_emb)

# Matrice de confusion
from pyspark.mllib.evaluation import MulticlassMetrics
preds_labels = predictions.select("prediction", "label") \
                           .rdd.map(lambda r: (r.prediction, float(r.label)))
metrics = MulticlassMetrics(preds_labels)
print("Matrice de confusion (lignes = vrai, colonnes = prédit) :")
print(metrics.confusionMatrix().toArray().astype(int))
categories = ["World", "Sports", "Business", "Sci/Tech"]
print("Catégories :", categories)

Question :

Comparez l’accuracy obtenue sur les données de test à une baseline triviale (classifieur aléatoire uniforme, classifieur de la classe majoritaire). Le modèle apporte-t-il un gain significatif ? Examinez les erreurs : pour quelle(s) paire(s) de catégories les confusions sont-elles les plus fréquentes ?

erreurs = predictions.filter(
    (F_col("label") != F_col("prediction"))
)

erreurs.select("label", "prediction", "text").show(truncate=130)

Question :

Examinez quelques articles mal classés. Voyez-vous une explication intuitive pour ces erreurs ? La représentation par centre de gravité des vecteurs GloVe peut-elle être responsable de certaines de ces erreurs ?

Question (bilan) :

Résumez le pipeline complet mis en œuvre dans cette séance, de la donnée originale (texte) jusqu’à la prédiction. Identifiez à chaque étape les choix qui influencent la qualité finale du modèle. Quels seraient vos premiers axes d’amélioration si vous aviez plus de temps ?