1use std::collections::HashMap;
19
20use itertools::Itertools;
21
22use super::{DataContentType, DataFile, PartitionSpecRef};
23use crate::error::invalid_data;
24use crate::spec::{ManifestContentType, ManifestFile, Operation, SchemaRef, Summary};
25use crate::{Error, ErrorKind, Result};
26
27const ADDED_DATA_FILES: &str = "added-data-files";
28const ADDED_DELETE_FILES: &str = "added-delete-files";
29const ADDED_EQUALITY_DELETES: &str = "added-equality-deletes";
30const ADDED_FILE_SIZE: &str = "added-files-size";
31const ADDED_POSITION_DELETES: &str = "added-position-deletes";
32const ADDED_POSITION_DELETE_FILES: &str = "added-position-delete-files";
33const ADDED_RECORDS: &str = "added-records";
34const DELETED_DATA_FILES: &str = "deleted-data-files";
35const DELETED_RECORDS: &str = "deleted-records";
36const ADDED_EQUALITY_DELETE_FILES: &str = "added-equality-delete-files";
37const REMOVED_DELETE_FILES: &str = "removed-delete-files";
38const REMOVED_EQUALITY_DELETES: &str = "removed-equality-deletes";
39const REMOVED_EQUALITY_DELETE_FILES: &str = "removed-equality-delete-files";
40const REMOVED_FILE_SIZE: &str = "removed-files-size";
41const REMOVED_POSITION_DELETES: &str = "removed-position-deletes";
42const REMOVED_POSITION_DELETE_FILES: &str = "removed-position-delete-files";
43const TOTAL_EQUALITY_DELETES: &str = "total-equality-deletes";
44const TOTAL_POSITION_DELETES: &str = "total-position-deletes";
45const TOTAL_DATA_FILES: &str = "total-data-files";
46const TOTAL_DELETE_FILES: &str = "total-delete-files";
47const TOTAL_RECORDS: &str = "total-records";
48const TOTAL_FILE_SIZE: &str = "total-files-size";
49const CHANGED_PARTITION_COUNT_PROP: &str = "changed-partition-count";
50const CHANGED_PARTITION_PREFIX: &str = "partitions.";
51
52#[derive(Default)]
56pub struct SnapshotSummaryCollector {
57 metrics: UpdateMetrics,
58 partition_metrics: HashMap<String, UpdateMetrics>,
59 max_changed_partitions_for_summaries: u64,
60 properties: HashMap<String, String>,
61 trust_partition_metrics: bool,
62}
63
64impl SnapshotSummaryCollector {
65 pub fn set(&mut self, key: &str, value: &str) {
67 self.properties.insert(key.to_string(), value.to_string());
68 }
69
70 pub fn set_partition_summary_limit(&mut self, limit: u64) {
73 self.max_changed_partitions_for_summaries = limit;
74 }
75
76 pub fn add_file(
78 &mut self,
79 data_file: &DataFile,
80 schema: SchemaRef,
81 partition_spec: PartitionSpecRef,
82 ) {
83 self.metrics.add_file(data_file);
84 if !data_file.partition.fields().is_empty() {
85 self.update_partition_metrics(schema, partition_spec, data_file, true);
86 }
87 }
88
89 pub fn remove_file(
91 &mut self,
92 data_file: &DataFile,
93 schema: SchemaRef,
94 partition_spec: PartitionSpecRef,
95 ) {
96 self.metrics.remove_file(data_file);
97 if !data_file.partition.fields().is_empty() {
98 self.update_partition_metrics(schema, partition_spec, data_file, false);
99 }
100 }
101
102 pub fn add_manifest(&mut self, manifest: &ManifestFile) {
104 self.trust_partition_metrics = false;
105 self.partition_metrics.clear();
106 self.metrics.add_manifest(manifest);
107 }
108
109 pub fn update_partition_metrics(
111 &mut self,
112 schema: SchemaRef,
113 partition_spec: PartitionSpecRef,
114 data_file: &DataFile,
115 is_add_file: bool,
116 ) {
117 let partition_path = partition_spec.partition_to_path(&data_file.partition, schema);
118 let metrics = self.partition_metrics.entry(partition_path).or_default();
119
120 if is_add_file {
121 metrics.add_file(data_file);
122 } else {
123 metrics.remove_file(data_file);
124 }
125 }
126
127 pub fn merge(&mut self, summary: SnapshotSummaryCollector) {
129 self.metrics.merge(&summary.metrics);
130 self.properties.extend(summary.properties);
131
132 if self.trust_partition_metrics && summary.trust_partition_metrics {
133 for (partition, partition_metric) in summary.partition_metrics.iter() {
134 self.partition_metrics
135 .entry(partition.to_string())
136 .or_default()
137 .merge(partition_metric);
138 }
139 } else {
140 self.partition_metrics.clear();
141 self.trust_partition_metrics = false;
142 }
143 }
144
145 pub fn build(&self) -> HashMap<String, String> {
147 let mut properties = self.metrics.to_map();
148 let changed_partitions_count = self.partition_metrics.len() as u64;
149 set_if_positive(
150 &mut properties,
151 changed_partitions_count,
152 CHANGED_PARTITION_COUNT_PROP,
153 );
154
155 if changed_partitions_count <= self.max_changed_partitions_for_summaries {
156 for (partition_path, update_metrics_partition) in &self.partition_metrics {
157 let property_key = format!("{CHANGED_PARTITION_PREFIX}{partition_path}");
158 let partition_summary = update_metrics_partition
159 .to_map()
160 .into_iter()
161 .map(|(property, value)| format!("{property}={value}"))
162 .join(",");
163
164 if !partition_summary.is_empty() {
165 properties.insert(property_key, partition_summary);
166 }
167 }
168 }
169 properties
170 }
171}
172
173#[derive(Debug, Default)]
174struct UpdateMetrics {
175 added_file_size: u64,
176 removed_file_size: u64,
177 added_data_files: u32,
178 removed_data_files: u32,
179 added_eq_delete_files: u64,
180 removed_eq_delete_files: u64,
181 added_pos_delete_files: u64,
182 removed_pos_delete_files: u64,
183 added_delete_files: u32,
184 removed_delete_files: u32,
185 added_records: u64,
186 deleted_records: u64,
187 added_pos_deletes: u64,
188 removed_pos_deletes: u64,
189 added_eq_deletes: u64,
190 removed_eq_deletes: u64,
191}
192
193impl UpdateMetrics {
194 fn add_file(&mut self, data_file: &DataFile) {
195 self.added_file_size += data_file.file_size_in_bytes;
196 match data_file.content_type() {
197 DataContentType::Data => {
198 self.added_data_files += 1;
199 self.added_records += data_file.record_count;
200 }
201 DataContentType::PositionDeletes => {
202 self.added_delete_files += 1;
203 self.added_pos_delete_files += 1;
204 self.added_pos_deletes += data_file.record_count;
205 }
206 DataContentType::EqualityDeletes => {
207 self.added_delete_files += 1;
208 self.added_eq_delete_files += 1;
209 self.added_eq_deletes += data_file.record_count;
210 }
211 }
212 }
213
214 fn remove_file(&mut self, data_file: &DataFile) {
215 self.removed_file_size += data_file.file_size_in_bytes;
216 match data_file.content_type() {
217 DataContentType::Data => {
218 self.removed_data_files += 1;
219 self.deleted_records += data_file.record_count;
220 }
221 DataContentType::PositionDeletes => {
222 self.removed_delete_files += 1;
223 self.removed_pos_delete_files += 1;
224 self.removed_pos_deletes += data_file.record_count;
225 }
226 DataContentType::EqualityDeletes => {
227 self.removed_delete_files += 1;
228 self.removed_eq_delete_files += 1;
229 self.removed_eq_deletes += data_file.record_count;
230 }
231 }
232 }
233
234 fn add_manifest(&mut self, manifest: &ManifestFile) {
235 match manifest.content {
236 ManifestContentType::Data => {
237 self.added_data_files += manifest.added_files_count.unwrap_or(0);
238 self.added_records += manifest.added_rows_count.unwrap_or(0);
239 self.removed_data_files += manifest.deleted_files_count.unwrap_or(0);
240 self.deleted_records += manifest.deleted_rows_count.unwrap_or(0);
241 }
242 ManifestContentType::Deletes => {
243 self.added_delete_files += manifest.added_files_count.unwrap_or(0);
244 self.removed_delete_files += manifest.deleted_files_count.unwrap_or(0);
245 }
246 }
247 }
248
249 fn to_map(&self) -> HashMap<String, String> {
250 let mut properties = HashMap::new();
251 set_if_positive(&mut properties, self.added_file_size, ADDED_FILE_SIZE);
252 set_if_positive(&mut properties, self.removed_file_size, REMOVED_FILE_SIZE);
253 set_if_positive(&mut properties, self.added_data_files, ADDED_DATA_FILES);
254 set_if_positive(&mut properties, self.removed_data_files, DELETED_DATA_FILES);
255 set_if_positive(
256 &mut properties,
257 self.added_eq_delete_files,
258 ADDED_EQUALITY_DELETE_FILES,
259 );
260 set_if_positive(
261 &mut properties,
262 self.removed_eq_delete_files,
263 REMOVED_EQUALITY_DELETE_FILES,
264 );
265 set_if_positive(
266 &mut properties,
267 self.added_pos_delete_files,
268 ADDED_POSITION_DELETE_FILES,
269 );
270 set_if_positive(
271 &mut properties,
272 self.removed_pos_delete_files,
273 REMOVED_POSITION_DELETE_FILES,
274 );
275 set_if_positive(&mut properties, self.added_delete_files, ADDED_DELETE_FILES);
276 set_if_positive(
277 &mut properties,
278 self.removed_delete_files,
279 REMOVED_DELETE_FILES,
280 );
281 set_if_positive(&mut properties, self.added_records, ADDED_RECORDS);
282 set_if_positive(&mut properties, self.deleted_records, DELETED_RECORDS);
283 set_if_positive(
284 &mut properties,
285 self.added_pos_deletes,
286 ADDED_POSITION_DELETES,
287 );
288 set_if_positive(
289 &mut properties,
290 self.removed_pos_deletes,
291 REMOVED_POSITION_DELETES,
292 );
293 set_if_positive(
294 &mut properties,
295 self.added_eq_deletes,
296 ADDED_EQUALITY_DELETES,
297 );
298 set_if_positive(
299 &mut properties,
300 self.removed_eq_deletes,
301 REMOVED_EQUALITY_DELETES,
302 );
303 properties
304 }
305
306 fn merge(&mut self, other: &UpdateMetrics) {
307 self.added_file_size += other.added_file_size;
308 self.removed_file_size += other.removed_file_size;
309 self.added_data_files += other.added_data_files;
310 self.removed_data_files += other.removed_data_files;
311 self.added_eq_delete_files += other.added_eq_delete_files;
312 self.removed_eq_delete_files += other.removed_eq_delete_files;
313 self.added_pos_delete_files += other.added_pos_delete_files;
314 self.removed_pos_delete_files += other.removed_pos_delete_files;
315 self.added_delete_files += other.added_delete_files;
316 self.removed_delete_files += other.removed_delete_files;
317 self.added_records += other.added_records;
318 self.deleted_records += other.deleted_records;
319 self.added_pos_deletes += other.added_pos_deletes;
320 self.removed_pos_deletes += other.removed_pos_deletes;
321 self.added_eq_deletes += other.added_eq_deletes;
322 self.removed_eq_deletes += other.removed_eq_deletes;
323 }
324}
325
326fn set_if_positive<T>(properties: &mut HashMap<String, String>, value: T, property_name: &str)
327where T: PartialOrd + Default + ToString {
328 if value > T::default() {
329 properties.insert(property_name.to_string(), value.to_string());
330 }
331}
332
333pub(crate) fn update_snapshot_summaries(
334 summary: Summary,
335 previous_summary: Option<&Summary>,
336 truncate_full_table: bool,
337) -> Result<Summary> {
338 if summary.operation != Operation::Append
340 && summary.operation != Operation::Overwrite
341 && summary.operation != Operation::Delete
342 {
343 return Err(invalid_data!("Operation is not supported."));
344 }
345
346 let mut summary = match previous_summary {
347 Some(prev_summary) if truncate_full_table && summary.operation == Operation::Overwrite => {
348 truncate_table_summary(summary, prev_summary).map_err(|err| {
349 Error::new(ErrorKind::Unexpected, "Failed to truncate table summary.")
350 .with_source(err)
351 })?
352 }
353 _ => summary,
354 };
355
356 update_totals(
357 &mut summary,
358 previous_summary,
359 TOTAL_DATA_FILES,
360 ADDED_DATA_FILES,
361 DELETED_DATA_FILES,
362 );
363
364 update_totals(
365 &mut summary,
366 previous_summary,
367 TOTAL_DELETE_FILES,
368 ADDED_DELETE_FILES,
369 REMOVED_DELETE_FILES,
370 );
371
372 update_totals(
373 &mut summary,
374 previous_summary,
375 TOTAL_RECORDS,
376 ADDED_RECORDS,
377 DELETED_RECORDS,
378 );
379
380 update_totals(
381 &mut summary,
382 previous_summary,
383 TOTAL_FILE_SIZE,
384 ADDED_FILE_SIZE,
385 REMOVED_FILE_SIZE,
386 );
387
388 update_totals(
389 &mut summary,
390 previous_summary,
391 TOTAL_POSITION_DELETES,
392 ADDED_POSITION_DELETES,
393 REMOVED_POSITION_DELETES,
394 );
395
396 update_totals(
397 &mut summary,
398 previous_summary,
399 TOTAL_EQUALITY_DELETES,
400 ADDED_EQUALITY_DELETES,
401 REMOVED_EQUALITY_DELETES,
402 );
403 Ok(summary)
404}
405
406fn get_prop(previous_summary: &Summary, prop: &str) -> Result<u64> {
407 let value_str = previous_summary
408 .additional_properties
409 .get(prop)
410 .map(String::as_str)
411 .unwrap_or("0");
412 value_str.parse::<u64>().map_err(|err| {
413 Error::new(
414 ErrorKind::Unexpected,
415 format!("Failed to parse summary property '{prop}' value '{value_str}' as u64."),
416 )
417 .with_source(err)
418 })
419}
420
421fn truncate_table_summary(mut summary: Summary, previous_summary: &Summary) -> Result<Summary> {
422 for prop in [
423 TOTAL_DATA_FILES,
424 TOTAL_DELETE_FILES,
425 TOTAL_RECORDS,
426 TOTAL_FILE_SIZE,
427 TOTAL_POSITION_DELETES,
428 TOTAL_EQUALITY_DELETES,
429 ] {
430 summary
431 .additional_properties
432 .insert(prop.to_string(), "0".to_string());
433 }
434
435 let value = get_prop(previous_summary, TOTAL_DATA_FILES)?;
436 if value != 0 {
437 summary
438 .additional_properties
439 .insert(DELETED_DATA_FILES.to_string(), value.to_string());
440 }
441 let value = get_prop(previous_summary, TOTAL_DELETE_FILES)?;
442 if value != 0 {
443 summary
444 .additional_properties
445 .insert(REMOVED_DELETE_FILES.to_string(), value.to_string());
446 }
447 let value = get_prop(previous_summary, TOTAL_RECORDS)?;
448 if value != 0 {
449 summary
450 .additional_properties
451 .insert(DELETED_RECORDS.to_string(), value.to_string());
452 }
453 let value = get_prop(previous_summary, TOTAL_FILE_SIZE)?;
454 if value != 0 {
455 summary
456 .additional_properties
457 .insert(REMOVED_FILE_SIZE.to_string(), value.to_string());
458 }
459
460 let value = get_prop(previous_summary, TOTAL_POSITION_DELETES)?;
461 if value != 0 {
462 summary
463 .additional_properties
464 .insert(REMOVED_POSITION_DELETES.to_string(), value.to_string());
465 }
466
467 let value = get_prop(previous_summary, TOTAL_EQUALITY_DELETES)?;
468 if value != 0 {
469 summary
470 .additional_properties
471 .insert(REMOVED_EQUALITY_DELETES.to_string(), value.to_string());
472 }
473
474 Ok(summary)
475}
476
477fn update_totals(
478 summary: &mut Summary,
479 previous_summary: Option<&Summary>,
480 total_property: &str,
481 added_property: &str,
482 removed_property: &str,
483) {
484 let previous_total = match previous_summary {
485 None => 0,
486 Some(prev_summary) => match prev_summary.additional_properties.get(total_property) {
487 Some(value_str) => match value_str.parse::<u64>() {
488 Ok(v) => v,
489 Err(parse_err) => {
490 tracing::warn!(
491 "Property '{total_property}' could not be parsed in the previous snapshot summary: {parse_err}. \
492 Skipping total computation.",
493 );
494 return;
495 }
496 },
497 None => {
498 tracing::debug!(
499 "Property '{total_property}' was not set in the previous snapshot summary. \
500 Skipping total computation."
501 );
502 return;
503 }
504 },
505 };
506
507 let parse_delta = |property: &str| -> Option<u64> {
513 match summary.additional_properties.get(property) {
514 None => Some(0),
515 Some(value) => match value.parse::<u64>() {
516 Ok(v) => Some(v),
517 Err(parse_err) => {
518 tracing::warn!(
519 "Property '{property}' could not be parsed when computing '{total_property}': {parse_err}. \
520 Skipping total computation.",
521 );
522 None
523 }
524 },
525 }
526 };
527
528 let (Some(added), Some(removed)) = (parse_delta(added_property), parse_delta(removed_property))
529 else {
530 return;
531 };
532
533 let new_total = previous_total + added - removed;
534 summary
535 .additional_properties
536 .insert(total_property.to_string(), new_total.to_string());
537}
538
539#[cfg(test)]
540mod tests {
541 use std::collections::HashMap;
542 use std::sync::Arc;
543
544 use super::*;
545 use crate::spec::{
546 DataFileFormat, Datum, Literal, NestedField, PartitionSpec, PrimitiveType, Schema, Struct,
547 Transform, Type, UnboundPartitionField,
548 };
549
550 #[test]
551 fn test_update_snapshot_summaries_append() {
552 let prev_props: HashMap<String, String> = [
553 (TOTAL_DATA_FILES.to_string(), "10".to_string()),
554 (TOTAL_DELETE_FILES.to_string(), "5".to_string()),
555 (TOTAL_RECORDS.to_string(), "100".to_string()),
556 (TOTAL_FILE_SIZE.to_string(), "1000".to_string()),
557 (TOTAL_POSITION_DELETES.to_string(), "3".to_string()),
558 (TOTAL_EQUALITY_DELETES.to_string(), "2".to_string()),
559 ]
560 .into_iter()
561 .collect();
562
563 let previous_summary = Summary {
564 operation: Operation::Append,
565 additional_properties: prev_props,
566 };
567
568 let new_props: HashMap<String, String> = [
569 (ADDED_DATA_FILES.to_string(), "4".to_string()),
570 (DELETED_DATA_FILES.to_string(), "1".to_string()),
571 (ADDED_DELETE_FILES.to_string(), "2".to_string()),
572 (REMOVED_DELETE_FILES.to_string(), "1".to_string()),
573 (ADDED_RECORDS.to_string(), "40".to_string()),
574 (DELETED_RECORDS.to_string(), "10".to_string()),
575 (ADDED_FILE_SIZE.to_string(), "400".to_string()),
576 (REMOVED_FILE_SIZE.to_string(), "100".to_string()),
577 (ADDED_POSITION_DELETES.to_string(), "5".to_string()),
578 (REMOVED_POSITION_DELETES.to_string(), "2".to_string()),
579 (ADDED_EQUALITY_DELETES.to_string(), "3".to_string()),
580 (REMOVED_EQUALITY_DELETES.to_string(), "1".to_string()),
581 ]
582 .into_iter()
583 .collect();
584
585 let summary = Summary {
586 operation: Operation::Append,
587 additional_properties: new_props,
588 };
589
590 let updated = update_snapshot_summaries(summary, Some(&previous_summary), false).unwrap();
591
592 assert_eq!(
593 updated.additional_properties.get(TOTAL_DATA_FILES).unwrap(),
594 "13"
595 );
596 assert_eq!(
597 updated
598 .additional_properties
599 .get(TOTAL_DELETE_FILES)
600 .unwrap(),
601 "6"
602 );
603 assert_eq!(
604 updated.additional_properties.get(TOTAL_RECORDS).unwrap(),
605 "130"
606 );
607 assert_eq!(
608 updated.additional_properties.get(TOTAL_FILE_SIZE).unwrap(),
609 "1300"
610 );
611 assert_eq!(
612 updated
613 .additional_properties
614 .get(TOTAL_POSITION_DELETES)
615 .unwrap(),
616 "6"
617 );
618 assert_eq!(
619 updated
620 .additional_properties
621 .get(TOTAL_EQUALITY_DELETES)
622 .unwrap(),
623 "4"
624 );
625 }
626
627 #[test]
628 fn test_truncate_table_summary() {
629 let prev_props: HashMap<String, String> = [
630 (TOTAL_DATA_FILES.to_string(), "10".to_string()),
631 (TOTAL_DELETE_FILES.to_string(), "5".to_string()),
632 (TOTAL_RECORDS.to_string(), "100".to_string()),
633 (TOTAL_FILE_SIZE.to_string(), "1000".to_string()),
634 (TOTAL_POSITION_DELETES.to_string(), "3".to_string()),
635 (TOTAL_EQUALITY_DELETES.to_string(), "2".to_string()),
636 ]
637 .into_iter()
638 .collect();
639
640 let previous_summary = Summary {
641 operation: Operation::Overwrite,
642 additional_properties: prev_props,
643 };
644
645 let mut new_props = HashMap::new();
646 new_props.insert("dummy".to_string(), "value".to_string());
647 let summary = Summary {
648 operation: Operation::Overwrite,
649 additional_properties: new_props,
650 };
651
652 let truncated = truncate_table_summary(summary, &previous_summary).unwrap();
653
654 assert_eq!(
655 truncated
656 .additional_properties
657 .get(TOTAL_DATA_FILES)
658 .unwrap(),
659 "0"
660 );
661 assert_eq!(
662 truncated
663 .additional_properties
664 .get(TOTAL_DELETE_FILES)
665 .unwrap(),
666 "0"
667 );
668 assert_eq!(
669 truncated.additional_properties.get(TOTAL_RECORDS).unwrap(),
670 "0"
671 );
672 assert_eq!(
673 truncated
674 .additional_properties
675 .get(TOTAL_FILE_SIZE)
676 .unwrap(),
677 "0"
678 );
679 assert_eq!(
680 truncated
681 .additional_properties
682 .get(TOTAL_POSITION_DELETES)
683 .unwrap(),
684 "0"
685 );
686 assert_eq!(
687 truncated
688 .additional_properties
689 .get(TOTAL_EQUALITY_DELETES)
690 .unwrap(),
691 "0"
692 );
693
694 assert_eq!(
695 truncated
696 .additional_properties
697 .get(DELETED_DATA_FILES)
698 .unwrap(),
699 "10"
700 );
701 assert_eq!(
702 truncated
703 .additional_properties
704 .get(REMOVED_DELETE_FILES)
705 .unwrap(),
706 "5"
707 );
708 assert_eq!(
709 truncated
710 .additional_properties
711 .get(DELETED_RECORDS)
712 .unwrap(),
713 "100"
714 );
715 assert_eq!(
716 truncated
717 .additional_properties
718 .get(REMOVED_FILE_SIZE)
719 .unwrap(),
720 "1000"
721 );
722 assert_eq!(
723 truncated
724 .additional_properties
725 .get(REMOVED_POSITION_DELETES)
726 .unwrap(),
727 "3"
728 );
729 assert_eq!(
730 truncated
731 .additional_properties
732 .get(REMOVED_EQUALITY_DELETES)
733 .unwrap(),
734 "2"
735 );
736 }
737
738 #[test]
739 fn test_update_snapshot_summaries_overwrite_truncate_handles_totals_above_i32_max() {
740 let big = (i32::MAX as u64 + 1).to_string(); let prev_props: HashMap<String, String> = [
745 (TOTAL_DATA_FILES.to_string(), big.clone()),
746 (TOTAL_DELETE_FILES.to_string(), "0".to_string()),
747 (TOTAL_RECORDS.to_string(), big.clone()),
748 (TOTAL_FILE_SIZE.to_string(), "0".to_string()),
749 (TOTAL_POSITION_DELETES.to_string(), "0".to_string()),
750 (TOTAL_EQUALITY_DELETES.to_string(), "0".to_string()),
751 ]
752 .into_iter()
753 .collect();
754
755 let previous_summary = Summary {
756 operation: Operation::Overwrite,
757 additional_properties: prev_props,
758 };
759
760 let summary = Summary {
761 operation: Operation::Overwrite,
762 additional_properties: HashMap::new(),
763 };
764
765 let updated = update_snapshot_summaries(summary, Some(&previous_summary), true)
766 .expect("overwrite truncation should accept totals above i32::MAX");
767 assert_eq!(
768 updated
769 .additional_properties
770 .get(DELETED_DATA_FILES)
771 .unwrap(),
772 &big
773 );
774 assert_eq!(
775 updated.additional_properties.get(DELETED_RECORDS).unwrap(),
776 &big
777 );
778 }
779
780 #[test]
781 fn test_update_snapshot_summaries_overwrite_truncate_returns_err_on_malformed_total() {
782 let prev_props: HashMap<String, String> = [
786 (TOTAL_DATA_FILES.to_string(), "not_a_number".to_string()),
787 (TOTAL_DELETE_FILES.to_string(), "0".to_string()),
788 (TOTAL_RECORDS.to_string(), "0".to_string()),
789 (TOTAL_FILE_SIZE.to_string(), "0".to_string()),
790 (TOTAL_POSITION_DELETES.to_string(), "0".to_string()),
791 (TOTAL_EQUALITY_DELETES.to_string(), "0".to_string()),
792 ]
793 .into_iter()
794 .collect();
795
796 let previous_summary = Summary {
797 operation: Operation::Overwrite,
798 additional_properties: prev_props,
799 };
800
801 let summary = Summary {
802 operation: Operation::Overwrite,
803 additional_properties: HashMap::new(),
804 };
805
806 let err = update_snapshot_summaries(summary, Some(&previous_summary), true)
807 .expect_err("malformed previous summary must produce an Err, not a panic");
808 assert!(
809 err.message().contains("truncate table summary"),
810 "expected wrapped 'Failed to truncate table summary' context, got: {}",
811 err.message()
812 );
813 }
814
815 #[test]
816 fn test_snapshot_summary_collector_build() {
817 let schema = Arc::new(
818 Schema::builder()
819 .with_fields(vec![
820 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
821 NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
822 ])
823 .build()
824 .unwrap(),
825 );
826
827 let partition_spec = Arc::new(
828 PartitionSpec::builder(schema.clone())
829 .add_unbound_fields(vec![
830 UnboundPartitionField::builder()
831 .source_ids(vec![2])
832 .name("year".to_string())
833 .transform(Transform::Identity)
834 .build()
835 .unwrap(),
836 ])
837 .unwrap()
838 .with_spec_id(1)
839 .build()
840 .unwrap(),
841 );
842
843 let mut collector = SnapshotSummaryCollector::default();
844 collector.set_partition_summary_limit(10);
845
846 let file1 = DataFile {
847 content: DataContentType::Data,
848 file_path: "s3://testbucket/path/to/file1.parquet".to_string(),
849 file_format: DataFileFormat::Parquet,
850 partition: Struct::from_iter(vec![]),
851 record_count: 10,
852 file_size_in_bytes: 100,
853 column_sizes: HashMap::from([(1, 46), (2, 48), (3, 48)]),
854 value_counts: HashMap::from([(1, 10), (2, 10), (3, 10)]),
855 null_value_counts: HashMap::from([(1, 0), (2, 0), (3, 0)]),
856 nan_value_counts: HashMap::new(),
857 lower_bounds: HashMap::from([
858 (1, Datum::long(1)),
859 (2, Datum::string("a")),
860 (3, Datum::string("x")),
861 ]),
862 upper_bounds: HashMap::from([
863 (1, Datum::long(1)),
864 (2, Datum::string("a")),
865 (3, Datum::string("x")),
866 ]),
867 key_metadata: None,
868 split_offsets: Some(vec![4]),
869 equality_ids: None,
870 sort_order_id: Some(0),
871 partition_spec_id: 0,
872 first_row_id: None,
873 referenced_data_file: None,
874 content_offset: None,
875 content_size_in_bytes: None,
876 };
877
878 let file2 = DataFile {
879 content: DataContentType::Data,
880 file_path: "s3://testbucket/path/to/file2.parquet".to_string(),
881 file_format: DataFileFormat::Parquet,
882 partition: Struct::from_iter(vec![Some(Literal::string("2025"))]),
883 record_count: 20,
884 file_size_in_bytes: 200,
885 column_sizes: HashMap::from([(1, 46), (2, 48), (3, 48)]),
886 value_counts: HashMap::from([(1, 20), (2, 20), (3, 20)]),
887 null_value_counts: HashMap::from([(1, 0), (2, 0), (3, 0)]),
888 nan_value_counts: HashMap::new(),
889 lower_bounds: HashMap::from([
890 (1, Datum::long(1)),
891 (2, Datum::string("a")),
892 (3, Datum::string("x")),
893 ]),
894 upper_bounds: HashMap::from([
895 (1, Datum::long(1)),
896 (2, Datum::string("a")),
897 (3, Datum::string("x")),
898 ]),
899 key_metadata: None,
900 split_offsets: Some(vec![4]),
901 equality_ids: None,
902 sort_order_id: Some(0),
903 partition_spec_id: 0,
904 first_row_id: None,
905 referenced_data_file: None,
906 content_offset: None,
907 content_size_in_bytes: None,
908 };
909
910 collector.add_file(&file1, schema.clone(), partition_spec.clone());
911 collector.add_file(&file2, schema.clone(), partition_spec.clone());
912
913 collector.remove_file(&file1, schema.clone(), partition_spec.clone());
914
915 let props = collector.build();
916
917 assert_eq!(props.get(ADDED_FILE_SIZE).unwrap(), "300");
918 assert_eq!(props.get(REMOVED_FILE_SIZE).unwrap(), "100");
919
920 let partition_key = format!("{}{}", CHANGED_PARTITION_PREFIX, "year=2025");
921
922 assert!(props.contains_key(&partition_key));
923
924 let partition_summary = props.get(&partition_key).unwrap();
925 assert!(partition_summary.contains(&format!("{ADDED_FILE_SIZE}=200")));
926 assert!(partition_summary.contains(&format!("{ADDED_DATA_FILES}=1")));
927 assert!(partition_summary.contains(&format!("{ADDED_RECORDS}=20")));
928 }
929
930 #[test]
931 fn test_snapshot_summary_collector_add_manifest() {
932 let mut collector = SnapshotSummaryCollector::default();
933 collector.set_partition_summary_limit(10);
934
935 let manifest = ManifestFile {
936 manifest_path: "file://dummy.manifest".to_string(),
937 manifest_length: 0,
938 partition_spec_id: 0,
939 content: ManifestContentType::Data,
940 sequence_number: 0,
941 min_sequence_number: 0,
942 added_snapshot_id: 0,
943 added_files_count: Some(3),
944 existing_files_count: Some(0),
945 deleted_files_count: Some(1),
946 added_rows_count: Some(100),
947 existing_rows_count: Some(0),
948 deleted_rows_count: Some(50),
949 partitions: Some(Vec::new()),
950 key_metadata: None,
951 first_row_id: None,
952 };
953
954 collector
955 .partition_metrics
956 .insert("dummy".to_string(), UpdateMetrics::default());
957 collector.add_manifest(&manifest);
958
959 let props = collector.build();
960 assert_eq!(props.get(ADDED_DATA_FILES).unwrap(), "3");
961 assert_eq!(props.get(DELETED_DATA_FILES).unwrap(), "1");
962 assert_eq!(props.get(ADDED_RECORDS).unwrap(), "100");
963 assert_eq!(props.get(DELETED_RECORDS).unwrap(), "50");
964 }
965
966 #[test]
967 fn test_snapshot_summary_collector_merge() {
968 let schema = Arc::new(
969 Schema::builder()
970 .with_fields(vec![
971 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
972 NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
973 ])
974 .build()
975 .unwrap(),
976 );
977
978 let partition_spec = Arc::new(
979 PartitionSpec::builder(schema.clone())
980 .add_unbound_fields(vec![
981 UnboundPartitionField::builder()
982 .source_ids(vec![2])
983 .name("year".to_string())
984 .transform(Transform::Identity)
985 .build()
986 .unwrap(),
987 ])
988 .unwrap()
989 .with_spec_id(1)
990 .build()
991 .unwrap(),
992 );
993
994 let mut summary_one = SnapshotSummaryCollector::default();
995 let mut summary_two = SnapshotSummaryCollector::default();
996
997 summary_one.add_file(
998 &DataFile {
999 content: DataContentType::Data,
1000 file_path: "test.parquet".into(),
1001 file_format: DataFileFormat::Parquet,
1002 partition: Struct::from_iter(vec![]),
1003 record_count: 10,
1004 file_size_in_bytes: 100,
1005 column_sizes: HashMap::new(),
1006 value_counts: HashMap::new(),
1007 null_value_counts: HashMap::new(),
1008 nan_value_counts: HashMap::new(),
1009 lower_bounds: HashMap::new(),
1010 upper_bounds: HashMap::new(),
1011 key_metadata: None,
1012 split_offsets: None,
1013 equality_ids: None,
1014 sort_order_id: None,
1015 partition_spec_id: 0,
1016 first_row_id: None,
1017 referenced_data_file: None,
1018 content_offset: None,
1019 content_size_in_bytes: None,
1020 },
1021 schema.clone(),
1022 partition_spec.clone(),
1023 );
1024
1025 summary_two.add_file(
1026 &DataFile {
1027 content: DataContentType::Data,
1028 file_path: "test.parquet".into(),
1029 file_format: DataFileFormat::Parquet,
1030 partition: Struct::from_iter(vec![]),
1031 record_count: 20,
1032 file_size_in_bytes: 200,
1033 column_sizes: HashMap::new(),
1034 value_counts: HashMap::new(),
1035 null_value_counts: HashMap::new(),
1036 nan_value_counts: HashMap::new(),
1037 lower_bounds: HashMap::new(),
1038 upper_bounds: HashMap::new(),
1039 key_metadata: None,
1040 split_offsets: None,
1041 equality_ids: None,
1042 sort_order_id: None,
1043 partition_spec_id: 0,
1044 first_row_id: None,
1045 referenced_data_file: None,
1046 content_offset: None,
1047 content_size_in_bytes: None,
1048 },
1049 schema.clone(),
1050 partition_spec.clone(),
1051 );
1052
1053 summary_one.merge(summary_two);
1054 let props = summary_one.build();
1055 assert_eq!(props.get(ADDED_DATA_FILES).unwrap(), "2");
1056 assert_eq!(props.get(ADDED_RECORDS).unwrap(), "30");
1057
1058 let mut summary_three = SnapshotSummaryCollector::default();
1059 let mut summary_four = SnapshotSummaryCollector::default();
1060
1061 summary_three.add_manifest(&ManifestFile {
1062 manifest_path: "test.manifest".to_string(),
1063 manifest_length: 0,
1064 partition_spec_id: 0,
1065 content: ManifestContentType::Data,
1066 sequence_number: 0,
1067 min_sequence_number: 0,
1068 added_snapshot_id: 0,
1069 added_files_count: Some(1),
1070 existing_files_count: Some(0),
1071 deleted_files_count: Some(0),
1072 added_rows_count: Some(5),
1073 existing_rows_count: Some(0),
1074 deleted_rows_count: Some(0),
1075 partitions: Some(Vec::new()),
1076 key_metadata: None,
1077 first_row_id: None,
1078 });
1079
1080 summary_four.add_file(
1081 &DataFile {
1082 content: DataContentType::Data,
1083 file_path: "test.parquet".into(),
1084 file_format: DataFileFormat::Parquet,
1085 partition: Struct::from_iter(vec![]),
1086 record_count: 1,
1087 file_size_in_bytes: 10,
1088 column_sizes: HashMap::new(),
1089 value_counts: HashMap::new(),
1090 null_value_counts: HashMap::new(),
1091 nan_value_counts: HashMap::new(),
1092 lower_bounds: HashMap::new(),
1093 upper_bounds: HashMap::new(),
1094 key_metadata: None,
1095 split_offsets: None,
1096 equality_ids: None,
1097 sort_order_id: None,
1098 partition_spec_id: 0,
1099 first_row_id: None,
1100 referenced_data_file: None,
1101 content_offset: None,
1102 content_size_in_bytes: None,
1103 },
1104 schema.clone(),
1105 partition_spec.clone(),
1106 );
1107
1108 summary_three.merge(summary_four);
1109 let props = summary_three.build();
1110
1111 assert_eq!(props.get(ADDED_DATA_FILES).unwrap(), "2");
1112 assert_eq!(props.get(ADDED_RECORDS).unwrap(), "6");
1113 assert!(
1114 props
1115 .iter()
1116 .all(|(k, _)| !k.starts_with(CHANGED_PARTITION_PREFIX))
1117 );
1118 }
1119
1120 #[test]
1121 fn test_update_totals_skipped_when_previous_summary_missing_totals() {
1122 let prev_props: HashMap<String, String> = [(TOTAL_DATA_FILES, "8")]
1123 .into_iter()
1124 .map(|(k, v)| (k.to_string(), v.to_string()))
1125 .collect();
1126
1127 let previous_summary = Summary {
1128 operation: Operation::Overwrite,
1129 additional_properties: prev_props,
1130 };
1131
1132 let new_props: HashMap<String, String> = [
1133 (ADDED_DATA_FILES, "4"),
1134 (ADDED_DELETE_FILES, "2"),
1135 (ADDED_RECORDS, "40"),
1136 (ADDED_FILE_SIZE, "400"),
1137 (ADDED_POSITION_DELETES, "5"),
1138 (ADDED_EQUALITY_DELETES, "3"),
1139 ]
1140 .into_iter()
1141 .map(|(k, v)| (k.to_string(), v.to_string()))
1142 .collect();
1143
1144 let summary = Summary {
1145 operation: Operation::Append,
1146 additional_properties: new_props,
1147 };
1148
1149 let updated = update_snapshot_summaries(summary, Some(&previous_summary), false).unwrap();
1150 let props = &updated.additional_properties;
1151
1152 assert_eq!(props.get(TOTAL_DATA_FILES).unwrap(), "12");
1153
1154 for total_field in [
1155 TOTAL_DELETE_FILES,
1156 TOTAL_RECORDS,
1157 TOTAL_FILE_SIZE,
1158 TOTAL_POSITION_DELETES,
1159 TOTAL_EQUALITY_DELETES,
1160 ] {
1161 assert!(
1162 !props.contains_key(total_field),
1163 "{total_field} should not be set when previous summary lacks it",
1164 );
1165 }
1166 }
1167
1168 #[test]
1169 fn test_update_totals_tolerates_unparsable_added_value() {
1170 let prev_props: HashMap<String, String> = [(TOTAL_DATA_FILES, "8"), (TOTAL_RECORDS, "80")]
1174 .into_iter()
1175 .map(|(k, v)| (k.to_string(), v.to_string()))
1176 .collect();
1177
1178 let previous_summary = Summary {
1179 operation: Operation::Append,
1180 additional_properties: prev_props,
1181 };
1182
1183 let new_props: HashMap<String, String> =
1184 [(ADDED_DATA_FILES, "not-a-number"), (ADDED_RECORDS, "40")]
1185 .into_iter()
1186 .map(|(k, v)| (k.to_string(), v.to_string()))
1187 .collect();
1188
1189 let summary = Summary {
1190 operation: Operation::Append,
1191 additional_properties: new_props,
1192 };
1193
1194 let updated = update_snapshot_summaries(summary, Some(&previous_summary), false).unwrap();
1196 let props = &updated.additional_properties;
1197
1198 assert!(
1200 !props.contains_key(TOTAL_DATA_FILES),
1201 "TOTAL_DATA_FILES should be skipped when its added value is unparsable",
1202 );
1203 assert_eq!(props.get(TOTAL_RECORDS).unwrap(), "120");
1205 }
1206
1207 #[test]
1208 fn test_update_totals_computed_when_no_previous_summary() {
1209 let new_props: HashMap<String, String> = [
1210 (ADDED_DATA_FILES, "4"),
1211 (ADDED_RECORDS, "40"),
1212 (ADDED_FILE_SIZE, "400"),
1213 ]
1214 .into_iter()
1215 .map(|(k, v)| (k.to_string(), v.to_string()))
1216 .collect();
1217
1218 let summary = Summary {
1219 operation: Operation::Append,
1220 additional_properties: new_props,
1221 };
1222
1223 let updated = update_snapshot_summaries(summary, None, false).unwrap();
1224 let props = &updated.additional_properties;
1225
1226 assert_eq!(props.get(TOTAL_DATA_FILES).unwrap(), "4");
1227 assert_eq!(props.get(TOTAL_RECORDS).unwrap(), "40");
1228 assert_eq!(props.get(TOTAL_FILE_SIZE).unwrap(), "400");
1229 }
1230
1231 #[test]
1232 fn test_update_totals_with_removes_only() {
1233 let prev_props: HashMap<String, String> = [
1234 (TOTAL_DATA_FILES, "10"),
1235 (TOTAL_DELETE_FILES, "5"),
1236 (TOTAL_RECORDS, "100"),
1237 (TOTAL_FILE_SIZE, "1000"),
1238 (TOTAL_POSITION_DELETES, "3"),
1239 (TOTAL_EQUALITY_DELETES, "2"),
1240 ]
1241 .into_iter()
1242 .map(|(k, v)| (k.to_string(), v.to_string()))
1243 .collect();
1244
1245 let previous_summary = Summary {
1246 operation: Operation::Overwrite,
1247 additional_properties: prev_props,
1248 };
1249
1250 let new_props: HashMap<String, String> = [
1251 (DELETED_DATA_FILES, "2"),
1252 (REMOVED_DELETE_FILES, "1"),
1253 (DELETED_RECORDS, "20"),
1254 (REMOVED_FILE_SIZE, "200"),
1255 (REMOVED_POSITION_DELETES, "1"),
1256 (REMOVED_EQUALITY_DELETES, "1"),
1257 ]
1258 .into_iter()
1259 .map(|(k, v)| (k.to_string(), v.to_string()))
1260 .collect();
1261
1262 let summary = Summary {
1263 operation: Operation::Delete,
1264 additional_properties: new_props,
1265 };
1266
1267 let updated = update_snapshot_summaries(summary, Some(&previous_summary), false).unwrap();
1268 let props = &updated.additional_properties;
1269
1270 assert_eq!(props.get(TOTAL_DATA_FILES).unwrap(), "8");
1271 assert_eq!(props.get(TOTAL_DELETE_FILES).unwrap(), "4");
1272 assert_eq!(props.get(TOTAL_RECORDS).unwrap(), "80");
1273 assert_eq!(props.get(TOTAL_FILE_SIZE).unwrap(), "800");
1274 assert_eq!(props.get(TOTAL_POSITION_DELETES).unwrap(), "2");
1275 assert_eq!(props.get(TOTAL_EQUALITY_DELETES).unwrap(), "1");
1276 }
1277}