(
&self,
partition: usize,
context: Arc<TaskContext>,
)
| 179 | } |
| 180 | |
| 181 | fn execute( |
| 182 | &self, |
| 183 | partition: usize, |
| 184 | context: Arc<TaskContext>, |
| 185 | ) -> Result<SendableRecordBatchStream> { |
| 186 | let mem_reservation = MemoryConsumer::new(format!("BufferExec[{partition}]")) |
| 187 | .register(context.memory_pool()); |
| 188 | let in_stream = self.input.execute(partition, context)?; |
| 189 | |
| 190 | // Set up the metrics for the stream. |
| 191 | let curr_mem_in = Arc::new(AtomicUsize::new(0)); |
| 192 | let curr_mem_out = Arc::clone(&curr_mem_in); |
| 193 | let mut max_mem_in = 0; |
| 194 | let max_mem = MetricBuilder::new(&self.metrics) |
| 195 | .with_category(MetricCategory::Bytes) |
| 196 | .gauge("max_mem_used", partition); |
| 197 | |
| 198 | let curr_queued_in = Arc::new(AtomicUsize::new(0)); |
| 199 | let curr_queued_out = Arc::clone(&curr_queued_in); |
| 200 | let mut max_queued_in = 0; |
| 201 | let max_queued = MetricBuilder::new(&self.metrics) |
| 202 | .with_category(MetricCategory::Rows) |
| 203 | .gauge("max_queued", partition); |
| 204 | |
| 205 | // Capture metrics when an element is queued on the stream. |
| 206 | let in_stream = in_stream.inspect_ok(move |v| { |
| 207 | let size = v.get_array_memory_size(); |
| 208 | let curr_size = curr_mem_in.fetch_add(size, Ordering::Relaxed) + size; |
| 209 | if curr_size > max_mem_in { |
| 210 | max_mem_in = curr_size; |
| 211 | max_mem.set(max_mem_in); |
| 212 | } |
| 213 | |
| 214 | let curr_queued = curr_queued_in.fetch_add(1, Ordering::Relaxed) + 1; |
| 215 | if curr_queued > max_queued_in { |
| 216 | max_queued_in = curr_queued; |
| 217 | max_queued.set(max_queued_in); |
| 218 | } |
| 219 | }); |
| 220 | // Buffer the input. |
| 221 | let out_stream = |
| 222 | MemoryBufferedStream::new(in_stream, self.capacity, mem_reservation); |
| 223 | // Update in the metrics that when an element gets out, some memory gets freed. |
| 224 | let out_stream = out_stream.inspect_ok(move |v| { |
| 225 | curr_mem_out.fetch_sub(v.get_array_memory_size(), Ordering::Relaxed); |
| 226 | curr_queued_out.fetch_sub(1, Ordering::Relaxed); |
| 227 | }); |
| 228 | |
| 229 | Ok(Box::pin(RecordBatchStreamAdapter::new( |
| 230 | self.schema(), |
| 231 | out_stream, |
| 232 | ))) |
| 233 | } |
| 234 | |
| 235 | fn metrics(&self) -> Option<MetricsSet> { |
| 236 | Some(self.metrics.clone_inner()) |
nothing calls this directly
no test coverage detected