1use std::{
4 fmt,
5 sync::{
6 Arc,
7 atomic::{AtomicU64, Ordering},
8 },
9};
10
11use super::options::ParquetReaderBackend;
12
13#[derive(Debug, Clone, PartialEq, Eq)]
15pub struct DeltaScanMetricsSnapshot {
16 pub snapshot_version: u64,
18 pub parquet_backend: ParquetReaderBackend,
20 pub scan_partitions_planned: u64,
22 pub files_planned: u64,
24 pub add_actions_excluded_during_planning: Option<u64>,
26 pub estimated_input_rows: Option<u64>,
28 pub estimated_input_bytes: Option<u64>,
30 pub scan_partitions_started: u64,
32 pub scan_partitions_completed: u64,
34 pub file_tasks_started: u64,
36 pub file_tasks_completed: u64,
38 pub scheduler_batches_emitted: u64,
40 pub scheduler_rows_emitted: u64,
42 pub deletion_vector_payloads_loaded: u64,
44 pub deletion_vectors_applied: u64,
46 pub deletion_vector_rows_deleted: u64,
48 pub deletion_vector_failures: u64,
50 pub deletion_vector_coordinate_rejections: u64,
52 pub parquet_data_file_range_get_operations: Option<u64>,
54 pub parquet_data_file_full_get_operations: Option<u64>,
56 pub parquet_data_file_bytes_received: Option<u64>,
58 pub estimated_parquet_task_bytes_admitted: Option<u64>,
60}
61
62#[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 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}