Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
Show all changes
32 commits
Select commit Hold shift + click to select a range
d1972fc
Merge remote-tracking branch 'refs/remotes/origin/master'
SemyonSinchenko Feb 4, 2025
c158815
wip
SemyonSinchenko Feb 7, 2025
ea11df6
wip
SemyonSinchenko Feb 8, 2025
0ddd5bd
wip
SemyonSinchenko Feb 8, 2025
da7eccc
wip
SemyonSinchenko Feb 8, 2025
fb784a3
The first working version
SemyonSinchenko Feb 8, 2025
9f8905f
Merge remote-tracking branch 'refs/remotes/graphframes/master'
SemyonSinchenko Feb 19, 2025
d58ed2a
Merge remote-tracking branch 'refs/remotes/graphframes/master'
SemyonSinchenko Feb 20, 2025
20f7575
WIP
SemyonSinchenko Feb 23, 2025
eee7b7b
Working version?
SemyonSinchenko Feb 23, 2025
130b12e
Merge remote-tracking branch 'refs/remotes/graphframes/master'
SemyonSinchenko Feb 23, 2025
7e325aa
Fix tests
SemyonSinchenko Feb 23, 2025
c47a57e
Fix tests
SemyonSinchenko Feb 23, 2025
fc8ebae
Fix CI typo
SemyonSinchenko Feb 23, 2025
f13c754
Fix typo in CI
SemyonSinchenko Feb 23, 2025
a21b5aa
Fix wget's verbose + GHA bug
SemyonSinchenko Feb 23, 2025
f4c91d6
Stop connect server
SemyonSinchenko Feb 23, 2025
e4d75f7
An attempt to fix a bug in GHA with a non-stopping tests
SemyonSinchenko Feb 23, 2025
0950cfd
Maybe https://github.com/grpc/grpc/issues/38290?
SemyonSinchenko Feb 23, 2025
8cb430c
Fix broken stop-cript
SemyonSinchenko Feb 23, 2025
19c1934
Ignore errors in clean-up
SemyonSinchenko Feb 24, 2025
1eef323
Verbosity in ci tests
SemyonSinchenko Feb 24, 2025
7fd1f23
Merge main
SemyonSinchenko Mar 10, 2025
cc60bcb
Typo
SemyonSinchenko Mar 10, 2025
5528d65
Fix merge-artifacts
SemyonSinchenko Mar 10, 2025
f88e19a
Fix merge artifacts
SemyonSinchenko Mar 10, 2025
59897fb
Apply pre-commit rules
SemyonSinchenko Mar 10, 2025
97054b0
Add the missing method
SemyonSinchenko Mar 10, 2025
90a326f
Restore accidently deleted part of CI
SemyonSinchenko Mar 10, 2025
c8bcf43
Typo
SemyonSinchenko Mar 11, 2025
9d7f714
Fixes from comments
SemyonSinchenko Mar 17, 2025
5a91659
Pin the pyspark version <4.0 and re-generate lock
SemyonSinchenko Mar 17, 2025
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Prev Previous commit
Next Next commit
Merge main
Classic tests are passed;
++some changes
  • Loading branch information
SemyonSinchenko committed Mar 10, 2025
commit 7fd1f231448c3bdfa11b3bc4eb3b74400e939222
4 changes: 2 additions & 2 deletions .github/workflows/python-ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -44,13 +44,13 @@ jobs:
- name: Test
working-directory: ./python
run: |
poetry run python -m unittest discover -s tests/ -v
poetry run python -m pytets

- name: Test SparkConnect
env:
SPARK_CONNECT_MODE_ENABLED: 1
working-directory: ./python
run: |
poetry run python dev/run_connect.py
Comment thread
rjurney marked this conversation as resolved.
poetry run python -m unittest discover -s tests/ -v
poetry run python -m pytest
poetry run python dev/stop_connect.py
20 changes: 5 additions & 15 deletions build.sbt
Original file line number Diff line number Diff line change
@@ -1,4 +1,6 @@
import ReleaseTransformations.*
import sbt.Credentials
import sbt.Keys.credentials

lazy val sparkVer = sys.props.getOrElse("spark.version", "3.5.4")
lazy val sparkBranch = sparkVer.substring(0, 3)
Expand Down Expand Up @@ -49,7 +51,8 @@ lazy val commonSetting = Seq(
"-XX:ReservedCodeCacheSize=384m",
"-XX:MaxMetaspaceSize=384m",
"--add-opens=java.base/sun.nio.ch=ALL-UNNAMED",
Comment thread
rjurney marked this conversation as resolved.
"--add-opens=java.base/java.lang=ALL-UNNAMED"))
"--add-opens=java.base/java.lang=ALL-UNNAMED"),
credentials += Credentials(Path.userHome / ".ivy2" / ".sbtcredentials"))

lazy val root = (project in file("."))
.settings(
Expand All @@ -60,7 +63,6 @@ lazy val root = (project in file("."))
// Global settings
Global / concurrentRestrictions := Seq(Tags.limitAll(1)),
autoAPIMappings := true,

coverageHighlighting := false,

// Release settings
Expand Down Expand Up @@ -91,16 +93,4 @@ lazy val connect = (project in file("graphframes-connect"))
Compile / PB.includePaths ++= Seq(file("src/main/protobuf")),
PB.protocVersion := "3.23.4", // Spark 3.5 branch
libraryDependencies ++= Seq(
"org.apache.spark" %% "spark-connect" % sparkVer % "provided" cross CrossVersion.for3Use2_13),

// Assembly and shading
assembly / test := {},
assembly / assemblyShadeRules := Seq(
ShadeRule.rename("com.google.protobuf.**" -> "org.sparkproject.connect.protobuf.@1").inAll),
assembly / assemblyMergeStrategy := {
case PathList("META-INF", xs @ _*) => MergeStrategy.discard
case x if x.endsWith("module-info.class") => MergeStrategy.discard
case x =>
val oldStrategy = (assembly / assemblyMergeStrategy).value
oldStrategy(x)
})
"org.apache.spark" %% "spark-connect" % sparkVer % "provided" cross CrossVersion.for3Use2_13))
42 changes: 20 additions & 22 deletions graphframes-connect/src/main/protobuf/graphframes.proto
Original file line number Diff line number Diff line change
Expand Up @@ -15,22 +15,20 @@ message GraphFramesAPI {
AggregateMessages aggregate_messages = 3;
BFS bfs = 4;
ConnectedComponents connected_components = 5;
Degrees degrees = 6;
DropIsolatedVertices drop_isolated_vertices = 7;
FilterEdges filter_edges = 8;
FilterVertices filter_vertices = 9;
Find find = 10;
InDegrees in_degrees = 11;
LabelPropagation label_propagation = 12;
OutDegrees out_degrees = 13;
PageRank page_rank = 14;
ParallelPersonalizedPageRank parallel_personalized_page_rank = 15;
Pregel pregel = 16;
ShortestPaths shortest_paths = 17;
StronglyConnectedComponents strongly_connected_components = 18;
SVDPlusPlus svd_plus_plus = 19;
TriangleCount triangle_count = 20;
Triplets triplets = 21;
DropIsolatedVertices drop_isolated_vertices = 6;
FilterEdges filter_edges = 7;
FilterVertices filter_vertices = 8;
Find find = 9;
LabelPropagation label_propagation = 10;
PageRank page_rank = 11;
ParallelPersonalizedPageRank parallel_personalized_page_rank = 12;
PowerIterationClustering power_iteration_clustering = 13;
Pregel pregel = 14;
ShortestPaths shortest_paths = 15;
StronglyConnectedComponents strongly_connected_components = 16;
SVDPlusPlus svd_plus_plus = 17;
TriangleCount triangle_count = 18;
Triplets triplets = 19;
}
}

Expand Down Expand Up @@ -67,8 +65,6 @@ message ConnectedComponents {
int32 broadcast_threshold = 3;
}

message Degrees {}

message DropIsolatedVertices {}

message FilterEdges {
Expand All @@ -83,14 +79,10 @@ message Find {
string pattern = 1;
}

message InDegrees {}

message LabelPropagation {
int32 max_iter = 1;
}

message OutDegrees {}

message PageRank {
double reset_probability = 1;
optional StringOrLongID source_id = 2;
Expand All @@ -104,6 +96,12 @@ message ParallelPersonalizedPageRank {
int32 max_iter = 3;
}

message PowerIterationClustering {
int32 k = 1;
int32 max_iter = 2;
optional string weight_col = 3;
}

message Pregel {
ColumnOrExpression agg_msgs = 1;
repeated ColumnOrExpression send_msg_to_dst = 2;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3,17 +3,14 @@
package org.apache.spark.sql.graphframes

import scala.jdk.CollectionConverters._

import org.graphframes.GraphFrame
import org.graphframes.{GraphFrame, GraphFramesUnreachableException}
import org.graphframes.connect.proto.{ColumnOrExpression, GraphFramesAPI, StringOrLongID}
import org.graphframes.connect.proto.ColumnOrExpression.ColOrExprCase
import org.graphframes.connect.proto.GraphFramesAPI.MethodCase
import org.graphframes.connect.proto.StringOrLongID.IdCase

import org.apache.spark.sql.{Column, DataFrame, Dataset}
import org.apache.spark.sql.connect.planner.SparkConnectPlanner
import org.apache.spark.sql.functions.{expr, lit}

import org.apache.spark.sql.functions.{col, expr, lit}
import com.google.protobuf.ByteString

object GraphFramesConnectUtils {
Expand Down Expand Up @@ -101,9 +98,6 @@ object GraphFramesConnectUtils {
.setBroadcastThreshold(cc.getBroadcastThreshold)
.run()
}
case MethodCase.DEGREES => {
graphFrame.degrees
}
case MethodCase.DROP_ISOLATED_VERTICES => {
graphFrame.dropIsolatedVertices().vertices
}
Expand All @@ -119,15 +113,9 @@ object GraphFramesConnectUtils {
case MethodCase.FIND => {
graphFrame.find(apiMessage.getFind.getPattern)
}
case MethodCase.IN_DEGREES => {
graphFrame.inDegrees
}
case MethodCase.LABEL_PROPAGATION => {
graphFrame.labelPropagation.maxIter(apiMessage.getLabelPropagation.getMaxIter).run()
}
case MethodCase.OUT_DEGREES => {
graphFrame.outDegrees
}
case MethodCase.PAGE_RANK => {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

There are two ways to use PageRank... one sets the reset probability and one sets the maxIter: https://graphframes.github.io/graphframes/docs/_site/user-guide.html#pagerank

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Take a look on the lines 123-127:

        if (pageRankProto.hasMaxIter) {
          pageRank.maxIter(pageRankProto.getMaxIter)
        } else {
          pageRank.tol(pageRankProto.getTol)
        }

val pageRankProto = apiMessage.getPageRank
val pageRank = graphFrame.pageRank.resetProbability(pageRankProto.getResetProbability)
Expand Down Expand Up @@ -160,6 +148,14 @@ object GraphFramesConnectUtils {
.run()
.vertices // See comment in the PageRank
}
case MethodCase.POWER_ITERATION_CLUSTERING => {
val pic = apiMessage.getPowerIterationClustering
if (pic.hasWeightCol) {
graphFrame.powerIterationClustering(pic.getK, pic.getMaxIter, Some(pic.getWeightCol))
} else {
graphFrame.powerIterationClustering(pic.getK, pic.getMaxIter, None)
}
}
case MethodCase.PREGEL => {
val pregelProto = apiMessage.getPregel
var pregel = graphFrame.pregel
Expand Down Expand Up @@ -211,6 +207,7 @@ object GraphFramesConnectUtils {
case MethodCase.TRIPLETS => {
graphFrame.triplets
}
case _ => throw new GraphFramesUnreachableException() // Unreachable
}
}
}
Loading
You are viewing a condensed version of this merge commit. You can view the full changes here.