Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
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
1 change: 1 addition & 0 deletions graphframes-connect/src/main/protobuf/graphframes.proto
Original file line number Diff line number Diff line change
Expand Up @@ -111,6 +111,7 @@ message Pregel {
string additional_col_name = 6;
ColumnOrExpression additional_col_initial = 7;
ColumnOrExpression additional_col_upd = 8;
optional bool early_stopping = 9;
}

message ShortestPaths {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -174,6 +174,10 @@ object GraphFramesConnectUtils {
.map(parseColumnOrExpression(_, planner))
.foldLeft(pregel)((p, col) => p.sendMsgToDst(col))

if (pregelProto.hasEarlyStopping) {
pregel = pregel.setEarlyStopping(pregelProto.getEarlyStopping)
}

pregel.run()
}
case MethodCase.SHORTEST_PATHS => {
Expand Down
9 changes: 9 additions & 0 deletions python/graphframes/connect/graphframe_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@ def __init__(self, graph: "GraphFrameConnect") -> None:
self._send_msg_to_src = []
self._send_msg_to_dst = []
self._agg_msg = None
self._early_stopping = False

def setMaxIter(self, value: int) -> Self:
self._max_iter = value
Expand All @@ -33,6 +34,10 @@ def setCheckpointInterval(self, value: int) -> Self:
self._checkpoint_interval = value
return self

def setEarlyStopping(self, value: bool) -> Self:
self._early_stopping = value
return self

def withVertexColumn(
self,
colName: str,
Expand Down Expand Up @@ -62,6 +67,7 @@ def __init__(
self,
max_iter: int,
checkpoint_interval: int,
early_stopping: bool,
vertex_col_name: str,
agg_msg: Column | str,
send2dst: list[Column | str],
Expand All @@ -74,6 +80,7 @@ def __init__(
super().__init__(None)
self.max_iter = max_iter
self.checkpoint_interval = checkpoint_interval
self.early_stopping = early_stopping
self.vertex_col_name = vertex_col_name
self.agg_msg = agg_msg
self.send2dst = send2dst
Expand All @@ -97,6 +104,7 @@ def plan(self, session: SparkConnectClient) -> proto.Relation:
additional_col_name=self.vertex_col_name,
additional_col_initial=make_column_or_expr(self.vertex_col_init, session),
additional_col_upd=make_column_or_expr(self.vertex_col_upd, session),
early_stopping=self.early_stopping,
)
pb_message = pb.GraphFramesAPI(
vertices=dataframe_to_proto(self.vertices, session),
Expand Down Expand Up @@ -129,6 +137,7 @@ def plan(self, session: SparkConnectClient) -> proto.Relation:
send2src=self._send_msg_to_src,
vertices=self.graph._vertices,
edges=self.graph._edges,
early_stopping=self._early_stopping,
),
session=self.graph._spark,
)
Expand Down
24 changes: 12 additions & 12 deletions python/graphframes/connect/proto/graphframes_pb2.py

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

4 changes: 4 additions & 0 deletions python/graphframes/connect/proto/graphframes_pb2.pyi
Original file line number Diff line number Diff line change
Expand Up @@ -251,6 +251,7 @@ class Pregel(_message.Message):
"additional_col_name",
"additional_col_initial",
"additional_col_upd",
"early_stopping",
)
AGG_MSGS_FIELD_NUMBER: _ClassVar[int]
SEND_MSG_TO_DST_FIELD_NUMBER: _ClassVar[int]
Expand All @@ -260,6 +261,7 @@ class Pregel(_message.Message):
ADDITIONAL_COL_NAME_FIELD_NUMBER: _ClassVar[int]
ADDITIONAL_COL_INITIAL_FIELD_NUMBER: _ClassVar[int]
ADDITIONAL_COL_UPD_FIELD_NUMBER: _ClassVar[int]
EARLY_STOPPING_FIELD_NUMBER: _ClassVar[int]
agg_msgs: ColumnOrExpression
send_msg_to_dst: _containers.RepeatedCompositeFieldContainer[ColumnOrExpression]
send_msg_to_src: _containers.RepeatedCompositeFieldContainer[ColumnOrExpression]
Expand All @@ -268,6 +270,7 @@ class Pregel(_message.Message):
additional_col_name: str
additional_col_initial: ColumnOrExpression
additional_col_upd: ColumnOrExpression
early_stopping: bool
def __init__(
self,
agg_msgs: _Optional[_Union[ColumnOrExpression, _Mapping]] = ...,
Expand All @@ -278,6 +281,7 @@ class Pregel(_message.Message):
additional_col_name: _Optional[str] = ...,
additional_col_initial: _Optional[_Union[ColumnOrExpression, _Mapping]] = ...,
additional_col_upd: _Optional[_Union[ColumnOrExpression, _Mapping]] = ...,
early_stopping: bool = ...,
) -> None: ...

class ShortestPaths(_message.Message):
Expand Down
Loading