Aller au contenu principal

Guide de l’utilisateur GraphFrames - Scala

Cet article présente des exemples du guide de l'utilisateur GraphFrames.

Scala
import org.apache.spark.sql._
import org.apache.spark.sql.functions._
import org.graphframes._

Création de GraphFrames

Vous pouvez créer des GraphFrames à partir de DataFrames de sommets et d'arêtes.

  • DataFrame de sommets : Un DataFrame de sommets doit contenir une colonne spéciale nommée id qui spécifie des ID uniques pour chaque sommet dans le Graphe.
  • DataFrame de périphérie : Un DataFrame de périphérie doit contenir deux colonnes spéciales : src (ID de sommet source de la périphérie) et dst (ID de sommet de destination de la périphérie).

Les deux DataFrames peuvent avoir d'autres colonnes arbitraires. Ces colonnes peuvent représenter les attributs de sommet et d'arête.

Créer les sommets et les arêtes

Scala
// Vertex DataFrame
val v = spark.createDataFrame(List(
("a", "Alice", 34),
("b", "Bob", 36),
("c", "Charlie", 30),
("d", "David", 29),
("e", "Esther", 32),
("f", "Fanny", 36),
("g", "Gabby", 60)
)).toDF("id", "name", "age")
// Edge DataFrame
val e = spark.createDataFrame(List(
("a", "b", "friend"),
("b", "c", "follow"),
("c", "b", "follow"),
("f", "c", "follow"),
("e", "f", "follow"),
("e", "d", "friend"),
("d", "a", "friend"),
("a", "e", "friend")
)).toDF("src", "dst", "relationship")

Créons un graphe à partir de ces sommets et de ces arêtes :

Scala
val g = GraphFrame(v, e)
Scala
// This example graph also comes with the GraphFrames package.
// val g = examples.Graphs.friends

Requêtes de Graphe et de DataFrame de base

GraphFrames fournit des requêtes Graphe simples, telles que le degré de nœud.

De plus, comme les GraphFrames représentent les graphes comme des paires de DataFrames de sommets et d’arêtes, il est facile d’effectuer des query puissantes directement sur les DataFrames de sommets et d’arêtes. Ces DataFrames sont disponibles en tant que champs de sommets et d'arêtes dans le GraphFrame.

Scala
display(g.vertices)
Scala
display(g.edges)

Le degré entrant des sommets :

Scala
display(g.inDegrees)

Le degré sortant des sommets :

Scala
display(g.outDegrees)

Le degré des sommets :

Scala
display(g.degrees)

Vous pouvez exécuter des query directement sur le DataFrame des sommets. Par exemple, nous pouvons trouver l'âge de la personne la plus jeune dans le graphe :

Scala
val youngest = g.vertices.groupBy().min("age")
display(youngest)

De même, vous pouvez exécuter des queries sur le DataFrame d’arêtes. Par exemple, comptons le nombre de relations « follow » dans le Graphe :

Scala
val numFollows = g.edges.filter("relationship = 'follow'").count()

Découverte de motif

Établissez des relations plus complexes impliquant des arêtes et des sommets à l'aide de motifs. La cellule suivante trouve les paires de sommets avec des arêtes dans les deux directions entre eux. Le résultat est un DataFrame, dans lequel les noms de colonne sont des clés de motif.

Consultez le Guide de l'utilisateur GraphFrame pour plus de détails sur l'API.

Scala
// Search for pairs of vertices with edges in both directions between them.
val motifs = g.find("(a)-[e]->(b); (b)-[e2]->(a)")
display(motifs)

Puisque le résultat est un DataFrame, vous pouvez élaborer des queries plus complexes à partir du motif. Cherchons toutes les relations réciproques dans lesquelles une personne a plus de 30 ans :

Scala
val filtered = motifs.filter("b.age > 30")
display(filtered)

queries avec état

La plupart des requêtes de motifs sont sans état et simples à exprimer, comme dans les exemples ci-dessus. Les exemples suivants illustrent des requêtes plus complexes qui transportent un état le long d'un chemin dans le motif. Exprimez ces queries en combinant la recherche de motifs GraphFrame avec des filtres sur le résultat, où les filtres utilisent des opérations de séquence pour construire une série de colonnes DataFrame.

Par exemple, supposons que vous vouliez identifier une chaîne de 4 sommets avec une propriété définie par une séquence de fonctions. C'est-à-dire, parmi les chaînes de 4 sommets a->b->c->d, identifiez le sous-ensemble de chaînes correspondant à ce filtre complexe :

  • Initialiser l'état sur le chemin d'accès.
  • Mettre à jour l'état en fonction du sommet a.
  • Mettre à jour l'état en fonction du sommet b.
  • etc. pour c et d.
  • Si l'état final correspond à une certaine condition, le filtre accepte alors la chaîne.

Les extraits de code suivants démontrent ce processus, où nous identifions des chaînes de 4 sommets de sorte qu’au moins 2 des 3 arêtes sont des relations « ami ». Dans cet exemple, l’état est le nombre actuel d’arêtes « ami » ; en général, il peut s’agir de n’importe quelle colonne DataFrame.

Scala
// Find chains of 4 vertices.
val chain4 = g.find("(a)-[ab]->(b); (b)-[bc]->(c); (c)-[cd]->(d)")

// Query on sequence, with state (cnt)
// (a) Define method for updating state given the next element of the motif.
def sumFriends(cnt: Column, relationship: Column): Column = {
when(relationship === "friend", cnt + 1).otherwise(cnt)
}
// (b) Use sequence operation to apply method to sequence of elements in motif.
// In this case, the elements are the 3 edges.
val condition = Seq("ab", "bc", "cd").
foldLeft(lit(0))((cnt, e) => sumFriends(cnt, col(e)("relationship")))
// (c) Apply filter to DataFrame.
val chainWith2Friends2 = chain4.where(condition >= 2)
display(chainWith2Friends2)

Sous-graphes

GraphFrames fournit des APIs pour construire des sous-graphes en filtrant sur les arêtes et les sommets. Ces filtres peuvent être combinés. Par exemple, le sous-graphe suivant ne contient que des personnes qui sont amies et qui ont plus de 30 ans.

Scala
// Select subgraph of users older than 30, and edges of type "friend"
val g2 = g
.filterEdges("relationship = 'friend'")
.filterVertices("age > 30")
.dropIsolatedVertices()

Filtres triplets complexes

L'exemple suivant montre comment sélectionner un sous-graphe basé sur des filtres de triplets qui opèrent sur une arête et ses sommets « src » et « dst ». Il est simple d'étendre cet exemple au-delà des triplets en utilisant des motifs plus complexes.

Scala
// Select subgraph based on edges "e" of type "follow"
// pointing from a younger user "a" to an older user "b".
val paths = g.find("(a)-[e]->(b)")
.filter("e.relationship = 'follow'")
.filter("a.age < b.age")
// "paths" contains vertex info. Extract the edges.
val e2 = paths.select("e.src", "e.dst", "e.relationship")
// In Spark 1.5+, the user may simplify this call:
// val e2 = paths.select("e.*")

// Construct the subgraph
val g2 = GraphFrame(g.vertices, e2)
Scala
display(g2.vertices)
Scala
display(g2.edges)

Algorithmes de graphe standard

Cette section décrit les algorithmes de graphe standard intégrés à GraphFrames.

Recherche en largeur (BFS)

Rechercher à partir d'« Esther » des utilisateurs de moins de 32 ans.

Scala
val paths: DataFrame = g.bfs.fromExpr("name = 'Esther'").toExpr("age < 32").run()
display(paths)

La recherche peut également limiter les filtres de bord et les longueurs de chemin maximales.

Scala
val filteredPaths = g.bfs.fromExpr("name = 'Esther'").toExpr("age < 32")
.edgeFilter("relationship != 'friend'")
.maxPathLength(3)
.run()
display(filteredPaths)

Composants connectés

Calculez l'appartenance de chaque sommet au composant connecté et renvoyez un graphe où chaque sommet se voit attribuer un ID de composant.

Scala
val result = g.connectedComponents.run() // doesn't work on Spark 1.4
display(result)

Composants fortement connectés

Calculez la composante fortement connexe (SCC) de chaque sommet et renvoyez un graphe avec chaque sommet affecté à la SCC contenant ce sommet.

Scala
val result = g.stronglyConnectedComponents.maxIter(10).run()
display(result.orderBy("component"))

Propagation des étiquettes

Exécutez un algorithme de propagation d'étiquettes statique pour détecter les communautés dans les réseaux.

Chaque nœud du réseau est initialement attribué à sa propre communauté. À chaque superpas, les nœuds envoient leur affiliation communautaire à tous les voisins et mettent à jour leur état à l'affiliation communautaire du mode des messages entrants.

Le LPA est un algorithme standard de détection de communautés pour les graphes. C'est peu coûteux en termes de calcul, bien que (1) la convergence ne soit pas garantie et (2) l'on puisse aboutir à des solutions triviales (tous les nœuds s'identifient en une seule communauté).

Scala
val result = g.labelPropagation.maxIter(5).run()
display(result.orderBy("label"))

PageRank

Identifiez les sommets importants dans un Graphe en fonction des connexions.

Scala
// Run PageRank until convergence to tolerance "tol".
val results = g.pageRank.resetProbability(0.15).tol(0.01).run()
display(results.vertices)
Scala
display(results.edges)
Scala
// Run PageRank for a fixed number of iterations.
val results2 = g.pageRank.resetProbability(0.15).maxIter(10).run()
display(results2.vertices)
Scala
// Run PageRank personalized for vertex "a"
val results3 = g.pageRank.resetProbability(0.15).maxIter(10).sourceId("a").run()
display(results3.vertices)

Chemins les plus courts

Calcule les chemins les plus courts vers l'ensemble donné de sommets de référence, où les sommets de référence sont spécifiés par l'ID de sommet.

Scala
val paths = g.shortestPaths.landmarks(Seq("a", "d")).run()
display(paths)

Comptage de triangles

Calcule le nombre de triangles passant par chaque sommet.

Scala
import org.graphframes.examples
val g: GraphFrame = examples.Graphs.friends // get example graph

val results = g.triangleCount.run()
results.select("id", "count").show()