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    /// Final execution partitions planned for the scan, including source repartitioning.
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 tasks that started, including independently read file ranges.
32    pub files_started: u64,
33    /// Data-file tasks that completed normally, including independently read file ranges.
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: AtomicU64,
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: AtomicU64::new(usize_to_u64_saturating(
112                    config.scan_partitions_planned,
113                )),
114                files_planned: usize_to_u64_saturating(config.files_planned),
115                files_filtered_during_planning: config.files_filtered_during_planning,
116                estimated_rows: config.estimated_rows,
117                estimated_bytes: config.estimated_bytes,
118                scan_partitions_started: AtomicU64::new(0),
119                scan_partitions_completed: AtomicU64::new(0),
120                files_started: AtomicU64::new(0),
121                files_completed: AtomicU64::new(0),
122                batches_produced: AtomicU64::new(0),
123                rows_produced: AtomicU64::new(0),
124                deletion_vector_payloads_loaded: AtomicU64::new(0),
125                deletion_vectors_applied: AtomicU64::new(0),
126                deletion_vector_rows_deleted: AtomicU64::new(0),
127                deletion_vector_failures: AtomicU64::new(0),
128                deletion_vector_rejections: AtomicU64::new(0),
129                parquet_data_file_range_get_operations: AtomicU64::new(0),
130                parquet_data_file_full_get_operations: AtomicU64::new(0),
131                parquet_data_file_bytes_received: AtomicU64::new(0),
132                parquet_data_file_opened_bytes: AtomicU64::new(0),
133            }),
134        }
135    }
136
137    /// Returns an immutable point-in-time copy of all scan metrics.
138    pub fn snapshot(&self) -> DeltaReadMetricsSnapshot {
139        let inner = self.inner.as_ref();
140        DeltaReadMetricsSnapshot {
141            snapshot_version: inner.snapshot_version,
142            reader_backend: inner.reader_backend,
143            scan_metadata_exhausted: inner.scan_metadata_exhausted,
144            scan_partitions_planned: load(&inner.scan_partitions_planned),
145            files_planned: inner.files_planned,
146            files_filtered_during_planning: inner.files_filtered_during_planning,
147            estimated_rows: inner.estimated_rows,
148            estimated_bytes: inner.estimated_bytes,
149            scan_partitions_started: load(&inner.scan_partitions_started),
150            scan_partitions_completed: load(&inner.scan_partitions_completed),
151            files_started: load(&inner.files_started),
152            files_completed: load(&inner.files_completed),
153            batches_produced: load(&inner.batches_produced),
154            rows_produced: load(&inner.rows_produced),
155            deletion_vector_payloads_loaded: load(&inner.deletion_vector_payloads_loaded),
156            deletion_vectors_applied: load(&inner.deletion_vectors_applied),
157            deletion_vector_rows_deleted: load(&inner.deletion_vector_rows_deleted),
158            deletion_vector_failures: load(&inner.deletion_vector_failures),
159            deletion_vector_rejections: load(&inner.deletion_vector_rejections),
160            parquet_data_file_range_get_operations: self
161                .parquet_metric(&inner.parquet_data_file_range_get_operations),
162            parquet_data_file_full_get_operations: self
163                .parquet_metric(&inner.parquet_data_file_full_get_operations),
164            parquet_data_file_bytes_received: self
165                .parquet_metric(&inner.parquet_data_file_bytes_received),
166            parquet_data_file_opened_bytes: self
167                .parquet_metric(&inner.parquet_data_file_opened_bytes),
168        }
169    }
170
171    fn parquet_metric(&self, counter: &AtomicU64) -> Option<u64> {
172        match self.inner.reader_backend {
173            DeltaReaderBackend::NativeAsync => Some(load(counter)),
174            DeltaReaderBackend::OfficialKernel => None,
175        }
176    }
177
178    #[allow(dead_code)]
179    pub(crate) fn record_scan_partitions_planned(&self, value: usize) {
180        self.inner
181            .scan_partitions_planned
182            .store(usize_to_u64_saturating(value), Ordering::Relaxed);
183    }
184
185    #[allow(dead_code)]
186    pub(crate) fn record_scan_partition_started(&self) {
187        saturating_fetch_add(&self.inner.scan_partitions_started, 1);
188    }
189
190    #[allow(dead_code)]
191    pub(crate) fn record_scan_partition_completed(&self) {
192        saturating_fetch_add(&self.inner.scan_partitions_completed, 1);
193    }
194
195    #[allow(dead_code)]
196    pub(crate) fn record_file_started(&self) {
197        saturating_fetch_add(&self.inner.files_started, 1);
198    }
199
200    #[allow(dead_code)]
201    pub(crate) fn record_file_completed(&self) {
202        saturating_fetch_add(&self.inner.files_completed, 1);
203    }
204
205    #[allow(dead_code)]
206    pub(crate) fn record_batch_produced(&self, rows: usize) {
207        saturating_fetch_add(&self.inner.batches_produced, 1);
208        saturating_fetch_add(&self.inner.rows_produced, usize_to_u64_saturating(rows));
209    }
210
211    #[allow(dead_code)]
212    pub(crate) fn record_deletion_vector_payload_loaded(&self) {
213        saturating_fetch_add(&self.inner.deletion_vector_payloads_loaded, 1);
214    }
215
216    #[allow(dead_code)]
217    pub(crate) fn record_deletion_vector_applied(&self) {
218        saturating_fetch_add(&self.inner.deletion_vectors_applied, 1);
219    }
220
221    #[allow(dead_code)]
222    pub(crate) fn record_deletion_vector_rows_deleted(&self, rows: usize) {
223        saturating_fetch_add(
224            &self.inner.deletion_vector_rows_deleted,
225            usize_to_u64_saturating(rows),
226        );
227    }
228
229    #[allow(dead_code)]
230    pub(crate) fn record_deletion_vector_failure(&self) {
231        saturating_fetch_add(&self.inner.deletion_vector_failures, 1);
232    }
233
234    #[allow(dead_code)]
235    pub(crate) fn record_deletion_vector_rejection(&self) {
236        saturating_fetch_add(&self.inner.deletion_vector_rejections, 1);
237    }
238
239    #[cfg(feature = "native-async")]
240    pub(crate) fn record_parquet_data_file_range_get_operation(&self) {
241        saturating_fetch_add(&self.inner.parquet_data_file_range_get_operations, 1);
242    }
243
244    #[cfg(feature = "native-async")]
245    pub(crate) fn record_parquet_data_file_full_get_operation(&self) {
246        saturating_fetch_add(&self.inner.parquet_data_file_full_get_operations, 1);
247    }
248
249    #[cfg(feature = "native-async")]
250    pub(crate) fn record_parquet_data_file_bytes_received(&self, bytes: usize) {
251        saturating_fetch_add(
252            &self.inner.parquet_data_file_bytes_received,
253            usize_to_u64_saturating(bytes),
254        );
255    }
256
257    #[cfg(feature = "native-async")]
258    pub(crate) fn record_parquet_data_file_opened_bytes(&self, bytes: u64) {
259        saturating_fetch_add(&self.inner.parquet_data_file_opened_bytes, bytes);
260    }
261}
262
263fn load(counter: &AtomicU64) -> u64 {
264    counter.load(Ordering::Relaxed)
265}
266
267#[allow(dead_code)]
268pub(crate) fn saturating_fetch_add(counter: &AtomicU64, value: u64) {
269    let _ = counter.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |current| {
270        Some(current.saturating_add(value))
271    });
272}
273
274fn usize_to_u64_saturating(value: usize) -> u64 {
275    u64::try_from(value).unwrap_or(u64::MAX)
276}
277
278#[cfg(test)]
279mod tests {
280    use std::{sync::atomic::Ordering, thread};
281
282    use super::{DeltaReadMetrics, DeltaReadMetricsConfig, saturating_fetch_add};
283    use crate::DeltaReaderBackend;
284
285    fn metrics(reader_backend: DeltaReaderBackend) -> DeltaReadMetrics {
286        DeltaReadMetrics::new(DeltaReadMetricsConfig {
287            snapshot_version: 7,
288            reader_backend,
289            scan_metadata_exhausted: Some(true),
290            scan_partitions_planned: 3,
291            files_planned: 5,
292            files_filtered_during_planning: Some(2),
293            estimated_rows: Some(99),
294            estimated_bytes: Some(42),
295        })
296    }
297
298    #[test]
299    fn snapshot_has_context_zeroes_and_backend_availability() {
300        let native = metrics(DeltaReaderBackend::NativeAsync).snapshot();
301        assert_eq!(native.snapshot_version, 7);
302        assert_eq!(native.reader_backend, DeltaReaderBackend::NativeAsync);
303        assert_eq!(native.scan_metadata_exhausted, Some(true));
304        assert_eq!(native.scan_partitions_planned, 3);
305        assert_eq!(native.files_planned, 5);
306        assert_eq!(native.files_filtered_during_planning, Some(2));
307        assert_eq!(native.estimated_rows, Some(99));
308        assert_eq!(native.estimated_bytes, Some(42));
309        assert_eq!(native.scan_partitions_started, 0);
310        assert_eq!(native.scan_partitions_completed, 0);
311        assert_eq!(native.files_started, 0);
312        assert_eq!(native.files_completed, 0);
313        assert_eq!(native.batches_produced, 0);
314        assert_eq!(native.rows_produced, 0);
315        assert_eq!(native.deletion_vector_payloads_loaded, 0);
316        assert_eq!(native.deletion_vectors_applied, 0);
317        assert_eq!(native.deletion_vector_rows_deleted, 0);
318        assert_eq!(native.deletion_vector_failures, 0);
319        assert_eq!(native.deletion_vector_rejections, 0);
320        assert_eq!(native.parquet_data_file_range_get_operations, Some(0));
321        assert_eq!(native.parquet_data_file_full_get_operations, Some(0));
322        assert_eq!(native.parquet_data_file_bytes_received, Some(0));
323        assert_eq!(native.parquet_data_file_opened_bytes, Some(0));
324
325        let official = metrics(DeltaReaderBackend::OfficialKernel).snapshot();
326        assert_eq!(official.parquet_data_file_range_get_operations, None);
327        assert_eq!(official.parquet_data_file_full_get_operations, None);
328        assert_eq!(official.parquet_data_file_bytes_received, None);
329        assert_eq!(official.parquet_data_file_opened_bytes, None);
330    }
331
332    #[test]
333    fn snapshot_maps_live_counters() {
334        let metrics = metrics(DeltaReaderBackend::NativeAsync);
335        metrics.record_scan_partitions_planned(16);
336        let counters = [
337            &metrics.inner.scan_partitions_started,
338            &metrics.inner.scan_partitions_completed,
339            &metrics.inner.files_started,
340            &metrics.inner.files_completed,
341            &metrics.inner.batches_produced,
342            &metrics.inner.rows_produced,
343            &metrics.inner.deletion_vector_payloads_loaded,
344            &metrics.inner.deletion_vectors_applied,
345            &metrics.inner.deletion_vector_rows_deleted,
346            &metrics.inner.deletion_vector_failures,
347            &metrics.inner.deletion_vector_rejections,
348            &metrics.inner.parquet_data_file_range_get_operations,
349            &metrics.inner.parquet_data_file_full_get_operations,
350            &metrics.inner.parquet_data_file_bytes_received,
351            &metrics.inner.parquet_data_file_opened_bytes,
352        ];
353        for (index, counter) in counters.into_iter().enumerate() {
354            saturating_fetch_add(counter, u64::try_from(index + 1).expect("small test value"));
355        }
356
357        let snapshot = metrics.snapshot();
358        assert_eq!(snapshot.scan_partitions_planned, 16);
359        assert_eq!(snapshot.scan_partitions_started, 1);
360        assert_eq!(snapshot.scan_partitions_completed, 2);
361        assert_eq!(snapshot.files_started, 3);
362        assert_eq!(snapshot.files_completed, 4);
363        assert_eq!(snapshot.batches_produced, 5);
364        assert_eq!(snapshot.rows_produced, 6);
365        assert_eq!(snapshot.deletion_vector_payloads_loaded, 7);
366        assert_eq!(snapshot.deletion_vectors_applied, 8);
367        assert_eq!(snapshot.deletion_vector_rows_deleted, 9);
368        assert_eq!(snapshot.deletion_vector_failures, 10);
369        assert_eq!(snapshot.deletion_vector_rejections, 11);
370        assert_eq!(snapshot.parquet_data_file_range_get_operations, Some(12));
371        assert_eq!(snapshot.parquet_data_file_full_get_operations, Some(13));
372        assert_eq!(snapshot.parquet_data_file_bytes_received, Some(14));
373        assert_eq!(snapshot.parquet_data_file_opened_bytes, Some(15));
374    }
375
376    #[test]
377    fn cloned_handles_saturate_under_concurrent_updates() -> Result<(), &'static str> {
378        let metrics = metrics(DeltaReaderBackend::NativeAsync);
379        metrics
380            .inner
381            .files_started
382            .store(u64::MAX - 1, Ordering::Relaxed);
383        let workers = (0..4)
384            .map(|_| {
385                let metrics = metrics.clone();
386                thread::spawn(move || {
387                    saturating_fetch_add(&metrics.inner.files_started, 1);
388                })
389            })
390            .collect::<Vec<_>>();
391
392        for worker in workers {
393            worker.join().map_err(|_| "metrics worker panicked")?;
394        }
395
396        assert_eq!(metrics.snapshot().files_started, u64::MAX);
397        Ok(())
398    }
399}