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: 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 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}