Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
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
Feedbacks + Doc update
Signed-off-by: joelrobin18 <joelrobin1818@gmail.com>
  • Loading branch information
joelrobin18 committed Dec 29, 2025
commit b77545704585f49584e861f09822fa9ad223c415
30 changes: 30 additions & 0 deletions docs/src/04-user-guide/10-pregel.md
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,36 @@ val edgeWeight = Pregel.edge("weight")

Under the hood, the passed name of the column will be resolved to get the corresponding element of the triplet structs.

### Memory Optimization for Triplets

By default, all vertex columns are included when constructing triplets. For algorithms with large per-vertex state (e.g., cycle detection storing sequences, random walks), this can create huge intermediate datasets in memory.

To reduce memory usage, you can specify only the columns that are actually needed using `requiredSrcColumns` and `requiredDstColumns`:

```scala
graph.pregel
.withVertexColumn("distances", ...)
.sendMsgToDst(Pregel.src("distances")) // Only needs "distances" from source
.requiredSrcColumns("distances") // Only include "distances" in src struct
.requiredDstColumns("distances") // Only include "distances" in dst struct
.aggMsgs(...)
.run()
```

In Python:

```python
graph.pregel \
.withVertexColumn("distances", ...) \
.sendMsgToDst(Pregel.src("distances")) \
.required_src_columns("distances") \
.required_dst_columns("distances") \
.aggMsgs(...) \
.run()
```

The `id` column and the active flag column (if used) are always included automatically, so you don't need to specify them.

### Sending Messages

GraphFrames Pregel API support arbitrary number of messages per vertex. Inside the Pregel API **graphs are always considered directed**. This means that if a vertex has an outgoing edge to another vertex, then the message will be sent to the destination vertex. To emulate the behavior of the undirected graph, the user can send the same message to both the source and the destination vertex.
Expand Down
16 changes: 8 additions & 8 deletions python/graphframes/connect/graphframes_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -119,30 +119,30 @@ def setIntermediateStorageLevel(self, storage_level: StorageLevel) -> Self:
self._storage_level = storage_level
return self

def requiredSrcColumns(self, colName: str, *colNames: str) -> Self:
def required_src_columns(self, col_name: str, *col_names: str) -> Self:
"""Specifies which source vertex columns are required when constructing triplets.

By default, all source vertex columns are included in triplets, which can create large
intermediate datasets for algorithms with significant state. Use this method to reduce
memory usage by specifying only the columns that are actually needed.

:param colName: the first required source vertex column name
:param colNames: additional required source vertex column names
:param col_name: the first required source vertex column name
:param col_names: additional required source vertex column names
"""
self._required_src_columns = [colName] + list(colNames)
self._required_src_columns = [col_name] + list(col_names)
return self

def requiredDstColumns(self, colName: str, *colNames: str) -> Self:
def required_dst_columns(self, col_name: str, *col_names: str) -> Self:
"""Specifies which destination vertex columns are required when constructing triplets.

By default, all destination vertex columns are included in triplets, which can create large
intermediate datasets for algorithms with significant state. Use this method to reduce
memory usage by specifying only the columns that are actually needed.

:param colName: the first required destination vertex column name
:param colNames: additional required destination vertex column names
:param col_name: the first required destination vertex column name
:param col_names: additional required destination vertex column names
"""
self._required_dst_columns = [colName] + list(colNames)
self._required_dst_columns = [col_name] + list(col_names)
return self

def run(self) -> DataFrame:
Expand Down
20 changes: 10 additions & 10 deletions python/graphframes/lib/pregel.py
Original file line number Diff line number Diff line change
Expand Up @@ -236,7 +236,7 @@ def setIntermediateStorageLevel(self, storage_level: StorageLevel) -> Self:
)
return self

def requiredSrcColumns(self, colName: str, *colNames: str) -> Self:
def required_src_columns(self, col_name: str, *col_names: str) -> Self:
"""Specifies which source vertex columns are required when constructing triplets.

By default, all source vertex columns are included in triplets, which can create large
Expand All @@ -246,17 +246,17 @@ def requiredSrcColumns(self, colName: str, *colNames: str) -> Self:

The ID column and the active flag column (if used) are always included automatically.

:param colName: the first required source vertex column name
:param colNames: additional required source vertex column names
:param col_name: the first required source vertex column name
:param col_names: additional required source vertex column names

See also :func:`requiredDstColumns`
See also :func:`required_dst_columns`
"""
self._java_obj.requiredSrcColumns(
colName, _to_seq(self.graph._spark.sparkContext, colNames)
col_name, _to_seq(self.graph._spark.sparkContext, col_names)
)
return self

def requiredDstColumns(self, colName: str, *colNames: str) -> Self:
def required_dst_columns(self, col_name: str, *col_names: str) -> Self:
"""Specifies which destination vertex columns are required when constructing triplets.

By default, all destination vertex columns are included in triplets, which can create large
Expand All @@ -266,13 +266,13 @@ def requiredDstColumns(self, colName: str, *colNames: str) -> Self:

The ID column and the active flag column (if used) are always included automatically.

:param colName: the first required destination vertex column name
:param colNames: additional required destination vertex column names
:param col_name: the first required destination vertex column name
:param col_names: additional required destination vertex column names

See also :func:`requiredSrcColumns`
See also :func:`required_src_columns`
"""
self._java_obj.requiredDstColumns(
colName, _to_seq(self.graph._spark.sparkContext, colNames)
col_name, _to_seq(self.graph._spark.sparkContext, col_names)
)
return self

Expand Down