Skip to main content

delta_arrow_reader/reader/
metrics.rs

1//! Delta scan planning and execution metrics.
2
3use std::{
4    fmt,
5    sync::{
6        Arc,
7        atomic::{AtomicU64, Ordering},
8    },
9};
10
11use super::options::ParquetReaderBackend;
12
13/// Immutable point-in-time metrics for one Delta scan.
14#[derive(Debug, Clone, PartialEq, Eq)]
15pub struct DeltaScanMetricsSnapshot {
16    /// Delta snapshot version selected for the scan.
17    pub snapshot_version: u64,
18    /// Parquet reader backend selected for the scan.
19    pub parquet_backend: ParquetReaderBackend,
20    /// Final execution partitions planned for the scan, including source repartitioning.
21    pub scan_partitions_planned: u64,
22    /// Data files selected during planning.
23    pub files_planned: u64,
24    /// Add actions excluded during planning, when known.
25    pub add_actions_excluded_during_planning: Option<u64>,
26    /// Estimated rows in selected input files before row filtering, when known for every file.
27    pub estimated_input_rows: Option<u64>,
28    /// Estimated bytes in selected input files, when every file reported a size.
29    pub estimated_input_bytes: Option<u64>,
30    /// Scan partitions whose execution started.
31    pub scan_partitions_started: u64,
32    /// Scan partitions that completed normally.
33    pub scan_partitions_completed: u64,
34    /// Data-file tasks that started, including independently read file ranges.
35    pub file_tasks_started: u64,
36    /// Data-file tasks that completed normally, including independently read file ranges.
37    pub file_tasks_completed: u64,
38    /// Batches emitted by the scheduler before final stream operations.
39    pub scheduler_batches_emitted: u64,
40    /// Rows emitted by the scheduler before final stream operations.
41    pub scheduler_rows_emitted: u64,
42    /// Deletion-vector payloads loaded.
43    pub deletion_vector_payloads_loaded: u64,
44    /// Deletion-vector masks applied.
45    pub deletion_vectors_applied: u64,
46    /// Rows removed by deletion-vector masks.
47    pub deletion_vector_rows_deleted: u64,
48    /// Deletion-vector read or masking failures.
49    pub deletion_vector_failures: u64,
50    /// Deletion-vector coordinate operations rejected by safety checks.
51    pub deletion_vector_coordinate_rejections: u64,
52    /// Direct Parquet ranged GET operations, or `None` for another backend.
53    pub parquet_data_file_range_get_operations: Option<u64>,
54    /// Direct Parquet full GET operations, or `None` for another backend.
55    pub parquet_data_file_full_get_operations: Option<u64>,
56    /// Direct Parquet payload bytes received, or `None` for another backend.
57    pub parquet_data_file_bytes_received: Option<u64>,
58    /// Estimated bytes admitted across direct Parquet tasks, or `None` for another backend.
59    pub estimated_parquet_task_bytes_admitted: Option<u64>,
60}
61
62/// Shared live metrics for one Delta scan.
63#[derive(Clone)]
64pub struct DeltaScanMetrics {
65    inner: Arc<DeltaScanMetricsInner>,
66}
67
68impl fmt::Debug for DeltaScanMetrics {
69    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
70        formatter
71            .debug_struct("DeltaScanMetrics")
72            .finish_non_exhaustive()
73    }
74}
75
76struct DeltaScanMetricsInner {
77    snapshot_version: u64,
78    parquet_backend: ParquetReaderBackend,
79    scan_partitions_planned: AtomicU64,
80    files_planned: u64,
81    add_actions_excluded_during_planning: Option<u64>,
82    estimated_input_rows: Option<u64>,
83    estimated_input_bytes: Option<u64>,
84    scan_partitions_started: AtomicU64,
85    scan_partitions_completed: AtomicU64,
86    file_tasks_started: AtomicU64,
87    file_tasks_completed: AtomicU64,
88    scheduler_batches_emitted: AtomicU64,
89    scheduler_rows_emitted: AtomicU64,
90    deletion_vector_payloads_loaded: AtomicU64,
91    deletion_vectors_applied: AtomicU64,
92    deletion_vector_rows_deleted: AtomicU64,
93    deletion_vector_failures: AtomicU64,
94    deletion_vector_coordinate_rejections: AtomicU64,
95    parquet_data_file_range_get_operations: AtomicU64,
96    parquet_data_file_full_get_operations: AtomicU64,
97    parquet_data_file_bytes_received: AtomicU64,
98    estimated_parquet_task_bytes_admitted: AtomicU64,
99}
100
101#[allow(dead_code)]
102pub(crate) struct DeltaScanMetricsConfig {
103    pub(crate) snapshot_version: u64,
104    pub(crate) parquet_backend: ParquetReaderBackend,
105    pub(crate) scan_partitions_planned: usize,
106    pub(crate) files_planned: usize,
107    pub(crate) add_actions_excluded_during_planning: Option<u64>,
108    pub(crate) estimated_input_rows: Option<u64>,
109    pub(crate) estimated_input_bytes: Option<u64>,
110}
111
112impl DeltaScanMetrics {
113    #[allow(dead_code)]
114    pub(crate) fn new(config: DeltaScanMetricsConfig) -> Self {
115        Self {
116            inner: Arc::new(DeltaScanMetricsInner {
117                snapshot_version: config.snapshot_version,
118                parquet_backend: config.parquet_backend,
119                scan_partitions_planned: AtomicU64::new(usize_to_u64_saturating(
120                    config.scan_partitions_planned,
121                )),
122                files_planned: usize_to_u64_saturating(config.files_planned),
123                add_actions_excluded_during_planning: config.add_actions_excluded_during_planning,
124                estimated_input_rows: config.estimated_input_rows,
125                estimated_input_bytes: config.estimated_input_bytes,
126                scan_partitions_started: AtomicU64::new(0),
127                scan_partitions_completed: AtomicU64::new(0),
128                file_tasks_started: AtomicU64::new(0),
129                file_tasks_completed: AtomicU64::new(0),
130                scheduler_batches_emitted: AtomicU64::new(0),
131                scheduler_rows_emitted: AtomicU64::new(0),
132                deletion_vector_payloads_loaded: AtomicU64::new(0),
133                deletion_vectors_applied: AtomicU64::new(0),
134                deletion_vector_rows_deleted: AtomicU64::new(0),
135                deletion_vector_failures: AtomicU64::new(0),
136                deletion_vector_coordinate_rejections: AtomicU64::new(0),
137                parquet_data_file_range_get_operations: AtomicU64::new(0),
138                parquet_data_file_full_get_operations: AtomicU64::new(0),
139                parquet_data_file_bytes_received: AtomicU64::new(0),
140                estimated_parquet_task_bytes_admitted: AtomicU64::new(0),
141            }),
142        }
143    }
144
145    /// Returns an immutable point-in-time copy of all scan metrics.
146    pub fn snapshot(&self) -> DeltaScanMetricsSnapshot {
147        let inner = self.inner.as_ref();
148        DeltaScanMetricsSnapshot {
149            snapshot_version: inner.snapshot_version,
150            parquet_backend: inner.parquet_backend,
151            scan_partitions_planned: load(&inner.scan_partitions_planned),
152            files_planned: inner.files_planned,
153            add_actions_excluded_during_planning: inner.add_actions_excluded_during_planning,
154            estimated_input_rows: inner.estimated_input_rows,
155            estimated_input_bytes: inner.estimated_input_bytes,
156            scan_partitions_started: load(&inner.scan_partitions_started),
157            scan_partitions_completed: load(&inner.scan_partitions_completed),
158            file_tasks_started: load(&inner.file_tasks_started),
159            file_tasks_completed: load(&inner.file_tasks_completed),
160            scheduler_batches_emitted: load(&inner.scheduler_batches_emitted),
161            scheduler_rows_emitted: load(&inner.scheduler_rows_emitted),
162            deletion_vector_payloads_loaded: load(&inner.deletion_vector_payloads_loaded),
163            deletion_vectors_applied: load(&inner.deletion_vectors_applied),
164            deletion_vector_rows_deleted: load(&inner.deletion_vector_rows_deleted),
165            deletion_vector_failures: load(&inner.deletion_vector_failures),
166            deletion_vector_coordinate_rejections: load(
167                &inner.deletion_vector_coordinate_rejections,
168            ),
169            parquet_data_file_range_get_operations: self
170                .parquet_metric(&inner.parquet_data_file_range_get_operations),
171            parquet_data_file_full_get_operations: self
172                .parquet_metric(&inner.parquet_data_file_full_get_operations),
173            parquet_data_file_bytes_received: self
174                .parquet_metric(&inner.parquet_data_file_bytes_received),
175            estimated_parquet_task_bytes_admitted: self
176                .parquet_metric(&inner.estimated_parquet_task_bytes_admitted),
177        }
178    }
179
180    fn parquet_metric(&self, counter: &AtomicU64) -> Option<u64> {
181        match self.inner.parquet_backend {
182            ParquetReaderBackend::Direct => Some(load(counter)),
183            ParquetReaderBackend::DeltaKernel => None,
184        }
185    }
186
187    #[allow(dead_code)]
188    pub(crate) fn record_scan_partitions_planned(&self, value: usize) {
189        self.inner
190            .scan_partitions_planned
191            .store(usize_to_u64_saturating(value), Ordering::Relaxed);
192    }
193
194    #[allow(dead_code)]
195    pub(crate) fn record_scan_partition_started(&self) {
196        saturating_fetch_add(&self.inner.scan_partitions_started, 1);
197    }
198
199    #[allow(dead_code)]
200    pub(crate) fn record_scan_partition_completed(&self) {
201        saturating_fetch_add(&self.inner.scan_partitions_completed, 1);
202    }
203
204    #[allow(dead_code)]
205    pub(crate) fn record_file_task_started(&self) {
206        saturating_fetch_add(&self.inner.file_tasks_started, 1);
207    }
208
209    #[allow(dead_code)]
210    pub(crate) fn record_file_task_completed(&self) {
211        saturating_fetch_add(&self.inner.file_tasks_completed, 1);
212    }
213
214    #[allow(dead_code)]
215    pub(crate) fn record_scheduler_batch_emitted(&self, rows: usize) {
216        saturating_fetch_add(&self.inner.scheduler_batches_emitted, 1);
217        saturating_fetch_add(
218            &self.inner.scheduler_rows_emitted,
219            usize_to_u64_saturating(rows),
220        );
221    }
222
223    #[allow(dead_code)]
224    pub(crate) fn record_deletion_vector_payload_loaded(&self) {
225        saturating_fetch_add(&self.inner.deletion_vector_payloads_loaded, 1);
226    }
227
228    #[allow(dead_code)]
229    pub(crate) fn record_deletion_vector_applied(&self) {
230        saturating_fetch_add(&self.inner.deletion_vectors_applied, 1);
231    }
232
233    #[allow(dead_code)]
234    pub(crate) fn record_deletion_vector_rows_deleted(&self, rows: usize) {
235        saturating_fetch_add(
236            &self.inner.deletion_vector_rows_deleted,
237            usize_to_u64_saturating(rows),
238        );
239    }
240
241    #[allow(dead_code)]
242    pub(crate) fn record_deletion_vector_failure(&self) {
243        saturating_fetch_add(&self.inner.deletion_vector_failures, 1);
244    }
245
246    #[allow(dead_code)]
247    pub(crate) fn record_deletion_vector_coordinate_rejection(&self) {
248        saturating_fetch_add(&self.inner.deletion_vector_coordinate_rejections, 1);
249    }
250
251    pub(crate) fn record_parquet_data_file_range_get_operation(&self) {
252        saturating_fetch_add(&self.inner.parquet_data_file_range_get_operations, 1);
253    }
254
255    pub(crate) fn record_parquet_data_file_full_get_operation(&self) {
256        saturating_fetch_add(&self.inner.parquet_data_file_full_get_operations, 1);
257    }
258
259    pub(crate) fn record_parquet_data_file_bytes_received(&self, bytes: usize) {
260        saturating_fetch_add(
261            &self.inner.parquet_data_file_bytes_received,
262            usize_to_u64_saturating(bytes),
263        );
264    }
265
266    pub(crate) fn record_estimated_parquet_task_bytes_admitted(&self, bytes: u64) {
267        saturating_fetch_add(&self.inner.estimated_parquet_task_bytes_admitted, bytes);
268    }
269}
270
271fn load(counter: &AtomicU64) -> u64 {
272    counter.load(Ordering::Relaxed)
273}
274
275#[allow(dead_code)]
276pub(crate) fn saturating_fetch_add(counter: &AtomicU64, value: u64) {
277    let _ = counter.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |current| {
278        Some(current.saturating_add(value))
279    });
280}
281
282fn usize_to_u64_saturating(value: usize) -> u64 {
283    u64::try_from(value).unwrap_or(u64::MAX)
284}
285
286#[cfg(test)]
287mod tests {
288    use std::{sync::atomic::Ordering, thread};
289
290    use super::{DeltaScanMetrics, DeltaScanMetricsConfig, saturating_fetch_add};
291    use crate::ParquetReaderBackend;
292
293    fn metrics(parquet_backend: ParquetReaderBackend) -> DeltaScanMetrics {
294        DeltaScanMetrics::new(DeltaScanMetricsConfig {
295            snapshot_version: 7,
296            parquet_backend,
297            scan_partitions_planned: 3,
298            files_planned: 5,
299            add_actions_excluded_during_planning: Some(2),
300            estimated_input_rows: Some(99),
301            estimated_input_bytes: Some(42),
302        })
303    }
304
305    #[test]
306    fn snapshot_has_context_zeroes_and_backend_availability() {
307        let direct = metrics(ParquetReaderBackend::Direct).snapshot();
308        assert_eq!(direct.snapshot_version, 7);
309        assert_eq!(direct.parquet_backend, ParquetReaderBackend::Direct);
310        assert_eq!(direct.scan_partitions_planned, 3);
311        assert_eq!(direct.files_planned, 5);
312        assert_eq!(direct.add_actions_excluded_during_planning, Some(2));
313        assert_eq!(direct.estimated_input_rows, Some(99));
314        assert_eq!(direct.estimated_input_bytes, Some(42));
315        assert_eq!(direct.scan_partitions_started, 0);
316        assert_eq!(direct.scan_partitions_completed, 0);
317        assert_eq!(direct.file_tasks_started, 0);
318        assert_eq!(direct.file_tasks_completed, 0);
319        assert_eq!(direct.scheduler_batches_emitted, 0);
320        assert_eq!(direct.scheduler_rows_emitted, 0);
321        assert_eq!(direct.deletion_vector_payloads_loaded, 0);
322        assert_eq!(direct.deletion_vectors_applied, 0);
323        assert_eq!(direct.deletion_vector_rows_deleted, 0);
324        assert_eq!(direct.deletion_vector_failures, 0);
325        assert_eq!(direct.deletion_vector_coordinate_rejections, 0);
326        assert_eq!(direct.parquet_data_file_range_get_operations, Some(0));
327        assert_eq!(direct.parquet_data_file_full_get_operations, Some(0));
328        assert_eq!(direct.parquet_data_file_bytes_received, Some(0));
329        assert_eq!(direct.estimated_parquet_task_bytes_admitted, Some(0));
330
331        let kernel = metrics(ParquetReaderBackend::DeltaKernel).snapshot();
332        assert_eq!(kernel.parquet_data_file_range_get_operations, None);
333        assert_eq!(kernel.parquet_data_file_full_get_operations, None);
334        assert_eq!(kernel.parquet_data_file_bytes_received, None);
335        assert_eq!(kernel.estimated_parquet_task_bytes_admitted, None);
336    }
337
338    #[test]
339    fn debug_output_is_safe_and_redacted() {
340        assert_eq!(
341            format!("{:?}", metrics(ParquetReaderBackend::Direct)),
342            "DeltaScanMetrics { .. }"
343        );
344    }
345
346    #[test]
347    fn snapshot_maps_live_counters() {
348        let metrics = metrics(ParquetReaderBackend::Direct);
349        metrics.record_scan_partitions_planned(16);
350        let counters = [
351            &metrics.inner.scan_partitions_started,
352            &metrics.inner.scan_partitions_completed,
353            &metrics.inner.file_tasks_started,
354            &metrics.inner.file_tasks_completed,
355            &metrics.inner.scheduler_batches_emitted,
356            &metrics.inner.scheduler_rows_emitted,
357            &metrics.inner.deletion_vector_payloads_loaded,
358            &metrics.inner.deletion_vectors_applied,
359            &metrics.inner.deletion_vector_rows_deleted,
360            &metrics.inner.deletion_vector_failures,
361            &metrics.inner.deletion_vector_coordinate_rejections,
362            &metrics.inner.parquet_data_file_range_get_operations,
363            &metrics.inner.parquet_data_file_full_get_operations,
364            &metrics.inner.parquet_data_file_bytes_received,
365            &metrics.inner.estimated_parquet_task_bytes_admitted,
366        ];
367        for (index, counter) in counters.into_iter().enumerate() {
368            saturating_fetch_add(counter, u64::try_from(index + 1).expect("small test value"));
369        }
370
371        let snapshot = metrics.snapshot();
372        assert_eq!(snapshot.scan_partitions_planned, 16);
373        assert_eq!(snapshot.scan_partitions_started, 1);
374        assert_eq!(snapshot.scan_partitions_completed, 2);
375        assert_eq!(snapshot.file_tasks_started, 3);
376        assert_eq!(snapshot.file_tasks_completed, 4);
377        assert_eq!(snapshot.scheduler_batches_emitted, 5);
378        assert_eq!(snapshot.scheduler_rows_emitted, 6);
379        assert_eq!(snapshot.deletion_vector_payloads_loaded, 7);
380        assert_eq!(snapshot.deletion_vectors_applied, 8);
381        assert_eq!(snapshot.deletion_vector_rows_deleted, 9);
382        assert_eq!(snapshot.deletion_vector_failures, 10);
383        assert_eq!(snapshot.deletion_vector_coordinate_rejections, 11);
384        assert_eq!(snapshot.parquet_data_file_range_get_operations, Some(12));
385        assert_eq!(snapshot.parquet_data_file_full_get_operations, Some(13));
386        assert_eq!(snapshot.parquet_data_file_bytes_received, Some(14));
387        assert_eq!(snapshot.estimated_parquet_task_bytes_admitted, Some(15));
388    }
389
390    #[test]
391    fn cloned_handles_saturate_under_concurrent_updates() -> Result<(), &'static str> {
392        let metrics = metrics(ParquetReaderBackend::Direct);
393        metrics
394            .inner
395            .file_tasks_started
396            .store(u64::MAX - 1, Ordering::Relaxed);
397        let workers = (0..4)
398            .map(|_| {
399                let metrics = metrics.clone();
400                thread::spawn(move || {
401                    saturating_fetch_add(&metrics.inner.file_tasks_started, 1);
402                })
403            })
404            .collect::<Vec<_>>();
405
406        for worker in workers {
407            worker.join().map_err(|_| "metrics worker panicked")?;
408        }
409
410        assert_eq!(metrics.snapshot().file_tasks_started, u64::MAX);
411        Ok(())
412    }
413}