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]