From db9ed41e0d6ef000cdc9f880344662794cece741 Mon Sep 17 00:00:00 2001 From: Raz Luvaton <16746759+rluvaton@users.noreply.github.com> Date: Sun, 19 Jul 2026 12:40:43 +0300 Subject: [PATCH] fix: track the time takes to init loser tree and build last in progress batches --- datafusion/physical-plan/src/sorts/merge.rs | 12 +++++------- 1 file changed, 5 insertions(+), 7 deletions(-) diff --git a/datafusion/physical-plan/src/sorts/merge.rs b/datafusion/physical-plan/src/sorts/merge.rs index 986da549f75c8..2090c16e0ca40 100644 --- a/datafusion/physical-plan/src/sorts/merge.rs +++ b/datafusion/physical-plan/src/sorts/merge.rs @@ -223,6 +223,9 @@ impl SortPreservingMergeStream { assert_eq!(uninitiated_partitions.len(), 0); + let elapsed_compute = self.metrics.elapsed_compute().clone(); + let mut timer = elapsed_compute.timer(); + // If there are no more uninitiated partitions, set up the loser tree and continue // to the next phase. @@ -230,11 +233,6 @@ impl SortPreservingMergeStream { drop(uninitiated_partitions); self.init_loser_tree(); - // NB timer records time taken on drop, so there are no - // calls to `timer.done()` below. - let elapsed_compute = self.metrics.elapsed_compute().clone(); - let mut timer = elapsed_compute.timer(); - loop { let stream_idx = self.loser_tree[0]; if !self.advance_cursors(stream_idx) { @@ -269,8 +267,6 @@ impl SortPreservingMergeStream { self.update_loser_tree(); } - drop(timer); - // When `build_record_batch()` hits an i32 offset overflow (e.g. // combined string offsets exceed 2 GB), it emits a partial batch // and keeps the remaining rows in `self.in_progress.indices`. @@ -279,7 +275,9 @@ impl SortPreservingMergeStream { // Repeated overflows are fine — each poll emits another partial // batch until `in_progress` is fully drained. while let Some(batch) = self.emit_in_progress_batch()? { + drop(timer); emitter.emit(batch).await; + timer = elapsed_compute.timer(); } Ok(()) })