EmilyMatt commented on code in PR #25372:
URL: https://github.com/apache/datafusion/pull/25372#discussion_r4120749747


##########
datafusion/physical-plan/src/sorts/stream.rs:
##########
@@ -94,25 +94,62 @@ impl FusedStreams {
     }
 }
 
-/// An `Arc<Rows>` that can be reused
+/// An `Arc<Rows>` that can be reused.
+///
+/// Owns the reservation covering every retained buffer for as long as it is
+/// retained, so the cache is visible to the pool rather than held off-book.
+/// A retained buffer keeps its capacity, so the cost is the high-water mark of
+/// each stream, not the size of the batch currently in flight.
 #[derive(Debug)]
 struct ReusableRows {
     inner: Vec<Option<Arc<Rows>>>,
+    reservation: MemoryReservation,
 }
 
 impl ReusableRows {
     // return a Rows for writing,
     // does not clone if the existing rows can be reused
-    fn take_next(&mut self, stream_idx: usize) -> Result<Rows> {
-        Arc::try_unwrap(self.inner[stream_idx].take().unwrap()).map_err(|_| {
-            internal_datafusion_err!(
-                "Rows from RowCursorStream is still in use by consumer"
-            )
-        })
+    fn take_next(&mut self, stream_idx: usize, converter: &RowConverter) -> 
Result<Rows> {
+        match self.inner[stream_idx].take() {
+            Some(rows) => Arc::try_unwrap(rows).map_err(|_| {
+                internal_datafusion_err!(
+                    "Rows from RowCursorStream is still in use by consumer"
+                )
+            }),
+            // Nothing retained yet, or already released, so start over.
+            None => Ok(converter.empty_rows(0, 0)),
+        }
     }
-    // save the Rows
-    fn save(&mut self, stream_idx: usize, rows: &Arc<Rows>) {
+
+    /// Account for a freshly built buffer, and retain it for reuse.
+    ///
+    /// The reservation is mandatory rather than best-effort. The buffer is 
live in the
+    /// cursor whether or not this slot keeps a handle to it, so declining to 
reserve
+    /// would hide it from the pool instead of avoiding it, and it frees 
nothing at this
+    /// point either, since the cursor holds the same `Arc`. Retention on top 
of the
+    /// reservation is free for the same reason. A pool that cannot cover the 
buffer
+    /// fails the query here, as it did before the buffer was cached at all.
+    fn save(&mut self, stream_idx: usize, rows: &Arc<Rows>) -> Result<()> {
         self.inner[stream_idx] = Some(Arc::clone(rows));
+        let retained = self.retained_size();
+        debug_assert!(retained >= self.reservation.size());
+        if let Err(e) = self.reservation.try_resize(retained) {
+            self.inner[stream_idx] = None;
+            return Err(e);
+        }
+        Ok(())
+    }
+
+    // drop whatever a finished stream was holding
+    fn release(&mut self, stream_idx: usize) {
+        debug_assert!(self.reservation.size() >= self.retained_size());
+        if let Some(rows) = self.inner[stream_idx].take() {
+            self.reservation.shrink(rows.size());
+        }
+    }
+
+    fn retained_size(&self) -> usize {

Review Comment:
   Yeah it's more of a design choice, I think it's negligible performance-wise, 
and makes for a cleaner, more stateless calculation



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


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

Reply via email to