Skip to content
Merged
Changes from 1 commit
Commits
Show all changes
31 commits
Select commit Hold shift + click to select a range
044a1c1
</think>
SemyonSinchenko Feb 28, 2026
54ebffe
refactor: extract ConnectedComponents implementation to TwoPhase
SemyonSinchenko Feb 28, 2026
87f43bb
fix: standardize connected components algorithm string and GraphX impl
SemyonSinchenko Feb 28, 2026
8057d93
feat: add getAlgorithm getter
SemyonSinchenko Feb 28, 2026
3f6c1d0
refactor: remove GraphX implementation and maxIter from TwoPhase
SemyonSinchenko Mar 5, 2026
b97f19c
feat: add runAQE method for AQE-based two-phase connected components
SemyonSinchenko Mar 5, 2026
44772b3
refactor: calcMinNbrSum object method, DecimalType(38,10), drop overflow
SemyonSinchenko Mar 5, 2026
44bfb4e
refactor: simplify calcMinNbrSum and remove unused variable
SemyonSinchenko Mar 5, 2026
354ec7d
refactor: extract shared output generation to buildOutput method
SemyonSinchenko Mar 5, 2026
6a8c0bf
feat: add intermediate storage level and unpersist in GraphX CC
SemyonSinchenko Mar 5, 2026
3d1972e
feat: use runAQE for two-phase CC when broadcastThreshold is -1
SemyonSinchenko Mar 5, 2026
d3d3230
refactor: remove unused spark variable in TwoPhase
SemyonSinchenko Mar 5, 2026
4cac662
feat: add isGraphPrepared to skip graph preparation in run and runAQE
SemyonSinchenko Mar 5, 2026
2321744
refactor: add isGraphPrepared and checkpoint params to algorithm runs
SemyonSinchenko Mar 5, 2026
2047d52
feat: add internal graph preparation control to ConnectedComponents
SemyonSinchenko Mar 5, 2026
4949d3c
refactor: rename _isGraphPrepared and update setter
SemyonSinchenko Mar 5, 2026
43982af
docs: clarify setIsGraphPrepared docstring on algorithm prep steps
SemyonSinchenko Mar 5, 2026
6d778b7
docs: replace Connected Components section with comprehensive algorit…
SemyonSinchenko Mar 5, 2026
075026b
feat: rename graphframes to two_phase and add checkpointing to rand_cont
SemyonSinchenko Mar 5, 2026
d2d2ecb
test: add local checkpointing to RandomizedContractionSuite
SemyonSinchenko Mar 5, 2026
5bfaac1
refactor: unpersist edges earlier and in finally block
SemyonSinchenko Mar 5, 2026
7813237
test: add System.gc() to make test more robust
SemyonSinchenko Mar 5, 2026
328f4fc
test: clean up Spark checkpoint directory in RandomizedContractionSuite
SemyonSinchenko Mar 5, 2026
cc3abbd
test(RandomizedContractionSuite): move checkpoint cleanup after unper…
SemyonSinchenko Mar 5, 2026
04ebe28
chore: check memory leaks for persisted only
SemyonSinchenko Mar 5, 2026
a4aee16
fix: make the test robust, not random
SemyonSinchenko Mar 5, 2026
8562ee6
chore: drop the redundant
SemyonSinchenko Mar 5, 2026
54a7ae9
chor: naming
SemyonSinchenko Mar 6, 2026
ab5aec2
Merge remote-tracking branch 'graphframes/main' into 775-refactor-cc-api
SemyonSinchenko Mar 11, 2026
6915dc2
fix: addressing James' comments
SemyonSinchenko Mar 11, 2026
ce062ff
fix: python docstrings
SemyonSinchenko Mar 11, 2026
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
feat: add internal graph preparation control to ConnectedComponents
  • Loading branch information
SemyonSinchenko committed Mar 5, 2026
commit 2047d520ae9f2d780167ba0914beddfbb52cd377
26 changes: 23 additions & 3 deletions core/src/main/scala/org/graphframes/lib/ConnectedComponents.scala
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,8 @@ class ConnectedComponents private[graphframes] (private val graph: GraphFrame)
private var algorithm: String = GraphFramesConf.getConnectedComponentsAlgorithm
.getOrElse(ALGO_TWO_PHASE)

private var _isGraphPrepared: Boolean = false

setCheckpointInterval(
GraphFramesConf.getConnectedComponentsCheckpointInterval.getOrElse(checkpointInterval))
setBroadcastThreshold(
Expand Down Expand Up @@ -99,6 +101,24 @@ class ConnectedComponents private[graphframes] (private val graph: GraphFrame)
*/
def getAlgorithm: String = algorithm

/**
* WARNING: INTERNAL API. For very experienced users only.
* Sets whether the graph has already been prepared (edges deduplicated and normalized).
* Only set this to `true` if you fully understand the internal graph preparation steps
* performed by ConnectedComponents and have done them yourself. Incorrect use WILL produce
* silently wrong results.
*
* @param value true if the graph is already prepared, false otherwise (default: false)
*/
@deprecated(
"INTERNAL API ONLY. This is an experimental internal option for advanced users who fully " +
"understand graph preparation internals. Misuse will produce silently wrong results.",
since = "0.9.3")
def setIsGraphPrepared(value: Boolean): this.type = {
_isGraphPrepared = value
this
}

/**
* Runs the algorithm.
*/
Expand All @@ -117,7 +137,7 @@ class ConnectedComponents private[graphframes] (private val graph: GraphFrame)
intermediateStorageLevel = intermediateStorageLevel,
useLabelsAsComponents = useLabelsAsComponents,
useLocalCheckpoints = useLocalCheckpoints,
isGraphPrepared = false)
isGraphPrepared = _isGraphPrepared)
} else {
TwoPhase.run(
graph,
Expand All @@ -126,7 +146,7 @@ class ConnectedComponents private[graphframes] (private val graph: GraphFrame)
intermediateStorageLevel = intermediateStorageLevel,
useLabelsAsComponents = useLabelsAsComponents,
useLocalCheckpoints = useLocalCheckpoints,
isGraphPrepared = false)
isGraphPrepared = _isGraphPrepared)
}
case ALGO_RANDOMIZED_CONTRACTION =>
RandomizedContraction.run(
Expand All @@ -135,7 +155,7 @@ class ConnectedComponents private[graphframes] (private val graph: GraphFrame)
intermediateStorageLevel = intermediateStorageLevel,
useLocalCheckpoints = useLocalCheckpoints,
checkpointInterval = checkpointInterval,
isGraphPrepared = false)
isGraphPrepared = _isGraphPrepared)
// the check is inside the setter
case _ => throw new GraphFramesUnreachableException()
}
Expand Down