Travaux pratiques - GraphFrames : chemins et centralités¶
Références externes utiles :
Objectif du TP¶
Ce TP introduit GraphFrames, la bibliothèque de traitement de graphes pour Apache Spark. Nous travaillons sur deux graphes :
un graphe de transport (villes des Pays-Bas et du Royaume-Uni, avec vraies distances kilométriques comme poids) ;
un graphe social (réseau de suivi entre utilisateurs, fictif).
La séance est organisée en trois parties :
Chemins : parcours en largeur (BFS), plus court chemin valué (Dijkstra), propriétés globales (diamètre, longueur moyenne).
Centralités locales : degré, intermédiarité, proximité.
PageRank et PageRank personnalisé.
Mise en place de l’environnement¶
import os, sys
os.environ['PYSPARK_PYTHON'] = sys.executable
os.environ['PYSPARK_DRIVER_PYTHON'] = sys.executable
os.environ['PYSPARK_SUBMIT_ARGS'] = '--packages io.graphframes:graphframes-spark3_2.12:0.12.2 pyspark-shell'
from graphframes import *
from pyspark.sql.types import *
import pyspark
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.config("spark.driver.memory", "4g") \
.getOrCreate()
# En environnement cluster :
#spark.sparkContext.setCheckpointDir("checkpoints/") # ou tout autre chemin accessible en écriture
Téléchargement des données :
dossier_tp = "tpCheminsCentralites/"
archive_zip = "tp-chemins-centralites.zip"
fichier_csv = dossier_tp + "transport-nodes.csv"
if not os.path.exists(dossier_tp + archive_zip):
os.system("mkdir -p " + dossier_tp)
sortie = "téléchargement OK" if os.system(
"wget -nv -nc http://cedric.cnam.fr/vertigo/Cours/RCP216/docs/"
+ archive_zip + " -P" + dossier_tp) == 0 else "problème de téléchargement"
print(sortie)
else:
print("archive présente dans " + dossier_tp)
if not os.path.exists(fichier_csv):
os.system("unzip " + dossier_tp + archive_zip + " -d " + dossier_tp)
Partie 1 : Chemins¶
Création du graphe de transport¶
node_fields = [
StructField("id", StringType(), True),
StructField("latitude", FloatType(), True),
StructField("longitude", FloatType(), True),
StructField("population", IntegerType(), True)
]
nodes = spark.read.csv(
"tpCheminsCentralites/tp-chemins-centralites/transport-nodes.csv",
header=True, schema=StructType(node_fields))
rels = spark.read.csv(
"tpCheminsCentralites/tp-chemins-centralites/transport-relationships.csv", header=True)
reversed_rels = (rels
.withColumn("newSrc", rels.dst)
.withColumn("newDst", rels.src)
.drop("dst", "src")
.withColumnRenamed("newSrc", "src")
.withColumnRenamed("newDst", "dst")
.select("src", "dst", "relationship", "cost"))
relationships = rels.union(reversed_rels)
g = GraphFrame(nodes, relationships)
Vérifions la structure du graphe :
print("Nombre de sommets :", g.vertices.count())
print("Nombre d'arêtes :", g.edges.count())
g.vertices.show()
Question
Pourquoi avoir créé reversed_rels et fait l’union avec rels ?
Quel effet cela a-t-il sur l’orientation du graphe résultant ?
Visualisation simple¶
import networkx as nx
import matplotlib.pyplot as plt
def plot_graph(edge_list, directed=False):
G = nx.DiGraph() if directed else nx.Graph()
for row in edge_list.select('src', 'dst').take(1000):
G.add_edge(row['src'], row['dst'])
plt.figure(figsize=(8, 6))
nx.draw(G, with_labels=True, font_weight='bold',
node_color='lightblue', node_size=800, font_size=9)
plt.show()
plot_graph(g.edges)
Parcours en largeur (BFS)¶
Le BFS explore le graphe en visitant d’abord tous les nœuds à distance 1 de
la source, puis à distance 2, etc. Dans GraphFrames, bfs trouve le plus
court chemin (en nombre de sauts) entre un nœud source et un ensemble de nœuds
cibles définis par un prédicat.
Trouvons d’abord les villes de taille moyenne (entre 100 000 et 300 000 hab.) :
g.vertices.filter("population > 100000 and population < 300000") \
.sort("population").show()
Puis le plus court chemin (en sauts) entre Den Haag et une ville de taille moyenne :
from_expr = "id='Den Haag'"
to_expr = "population > 100000 and population < 300000 and id <> 'Den Haag'"
result = g.bfs(from_expr, to_expr)
# On ne conserve que les colonnes de nœuds (pas les arêtes)
colonnes_noeuds = [c for c in result.columns if not c.startswith("e")]
result.select(colonnes_noeuds).show(truncate=False)
Question
Pourquoi le résultat ne contient-il qu’Ipswich, alors qu’il existe d’autres villes dans l’intervalle de population (Colchester, par exemple) ?
Exercice
Modifiez la requête pour trouver le plus court chemin (en sauts) entre Amsterdam et Colchester. Quelle est la longueur de ce chemin ?
Plus court chemin valué (Dijkstra)¶
Le BFS minimise le nombre de sauts. Pour minimiser la somme des coûts
(distances kilométriques ici), il faut employer l’algorithme de Dijkstra. Comme
GraphFrames ne l’implémente pas nativement, on le reconstruit via le passage de
messages (AggregateMessages).
from graphframes.lib import AggregateMessages as AM
from pyspark.sql import functions as F
add_path_udf = F.udf(lambda path, id: path + [id], ArrayType(StringType()))
def get_cached_dataframe(df):
"""Remplace AM.getCachedDataFrame (absent de graphframes-py récent).
Cette méthode servait à mettre en cache le DataFrame à chaque itération
pour éviter que la lignée Spark (lineage) ne s'allonge indéfiniment dans
la boucle. On peut la remplacer par un localCheckpoint qui matérialise le
DataFrame en mémoire et sur disque ET coupe la lignée logique."""
return df.localCheckpoint(eager=True)
En environnement cluster, pour la tolérance aux pannes préférer :
def get_cached_dataframe(df):
df = df.cache()
df.checkpoint(eager=True)
return df
def shortest_path(g, origin, destination, column_name="cost"):
if g.vertices.filter(g.vertices.id == destination).count() == 0:
return (spark.createDataFrame(sc.emptyRDD(), g.vertices.schema)
.withColumn("path", F.array()))
vertices = (g.vertices
.withColumn("visited", F.lit(False))
.withColumn("distance", F.when(g.vertices["id"] == origin, 0)
.otherwise(float("inf")))
.withColumn("path", F.array()))
cached_vertices = get_cached_dataframe(vertices)
g2 = GraphFrame(cached_vertices, g.edges)
while g2.vertices.filter('visited == False').first():
current_id = g2.vertices.filter('visited == False').sort("distance").first().id
msg_distance = AM.edge[column_name] + AM.src['distance']
msg_path = add_path_udf(AM.src["path"], AM.src["id"])
msg_for_dst = F.when(AM.src['id'] == current_id,
F.struct(msg_distance, msg_path))
new_distances = g2.aggregateMessages(F.min(AM.msg).alias("aggMess"),
sendToDst=msg_for_dst)
new_visited = F.when(g2.vertices.visited | (g2.vertices.id == current_id),
True).otherwise(False)
new_distance = F.when(new_distances["aggMess"].isNotNull() &
(new_distances.aggMess["col1"] < g2.vertices.distance),
new_distances.aggMess["col1"]) \
.otherwise(g2.vertices.distance)
new_path = F.when(new_distances["aggMess"].isNotNull() &
(new_distances.aggMess["col1"] < g2.vertices.distance),
new_distances.aggMess["col2"].cast("array<string>")) \
.otherwise(g2.vertices.path)
new_vertices = (g2.vertices.join(new_distances, on="id", how="left_outer")
.drop(new_distances["id"])
.withColumn("visited", new_visited)
.withColumn("newDistance", new_distance)
.withColumn("newPath", new_path)
.drop("aggMess", "distance", "path")
.withColumnRenamed("newDistance", "distance")
.withColumnRenamed("newPath", "path"))
cached_new_vertices = get_cached_dataframe(new_vertices)
g2 = GraphFrame(cached_new_vertices, g2.edges)
if g2.vertices.filter(g2.vertices.id == destination).first().visited:
return (g2.vertices.filter(g2.vertices.id == destination)
.withColumn("newPath", add_path_udf("path", "id"))
.drop("visited", "path")
.withColumnRenamed("newPath", "path"))
return (spark.createDataFrame(sc.emptyRDD(), g.vertices.schema)
.withColumn("path", F.array()))
Comparons maintenant les deux chemins Amsterdam → Colchester :
# Chemin le plus court en coût (km)
result_cout = shortest_path(g, "Amsterdam", "Colchester", "cost")
result_cout.select("id", "distance", "path").show(truncate=False)
# Chemin le plus court en nombre de sauts (BFS)
from_expr = "id='Amsterdam'"
to_expr = "id='Colchester'"
result_sauts = g.bfs(from_expr, to_expr)
colonnes = [c for c in result_sauts.columns if not c.startswith("e")]
result_sauts.select(colonnes).show(truncate=False)
Question
Les deux chemins sont-ils identiques ? Que concluez-vous sur le choix de la mesure de distance dans un graphe valué ?
Propriétés globales du graphe¶
Pour caractériser la structure d’ensemble, on calcule la longueur moyenne des
plus courts chemins et le diamètre. Comme GraphFrames ne les calcule pas
directement, on utilise la variante BFS shortestPaths qui renvoie, pour
chaque nœud, sa distance à un ensemble de nœuds cibles. Attention, par distance
on entend ici le nombre de liens traversés. Afin d’obtenir les vraies distances
(ou plus généralement coût de la traversée des liens) il faudrait appliquer
la fonction shortest_path définie plus haut à toutes les paires de sommets,
ce qui est très coûteux.
# Distances de chaque nœud vers tous les autres (landmarks = liste de tous les ids)
all_ids = [row.id for row in g.vertices.select("id").collect()]
result_sp = g.shortestPaths(landmarks=all_ids)
result_sp.select("id", "distances").show(truncate=False)
# Calcul du diamètre et de la longueur moyenne à partir du résultat
from pyspark.sql.functions import col, explode, map_values
# Extraire toutes les distances (valeurs du dictionnaire "distances")
distances_df = result_sp.select(
col("id").alias("src"),
explode(col("distances")).alias("dst", "dist")
).filter("src != dst")
diametre = distances_df.agg({"dist": "max"}).collect()[0][0]
longueur_moy = distances_df.agg({"dist": "avg"}).collect()[0][0]
print(f"Diamètre : {diametre}")
print(f"Longueur moy. chemins: {longueur_moy:.2f}")
Partie 2 : Centralités¶
Centralité de degré¶
La centralité de degré mesure simplement le nombre de liens directs d’un nœud. Dans un graphe orienté, on distingue le degré entrant (« indice de popularité ») et le degré sortant (« indice d’activité »).
Pour le graphe de transport :
total_degree = g.degrees
in_degree = g.inDegrees
out_degree = g.outDegrees
(total_degree.join(in_degree, "id", how="left")
.join(out_degree, "id", how="left")
.fillna(0)
.sort("inDegree", ascending=False)
.show())
Pour le graphe social :
social_total = social_g.degrees
social_in = social_g.inDegrees
social_out = social_g.outDegrees
(social_total.join(social_in, "id", how="left")
.join(social_out, "id", how="left")
.fillna(0)
.sort("inDegree", ascending=False)
.show())
Question
Dans les deux graphes, observez-vous que inDegree == outDegree pour
tous les nœuds ? Expliquez pourquoi c’est le cas ou non selon le graphe.
Exercice
Calculez la distribution de degrés du graphe social : pour chaque valeur de degré total, combien de nœuds ont ce degré ? Affichez le résultat sous forme de tableau trié par degré.
Centralité d’intermédiarité¶
La centralité d’intermédiarité (betweenness centrality) d’un nœud \(v\) mesure la fraction des plus courts chemins entre toutes paires de nœuds qui passent par \(v\). Un nœud avec une intermédiarité élevée est un pont structurel : sa suppression peut déconnecter des parties du réseau.
Comme GraphFrames ne propose pas cette centralité en standard, nous l’obtenons ici (sous forme approximative) avec NetworkX, en travaillant sur le graphe chargé en mémoire locale (acceptable pour notre petit graphe social).
# Reconstruction du graphe dans NetworkX à partir des données Spark
G_social = nx.DiGraph()
for row in social_e.select('src', 'dst').collect():
G_social.add_edge(row['src'], row['dst'])
# Calcul de la centralité d'intermédiarité (normalisée)
betweenness = nx.betweenness_centrality(G_social, normalized=True)
sorted(betweenness.items(), key=lambda x: x[1], reverse=True)
Question
Quel(s) nœud(s) ont la plus forte centralité d’intermédiarité ? Retrouvez-les sur la visualisation du graphe social. Leur position vous semble-t-elle cohérente avec un rôle de « pont » ou « courtier » ?
Exercice
Faites la même chose pour le graphe de transport. Quelles villes jouent le rôle de « points de passage obligés » dans le réseau ?
Centralité de proximité¶
La centralité de proximité (closeness centrality) d’un nœud mesure l’inverse de la somme de ses distances à tous les autres nœuds : d’un nœud plus central on peut atteindre plus rapidement l’ensemble du réseau.
from graphframes.lib import AggregateMessages as AM
from pyspark.sql import functions as F
from pyspark.sql.types import *
from operator import itemgetter
paths_type = ArrayType(
StructType([StructField("id", StringType()),
StructField("distance", IntegerType())]))
def new_paths(paths, id):
paths = [{"id": c1, "distance": c2 + 1} for c1, c2 in paths if c1 != id]
paths.append({"id": id, "distance": 1})
return paths
def flatten(ids):
flat = [item for sub in ids for item in sub]
return list(dict(sorted(flat, key=itemgetter(0))).items())
def merge_paths(ids, new_ids, id):
joined = ids + (new_ids if new_ids else [])
merged = [(c1, c2) for c1, c2 in joined if c1 != id]
best = dict(sorted(merged, key=itemgetter(1), reverse=True))
return [{"id": c1, "distance": c2} for c1, c2 in best.items()]
def calculate_closeness(ids):
n = len(ids)
total = sum(c2 for c1, c2 in ids)
return 0.0 if total == 0 else n / total
new_paths_udf = F.udf(new_paths, paths_type)
flatten_udf = F.udf(flatten, paths_type)
merge_paths_udf = F.udf(merge_paths, paths_type)
closeness_udf = F.udf(calculate_closeness, DoubleType())
Pour le graphe de transport :
vertices = g.vertices.withColumn("ids", F.array())
cached_v = get_cached_dataframe(vertices)
g2 = GraphFrame(cached_v, g.edges)
for _ in range(g2.vertices.count()):
msg_dst = new_paths_udf(AM.src["ids"], AM.src["id"])
msg_src = new_paths_udf(AM.dst["ids"], AM.dst["id"])
agg = g2.aggregateMessages(F.collect_set(AM.msg).alias("agg"),
sendToSrc=msg_src, sendToDst=msg_dst)
res = agg.withColumn("newIds", flatten_udf("agg")).drop("agg")
new_v = (g2.vertices.join(res, on="id", how="left_outer")
.withColumn("mergedIds", merge_paths_udf("ids", "newIds", "id"))
.drop("ids", "newIds")
.withColumnRenamed("mergedIds", "ids"))
g2 = GraphFrame(get_cached_dataframe(new_v), g2.edges)
(g2.vertices
.withColumn("closeness", closeness_udf("ids"))
.sort("closeness", ascending=False)
.select("id", "closeness")
.show(truncate=False))
Exercice
Refaites ce calcul pour le graphe social. Les mêmes nœuds sont les plus centraux selon la proximité, le degré ou l’intermédiarité ?
Partie 3 : PageRank¶
PageRank standard¶
Le PageRank modélise un marcheur aléatoire sur le graphe : à chaque étape, il suit un lien au hasard avec probabilité \(d\) (ici 0,85), ou « téléporte » vers un nœud aléatoire avec probabilité \(1-d\). Le PageRank d’un nœud est la fraction du temps que le marcheur y passe à l’équilibre.
Deux variantes sont disponibles dans GraphFrames :
Nombre d’itérations fixé (arrêt mécanique) :
results = g.pageRank(resetProbability=0.15, maxIter=20)
results.vertices.sort("pagerank", ascending=False).show()
social_results = social_g.pageRank(resetProbability=0.15, maxIter=20)
social_results.vertices.sort("pagerank", ascending=False).show()
Jusqu’à convergence (arrêt dès que les valeurs varient de moins de tol) :
results2 = g.pageRank(resetProbability=0.15, tol=0.01)
results2.vertices.sort("pagerank", ascending=False).show()
Le calcul est très long pour ce (petit) graphe social, évitons donc de le faire.
Question
Comparez les classements obtenus par PageRank et par centralité de degré entrant pour le graphe social. Dans quel cas les deux diffèrent-ils le plus ? Pourquoi ?
Correction
PageRank calcule, pour chaque nœud X, une sorte de degré entrant (« popularité ») pondéré par le PageRank (la « popularité ») des nœuds qui « suivent » le nœud X. Dans l’exemple de réseau social, le degré entrant et le PageRank diffèrent le plus pour Mark, qui est suivi seulement par Doug (donc
inDegree = 1), mais Doug a le PageRank le plus élevé, ce qui élève le PageRank de Mark.
Question
Dans quel cas est-il préférable d’utiliser la variante à nombre d’itérations fixé plutôt que la variante à convergence ?
PageRank personnalisé¶
Le PageRank personnalisé (Personalized PageRank, PPR) modifie le mécanisme de téléportation : au lieu de sauter vers un nœud aléatoire quelconque, le marcheur revient toujours vers un nœud (ou un ensemble de nœuds) spécifique. Le PPR depuis un nœud \(u\) mesure la proximité structurelle de chaque nœud du graphe par rapport à \(u\). Il est utilisé pour la recommandation (« qui devrait suivre Doug ? »).
me = "Doug"
ppr_results = social_g.pageRank(resetProbability=0.15, maxIter=20, sourceId=me)
people_to_follow = ppr_results.vertices.sort("pagerank", ascending=False)
already_follows = list(social_g.edges.filter(f"src = '{me}'").toPandas()["dst"])
exclude = already_follows + [me]
people_to_follow[~people_to_follow.id.isin(exclude)].show()
Question
En vous aidant de la visualisation du graphe social, vérifiez si le résultat est intuitif. Qui sont les « amis des amis » de Doug qui lui sont suggérés en premier ?
Synthèse¶
À l’issue de ce TP, vous êtes en mesure de :
créer un graphe GraphFrames à partir de DataFrames Spark et (si sa taille le permet) le visualiser avec NetworkX ;
calculer les plus courts chemins (BFS en nombre de sauts, alg. de Dijkstra en coût) et des propriétés globales (diamètre, longueur moyenne) ;
calculer et comparer plusieurs mesures de centralité (degré, intermédiarité, proximité, PageRank) ;
utiliser le PageRank personnalisé pour la recommandation.
Les mêmes algorithmes, exécutés ici sur de petits graphes de démonstration, passent à l’échelle sur de très grands graphes grâce à la parallélisation Spark.
Références¶
Ce TP s’appuie en partie sur le livre Graph Algorithms: Practical Examples in Apache Spark and Neo4j de Mark Needham et Amy E. Hodler (O’Reilly, 2019).