Guide de l’utilisateur GraphFrames - Scala
Cet article présente des exemples du guide de l'utilisateur GraphFrames.
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
idqui 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) etdst(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
// 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 :
val g = GraphFrame(v, e)
// 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.
display(g.vertices)
display(g.edges)
Le degré entrant des sommets :
display(g.inDegrees)
Le degré sortant des sommets :
display(g.outDegrees)
Le degré des sommets :
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 :
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 :
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.
// 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 :
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.
// 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.
// 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.
// 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)
display(g2.vertices)
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.
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.
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.
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.
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é).
val result = g.labelPropagation.maxIter(5).run()
display(result.orderBy("label"))
PageRank
Identifiez les sommets importants dans un Graphe en fonction des connexions.
// Run PageRank until convergence to tolerance "tol".
val results = g.pageRank.resetProbability(0.15).tol(0.01).run()
display(results.vertices)
display(results.edges)
// Run PageRank for a fixed number of iterations.
val results2 = g.pageRank.resetProbability(0.15).maxIter(10).run()
display(results2.vertices)
// 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.
val paths = g.shortestPaths.landmarks(Seq("a", "d")).run()
display(paths)
Comptage de triangles
Calcule le nombre de triangles passant par chaque sommet.
import org.graphframes.examples
val g: GraphFrame = examples.Graphs.friends // get example graph
val results = g.triangleCount.run()
results.select("id", "count").show()