MCPcopy Create free account
hub / github.com/apache/datafusion / execute

Method execute

datafusion/physical-plan/src/buffer.rs:181–233  ·  view source on GitHub ↗
(
        &self,
        partition: usize,
        context: Arc<TaskContext>,
    )

Source from the content-addressed store, hash-verified

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())

Callers

nothing calls this directly

Calls 7

newFunction · 0.85
memory_poolMethod · 0.80
gaugeMethod · 0.80
registerMethod · 0.45
with_categoryMethod · 0.45
setMethod · 0.45
schemaMethod · 0.45

Tested by

no test coverage detected