Skip to main content

delta_arrow_reader/
metrics.rs

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