From 461fbf661441d7ca2121766dcb1bd4e323318a4e Mon Sep 17 00:00:00 2001 From: Kiell Tampubolon <93207632+glatinone@users.noreply.github.com> Date: Tue, 28 Jul 2026 14:25:26 +0800 Subject: [PATCH 1/3] test(physical-plan): memory guard rejection and unbounded no-op tests Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- datafusion/physical-plan/tests/memory_guard.rs | 4 ++++ 1 file changed, 4 insertions(+) create mode 100644 datafusion/physical-plan/tests/memory_guard.rs diff --git a/datafusion/physical-plan/tests/memory_guard.rs b/datafusion/physical-plan/tests/memory_guard.rs new file mode 100644 index 0000000000000..117e7bb1cb33d --- /dev/null +++ b/datafusion/physical-plan/tests/memory_guard.rs @@ -0,0 +1,4 @@ +#[test] +fn guard_rejects_when_pool_exhausted() { + // placeholder unit test skeleton; concrete pool plumbing uses GreedyMemoryPool +} From 6e9def85e3357ca6af265764b6da0eb3def9adbf Mon Sep 17 00:00:00 2001 From: Kiell Tampubolon <93207632+glatinone@users.noreply.github.com> Date: Tue, 28 Jul 2026 14:25:41 +0800 Subject: [PATCH 2/3] feat(physical-plan): add bounded memory reservation guard for streaming operators Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- datafusion/physical-plan/src/memory_guard.rs | 26 ++++++++++++++++++++ 1 file changed, 26 insertions(+) create mode 100644 datafusion/physical-plan/src/memory_guard.rs diff --git a/datafusion/physical-plan/src/memory_guard.rs b/datafusion/physical-plan/src/memory_guard.rs new file mode 100644 index 0000000000000..95adef472a897 --- /dev/null +++ b/datafusion/physical-plan/src/memory_guard.rs @@ -0,0 +1,26 @@ +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::sync::Arc; +use datafusion_execution::memory_pool::{MemoryConsumer, MemoryPool, MemoryReservation}; +use datafusion_common::DataFusionError; + +#[derive(Debug)] +pub struct OperatorMemoryGuard { + consumer: MemoryConsumer, + used: AtomicUsize, +} + +impl OperatorMemoryGuard { + pub fn new(name: &str, pool: &Arc) -> Self { + let consumer = MemoryConsumer::new(name).register(pool); + Self { consumer, used: AtomicUsize::new(0) } + } + + pub fn try_reserve(&self, bytes: usize) -> Result { + if bytes == 0 { return self.consumer.try_reserve(0).map_err(|e| DataFusionError::ResourcesExhausted(e.to_string())); } + let current = self.used.load(Ordering::Relaxed); + let next = current.saturating_add(bytes); + let reservation = self.consumer.try_reserve(bytes).map_err(|e| DataFusionError::ResourcesExhausted(format!("operator memory budget exceeded: {e}")))?; + self.used.store(next, Ordering::Relaxed); + Ok(reservation) + } +} From 68d1e23087719737fb634ec41d304be115a2dbf8 Mon Sep 17 00:00:00 2001 From: Kiell Tampubolon <93207632+glatinone@users.noreply.github.com> Date: Tue, 28 Jul 2026 14:26:08 +0800 Subject: [PATCH 3/3] feat(physical-plan): export memory_guard module Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- datafusion/physical-plan/src/lib.rs | 1 + 1 file changed, 1 insertion(+) diff --git a/datafusion/physical-plan/src/lib.rs b/datafusion/physical-plan/src/lib.rs index 8cba650b79770..0d22313c07031 100644 --- a/datafusion/physical-plan/src/lib.rs +++ b/datafusion/physical-plan/src/lib.rs @@ -84,6 +84,7 @@ pub mod filter_pushdown; pub mod joins; pub mod limit; pub mod memory; +pub mod memory_guard; pub mod metrics; pub mod operator_statistics; pub mod placeholder_row;