Repository navigation
feat: SparkConnect support #506
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from 1 commit
d1972fc
c158815
ea11df6
0ddd5bd
da7eccc
fb784a3
9f8905f
d58ed2a
20f7575
eee7b7b
130b12e
7e325aa
c47a57e
fc8ebae
f13c754
a21b5aa
f4c91d6
e4d75f7
0950cfd
8cb430c
19c1934
1eef323
7fd1f23
cc60bcb
5528d65
f88e19a
59897fb
97054b0
90a326f
c8bcf43
9d7f714
5a91659
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
Classic tests are passed; ++some changes
- Loading branch information
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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 { | ||
|
|
@@ -101,9 +98,6 @@ object GraphFramesConnectUtils { | |
| .setBroadcastThreshold(cc.getBroadcastThreshold) | ||
| .run() | ||
| } | ||
| case MethodCase.DEGREES => { | ||
| graphFrame.degrees | ||
| } | ||
| case MethodCase.DROP_ISOLATED_VERTICES => { | ||
| graphFrame.dropIsolatedVertices().vertices | ||
| } | ||
|
|
@@ -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 => { | ||
|
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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
Collaborator
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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) | ||
|
|
@@ -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 | ||
|
|
@@ -211,6 +207,7 @@ object GraphFramesConnectUtils { | |
| case MethodCase.TRIPLETS => { | ||
| graphFrame.triplets | ||
| } | ||
| case _ => throw new GraphFramesUnreachableException() // Unreachable | ||
| } | ||
| } | ||
| } | ||
Uh oh!
There was an error while loading. Please reload this page.