From 2b54d190b0c5d1c172f3a1c0baedab61839ec6b9 Mon Sep 17 00:00:00 2001 From: semyonsinchenko Date: Fri, 24 Oct 2025 11:25:44 +0200 Subject: [PATCH 1/2] maximal independent set --- connect/src/main/protobuf/graphframes.proto | 8 + .../graphframes/GraphFramesConnectUtils.scala | 13 + .../scala/org/graphframes/GraphFrame.scala | 9 + .../lib/MaximalIndependentSet.scala | 225 ++++++++++++++++++ .../lib/MaximalIndependentSetSuite.scala | 121 ++++++++++ python/graphframes/classic/graphframe.py | 15 ++ .../graphframes/connect/graphframes_client.py | 55 +++++ .../connect/proto/graphframes_pb2.py | 90 +++---- .../connect/proto/graphframes_pb2.pyi | 136 ++++------- python/graphframes/graphframe.py | 42 ++++ python/tests/test_graphframes.py | 22 ++ 11 files changed, 597 insertions(+), 139 deletions(-) create mode 100644 core/src/main/scala/org/graphframes/lib/MaximalIndependentSet.scala create mode 100644 core/src/test/scala/org/graphframes/lib/MaximalIndependentSetSuite.scala diff --git a/connect/src/main/protobuf/graphframes.proto b/connect/src/main/protobuf/graphframes.proto index c22223941..42a0b3c27 100644 --- a/connect/src/main/protobuf/graphframes.proto +++ b/connect/src/main/protobuf/graphframes.proto @@ -35,6 +35,7 @@ message GraphFramesAPI { SVDPlusPlus svd_plus_plus = 18; TriangleCount triangle_count = 19; Triplets triplets = 20; + MaximalIndependentSet mis = 22; } } @@ -186,3 +187,10 @@ message TriangleCount { } message Triplets {} + +message MaximalIndependentSet { + int32 checkpoint_interval = 1; + optional StorageLevel storage_level = 2; + bool use_local_checkpoints = 3; + int64 seed = 4; +} diff --git a/connect/src/main/scala/org/apache/spark/sql/graphframes/GraphFramesConnectUtils.scala b/connect/src/main/scala/org/apache/spark/sql/graphframes/GraphFramesConnectUtils.scala index ea63ae281..bf7836925 100644 --- a/connect/src/main/scala/org/apache/spark/sql/graphframes/GraphFramesConnectUtils.scala +++ b/connect/src/main/scala/org/apache/spark/sql/graphframes/GraphFramesConnectUtils.scala @@ -399,6 +399,19 @@ object GraphFramesConnectUtils { case proto.GraphFramesAPI.MethodCase.TRIPLETS => { graphFrame.triplets } + case proto.GraphFramesAPI.MethodCase.MIS => { + val mis = graphFrame.maximalIndependentSet + .setCheckpointInterval(apiMessage.getMis.getCheckpointInterval) + .setUseLocalCheckpoints(apiMessage.getMis.getUseLocalCheckpoints) + + if (apiMessage.getMis.hasStorageLevel) { + mis + .setIntermediateStorageLevel(parseStorageLevel(apiMessage.getMis.getStorageLevel)) + .run(apiMessage.getMis.getSeed) + } else { + mis.run(apiMessage.getMis.getSeed) + } + } case _ => throw new GraphFramesUnreachableException() // Unreachable } } diff --git a/core/src/main/scala/org/graphframes/GraphFrame.scala b/core/src/main/scala/org/graphframes/GraphFrame.scala index af2a1ae1e..11df5018c 100644 --- a/core/src/main/scala/org/graphframes/GraphFrame.scala +++ b/core/src/main/scala/org/graphframes/GraphFrame.scala @@ -701,6 +701,15 @@ class GraphFrame private ( */ def detectingCycles: DetectingCycles = new DetectingCycles(this) + /** + * Maximal Independent Set algorithm. + * + * See [[org.graphframes.lib.MaximalIndependentSet]] for more details. + * + * @group stdlib + */ + def maximalIndependentSet: MaximalIndependentSet = new MaximalIndependentSet(this) + /** * Converts the directed graph into an undirected graph by ensuring that all directed edges are * bidirectional. For every directed edge (src, dst), a corresponding edge (dst, src) is added. diff --git a/core/src/main/scala/org/graphframes/lib/MaximalIndependentSet.scala b/core/src/main/scala/org/graphframes/lib/MaximalIndependentSet.scala new file mode 100644 index 000000000..41cca2d0c --- /dev/null +++ b/core/src/main/scala/org/graphframes/lib/MaximalIndependentSet.scala @@ -0,0 +1,225 @@ +package org.graphframes.lib + +import org.apache.spark.sql.DataFrame +import org.apache.spark.sql.functions.* +import org.apache.spark.sql.types.DoubleType +import org.apache.spark.storage.StorageLevel +import org.graphframes.GraphFrame +import org.graphframes.Logging +import org.graphframes.WithCheckpointInterval +import org.graphframes.WithIntermediateStorageLevel +import org.graphframes.WithLocalCheckpoints + +import java.io.IOException + +/** + * This class implements a distributed algorithm for finding a Maximal Independent Set (MIS) in a + * graph. + * + * An MIS is a set of vertices such that no two vertices in the set are adjacent (i.e., there is + * no edge between any two vertices in the set), and the set is maximal, meaning that adding any + * other vertex to the set would violate the independence property. Note that this implementation + * finds a maximal (but not necessarily maximum) independent set; that is, it ensures no more + * vertices can be added to the set, but does not guarantee that the set has the largest possible + * number of vertices among all possible independent sets in the graph. + * + * The algorithm implemented here is based on the paper: Ghaffari, Mohsen. "An improved + * distributed algorithm for maximal independent set." Proceedings of the twenty-seventh annual + * ACM-SIAM symposium on Discrete algorithms. Society for Industrial and Applied Mathematics, + * 2016. + * + * Note: This is a randomized, non-deterministic algorithm. The result may vary between runs even + * if a fixed random seed is provided because how Apache Spark works. + * + * @param graph + */ +class MaximalIndependentSet private[graphframes] (private val graph: GraphFrame) + extends Serializable + with WithIntermediateStorageLevel + with WithCheckpointInterval + with WithLocalCheckpoints { + def run(seed: Long): DataFrame = { + MaximalIndependentSet.run( + graph, + checkpointInterval, + useLocalCheckpoints, + intermediateStorageLevel, + seed) + } +} + +object MaximalIndependentSet extends Serializable with Logging { + private val probCol = "prob" + private val degCol = "effectiveDegree" + private val isNominated = "isNominated" + private val notJoinedMISCol = "notJoinMIS" + private val isMIS = "isMIS" + + private def run( + graph: GraphFrame, + checkpointInterval: Int, + useLocalCheckpoints: Boolean, + storageLevel: StorageLevel, + seed: Long): DataFrame = { + // initial p = 1/2 + var vertices = + graph.vertices + .select(col(GraphFrame.ID), lit(0.5).cast(DoubleType).alias(probCol)) + .persist(storageLevel) + + // make edges undirected and de-duplicate + // persist() for future usage + val edges = graph.edges + .select(GraphFrame.SRC, GraphFrame.DST) + .union( + graph.edges.select( + col(GraphFrame.DST).alias(GraphFrame.SRC), + col(GraphFrame.SRC).alias(GraphFrame.DST))) + .filter(col(GraphFrame.SRC) =!= col(GraphFrame.DST)) + .distinct() + .persist(storageLevel) + + var misDF = graph.vertices.select(col(GraphFrame.ID), lit(false).alias(isMIS)) + + var i = 0 + var converged = false + val spark = graph.vertices.sparkSession + + val shouldCheckpoint = checkpointInterval > 0 + if (!useLocalCheckpoints && spark.sparkContext.getCheckpointDir.isEmpty) { + // Spark-Connect workaround + spark.sparkContext + .setCheckpointDir(spark.conf + .getOption("spark.checkpoint.dir") match { + case Some(d) => d + case None => + throw new IOException( + "Checkpoint directory is not set. Please set it first using sc.setCheckpointDir()" + + "or by specifying the conf 'spark.checkpoint.dir'.") + }) + } + + val rng = new util.Random(seed) + + // randomized algorithms are not working with AQE well + val originalAQE = spark.conf.get("spark.sql.adaptive.enabled") + try { + spark.conf.set("spark.sql.adaptive.enabled", "false") + + while (!converged) { + val iterSeed = rng.nextLong() + // compute effective degree as a sum of nbrs p + val effectiveDegrees = + edges + .join(vertices, col(GraphFrame.ID) === col(GraphFrame.DST)) + .groupBy(GraphFrame.SRC) + .agg(sum(col(probCol)).alias(degCol)) + + // update p per vertex by condition: + // if effective degree >= 2 then p / 2 + // else min(2p, 1/2) + // + // + mark vertices based on p + val probs = vertices + .join(effectiveDegrees, col(GraphFrame.ID) === col(GraphFrame.SRC)) + .drop(GraphFrame.SRC) + .withColumn( + probCol, + when(col(degCol) >= lit(2), col(probCol) / lit(2.0)).otherwise( + when(lit(2) * col(probCol) <= lit(0.5), lit(2) * col(probCol)).otherwise(lit(0.5)))) + .withColumn(isNominated, col(probCol) >= rand(iterSeed)) + .select(GraphFrame.ID, isNominated, probCol) + .persist(storageLevel) + + val isolatedVertices = + vertices + .join(probs.select(col(GraphFrame.ID)), Seq(GraphFrame.ID), "left_anti") + .select(GraphFrame.ID) + + // if no nbr of v is marked and v is marked, + // v is joined MIS and removed with all it's nbrs + val isJoinedMIS = probs + .join( + edges + .join(probs, col(GraphFrame.ID) === col(GraphFrame.DST)) + .groupBy(GraphFrame.SRC) + .agg(bool_or(col(isNominated)).alias(notJoinedMISCol)), + col(GraphFrame.SRC) === col(GraphFrame.ID)) + .select(GraphFrame.ID, probCol, isNominated, notJoinedMISCol) + + val joinedMIS = + isJoinedMIS.filter((!col(notJoinedMISCol)) && col(isNominated)).select(GraphFrame.ID) + + // update curent MIS + val updatedMIS = misDF + .join( + isolatedVertices.select(col(GraphFrame.ID), lit(true).alias("f")), + Seq(GraphFrame.ID), + "left") + .select(col(GraphFrame.ID), (col(isMIS) || col("f")).alias(isMIS)) + .join( + joinedMIS.select(col(GraphFrame.ID), lit(true).alias("f")), + Seq(GraphFrame.ID), + "left") + .select(col(GraphFrame.ID), (col(isMIS) || col("f")).alias(isMIS)) + .persist(storageLevel) + + // We cannot not checkpoint current MIS, otherwise it is almost not working. + if (useLocalCheckpoints) { + val newMis = updatedMIS.localCheckpoint(eager = true) + newMis.count() + misDF.unpersist() + misDF = newMis + } else { + val newMis = updatedMIS.checkpoint(eager = true) + newMis.count() + misDF.unpersist() + misDF = newMis + } + + val neighborsOfMIS = edges + .join(joinedMIS, col(GraphFrame.ID) === col(GraphFrame.DST)) + .select(col(GraphFrame.SRC)) + + val updatedVertices = probs + .join(joinedMIS, Seq(GraphFrame.ID), "left_anti") + .join(neighborsOfMIS, col(GraphFrame.ID) === col(GraphFrame.SRC), "left_anti") + .select(GraphFrame.ID, probCol) + + // checkpointing of vertices + if (shouldCheckpoint && (i % checkpointInterval == 0)) { + if (useLocalCheckpoints) { + vertices = updatedVertices.localCheckpoint(eager = true) + } else { + vertices = updatedVertices.checkpoint(eager = true) + } + } else { + vertices = updatedVertices + } + + // algorithm stops if no more vertex left + converged = vertices.isEmpty + + updatedVertices.unpersist() + probs.unpersist() + + logInfo(s"iteration $i finished, vertices left: ${vertices.count()}") + i += 1 + } + + vertices.unpersist(true) + edges.unpersist(true) + + val mis = misDF.filter(col(isMIS)).select(GraphFrame.ID).persist(storageLevel) + // materialize + mis.count() + resultIsPersistent() + misDF.unpersist(true) + + mis + } finally { + // Restore original AQE setting + spark.conf.set("spark.sql.adaptive.enabled", originalAQE) + } + } +} diff --git a/core/src/test/scala/org/graphframes/lib/MaximalIndependentSetSuite.scala b/core/src/test/scala/org/graphframes/lib/MaximalIndependentSetSuite.scala new file mode 100644 index 000000000..80c5416d4 --- /dev/null +++ b/core/src/test/scala/org/graphframes/lib/MaximalIndependentSetSuite.scala @@ -0,0 +1,121 @@ +package org.graphframes.lib + +import org.apache.spark.sql.DataFrame +import org.apache.spark.sql.functions.col +import org.graphframes.* +import org.graphframes.examples.Graphs + +class MaximalIndependentSetSuite extends SparkFunSuite with GraphFrameTestSparkContext { + test("isolated vertices should be included in MIS") { + // Create a graph with isolated vertices + val vertices = + spark.createDataFrame(Seq((0L, "a"), (1L, "b"), (2L, "c"), (3L, "d"))).toDF("id", "name") + + // Only connect vertices 0 and 1 + val edges = spark.createDataFrame(Seq((0L, 1L, "edge1"))).toDF("src", "dst", "name") + + val graph = GraphFrame(vertices, edges) + val mis = graph.maximalIndependentSet.run(seed = 12345L) + + // Check that all vertices are in the MIS (since 2 and 3 are isolated) + val misIds = mis.select("id").collect().map(_.getLong(0)).toSet + assert(misIds.size == 3, "MIS should contain 2 isolated vertices and one of linked") + assert(misIds.contains(2L), "Isolated vertex 2 should be in MIS") + assert(misIds.contains(3L), "Isolated vertex 3 should be in MIS") + + mis.unpersist() + } + + def isIndependent(graph: GraphFrame, mis: DataFrame): Boolean = { + graph.edges + .join(mis, col(GraphFrame.SRC) === col(GraphFrame.ID)) + .select(col(GraphFrame.DST)) + .join(mis, col(GraphFrame.DST) === col(GraphFrame.ID)) + .count() == 0 + } + + def isMaximal(graph: GraphFrame, mis: DataFrame): Boolean = { + val undirectedG = graph.asUndirected() + val verticesNotInMIS = undirectedG.vertices.join(mis, Seq(GraphFrame.ID), "left_anti") + + val verticesWithEdgesToMIS = undirectedG.edges + .join(mis, col(GraphFrame.ID) === col(GraphFrame.DST)) + .select(GraphFrame.SRC) + .distinct() + + val countVerticesNotInMIS = verticesNotInMIS.count() + val countVerticesWithEdgesToMIS = verticesWithEdgesToMIS.count() + + countVerticesNotInMIS == countVerticesWithEdgesToMIS + } + + test("correct MIS, seed 12345") { + val graph = Graphs.friends + + val mis = graph.maximalIndependentSet.run(seed = 12345L) + + assert(isIndependent(graph, mis)) + assert(isMaximal(graph, mis)) + + mis.unpersist() + } + + test("correct MIS, seed 23456") { + val graph = Graphs.friends + + val mis = graph.maximalIndependentSet.run(seed = 23456L) + + assert(isIndependent(graph, mis)) + assert(isMaximal(graph, mis)) + + mis.unpersist() + } + + test("MIS on empty graph") { + val emptyGraph = Graphs.empty[Long] + val mis = emptyGraph.maximalIndependentSet.run(seed = 12345L) + assert(mis.count() == 0, "MIS of empty graph should be empty") + mis.unpersist() + } + + test("MIS on single vertex graph") { + val vertices = spark.createDataFrame(Seq((0L, "vertex"))).toDF("id", "name") + val edges = spark.createDataFrame(Seq.empty[(Long, Long)]).toDF("src", "dst") + val graph = GraphFrame(vertices, edges) + + val mis = graph.maximalIndependentSet.run(seed = 12345L) + assert(mis.count() == 1, "MIS of single vertex graph should contain one vertex") + + val misId = mis.select("id").collect()(0).getLong(0) + assert(misId == 0L, "MIS should contain vertex with id 0") + + mis.unpersist() + } + + test("MIS on disconnected vertices") { + val n = 5L + val vertices = spark.range(n).toDF("id") + val edges = spark.createDataFrame(Seq.empty[(Long, Long)]).toDF("src", "dst") + val graph = GraphFrame(vertices, edges) + + val mis = graph.maximalIndependentSet.run(seed = 12345L) + assert(mis.count() == n, s"MIS should contain all $n vertices for disconnected graph") + + mis.unpersist() + } + + test("MIS on complete graph of 5 vertices") { + val vertices = spark.range(5).toDF("id") + val edges = for { + i <- 0L until 5L + j <- (i + 1L) until 5L + } yield (i, j) + val edgeDF = spark.createDataFrame(edges).toDF("src", "dst") + val graph = GraphFrame(vertices, edgeDF) + + val mis = graph.maximalIndependentSet.run(seed = 12345L) + assert(mis.count() == 1, "MIS of complete graph should contain exactly one vertex") + + mis.unpersist() + } +} diff --git a/python/graphframes/classic/graphframe.py b/python/graphframes/classic/graphframe.py index ab9618906..b9a7fcf1e 100644 --- a/python/graphframes/classic/graphframe.py +++ b/python/graphframes/classic/graphframe.py @@ -341,3 +341,18 @@ def powerIterationClustering( weightCol = self._spark._jvm.scala.Option.empty() jdf = self._jvm_graph.powerIterationClustering(k, maxIter, weightCol) return DataFrame(jdf, self._spark) + + def maximal_independent_set( + self, + checkpoint_interval: int, + storage_level: StorageLevel, + use_local_checkpoints: bool, + seed: int, + ) -> DataFrame: + builder = self._jvm_graph.maximalIndependentSet() + builder.setCheckpointInterval(checkpoint_interval) + builder.setIntermediateStorageLevel(storage_level_to_jvm(storage_level, self._spark)) + builder.setUseLocalCheckpoints(use_local_checkpoints) + + jdf = builder.run(seed) + return DataFrame(jdf, self._spark) diff --git a/python/graphframes/connect/graphframes_client.py b/python/graphframes/connect/graphframes_client.py index 2e64d7a41..2ea77cee8 100644 --- a/python/graphframes/connect/graphframes_client.py +++ b/python/graphframes/connect/graphframes_client.py @@ -1064,3 +1064,58 @@ def plan(self, session: SparkConnectClient) -> proto.Relation: return _dataframe_from_plan( TriangleCount(self._vertices, self._edges, storage_level), self._spark ) + + def maximal_independent_set( + self, + checkpoint_interval: int, + storage_level: StorageLevel, + use_local_checkpoints: bool, + seed: int, + ) -> DataFrame: + @final + class MaximalIndependentSet(LogicalPlan): + def __init__( + self, + v: DataFrame, + e: DataFrame, + checkpoint_interval: int, + storage_level: StorageLevel, + use_local_checkpoints: bool, + seed: int, + ) -> None: + super().__init__(None) + self.v = v + self.e = e + self.checkpoint_interval = checkpoint_interval + self.storage_level = (storage_level,) + self.use_local_checkpoints = use_local_checkpoints + self.seed = seed + + @override + def plan(self, session: SparkConnectClient) -> proto.Relation: + graphframes_api_call = GraphFrameConnect._get_pb_api_message( + self.v, self.e, session + ) + graphframes_api_call.mis.CopyFrom( + pb.MaximalIndependentSet( + checkpoint_interval=self.checkpoint_interval, + storage_level=storage_level_to_proto(self.storage_level), + use_local_checkpoints=self.use_local_checkpoints, + seed=self.seed, + ) + ) + plan = self._create_proto_relation() + plan.extension.Pack(graphframes_api_call) + return plan + + return _dataframe_from_plan( + MaximalIndependentSet( + self._vertices, + self._edges, + checkpoint_interval, + storage_level, + use_local_checkpoints, + seed, + ), + self._spark, + ) diff --git a/python/graphframes/connect/proto/graphframes_pb2.py b/python/graphframes/connect/proto/graphframes_pb2.py index a34f42a42..cf3731ba9 100644 --- a/python/graphframes/connect/proto/graphframes_pb2.py +++ b/python/graphframes/connect/proto/graphframes_pb2.py @@ -19,7 +19,7 @@ DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile( - b'\n\x11graphframes.proto\x12\x1dorg.graphframes.connect.proto"\xb3\r\n\x0eGraphFramesAPI\x12\x1a\n\x08vertices\x18\x01 \x01(\x0cR\x08vertices\x12\x14\n\x05\x65\x64ges\x18\x02 \x01(\x0cR\x05\x65\x64ges\x12\x61\n\x12\x61ggregate_messages\x18\x03 \x01(\x0b\x32\x30.org.graphframes.connect.proto.AggregateMessagesH\x00R\x11\x61ggregateMessages\x12\x36\n\x03\x62\x66s\x18\x04 \x01(\x0b\x32".org.graphframes.connect.proto.BFSH\x00R\x03\x62\x66s\x12g\n\x14\x63onnected_components\x18\x05 \x01(\x0b\x32\x32.org.graphframes.connect.proto.ConnectedComponentsH\x00R\x13\x63onnectedComponents\x12k\n\x16\x64rop_isolated_vertices\x18\x06 \x01(\x0b\x32\x33.org.graphframes.connect.proto.DropIsolatedVerticesH\x00R\x14\x64ropIsolatedVertices\x12[\n\x10\x64\x65tecting_cycles\x18\x07 \x01(\x0b\x32..org.graphframes.connect.proto.DetectingCyclesH\x00R\x0f\x64\x65tectingCycles\x12O\n\x0c\x66ilter_edges\x18\x08 \x01(\x0b\x32*.org.graphframes.connect.proto.FilterEdgesH\x00R\x0b\x66ilterEdges\x12X\n\x0f\x66ilter_vertices\x18\t \x01(\x0b\x32-.org.graphframes.connect.proto.FilterVerticesH\x00R\x0e\x66ilterVertices\x12\x39\n\x04\x66ind\x18\n \x01(\x0b\x32#.org.graphframes.connect.proto.FindH\x00R\x04\x66ind\x12^\n\x11label_propagation\x18\x0b \x01(\x0b\x32/.org.graphframes.connect.proto.LabelPropagationH\x00R\x10labelPropagation\x12\x46\n\tpage_rank\x18\x0c \x01(\x0b\x32\'.org.graphframes.connect.proto.PageRankH\x00R\x08pageRank\x12\x84\x01\n\x1fparallel_personalized_page_rank\x18\r \x01(\x0b\x32;.org.graphframes.connect.proto.ParallelPersonalizedPageRankH\x00R\x1cparallelPersonalizedPageRank\x12w\n\x1apower_iteration_clustering\x18\x0e \x01(\x0b\x32\x37.org.graphframes.connect.proto.PowerIterationClusteringH\x00R\x18powerIterationClustering\x12?\n\x06pregel\x18\x0f \x01(\x0b\x32%.org.graphframes.connect.proto.PregelH\x00R\x06pregel\x12U\n\x0eshortest_paths\x18\x10 \x01(\x0b\x32,.org.graphframes.connect.proto.ShortestPathsH\x00R\rshortestPaths\x12\x80\x01\n\x1dstrongly_connected_components\x18\x11 \x01(\x0b\x32:.org.graphframes.connect.proto.StronglyConnectedComponentsH\x00R\x1bstronglyConnectedComponents\x12P\n\rsvd_plus_plus\x18\x12 \x01(\x0b\x32*.org.graphframes.connect.proto.SVDPlusPlusH\x00R\x0bsvdPlusPlus\x12U\n\x0etriangle_count\x18\x13 \x01(\x0b\x32,.org.graphframes.connect.proto.TriangleCountH\x00R\rtriangleCount\x12\x45\n\x08triplets\x18\x14 \x01(\x0b\x32\'.org.graphframes.connect.proto.TripletsH\x00R\x08tripletsB\x08\n\x06method"\xd7\x02\n\x0cStorageLevel\x12\x1d\n\tdisk_only\x18\x01 \x01(\x08H\x00R\x08\x64iskOnly\x12 \n\x0b\x64isk_only_2\x18\x02 \x01(\x08H\x00R\tdiskOnly2\x12 \n\x0b\x64isk_only_3\x18\x03 \x01(\x08H\x00R\tdiskOnly3\x12(\n\x0fmemory_and_disk\x18\x04 \x01(\x08H\x00R\rmemoryAndDisk\x12+\n\x11memory_and_disk_2\x18\x05 \x01(\x08H\x00R\x0ememoryAndDisk2\x12\x33\n\x15memory_and_disk_deser\x18\x06 \x01(\x08H\x00R\x12memoryAndDiskDeser\x12!\n\x0bmemory_only\x18\x07 \x01(\x08H\x00R\nmemoryOnly\x12$\n\rmemory_only_2\x18\x08 \x01(\x08H\x00R\x0bmemoryOnly2B\x0f\n\rstorage_level"M\n\x12\x43olumnOrExpression\x12\x12\n\x03\x63ol\x18\x01 \x01(\x0cH\x00R\x03\x63ol\x12\x14\n\x04\x65xpr\x18\x02 \x01(\tH\x00R\x04\x65xprB\r\n\x0b\x63ol_or_expr"P\n\x0eStringOrLongID\x12\x19\n\x07long_id\x18\x01 \x01(\x03H\x00R\x06longId\x12\x1d\n\tstring_id\x18\x02 \x01(\tH\x00R\x08stringIdB\x04\n\x02id"\xee\x02\n\x11\x41ggregateMessages\x12J\n\x07\x61gg_col\x18\x01 \x03(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionR\x06\x61ggCol\x12Q\n\x0bsend_to_src\x18\x02 \x03(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionR\tsendToSrc\x12Q\n\x0bsend_to_dst\x18\x03 \x03(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionR\tsendToDst\x12U\n\rstorage_level\x18\x04 \x01(\x0b\x32+.org.graphframes.connect.proto.StorageLevelH\x00R\x0cstorageLevel\x88\x01\x01\x42\x10\n\x0e_storage_level"\x9d\x02\n\x03\x42\x46S\x12N\n\tfrom_expr\x18\x01 \x01(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionR\x08\x66romExpr\x12J\n\x07to_expr\x18\x02 \x01(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionR\x06toExpr\x12R\n\x0b\x65\x64ge_filter\x18\x03 \x01(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionR\nedgeFilter\x12&\n\x0fmax_path_length\x18\x04 \x01(\x05R\rmaxPathLength"\x86\x03\n\x13\x43onnectedComponents\x12\x1c\n\talgorithm\x18\x01 \x01(\tR\talgorithm\x12/\n\x13\x63heckpoint_interval\x18\x02 \x01(\x05R\x12\x63heckpointInterval\x12/\n\x13\x62roadcast_threshold\x18\x03 \x01(\x05R\x12\x62roadcastThreshold\x12\x37\n\x18use_labels_as_components\x18\x04 \x01(\x08R\x15useLabelsAsComponents\x12\x32\n\x15use_local_checkpoints\x18\x05 \x01(\x08R\x13useLocalCheckpoints\x12\x19\n\x08max_iter\x18\x06 \x01(\x05R\x07maxIter\x12U\n\rstorage_level\x18\x07 \x01(\x0b\x32+.org.graphframes.connect.proto.StorageLevelH\x00R\x0cstorageLevel\x88\x01\x01\x42\x10\n\x0e_storage_level"\xdf\x01\n\x0f\x44\x65tectingCycles\x12\x32\n\x15use_local_checkpoints\x18\x01 \x01(\x08R\x13useLocalCheckpoints\x12/\n\x13\x63heckpoint_interval\x18\x02 \x01(\x05R\x12\x63heckpointInterval\x12U\n\rstorage_level\x18\x03 \x01(\x0b\x32+.org.graphframes.connect.proto.StorageLevelH\x00R\x0cstorageLevel\x88\x01\x01\x42\x10\n\x0e_storage_level"\x16\n\x14\x44ropIsolatedVertices"^\n\x0b\x46ilterEdges\x12O\n\tcondition\x18\x01 \x01(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionR\tcondition"a\n\x0e\x46ilterVertices\x12O\n\tcondition\x18\x02 \x01(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionR\tcondition" \n\x04\x46ind\x12\x18\n\x07pattern\x18\x01 \x01(\tR\x07pattern"\x99\x02\n\x10LabelPropagation\x12\x1c\n\talgorithm\x18\x01 \x01(\tR\talgorithm\x12\x19\n\x08max_iter\x18\x02 \x01(\x05R\x07maxIter\x12\x32\n\x15use_local_checkpoints\x18\x03 \x01(\x08R\x13useLocalCheckpoints\x12/\n\x13\x63heckpoint_interval\x18\x04 \x01(\x05R\x12\x63heckpointInterval\x12U\n\rstorage_level\x18\x05 \x01(\x0b\x32+.org.graphframes.connect.proto.StorageLevelH\x00R\x0cstorageLevel\x88\x01\x01\x42\x10\n\x0e_storage_level"\xe2\x01\n\x08PageRank\x12+\n\x11reset_probability\x18\x01 \x01(\x01R\x10resetProbability\x12O\n\tsource_id\x18\x02 \x01(\x0b\x32-.org.graphframes.connect.proto.StringOrLongIDH\x00R\x08sourceId\x88\x01\x01\x12\x1e\n\x08max_iter\x18\x03 \x01(\x05H\x01R\x07maxIter\x88\x01\x01\x12\x15\n\x03tol\x18\x04 \x01(\x01H\x02R\x03tol\x88\x01\x01\x42\x0c\n\n_source_idB\x0b\n\t_max_iterB\x06\n\x04_tol"\xb4\x01\n\x1cParallelPersonalizedPageRank\x12+\n\x11reset_probability\x18\x01 \x01(\x01R\x10resetProbability\x12L\n\nsource_ids\x18\x02 \x03(\x0b\x32-.org.graphframes.connect.proto.StringOrLongIDR\tsourceIds\x12\x19\n\x08max_iter\x18\x03 \x01(\x05R\x07maxIter"v\n\x18PowerIterationClustering\x12\x0c\n\x01k\x18\x01 \x01(\x05R\x01k\x12\x19\n\x08max_iter\x18\x02 \x01(\x05R\x07maxIter\x12"\n\nweight_col\x18\x03 \x01(\tH\x00R\tweightCol\x88\x01\x01\x42\r\n\x0b_weight_col"\xe6\t\n\x06Pregel\x12L\n\x08\x61gg_msgs\x18\x01 \x01(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionR\x07\x61ggMsgs\x12X\n\x0fsend_msg_to_dst\x18\x02 \x03(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionR\x0csendMsgToDst\x12X\n\x0fsend_msg_to_src\x18\x03 \x03(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionR\x0csendMsgToSrc\x12/\n\x13\x63heckpoint_interval\x18\x04 \x01(\x05R\x12\x63heckpointInterval\x12\x19\n\x08max_iter\x18\x05 \x01(\x05R\x07maxIter\x12.\n\x13\x61\x64\x64itional_col_name\x18\x06 \x01(\tR\x11\x61\x64\x64itionalColName\x12g\n\x16\x61\x64\x64itional_col_initial\x18\x07 \x01(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionR\x14\x61\x64\x64itionalColInitial\x12_\n\x12\x61\x64\x64itional_col_upd\x18\x08 \x01(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionR\x10\x61\x64\x64itionalColUpd\x12*\n\x0e\x65\x61rly_stopping\x18\t \x01(\x08H\x00R\rearlyStopping\x88\x01\x01\x12\x32\n\x15use_local_checkpoints\x18\n \x01(\x08R\x13useLocalCheckpoints\x12U\n\rstorage_level\x18\x0b \x01(\x0b\x32+.org.graphframes.connect.proto.StorageLevelH\x01R\x0cstorageLevel\x88\x01\x01\x12\x37\n\x16stop_if_all_non_active\x18\x0c \x01(\x08H\x02R\x12stopIfAllNonActive\x88\x01\x01\x12\x66\n\x13initial_active_expr\x18\r \x01(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionH\x03R\x11initialActiveExpr\x88\x01\x01\x12\x64\n\x12update_active_expr\x18\x0e \x01(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionH\x04R\x10updateActiveExpr\x88\x01\x01\x12\x45\n\x1dskip_messages_from_non_active\x18\x0f \x01(\x08H\x05R\x19skipMessagesFromNonActive\x88\x01\x01\x42\x11\n\x0f_early_stoppingB\x10\n\x0e_storage_levelB\x19\n\x17_stop_if_all_non_activeB\x16\n\x14_initial_active_exprB\x15\n\x13_update_active_exprB \n\x1e_skip_messages_from_non_active"\xc8\x02\n\rShortestPaths\x12K\n\tlandmarks\x18\x01 \x03(\x0b\x32-.org.graphframes.connect.proto.StringOrLongIDR\tlandmarks\x12\x1c\n\talgorithm\x18\x02 \x01(\tR\talgorithm\x12\x32\n\x15use_local_checkpoints\x18\x03 \x01(\x08R\x13useLocalCheckpoints\x12/\n\x13\x63heckpoint_interval\x18\x04 \x01(\x05R\x12\x63heckpointInterval\x12U\n\rstorage_level\x18\x05 \x01(\x0b\x32+.org.graphframes.connect.proto.StorageLevelH\x00R\x0cstorageLevel\x88\x01\x01\x42\x10\n\x0e_storage_level"8\n\x1bStronglyConnectedComponents\x12\x19\n\x08max_iter\x18\x01 \x01(\x05R\x07maxIter"\xd6\x01\n\x0bSVDPlusPlus\x12\x12\n\x04rank\x18\x01 \x01(\x05R\x04rank\x12\x19\n\x08max_iter\x18\x02 \x01(\x05R\x07maxIter\x12\x1b\n\tmin_value\x18\x03 \x01(\x01R\x08minValue\x12\x1b\n\tmax_value\x18\x04 \x01(\x01R\x08maxValue\x12\x16\n\x06gamma1\x18\x05 \x01(\x01R\x06gamma1\x12\x16\n\x06gamma2\x18\x06 \x01(\x01R\x06gamma2\x12\x16\n\x06gamma6\x18\x07 \x01(\x01R\x06gamma6\x12\x16\n\x06gamma7\x18\x08 \x01(\x01R\x06gamma7"x\n\rTriangleCount\x12U\n\rstorage_level\x18\x01 \x01(\x0b\x32+.org.graphframes.connect.proto.StorageLevelH\x00R\x0cstorageLevel\x88\x01\x01\x42\x10\n\x0e_storage_level"\n\n\x08TripletsB\xd2\x01\n!com.org.graphframes.connect.protoB\x10GraphframesProtoH\x01P\x01\xa0\x01\x01\xa2\x02\x04OGCP\xaa\x02\x1dOrg.Graphframes.Connect.Proto\xca\x02\x1dOrg\\Graphframes\\Connect\\Proto\xe2\x02)Org\\Graphframes\\Connect\\Proto\\GPBMetadata\xea\x02 Org::Graphframes::Connect::Protob\x06proto3' + b'\n\x11graphframes.proto\x12\x1dorg.graphframes.connect.proto"\xfd\r\n\x0eGraphFramesAPI\x12\x1a\n\x08vertices\x18\x01 \x01(\x0cR\x08vertices\x12\x14\n\x05\x65\x64ges\x18\x02 \x01(\x0cR\x05\x65\x64ges\x12\x61\n\x12\x61ggregate_messages\x18\x03 \x01(\x0b\x32\x30.org.graphframes.connect.proto.AggregateMessagesH\x00R\x11\x61ggregateMessages\x12\x36\n\x03\x62\x66s\x18\x04 \x01(\x0b\x32".org.graphframes.connect.proto.BFSH\x00R\x03\x62\x66s\x12g\n\x14\x63onnected_components\x18\x05 \x01(\x0b\x32\x32.org.graphframes.connect.proto.ConnectedComponentsH\x00R\x13\x63onnectedComponents\x12k\n\x16\x64rop_isolated_vertices\x18\x06 \x01(\x0b\x32\x33.org.graphframes.connect.proto.DropIsolatedVerticesH\x00R\x14\x64ropIsolatedVertices\x12[\n\x10\x64\x65tecting_cycles\x18\x07 \x01(\x0b\x32..org.graphframes.connect.proto.DetectingCyclesH\x00R\x0f\x64\x65tectingCycles\x12O\n\x0c\x66ilter_edges\x18\x08 \x01(\x0b\x32*.org.graphframes.connect.proto.FilterEdgesH\x00R\x0b\x66ilterEdges\x12X\n\x0f\x66ilter_vertices\x18\t \x01(\x0b\x32-.org.graphframes.connect.proto.FilterVerticesH\x00R\x0e\x66ilterVertices\x12\x39\n\x04\x66ind\x18\n \x01(\x0b\x32#.org.graphframes.connect.proto.FindH\x00R\x04\x66ind\x12^\n\x11label_propagation\x18\x0b \x01(\x0b\x32/.org.graphframes.connect.proto.LabelPropagationH\x00R\x10labelPropagation\x12\x46\n\tpage_rank\x18\x0c \x01(\x0b\x32\'.org.graphframes.connect.proto.PageRankH\x00R\x08pageRank\x12\x84\x01\n\x1fparallel_personalized_page_rank\x18\r \x01(\x0b\x32;.org.graphframes.connect.proto.ParallelPersonalizedPageRankH\x00R\x1cparallelPersonalizedPageRank\x12w\n\x1apower_iteration_clustering\x18\x0e \x01(\x0b\x32\x37.org.graphframes.connect.proto.PowerIterationClusteringH\x00R\x18powerIterationClustering\x12?\n\x06pregel\x18\x0f \x01(\x0b\x32%.org.graphframes.connect.proto.PregelH\x00R\x06pregel\x12U\n\x0eshortest_paths\x18\x10 \x01(\x0b\x32,.org.graphframes.connect.proto.ShortestPathsH\x00R\rshortestPaths\x12\x80\x01\n\x1dstrongly_connected_components\x18\x11 \x01(\x0b\x32:.org.graphframes.connect.proto.StronglyConnectedComponentsH\x00R\x1bstronglyConnectedComponents\x12P\n\rsvd_plus_plus\x18\x12 \x01(\x0b\x32*.org.graphframes.connect.proto.SVDPlusPlusH\x00R\x0bsvdPlusPlus\x12U\n\x0etriangle_count\x18\x13 \x01(\x0b\x32,.org.graphframes.connect.proto.TriangleCountH\x00R\rtriangleCount\x12\x45\n\x08triplets\x18\x14 \x01(\x0b\x32\'.org.graphframes.connect.proto.TripletsH\x00R\x08triplets\x12H\n\x03mis\x18\x16 \x01(\x0b\x32\x34.org.graphframes.connect.proto.MaximalIndependentSetH\x00R\x03misB\x08\n\x06method"\xd7\x02\n\x0cStorageLevel\x12\x1d\n\tdisk_only\x18\x01 \x01(\x08H\x00R\x08\x64iskOnly\x12 \n\x0b\x64isk_only_2\x18\x02 \x01(\x08H\x00R\tdiskOnly2\x12 \n\x0b\x64isk_only_3\x18\x03 \x01(\x08H\x00R\tdiskOnly3\x12(\n\x0fmemory_and_disk\x18\x04 \x01(\x08H\x00R\rmemoryAndDisk\x12+\n\x11memory_and_disk_2\x18\x05 \x01(\x08H\x00R\x0ememoryAndDisk2\x12\x33\n\x15memory_and_disk_deser\x18\x06 \x01(\x08H\x00R\x12memoryAndDiskDeser\x12!\n\x0bmemory_only\x18\x07 \x01(\x08H\x00R\nmemoryOnly\x12$\n\rmemory_only_2\x18\x08 \x01(\x08H\x00R\x0bmemoryOnly2B\x0f\n\rstorage_level"M\n\x12\x43olumnOrExpression\x12\x12\n\x03\x63ol\x18\x01 \x01(\x0cH\x00R\x03\x63ol\x12\x14\n\x04\x65xpr\x18\x02 \x01(\tH\x00R\x04\x65xprB\r\n\x0b\x63ol_or_expr"P\n\x0eStringOrLongID\x12\x19\n\x07long_id\x18\x01 \x01(\x03H\x00R\x06longId\x12\x1d\n\tstring_id\x18\x02 \x01(\tH\x00R\x08stringIdB\x04\n\x02id"\xee\x02\n\x11\x41ggregateMessages\x12J\n\x07\x61gg_col\x18\x01 \x03(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionR\x06\x61ggCol\x12Q\n\x0bsend_to_src\x18\x02 \x03(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionR\tsendToSrc\x12Q\n\x0bsend_to_dst\x18\x03 \x03(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionR\tsendToDst\x12U\n\rstorage_level\x18\x04 \x01(\x0b\x32+.org.graphframes.connect.proto.StorageLevelH\x00R\x0cstorageLevel\x88\x01\x01\x42\x10\n\x0e_storage_level"\x9d\x02\n\x03\x42\x46S\x12N\n\tfrom_expr\x18\x01 \x01(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionR\x08\x66romExpr\x12J\n\x07to_expr\x18\x02 \x01(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionR\x06toExpr\x12R\n\x0b\x65\x64ge_filter\x18\x03 \x01(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionR\nedgeFilter\x12&\n\x0fmax_path_length\x18\x04 \x01(\x05R\rmaxPathLength"\x86\x03\n\x13\x43onnectedComponents\x12\x1c\n\talgorithm\x18\x01 \x01(\tR\talgorithm\x12/\n\x13\x63heckpoint_interval\x18\x02 \x01(\x05R\x12\x63heckpointInterval\x12/\n\x13\x62roadcast_threshold\x18\x03 \x01(\x05R\x12\x62roadcastThreshold\x12\x37\n\x18use_labels_as_components\x18\x04 \x01(\x08R\x15useLabelsAsComponents\x12\x32\n\x15use_local_checkpoints\x18\x05 \x01(\x08R\x13useLocalCheckpoints\x12\x19\n\x08max_iter\x18\x06 \x01(\x05R\x07maxIter\x12U\n\rstorage_level\x18\x07 \x01(\x0b\x32+.org.graphframes.connect.proto.StorageLevelH\x00R\x0cstorageLevel\x88\x01\x01\x42\x10\n\x0e_storage_level"\xdf\x01\n\x0f\x44\x65tectingCycles\x12\x32\n\x15use_local_checkpoints\x18\x01 \x01(\x08R\x13useLocalCheckpoints\x12/\n\x13\x63heckpoint_interval\x18\x02 \x01(\x05R\x12\x63heckpointInterval\x12U\n\rstorage_level\x18\x03 \x01(\x0b\x32+.org.graphframes.connect.proto.StorageLevelH\x00R\x0cstorageLevel\x88\x01\x01\x42\x10\n\x0e_storage_level"\x16\n\x14\x44ropIsolatedVertices"^\n\x0b\x46ilterEdges\x12O\n\tcondition\x18\x01 \x01(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionR\tcondition"a\n\x0e\x46ilterVertices\x12O\n\tcondition\x18\x02 \x01(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionR\tcondition" \n\x04\x46ind\x12\x18\n\x07pattern\x18\x01 \x01(\tR\x07pattern"\x99\x02\n\x10LabelPropagation\x12\x1c\n\talgorithm\x18\x01 \x01(\tR\talgorithm\x12\x19\n\x08max_iter\x18\x02 \x01(\x05R\x07maxIter\x12\x32\n\x15use_local_checkpoints\x18\x03 \x01(\x08R\x13useLocalCheckpoints\x12/\n\x13\x63heckpoint_interval\x18\x04 \x01(\x05R\x12\x63heckpointInterval\x12U\n\rstorage_level\x18\x05 \x01(\x0b\x32+.org.graphframes.connect.proto.StorageLevelH\x00R\x0cstorageLevel\x88\x01\x01\x42\x10\n\x0e_storage_level"\xe2\x01\n\x08PageRank\x12+\n\x11reset_probability\x18\x01 \x01(\x01R\x10resetProbability\x12O\n\tsource_id\x18\x02 \x01(\x0b\x32-.org.graphframes.connect.proto.StringOrLongIDH\x00R\x08sourceId\x88\x01\x01\x12\x1e\n\x08max_iter\x18\x03 \x01(\x05H\x01R\x07maxIter\x88\x01\x01\x12\x15\n\x03tol\x18\x04 \x01(\x01H\x02R\x03tol\x88\x01\x01\x42\x0c\n\n_source_idB\x0b\n\t_max_iterB\x06\n\x04_tol"\xb4\x01\n\x1cParallelPersonalizedPageRank\x12+\n\x11reset_probability\x18\x01 \x01(\x01R\x10resetProbability\x12L\n\nsource_ids\x18\x02 \x03(\x0b\x32-.org.graphframes.connect.proto.StringOrLongIDR\tsourceIds\x12\x19\n\x08max_iter\x18\x03 \x01(\x05R\x07maxIter"v\n\x18PowerIterationClustering\x12\x0c\n\x01k\x18\x01 \x01(\x05R\x01k\x12\x19\n\x08max_iter\x18\x02 \x01(\x05R\x07maxIter\x12"\n\nweight_col\x18\x03 \x01(\tH\x00R\tweightCol\x88\x01\x01\x42\r\n\x0b_weight_col"\xe6\t\n\x06Pregel\x12L\n\x08\x61gg_msgs\x18\x01 \x01(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionR\x07\x61ggMsgs\x12X\n\x0fsend_msg_to_dst\x18\x02 \x03(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionR\x0csendMsgToDst\x12X\n\x0fsend_msg_to_src\x18\x03 \x03(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionR\x0csendMsgToSrc\x12/\n\x13\x63heckpoint_interval\x18\x04 \x01(\x05R\x12\x63heckpointInterval\x12\x19\n\x08max_iter\x18\x05 \x01(\x05R\x07maxIter\x12.\n\x13\x61\x64\x64itional_col_name\x18\x06 \x01(\tR\x11\x61\x64\x64itionalColName\x12g\n\x16\x61\x64\x64itional_col_initial\x18\x07 \x01(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionR\x14\x61\x64\x64itionalColInitial\x12_\n\x12\x61\x64\x64itional_col_upd\x18\x08 \x01(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionR\x10\x61\x64\x64itionalColUpd\x12*\n\x0e\x65\x61rly_stopping\x18\t \x01(\x08H\x00R\rearlyStopping\x88\x01\x01\x12\x32\n\x15use_local_checkpoints\x18\n \x01(\x08R\x13useLocalCheckpoints\x12U\n\rstorage_level\x18\x0b \x01(\x0b\x32+.org.graphframes.connect.proto.StorageLevelH\x01R\x0cstorageLevel\x88\x01\x01\x12\x37\n\x16stop_if_all_non_active\x18\x0c \x01(\x08H\x02R\x12stopIfAllNonActive\x88\x01\x01\x12\x66\n\x13initial_active_expr\x18\r \x01(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionH\x03R\x11initialActiveExpr\x88\x01\x01\x12\x64\n\x12update_active_expr\x18\x0e \x01(\x0b\x32\x31.org.graphframes.connect.proto.ColumnOrExpressionH\x04R\x10updateActiveExpr\x88\x01\x01\x12\x45\n\x1dskip_messages_from_non_active\x18\x0f \x01(\x08H\x05R\x19skipMessagesFromNonActive\x88\x01\x01\x42\x11\n\x0f_early_stoppingB\x10\n\x0e_storage_levelB\x19\n\x17_stop_if_all_non_activeB\x16\n\x14_initial_active_exprB\x15\n\x13_update_active_exprB \n\x1e_skip_messages_from_non_active"\xc8\x02\n\rShortestPaths\x12K\n\tlandmarks\x18\x01 \x03(\x0b\x32-.org.graphframes.connect.proto.StringOrLongIDR\tlandmarks\x12\x1c\n\talgorithm\x18\x02 \x01(\tR\talgorithm\x12\x32\n\x15use_local_checkpoints\x18\x03 \x01(\x08R\x13useLocalCheckpoints\x12/\n\x13\x63heckpoint_interval\x18\x04 \x01(\x05R\x12\x63heckpointInterval\x12U\n\rstorage_level\x18\x05 \x01(\x0b\x32+.org.graphframes.connect.proto.StorageLevelH\x00R\x0cstorageLevel\x88\x01\x01\x42\x10\n\x0e_storage_level"8\n\x1bStronglyConnectedComponents\x12\x19\n\x08max_iter\x18\x01 \x01(\x05R\x07maxIter"\xd6\x01\n\x0bSVDPlusPlus\x12\x12\n\x04rank\x18\x01 \x01(\x05R\x04rank\x12\x19\n\x08max_iter\x18\x02 \x01(\x05R\x07maxIter\x12\x1b\n\tmin_value\x18\x03 \x01(\x01R\x08minValue\x12\x1b\n\tmax_value\x18\x04 \x01(\x01R\x08maxValue\x12\x16\n\x06gamma1\x18\x05 \x01(\x01R\x06gamma1\x12\x16\n\x06gamma2\x18\x06 \x01(\x01R\x06gamma2\x12\x16\n\x06gamma6\x18\x07 \x01(\x01R\x06gamma6\x12\x16\n\x06gamma7\x18\x08 \x01(\x01R\x06gamma7"x\n\rTriangleCount\x12U\n\rstorage_level\x18\x01 \x01(\x0b\x32+.org.graphframes.connect.proto.StorageLevelH\x00R\x0cstorageLevel\x88\x01\x01\x42\x10\n\x0e_storage_level"\n\n\x08Triplets"\xf9\x01\n\x15MaximalIndependentSet\x12/\n\x13\x63heckpoint_interval\x18\x01 \x01(\x05R\x12\x63heckpointInterval\x12U\n\rstorage_level\x18\x02 \x01(\x0b\x32+.org.graphframes.connect.proto.StorageLevelH\x00R\x0cstorageLevel\x88\x01\x01\x12\x32\n\x15use_local_checkpoints\x18\x03 \x01(\x08R\x13useLocalCheckpoints\x12\x12\n\x04seed\x18\x04 \x01(\x03R\x04seedB\x10\n\x0e_storage_levelB\xd2\x01\n!com.org.graphframes.connect.protoB\x10GraphframesProtoH\x01P\x01\xa0\x01\x01\xa2\x02\x04OGCP\xaa\x02\x1dOrg.Graphframes.Connect.Proto\xca\x02\x1dOrg\\Graphframes\\Connect\\Proto\xe2\x02)Org\\Graphframes\\Connect\\Proto\\GPBMetadata\xea\x02 Org::Graphframes::Connect::Protob\x06proto3' ) _globals = globals() @@ -31,47 +31,49 @@ "DESCRIPTOR" ]._serialized_options = b"\n!com.org.graphframes.connect.protoB\020GraphframesProtoH\001P\001\240\001\001\242\002\004OGCP\252\002\035Org.Graphframes.Connect.Proto\312\002\035Org\\Graphframes\\Connect\\Proto\342\002)Org\\Graphframes\\Connect\\Proto\\GPBMetadata\352\002 Org::Graphframes::Connect::Proto" _globals["_GRAPHFRAMESAPI"]._serialized_start = 53 - _globals["_GRAPHFRAMESAPI"]._serialized_end = 1768 - _globals["_STORAGELEVEL"]._serialized_start = 1771 - _globals["_STORAGELEVEL"]._serialized_end = 2114 - _globals["_COLUMNOREXPRESSION"]._serialized_start = 2116 - _globals["_COLUMNOREXPRESSION"]._serialized_end = 2193 - _globals["_STRINGORLONGID"]._serialized_start = 2195 - _globals["_STRINGORLONGID"]._serialized_end = 2275 - _globals["_AGGREGATEMESSAGES"]._serialized_start = 2278 - _globals["_AGGREGATEMESSAGES"]._serialized_end = 2644 - _globals["_BFS"]._serialized_start = 2647 - _globals["_BFS"]._serialized_end = 2932 - _globals["_CONNECTEDCOMPONENTS"]._serialized_start = 2935 - _globals["_CONNECTEDCOMPONENTS"]._serialized_end = 3325 - _globals["_DETECTINGCYCLES"]._serialized_start = 3328 - _globals["_DETECTINGCYCLES"]._serialized_end = 3551 - _globals["_DROPISOLATEDVERTICES"]._serialized_start = 3553 - _globals["_DROPISOLATEDVERTICES"]._serialized_end = 3575 - _globals["_FILTEREDGES"]._serialized_start = 3577 - _globals["_FILTEREDGES"]._serialized_end = 3671 - _globals["_FILTERVERTICES"]._serialized_start = 3673 - _globals["_FILTERVERTICES"]._serialized_end = 3770 - _globals["_FIND"]._serialized_start = 3772 - _globals["_FIND"]._serialized_end = 3804 - _globals["_LABELPROPAGATION"]._serialized_start = 3807 - _globals["_LABELPROPAGATION"]._serialized_end = 4088 - _globals["_PAGERANK"]._serialized_start = 4091 - _globals["_PAGERANK"]._serialized_end = 4317 - _globals["_PARALLELPERSONALIZEDPAGERANK"]._serialized_start = 4320 - _globals["_PARALLELPERSONALIZEDPAGERANK"]._serialized_end = 4500 - _globals["_POWERITERATIONCLUSTERING"]._serialized_start = 4502 - _globals["_POWERITERATIONCLUSTERING"]._serialized_end = 4620 - _globals["_PREGEL"]._serialized_start = 4623 - _globals["_PREGEL"]._serialized_end = 5877 - _globals["_SHORTESTPATHS"]._serialized_start = 5880 - _globals["_SHORTESTPATHS"]._serialized_end = 6208 - _globals["_STRONGLYCONNECTEDCOMPONENTS"]._serialized_start = 6210 - _globals["_STRONGLYCONNECTEDCOMPONENTS"]._serialized_end = 6266 - _globals["_SVDPLUSPLUS"]._serialized_start = 6269 - _globals["_SVDPLUSPLUS"]._serialized_end = 6483 - _globals["_TRIANGLECOUNT"]._serialized_start = 6485 - _globals["_TRIANGLECOUNT"]._serialized_end = 6605 - _globals["_TRIPLETS"]._serialized_start = 6607 - _globals["_TRIPLETS"]._serialized_end = 6617 + _globals["_GRAPHFRAMESAPI"]._serialized_end = 1842 + _globals["_STORAGELEVEL"]._serialized_start = 1845 + _globals["_STORAGELEVEL"]._serialized_end = 2188 + _globals["_COLUMNOREXPRESSION"]._serialized_start = 2190 + _globals["_COLUMNOREXPRESSION"]._serialized_end = 2267 + _globals["_STRINGORLONGID"]._serialized_start = 2269 + _globals["_STRINGORLONGID"]._serialized_end = 2349 + _globals["_AGGREGATEMESSAGES"]._serialized_start = 2352 + _globals["_AGGREGATEMESSAGES"]._serialized_end = 2718 + _globals["_BFS"]._serialized_start = 2721 + _globals["_BFS"]._serialized_end = 3006 + _globals["_CONNECTEDCOMPONENTS"]._serialized_start = 3009 + _globals["_CONNECTEDCOMPONENTS"]._serialized_end = 3399 + _globals["_DETECTINGCYCLES"]._serialized_start = 3402 + _globals["_DETECTINGCYCLES"]._serialized_end = 3625 + _globals["_DROPISOLATEDVERTICES"]._serialized_start = 3627 + _globals["_DROPISOLATEDVERTICES"]._serialized_end = 3649 + _globals["_FILTEREDGES"]._serialized_start = 3651 + _globals["_FILTEREDGES"]._serialized_end = 3745 + _globals["_FILTERVERTICES"]._serialized_start = 3747 + _globals["_FILTERVERTICES"]._serialized_end = 3844 + _globals["_FIND"]._serialized_start = 3846 + _globals["_FIND"]._serialized_end = 3878 + _globals["_LABELPROPAGATION"]._serialized_start = 3881 + _globals["_LABELPROPAGATION"]._serialized_end = 4162 + _globals["_PAGERANK"]._serialized_start = 4165 + _globals["_PAGERANK"]._serialized_end = 4391 + _globals["_PARALLELPERSONALIZEDPAGERANK"]._serialized_start = 4394 + _globals["_PARALLELPERSONALIZEDPAGERANK"]._serialized_end = 4574 + _globals["_POWERITERATIONCLUSTERING"]._serialized_start = 4576 + _globals["_POWERITERATIONCLUSTERING"]._serialized_end = 4694 + _globals["_PREGEL"]._serialized_start = 4697 + _globals["_PREGEL"]._serialized_end = 5951 + _globals["_SHORTESTPATHS"]._serialized_start = 5954 + _globals["_SHORTESTPATHS"]._serialized_end = 6282 + _globals["_STRONGLYCONNECTEDCOMPONENTS"]._serialized_start = 6284 + _globals["_STRONGLYCONNECTEDCOMPONENTS"]._serialized_end = 6340 + _globals["_SVDPLUSPLUS"]._serialized_start = 6343 + _globals["_SVDPLUSPLUS"]._serialized_end = 6557 + _globals["_TRIANGLECOUNT"]._serialized_start = 6559 + _globals["_TRIANGLECOUNT"]._serialized_end = 6679 + _globals["_TRIPLETS"]._serialized_start = 6681 + _globals["_TRIPLETS"]._serialized_end = 6691 + _globals["_MAXIMALINDEPENDENTSET"]._serialized_start = 6694 + _globals["_MAXIMALINDEPENDENTSET"]._serialized_end = 6943 # @@protoc_insertion_point(module_scope) diff --git a/python/graphframes/connect/proto/graphframes_pb2.pyi b/python/graphframes/connect/proto/graphframes_pb2.pyi index ffe59932d..6544824da 100644 --- a/python/graphframes/connect/proto/graphframes_pb2.pyi +++ b/python/graphframes/connect/proto/graphframes_pb2.pyi @@ -11,28 +11,7 @@ from google.protobuf.internal import containers as _containers DESCRIPTOR: _descriptor.FileDescriptor class GraphFramesAPI(_message.Message): - __slots__ = ( - "vertices", - "edges", - "aggregate_messages", - "bfs", - "connected_components", - "drop_isolated_vertices", - "detecting_cycles", - "filter_edges", - "filter_vertices", - "find", - "label_propagation", - "page_rank", - "parallel_personalized_page_rank", - "power_iteration_clustering", - "pregel", - "shortest_paths", - "strongly_connected_components", - "svd_plus_plus", - "triangle_count", - "triplets", - ) + __slots__ = () VERTICES_FIELD_NUMBER: _ClassVar[int] EDGES_FIELD_NUMBER: _ClassVar[int] AGGREGATE_MESSAGES_FIELD_NUMBER: _ClassVar[int] @@ -53,6 +32,7 @@ class GraphFramesAPI(_message.Message): SVD_PLUS_PLUS_FIELD_NUMBER: _ClassVar[int] TRIANGLE_COUNT_FIELD_NUMBER: _ClassVar[int] TRIPLETS_FIELD_NUMBER: _ClassVar[int] + MIS_FIELD_NUMBER: _ClassVar[int] vertices: bytes edges: bytes aggregate_messages: AggregateMessages @@ -73,6 +53,7 @@ class GraphFramesAPI(_message.Message): svd_plus_plus: SVDPlusPlus triangle_count: TriangleCount triplets: Triplets + mis: MaximalIndependentSet def __init__( self, vertices: _Optional[bytes] = ..., @@ -99,19 +80,11 @@ class GraphFramesAPI(_message.Message): svd_plus_plus: _Optional[_Union[SVDPlusPlus, _Mapping]] = ..., triangle_count: _Optional[_Union[TriangleCount, _Mapping]] = ..., triplets: _Optional[_Union[Triplets, _Mapping]] = ..., + mis: _Optional[_Union[MaximalIndependentSet, _Mapping]] = ..., ) -> None: ... class StorageLevel(_message.Message): - __slots__ = ( - "disk_only", - "disk_only_2", - "disk_only_3", - "memory_and_disk", - "memory_and_disk_2", - "memory_and_disk_deser", - "memory_only", - "memory_only_2", - ) + __slots__ = () DISK_ONLY_FIELD_NUMBER: _ClassVar[int] DISK_ONLY_2_FIELD_NUMBER: _ClassVar[int] DISK_ONLY_3_FIELD_NUMBER: _ClassVar[int] @@ -141,7 +114,7 @@ class StorageLevel(_message.Message): ) -> None: ... class ColumnOrExpression(_message.Message): - __slots__ = ("col", "expr") + __slots__ = () COL_FIELD_NUMBER: _ClassVar[int] EXPR_FIELD_NUMBER: _ClassVar[int] col: bytes @@ -149,7 +122,7 @@ class ColumnOrExpression(_message.Message): def __init__(self, col: _Optional[bytes] = ..., expr: _Optional[str] = ...) -> None: ... class StringOrLongID(_message.Message): - __slots__ = ("long_id", "string_id") + __slots__ = () LONG_ID_FIELD_NUMBER: _ClassVar[int] STRING_ID_FIELD_NUMBER: _ClassVar[int] long_id: int @@ -157,7 +130,7 @@ class StringOrLongID(_message.Message): def __init__(self, long_id: _Optional[int] = ..., string_id: _Optional[str] = ...) -> None: ... class AggregateMessages(_message.Message): - __slots__ = ("agg_col", "send_to_src", "send_to_dst", "storage_level") + __slots__ = () AGG_COL_FIELD_NUMBER: _ClassVar[int] SEND_TO_SRC_FIELD_NUMBER: _ClassVar[int] SEND_TO_DST_FIELD_NUMBER: _ClassVar[int] @@ -175,7 +148,7 @@ class AggregateMessages(_message.Message): ) -> None: ... class BFS(_message.Message): - __slots__ = ("from_expr", "to_expr", "edge_filter", "max_path_length") + __slots__ = () FROM_EXPR_FIELD_NUMBER: _ClassVar[int] TO_EXPR_FIELD_NUMBER: _ClassVar[int] EDGE_FILTER_FIELD_NUMBER: _ClassVar[int] @@ -193,15 +166,7 @@ class BFS(_message.Message): ) -> None: ... class ConnectedComponents(_message.Message): - __slots__ = ( - "algorithm", - "checkpoint_interval", - "broadcast_threshold", - "use_labels_as_components", - "use_local_checkpoints", - "max_iter", - "storage_level", - ) + __slots__ = () ALGORITHM_FIELD_NUMBER: _ClassVar[int] CHECKPOINT_INTERVAL_FIELD_NUMBER: _ClassVar[int] BROADCAST_THRESHOLD_FIELD_NUMBER: _ClassVar[int] @@ -228,7 +193,7 @@ class ConnectedComponents(_message.Message): ) -> None: ... class DetectingCycles(_message.Message): - __slots__ = ("use_local_checkpoints", "checkpoint_interval", "storage_level") + __slots__ = () USE_LOCAL_CHECKPOINTS_FIELD_NUMBER: _ClassVar[int] CHECKPOINT_INTERVAL_FIELD_NUMBER: _ClassVar[int] STORAGE_LEVEL_FIELD_NUMBER: _ClassVar[int] @@ -247,7 +212,7 @@ class DropIsolatedVertices(_message.Message): def __init__(self) -> None: ... class FilterEdges(_message.Message): - __slots__ = ("condition",) + __slots__ = () CONDITION_FIELD_NUMBER: _ClassVar[int] condition: ColumnOrExpression def __init__( @@ -255,7 +220,7 @@ class FilterEdges(_message.Message): ) -> None: ... class FilterVertices(_message.Message): - __slots__ = ("condition",) + __slots__ = () CONDITION_FIELD_NUMBER: _ClassVar[int] condition: ColumnOrExpression def __init__( @@ -263,19 +228,13 @@ class FilterVertices(_message.Message): ) -> None: ... class Find(_message.Message): - __slots__ = ("pattern",) + __slots__ = () PATTERN_FIELD_NUMBER: _ClassVar[int] pattern: str def __init__(self, pattern: _Optional[str] = ...) -> None: ... class LabelPropagation(_message.Message): - __slots__ = ( - "algorithm", - "max_iter", - "use_local_checkpoints", - "checkpoint_interval", - "storage_level", - ) + __slots__ = () ALGORITHM_FIELD_NUMBER: _ClassVar[int] MAX_ITER_FIELD_NUMBER: _ClassVar[int] USE_LOCAL_CHECKPOINTS_FIELD_NUMBER: _ClassVar[int] @@ -296,7 +255,7 @@ class LabelPropagation(_message.Message): ) -> None: ... class PageRank(_message.Message): - __slots__ = ("reset_probability", "source_id", "max_iter", "tol") + __slots__ = () RESET_PROBABILITY_FIELD_NUMBER: _ClassVar[int] SOURCE_ID_FIELD_NUMBER: _ClassVar[int] MAX_ITER_FIELD_NUMBER: _ClassVar[int] @@ -314,7 +273,7 @@ class PageRank(_message.Message): ) -> None: ... class ParallelPersonalizedPageRank(_message.Message): - __slots__ = ("reset_probability", "source_ids", "max_iter") + __slots__ = () RESET_PROBABILITY_FIELD_NUMBER: _ClassVar[int] SOURCE_IDS_FIELD_NUMBER: _ClassVar[int] MAX_ITER_FIELD_NUMBER: _ClassVar[int] @@ -329,7 +288,7 @@ class ParallelPersonalizedPageRank(_message.Message): ) -> None: ... class PowerIterationClustering(_message.Message): - __slots__ = ("k", "max_iter", "weight_col") + __slots__ = () K_FIELD_NUMBER: _ClassVar[int] MAX_ITER_FIELD_NUMBER: _ClassVar[int] WEIGHT_COL_FIELD_NUMBER: _ClassVar[int] @@ -344,23 +303,7 @@ class PowerIterationClustering(_message.Message): ) -> None: ... class Pregel(_message.Message): - __slots__ = ( - "agg_msgs", - "send_msg_to_dst", - "send_msg_to_src", - "checkpoint_interval", - "max_iter", - "additional_col_name", - "additional_col_initial", - "additional_col_upd", - "early_stopping", - "use_local_checkpoints", - "storage_level", - "stop_if_all_non_active", - "initial_active_expr", - "update_active_expr", - "skip_messages_from_non_active", - ) + __slots__ = () AGG_MSGS_FIELD_NUMBER: _ClassVar[int] SEND_MSG_TO_DST_FIELD_NUMBER: _ClassVar[int] SEND_MSG_TO_SRC_FIELD_NUMBER: _ClassVar[int] @@ -411,13 +354,7 @@ class Pregel(_message.Message): ) -> None: ... class ShortestPaths(_message.Message): - __slots__ = ( - "landmarks", - "algorithm", - "use_local_checkpoints", - "checkpoint_interval", - "storage_level", - ) + __slots__ = () LANDMARKS_FIELD_NUMBER: _ClassVar[int] ALGORITHM_FIELD_NUMBER: _ClassVar[int] USE_LOCAL_CHECKPOINTS_FIELD_NUMBER: _ClassVar[int] @@ -438,22 +375,13 @@ class ShortestPaths(_message.Message): ) -> None: ... class StronglyConnectedComponents(_message.Message): - __slots__ = ("max_iter",) + __slots__ = () MAX_ITER_FIELD_NUMBER: _ClassVar[int] max_iter: int def __init__(self, max_iter: _Optional[int] = ...) -> None: ... class SVDPlusPlus(_message.Message): - __slots__ = ( - "rank", - "max_iter", - "min_value", - "max_value", - "gamma1", - "gamma2", - "gamma6", - "gamma7", - ) + __slots__ = () RANK_FIELD_NUMBER: _ClassVar[int] MAX_ITER_FIELD_NUMBER: _ClassVar[int] MIN_VALUE_FIELD_NUMBER: _ClassVar[int] @@ -483,7 +411,7 @@ class SVDPlusPlus(_message.Message): ) -> None: ... class TriangleCount(_message.Message): - __slots__ = ("storage_level",) + __slots__ = () STORAGE_LEVEL_FIELD_NUMBER: _ClassVar[int] storage_level: StorageLevel def __init__(self, storage_level: _Optional[_Union[StorageLevel, _Mapping]] = ...) -> None: ... @@ -491,3 +419,21 @@ class TriangleCount(_message.Message): class Triplets(_message.Message): __slots__ = () def __init__(self) -> None: ... + +class MaximalIndependentSet(_message.Message): + __slots__ = () + CHECKPOINT_INTERVAL_FIELD_NUMBER: _ClassVar[int] + STORAGE_LEVEL_FIELD_NUMBER: _ClassVar[int] + USE_LOCAL_CHECKPOINTS_FIELD_NUMBER: _ClassVar[int] + SEED_FIELD_NUMBER: _ClassVar[int] + checkpoint_interval: int + storage_level: StorageLevel + use_local_checkpoints: bool + seed: int + def __init__( + self, + checkpoint_interval: _Optional[int] = ..., + storage_level: _Optional[_Union[StorageLevel, _Mapping]] = ..., + use_local_checkpoints: _Optional[bool] = ..., + seed: _Optional[int] = ..., + ) -> None: ... diff --git a/python/graphframes/graphframe.py b/python/graphframes/graphframe.py index d2cae7b9e..18d65886c 100644 --- a/python/graphframes/graphframe.py +++ b/python/graphframes/graphframe.py @@ -456,6 +456,48 @@ def connectedComponents( storage_level=storage_level, ) + def maximal_independent_set( + self, + seed: int = 42, + checkpoint_interval: int = 2, + use_local_checkpoints: bool = False, + storage_level: StorageLevel = StorageLevel.MEMORY_AND_DISK_DESER, + ) -> DataFrame: + """ + This method implements a distributed algorithm for finding a Maximal Independent Set (MIS) + in a graph. + + An MIS is a set of vertices such that no two vertices in the set are adjacent (i.e., there + is no edge between any two vertices in the set), and the set is maximal, meaning that adding + any other vertex to the set would violate the independence property. Note that this + implementation finds a maximal (but not necessarily maximum) independent set; that is, it + ensures no more vertices can be added to the set, but does not guarantee that the set has + the largest possible number of vertices among all possible independent sets in the graph. + + The algorithm implemented here is based on the paper: Ghaffari, Mohsen. "An improved + distributed algorithm for maximal independent set." Proceedings of the twenty-seventh annual + ACM-SIAM symposium on Discrete algorithms. Society for Industrial and Applied Mathematics, + 2016. + + Note: This is a randomized, non-deterministic algorithm. The result may vary between runs + even if a fixed random seed is provided because of how Apache Spark works. + + :param seed: random seed used for tie-breaking in the algorithm (default: 42) + :param checkpoint_interval: checkpoint interval in terms of number of iterations (default: 2) + :param use_local_checkpoints: whether to use local checkpoints (default: False); + local checkpoints are faster and do not require setting + a persistent checkpoint directory; however, they are less + reliable and require executors to have sufficient local disk space. + :param storage_level: storage level for both intermediate and final DataFrames + (default: MEMORY_AND_DISK_DESER) + + :return: DataFrame with new vertex column "selected", where "true" indicates the vertex + is part of the Maximal Independent Set + """ # noqa: E501 + return self._impl.maximal_independent_set( + checkpoint_interval, storage_level, use_local_checkpoints, seed + ) + def labelPropagation( self, maxIter: int, diff --git a/python/tests/test_graphframes.py b/python/tests/test_graphframes.py index f5aa0eace..73eadf856 100644 --- a/python/tests/test_graphframes.py +++ b/python/tests/test_graphframes.py @@ -487,6 +487,28 @@ def test_cycles_finding(spark: SparkSession, args: PregelArguments) -> None: _ = res.unpersist() +@pytest.mark.parametrize("storage_level", STORAGE_LEVELS, ids=STORAGE_LEVELS_IDS) +def test_mis(spark: SparkSession, storage_level: StorageLevel) -> None: + # Create a graph with isolated vertices + vertices = spark.createDataFrame( + [(0, "a"), (1, "b"), (2, "c"), (3, "d")], ["id", "name"] + ) + + # Only connect vertices 0 and 1 + edges = spark.createDataFrame([(0, 1, "edge1")], ["src", "dst", "name"]) + + graph = GraphFrame(vertices, edges) + mis = graph.maximal_independent_set(storage_level=storage_level, seed=12345) + + # Check that all vertices are in the MIS (since 2 and 3 are isolated) + mis_ids = set(row[0] for row in mis.select("id").collect()) + assert len(mis_ids) == 3, "MIS should contain 2 isolated vertices and one of linked" + assert 2 in mis_ids, "Isolated vertex 2 should be in MIS" + assert 3 in mis_ids, "Isolated vertex 3 should be in MIS" + + _ = mis.unpersist() + + @pytest.mark.skipif(is_remote(), reason="DISABLE FOR CONNECT") def test_svd_plus_plus(examples, spark: SparkSession): g = _from_java_gf(getattr(examples, "ALSSyntheticData")(), spark) From 0bf76f27c35fb1c4ae81664da944c860ea7ab256 Mon Sep 17 00:00:00 2001 From: semyonsinchenko Date: Fri, 24 Oct 2025 14:30:06 +0200 Subject: [PATCH 2/2] fix connect --- python/graphframes/connect/graphframes_client.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/python/graphframes/connect/graphframes_client.py b/python/graphframes/connect/graphframes_client.py index 2ea77cee8..0a7078f45 100644 --- a/python/graphframes/connect/graphframes_client.py +++ b/python/graphframes/connect/graphframes_client.py @@ -1087,7 +1087,7 @@ def __init__( self.v = v self.e = e self.checkpoint_interval = checkpoint_interval - self.storage_level = (storage_level,) + self.storage_level = storage_level self.use_local_checkpoints = use_local_checkpoints self.seed = seed