This is an automated email from the ASF dual-hosted git repository.

andygrove pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/datafusion-ballista.git


The following commit(s) were added to refs/heads/main by this push:
     new 2959cb7f9 perf: avoid redundant TaskStatus clones in scheduler 
update_task_status (#1983)
2959cb7f9 is described below

commit 2959cb7f9cd6a6515bb14cdde9f67302cacea812
Author: Andy Grove <[email protected]>
AuthorDate: Thu Jul 9 21:18:58 2026 -0600

    perf: avoid redundant TaskStatus clones in scheduler update_task_status 
(#1983)
---
 ballista/scheduler/src/state/aqe/mod.rs         | 19 ++++++++++---------
 ballista/scheduler/src/state/execution_graph.rs | 22 +++++++++++-----------
 ballista/scheduler/src/state/execution_stage.rs | 10 +++++-----
 3 files changed, 26 insertions(+), 25 deletions(-)

diff --git a/ballista/scheduler/src/state/aqe/mod.rs 
b/ballista/scheduler/src/state/aqe/mod.rs
index 74cedb1f4..7a43d1d88 100644
--- a/ballista/scheduler/src/state/aqe/mod.rs
+++ b/ballista/scheduler/src/state/aqe/mod.rs
@@ -656,7 +656,7 @@ impl ExecutionGraph for AdaptiveExecutionGraph {
                             );
                             continue;
                         }
-                        let partition_id = task_status.clone().partition_id as 
usize;
+                        let partition_id = task_status.partition_id as usize;
                         let task_identity = format!(
                             "TID {} {}/{}.{}/{}",
                             task_status.task_id,
@@ -665,19 +665,20 @@ impl ExecutionGraph for AdaptiveExecutionGraph {
                             task_stage_attempt_num,
                             partition_id
                         );
-                        let operator_metrics = task_status.metrics.clone();
 
-                        if !running_stage
-                            .update_task_info(partition_id, 
task_status.clone())
-                        {
+                        if !running_stage.update_task_info(partition_id, 
&task_status) {
                             continue;
                         }
+
+                        let TaskStatus {
+                            status,
+                            metrics: operator_metrics,
+                            ..
+                        } = task_status;
                         //
                         // handle task failure
                         //
-                        if let Some(task_status::Status::Failed(failed_task)) =
-                            task_status.status
-                        {
+                        if let Some(task_status::Status::Failed(failed_task)) 
= status {
                             let failed_reason = failed_task.failed_reason;
 
                             match failed_reason {
@@ -786,7 +787,7 @@ impl ExecutionGraph for AdaptiveExecutionGraph {
                         //
                         else if let Some(task_status::Status::Successful(
                             successful_task,
-                        )) = task_status.status
+                        )) = status
                         {
                             // update task metrics for successfu task
                             running_stage
diff --git a/ballista/scheduler/src/state/execution_graph.rs 
b/ballista/scheduler/src/state/execution_graph.rs
index 79cbf9c70..879a17296 100644
--- a/ballista/scheduler/src/state/execution_graph.rs
+++ b/ballista/scheduler/src/state/execution_graph.rs
@@ -791,7 +791,7 @@ impl ExecutionGraph for StaticExecutionGraph {
                             );
                             continue;
                         }
-                        let partition_id = task_status.clone().partition_id as 
usize;
+                        let partition_id = task_status.partition_id as usize;
                         let task_identity = format!(
                             "TID {} {}/{}.{}/{}",
                             task_status.task_id,
@@ -800,17 +800,17 @@ impl ExecutionGraph for StaticExecutionGraph {
                             task_stage_attempt_num,
                             partition_id
                         );
-                        let operator_metrics = task_status.metrics.clone();
-
-                        if !running_stage
-                            .update_task_info(partition_id, 
task_status.clone())
-                        {
+                        if !running_stage.update_task_info(partition_id, 
&task_status) {
                             continue;
                         }
 
-                        if let Some(task_status::Status::Failed(failed_task)) =
-                            task_status.status
-                        {
+                        let TaskStatus {
+                            status,
+                            metrics: operator_metrics,
+                            ..
+                        } = task_status;
+
+                        if let Some(task_status::Status::Failed(failed_task)) 
= status {
                             let failed_reason = failed_task.failed_reason;
 
                             match failed_reason {
@@ -915,7 +915,7 @@ impl ExecutionGraph for StaticExecutionGraph {
                             }
                         } else if let Some(task_status::Status::Successful(
                             successful_task,
-                        )) = task_status.status
+                        )) = status
                         {
                             // update task metrics for successfu task
                             running_stage
@@ -966,7 +966,7 @@ impl ExecutionGraph for StaticExecutionGraph {
                     for task_status in stage_task_statuses.into_iter() {
                         let task_stage_attempt_num =
                             task_status.stage_attempt_num as usize;
-                        let partition_id = task_status.clone().partition_id as 
usize;
+                        let partition_id = task_status.partition_id as usize;
                         let task_identity = format!(
                             "TID {} {}/{}.{}/{}",
                             task_status.task_id,
diff --git a/ballista/scheduler/src/state/execution_stage.rs 
b/ballista/scheduler/src/state/execution_stage.rs
index 698d95182..10440f357 100644
--- a/ballista/scheduler/src/state/execution_stage.rs
+++ b/ballista/scheduler/src/state/execution_stage.rs
@@ -669,7 +669,7 @@ impl RunningStage {
     }
 
     /// Update the TaskInfo for task partition
-    pub fn update_task_info(&mut self, partition_id: usize, status: 
TaskStatus) -> bool {
+    pub fn update_task_info(&mut self, partition_id: usize, status: 
&TaskStatus) -> bool {
         debug!("Updating TaskInfo for partition {partition_id}");
         let Some(task_info) = self.task_infos[partition_id].as_ref() else {
             warn!(
@@ -686,7 +686,7 @@ impl RunningStage {
             return false;
         }
         let scheduled_time = task_info.scheduled_time;
-        let task_status = status.status.unwrap();
+        let task_status = status.status.as_ref().unwrap();
         let updated_task_info = TaskInfo {
             task_id,
             scheduled_time,
@@ -1178,7 +1178,7 @@ mod tests {
         // Simulates receiving a status update for a task that was already
         // reset (e.g., executor heartbeat timed out).
         let status = make_task_status(0, 0);
-        let result = stage.update_task_info(0, status);
+        let result = stage.update_task_info(0, &status);
 
         // Should return false (update rejected), not panic.
         assert!(!result);
@@ -1203,7 +1203,7 @@ mod tests {
         });
 
         let status = make_task_status(0, 0);
-        let result = stage.update_task_info(0, status);
+        let result = stage.update_task_info(0, &status);
 
         assert!(result);
         assert!(matches!(
@@ -1240,7 +1240,7 @@ mod tests {
 
         // Executor sends a late status update for partition 0.
         let status = make_task_status(0, 0);
-        let result = stage.update_task_info(0, status);
+        let result = stage.update_task_info(0, &status);
 
         // Should gracefully reject the update, not panic.
         assert!(!result);


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to