1use std::sync::{
2 Arc,
3 atomic::{AtomicU64, Ordering},
4};
5
6use crate::DeltaReaderBackend;
7
8#[derive(Debug, Clone, PartialEq, Eq)]
10pub struct DeltaReadMetricsSnapshot {
11 pub snapshot_version: u64,
13 pub reader_backend: DeltaReaderBackend,
15 pub scan_metadata_exhausted: Option<bool>,
17 pub scan_partitions_planned: u64,
19 pub files_planned: u64,
21 pub files_filtered_during_planning: Option<u64>,
23 pub estimated_rows: Option<u64>,
25 pub estimated_bytes: Option<u64>,
27 pub scan_partitions_started: u64,
29 pub scan_partitions_completed: u64,
31 pub files_started: u64,
33 pub files_completed: u64,
35 pub batches_produced: u64,
37 pub rows_produced: u64,
39 pub deletion_vector_payloads_loaded: u64,
41 pub deletion_vectors_applied: u64,
43 pub deletion_vector_rows_deleted: u64,
45 pub deletion_vector_failures: u64,
47 pub deletion_vector_rejections: u64,
49 pub parquet_data_file_range_get_operations: Option<u64>,
51 pub parquet_data_file_full_get_operations: Option<u64>,
53 pub parquet_data_file_bytes_received: Option<u64>,
55 pub parquet_data_file_opened_bytes: Option<u64>,
57}
58
59#[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 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}