1use std::collections::{HashMap, HashSet};
19use std::sync::Arc;
20
21use async_trait::async_trait;
22use uuid::Uuid;
23
24use crate::error::Result;
25use crate::spec::{DataFile, ManifestEntry, ManifestFile, Operation};
26use crate::table::Table;
27use crate::transaction::snapshot::{
28 DefaultManifestProcess, SnapshotProduceOperation, SnapshotProducer,
29};
30use crate::transaction::{ActionCommit, TransactionAction};
31
32pub struct FastAppendAction {
34 check_duplicate: bool,
35 commit_uuid: Option<Uuid>,
37 snapshot_properties: HashMap<String, String>,
38 added_data_files: Vec<DataFile>,
39}
40
41impl FastAppendAction {
42 pub(crate) fn new() -> Self {
43 Self {
44 check_duplicate: true,
45 commit_uuid: None,
46 snapshot_properties: HashMap::default(),
47 added_data_files: vec![],
48 }
49 }
50
51 pub fn with_check_duplicate(mut self, v: bool) -> Self {
53 self.check_duplicate = v;
54 self
55 }
56
57 pub fn add_data_files(mut self, data_files: impl IntoIterator<Item = DataFile>) -> Self {
59 self.added_data_files.extend(data_files);
60 self
61 }
62
63 pub fn set_commit_uuid(mut self, commit_uuid: Uuid) -> Self {
65 self.commit_uuid = Some(commit_uuid);
66 self
67 }
68
69 pub fn set_snapshot_properties(mut self, snapshot_properties: HashMap<String, String>) -> Self {
71 self.snapshot_properties = snapshot_properties;
72 self
73 }
74
75 fn dedupe_added_files(&self) -> Vec<DataFile> {
79 let mut seen = HashSet::with_capacity(self.added_data_files.len());
80 self.added_data_files
81 .iter()
82 .filter(|data_file| seen.insert(data_file.file_path.as_str()))
83 .cloned()
84 .collect()
85 }
86}
87
88#[async_trait]
89impl TransactionAction for FastAppendAction {
90 async fn commit(self: Arc<Self>, table: &Table) -> Result<ActionCommit> {
91 let snapshot_producer = SnapshotProducer::new(
92 table,
93 self.commit_uuid.unwrap_or_else(Uuid::now_v7),
94 self.snapshot_properties.clone(),
95 self.dedupe_added_files(),
96 );
97
98 snapshot_producer.validate_added_data_files()?;
100
101 if self.check_duplicate {
103 snapshot_producer.validate_duplicate_files().await?;
104 }
105
106 snapshot_producer
107 .commit(FastAppendOperation, DefaultManifestProcess)
108 .await
109 }
110}
111
112struct FastAppendOperation;
113
114impl SnapshotProduceOperation for FastAppendOperation {
115 fn operation(&self) -> Operation {
116 Operation::Append
117 }
118
119 async fn delete_entries(
120 &self,
121 _snapshot_produce: &SnapshotProducer<'_>,
122 ) -> Result<Vec<ManifestEntry>> {
123 Ok(vec![])
124 }
125
126 async fn existing_manifest(
127 &self,
128 snapshot_produce: &SnapshotProducer<'_>,
129 ) -> Result<Vec<ManifestFile>> {
130 let Some(snapshot) = snapshot_produce.table.metadata().current_snapshot() else {
131 return Ok(vec![]);
132 };
133
134 let manifest_list = snapshot_produce
135 .table
136 .manifest_list_reader(snapshot)
137 .load()
138 .await?;
139
140 Ok(manifest_list
141 .entries()
142 .iter()
143 .filter(|entry| {
144 entry.has_added_files() || entry.has_existing_files() || entry.has_deleted_files()
148 })
149 .cloned()
150 .collect())
151 }
152}
153
154#[cfg(test)]
155mod tests {
156 use std::collections::HashMap;
157 use std::fs;
158 use std::sync::Arc;
159
160 use minijinja::{AutoEscape, Environment, Value, context};
161 use tempfile::TempDir;
162 use uuid::Uuid;
163
164 use crate::encryption::kms::MemoryKeyManagementClient;
165 use crate::encryption::{SensitiveBytes, StandardKeyMetadata};
166 use crate::io::FileIO;
167 use crate::spec::{
168 DataContentType, DataFile, DataFileBuilder, DataFileFormat, Literal, MAIN_BRANCH,
169 ManifestEntry, ManifestListWriter, ManifestStatus, ManifestWriterBuilder, SnapshotRef,
170 Struct, TableMetadata,
171 };
172 use crate::table::Table;
173 use crate::test_utils::{make_encrypted_table, test_runtime};
174 use crate::transaction::tests::make_v2_minimal_table;
175 use crate::transaction::{Transaction, TransactionAction};
176 use crate::{TableIdent, TableRequirement, TableUpdate};
177
178 fn render_template(template: &str, ctx: Value) -> String {
179 let mut env = Environment::new();
180 env.set_auto_escape_callback(|_| AutoEscape::None);
181 env.render_str(template, ctx).unwrap()
182 }
183
184 async fn make_table_with_delete_only_manifest() -> (Table, TempDir, String) {
191 let tmp_dir = TempDir::new().unwrap();
192 let table_location = tmp_dir.path().join("table1");
193 let manifest_list_location = table_location.join("metadata/manifests_list_1.avro");
194 let table_metadata_location = table_location.join("metadata/v1.json");
195
196 let file_io = FileIO::new_with_fs();
197
198 let template = fs::read_to_string(format!(
199 "{}/testdata/example_table_metadata_v2.json",
200 env!("CARGO_MANIFEST_DIR")
201 ))
202 .unwrap();
203 let metadata_json = render_template(&template, context! {
205 table_location => &table_location,
206 manifest_list_1_location => &manifest_list_location,
207 manifest_list_2_location => &manifest_list_location,
208 table_metadata_1_location => &table_metadata_location,
209 });
210 let table_metadata = serde_json::from_str::<TableMetadata>(&metadata_json).unwrap();
211
212 let table = Table::builder()
213 .metadata(table_metadata)
214 .identifier(TableIdent::from_strs(["db", "table1"]).unwrap())
215 .file_io(file_io)
216 .metadata_location(table_metadata_location.to_str().unwrap())
217 .runtime(test_runtime())
218 .build()
219 .unwrap();
220
221 let current_snapshot = table.metadata().current_snapshot().unwrap();
222 let schema = current_snapshot.schema(table.metadata()).unwrap();
223 let partition_spec = table.metadata().default_partition_spec();
224
225 let next_manifest_file = |location: &str| {
226 table
227 .file_io()
228 .new_output(format!(
229 "{}/metadata/manifest_{}.avro",
230 location,
231 Uuid::new_v4()
232 ))
233 .unwrap()
234 };
235 let table_location_str = table_location.to_str().unwrap().to_string();
236
237 let mut data_writer = ManifestWriterBuilder::new(
239 next_manifest_file(&table_location_str),
240 Some(current_snapshot.snapshot_id()),
241 schema.clone(),
242 partition_spec.as_ref().clone(),
243 )
244 .build_v2_data();
245 data_writer
246 .add_entry(
247 ManifestEntry::builder()
248 .status(ManifestStatus::Added)
249 .data_file(
250 DataFileBuilder::default()
251 .partition_spec_id(0)
252 .content(DataContentType::Data)
253 .file_path(format!("{table_location_str}/data.parquet"))
254 .file_format(DataFileFormat::Parquet)
255 .file_size_in_bytes(100)
256 .record_count(1)
257 .partition(Struct::from_iter([Some(Literal::long(100))]))
258 .build()
259 .unwrap(),
260 )
261 .build(),
262 )
263 .unwrap();
264 let data_manifest = data_writer.write_manifest_file().await.unwrap();
265
266 let mut delete_writer = ManifestWriterBuilder::new(
268 next_manifest_file(&table_location_str),
269 Some(current_snapshot.snapshot_id()),
270 schema.clone(),
271 partition_spec.as_ref().clone(),
272 )
273 .build_v2_data();
274 delete_writer
275 .add_delete_entry(
276 ManifestEntry::builder()
277 .status(ManifestStatus::Deleted)
278 .sequence_number(0)
279 .file_sequence_number(0)
280 .data_file(
281 DataFileBuilder::default()
282 .partition_spec_id(0)
283 .content(DataContentType::Data)
284 .file_path(format!("{table_location_str}/removed.parquet"))
285 .file_format(DataFileFormat::Parquet)
286 .file_size_in_bytes(100)
287 .record_count(1)
288 .partition(Struct::from_iter([Some(Literal::long(100))]))
289 .build()
290 .unwrap(),
291 )
292 .build(),
293 )
294 .unwrap();
295 let delete_manifest = delete_writer.write_manifest_file().await.unwrap();
296 let delete_manifest_path = delete_manifest.manifest_path.clone();
297
298 assert!(delete_manifest.has_deleted_files());
300 assert!(!delete_manifest.has_added_files());
301 assert!(!delete_manifest.has_existing_files());
302
303 let mut manifest_list_writer = ManifestListWriter::v2(
304 table
305 .file_io()
306 .new_output(current_snapshot.manifest_list())
307 .unwrap()
308 .writer()
309 .await
310 .unwrap(),
311 current_snapshot.snapshot_id(),
312 current_snapshot.parent_snapshot_id(),
313 current_snapshot.sequence_number(),
314 );
315 manifest_list_writer
316 .add_manifests(vec![data_manifest, delete_manifest].into_iter())
317 .unwrap();
318 manifest_list_writer.close().await.unwrap();
319
320 (table, tmp_dir, delete_manifest_path)
321 }
322
323 #[tokio::test]
327 async fn test_fast_append_preserves_delete_only_manifest() {
328 let (table, _tmp_dir, delete_manifest_path) = make_table_with_delete_only_manifest().await;
329
330 let new_file = DataFileBuilder::default()
332 .content(DataContentType::Data)
333 .file_path(format!("{}/appended.parquet", table.metadata().location()))
334 .file_format(DataFileFormat::Parquet)
335 .file_size_in_bytes(100)
336 .record_count(1)
337 .partition_spec_id(table.metadata().default_partition_spec_id())
338 .partition(Struct::from_iter([Some(Literal::long(100))]))
339 .build()
340 .unwrap();
341
342 let tx = Transaction::new(&table);
343 let action = tx.fast_append().add_data_files(vec![new_file]);
344 let mut action_commit = Arc::new(action).commit(&table).await.unwrap();
345 let updates = action_commit.take_updates();
346
347 let new_snapshot: SnapshotRef = if let TableUpdate::AddSnapshot { snapshot } = &updates[0] {
348 SnapshotRef::new(snapshot.clone())
349 } else {
350 unreachable!("first update of a fast append should be AddSnapshot")
351 };
352
353 let manifest_list = table
354 .manifest_list_reader(&new_snapshot)
355 .load()
356 .await
357 .unwrap();
358
359 assert!(
360 manifest_list
361 .entries()
362 .iter()
363 .any(|m| m.manifest_path == delete_manifest_path),
364 "delete-only manifest {delete_manifest_path} was dropped from the new snapshot's \
365 manifest list; the files it removed would reappear as live data"
366 );
367 }
368
369 async fn committed_data_files(table: &Table, updates: &[TableUpdate]) -> Vec<DataFile> {
371 let TableUpdate::AddSnapshot { snapshot } = &updates[0] else {
372 unreachable!("first update is always AddSnapshot")
373 };
374 let manifest_list = table
375 .manifest_list_reader(&SnapshotRef::new(snapshot.clone()))
376 .load()
377 .await
378 .unwrap();
379 assert_eq!(1, manifest_list.entries().len());
380 table
381 .manifest_reader()
382 .read(&manifest_list.entries()[0])
383 .await
384 .unwrap()
385 .entries()
386 .iter()
387 .map(|entry| entry.data_file().clone())
388 .collect()
389 }
390
391 #[tokio::test]
392 async fn test_fast_append_writes_encrypted_manifest() {
393 let table = make_encrypted_table().await;
394 assert!(
395 table.encryption_manager().is_some(),
396 "fixture table should have an EncryptionManager"
397 );
398
399 let new_file = DataFileBuilder::default()
400 .content(DataContentType::Data)
401 .file_path("memory:///table/data/00000.parquet".to_string())
402 .file_format(DataFileFormat::Parquet)
403 .partition(Struct::empty())
404 .record_count(100)
405 .file_size_in_bytes(4096)
406 .partition_spec_id(table.metadata().default_partition_spec_id())
407 .build()
408 .unwrap();
409
410 let tx = Transaction::new(&table);
411 let action = tx.fast_append().add_data_files(vec![new_file]);
412 let mut action_commit = Arc::new(action).commit(&table).await.unwrap();
413 let updates = action_commit.take_updates();
414
415 let new_snapshot: SnapshotRef = updates
416 .iter()
417 .find_map(|u| match u {
418 TableUpdate::AddSnapshot { snapshot } => Some(SnapshotRef::new(snapshot.clone())),
419 _ => None,
420 })
421 .expect("a fast append should emit an AddSnapshot update");
422
423 let manifest_list_key_metadata = table
424 .encryption_manager()
425 .unwrap()
426 .decrypt_manifest_list_key_metadata(new_snapshot.encryption_key_id().unwrap())
427 .await
428 .unwrap();
429 let manifest_list_size = table
430 .file_io()
431 .new_input(new_snapshot.manifest_list())
432 .unwrap()
433 .metadata()
434 .await
435 .unwrap()
436 .size;
437 assert_eq!(
438 manifest_list_key_metadata.file_length(),
439 Some(manifest_list_size)
440 );
441
442 let manifest_list = table
443 .manifest_list_reader(&new_snapshot)
444 .load()
445 .await
446 .unwrap();
447 let manifest_file = manifest_list
448 .entries()
449 .iter()
450 .find(|m| m.added_files_count.unwrap_or(0) > 0)
451 .expect("new snapshot should carry the appended data manifest");
452
453 let key_metadata_bytes = manifest_file
455 .key_metadata
456 .as_ref()
457 .expect("encrypted manifest must record key metadata");
458 let key_metadata = StandardKeyMetadata::decode(key_metadata_bytes)
459 .expect("recorded key metadata must decode as StandardKeyMetadata");
460 let manifest_size = table
461 .file_io()
462 .new_input(&manifest_file.manifest_path)
463 .unwrap()
464 .metadata()
465 .await
466 .unwrap()
467 .size;
468 assert_eq!(key_metadata.file_length(), Some(manifest_size));
469 assert_eq!(manifest_file.manifest_length, manifest_size as i64);
470
471 let manifest = table.manifest_reader().read(manifest_file).await.unwrap();
476 assert_eq!(manifest.entries().len(), 1);
477 assert_eq!(
478 manifest.entries()[0].data_file().file_path(),
479 "memory:///table/data/00000.parquet"
480 );
481 }
482
483 #[tokio::test]
484 async fn test_empty_data_append_action() {
485 let table = make_v2_minimal_table();
486 let tx = Transaction::new(&table);
487 let action = tx.fast_append().add_data_files(vec![]);
488 assert!(Arc::new(action).commit(&table).await.is_err());
489 }
490
491 #[tokio::test]
495 async fn test_fast_append_honors_write_metadata_path() {
496 let base = make_v2_minimal_table();
497 let metadata_root = format!("{}/custom-meta", base.metadata().location());
498 let metadata = base
499 .metadata()
500 .clone()
501 .into_builder(None)
502 .set_properties(HashMap::from([(
503 "write.metadata.path".to_string(),
504 metadata_root.clone(),
505 )]))
506 .unwrap()
507 .build()
508 .unwrap()
509 .metadata;
510 let table = base.with_metadata(Arc::new(metadata));
511
512 let data_file = DataFileBuilder::default()
513 .content(DataContentType::Data)
514 .file_path(format!("{}/data/1.parquet", table.metadata().location()))
515 .file_format(DataFileFormat::Parquet)
516 .file_size_in_bytes(100)
517 .record_count(1)
518 .partition_spec_id(table.metadata().default_partition_spec_id())
519 .partition(Struct::from_iter([Some(Literal::long(300))]))
520 .build()
521 .unwrap();
522
523 let tx = Transaction::new(&table);
524 let action = tx.fast_append().add_data_files(vec![data_file]);
525 let mut action_commit = Arc::new(action).commit(&table).await.unwrap();
526 let updates = action_commit.take_updates();
527
528 let new_snapshot: SnapshotRef = if let TableUpdate::AddSnapshot { snapshot } = &updates[0] {
529 SnapshotRef::new(snapshot.clone())
530 } else {
531 unreachable!("first update of a fast append should be AddSnapshot")
532 };
533
534 let prefix = format!("{metadata_root}/");
535
536 assert!(
538 new_snapshot.manifest_list().starts_with(prefix.as_str()),
539 "manifest list {} not under configured write.metadata.path {metadata_root}",
540 new_snapshot.manifest_list()
541 );
542
543 let manifest_list = table
545 .manifest_list_reader(&new_snapshot)
546 .load()
547 .await
548 .unwrap();
549 assert!(
550 !manifest_list.entries().is_empty(),
551 "expected at least one manifest entry"
552 );
553 for entry in manifest_list.entries() {
554 assert!(
555 entry.manifest_path.starts_with(prefix.as_str()),
556 "manifest {} not under configured write.metadata.path {metadata_root}",
557 entry.manifest_path
558 );
559 }
560 }
561
562 #[tokio::test]
563 async fn test_set_snapshot_properties() {
564 let table = make_v2_minimal_table();
565 let tx = Transaction::new(&table);
566
567 let mut snapshot_properties = HashMap::new();
568 snapshot_properties.insert("key".to_string(), "val".to_string());
569
570 let data_file = DataFileBuilder::default()
571 .content(DataContentType::Data)
572 .file_path("test/1.parquet".to_string())
573 .file_format(DataFileFormat::Parquet)
574 .file_size_in_bytes(100)
575 .record_count(1)
576 .partition_spec_id(table.metadata().default_partition_spec_id())
577 .partition(Struct::from_iter([Some(Literal::long(300))]))
578 .build()
579 .unwrap();
580
581 let action = tx
582 .fast_append()
583 .set_snapshot_properties(snapshot_properties)
584 .add_data_files(vec![data_file]);
585 let mut action_commit = Arc::new(action).commit(&table).await.unwrap();
586 let updates = action_commit.take_updates();
587
588 let new_snapshot = if let TableUpdate::AddSnapshot { snapshot } = &updates[0] {
590 snapshot
591 } else {
592 unreachable!()
593 };
594 assert_eq!(
595 new_snapshot
596 .summary()
597 .additional_properties
598 .get("key")
599 .unwrap(),
600 "val"
601 );
602 }
603
604 const FIXTURE_MASTER_KEY_ID: &str = "master-1";
606 const FIXTURE_MASTER_KEY_BYTES: [u8; 16] = [
607 0x00, 0x01, 0x02, 0x03, 0x04, 0x05, 0x06, 0x07, 0x08, 0x09, 0x0a, 0x0b, 0x0c, 0x0d, 0x0e,
608 0x0f,
609 ];
610
611 async fn make_v3_encrypted_table() -> Table {
612 let json = fs::read_to_string(format!(
613 "{}/testdata/table_metadata/TableMetadataV3ValidEncryption.json",
614 env!("CARGO_MANIFEST_DIR")
615 ))
616 .unwrap();
617 let metadata = serde_json::from_str::<TableMetadata>(&json).unwrap();
618
619 let kms = MemoryKeyManagementClient::new();
620 kms.add_master_key_bytes(
621 FIXTURE_MASTER_KEY_ID,
622 SensitiveBytes::new(FIXTURE_MASTER_KEY_BYTES),
623 )
624 .unwrap();
625
626 let file_io = FileIO::new_with_memory();
627
628 let manifest_list_bytes = fs::read(format!(
629 "{}/testdata/manifests_lists/manifest-list-v3-encrypted.avro",
630 env!("CARGO_MANIFEST_DIR")
631 ))
632 .unwrap();
633 let parent_manifest_list = metadata.current_snapshot().unwrap().manifest_list();
634 file_io
635 .new_output(parent_manifest_list)
636 .unwrap()
637 .write(manifest_list_bytes.into())
638 .await
639 .unwrap();
640
641 Table::builder()
642 .metadata(metadata)
643 .metadata_location("memory:///table/metadata/v1.json")
644 .identifier(TableIdent::from_strs(["ns1", "enc"]).unwrap())
645 .file_io(file_io)
646 .kms_client(Arc::new(kms))
647 .runtime(test_runtime())
648 .build()
649 .unwrap()
650 }
651
652 #[tokio::test]
653 async fn test_commit_with_encryption_adds_keys_and_records_snapshot_key_id() {
654 let table = make_v3_encrypted_table().await;
655
656 let data_file = DataFileBuilder::default()
657 .content(DataContentType::Data)
658 .file_path("test/1.parquet".to_string())
659 .file_format(DataFileFormat::Parquet)
660 .file_size_in_bytes(100)
661 .record_count(1)
662 .partition_spec_id(table.metadata().default_partition_spec_id())
663 .partition(Struct::empty())
664 .build()
665 .unwrap();
666
667 let tx = Transaction::new(&table);
668 let action = tx.fast_append().add_data_files(vec![data_file]);
669 let mut action_commit = Arc::new(action).commit(&table).await.unwrap();
670 let updates = action_commit.take_updates();
671
672 let added_key_ids: Vec<String> = updates
673 .iter()
674 .filter_map(|u| match u {
675 TableUpdate::AddEncryptionKey { encryption_key } => {
676 Some(encryption_key.key_id().to_string())
677 }
678 _ => None,
679 })
680 .collect();
681 assert!(!added_key_ids.is_empty(), "got {updates:?}");
682
683 let new_snapshot = updates
685 .iter()
686 .find_map(|u| match u {
687 TableUpdate::AddSnapshot { snapshot } => Some(snapshot),
688 _ => None,
689 })
690 .expect("commit should add a snapshot");
691
692 let snapshot_key_id = new_snapshot
693 .encryption_key_id()
694 .expect("encrypted snapshot should record its manifest-list key id");
695 assert!(
696 added_key_ids.iter().any(|id| id == snapshot_key_id),
697 "snapshot key id {snapshot_key_id} not in added keys {added_key_ids:?}"
698 );
699
700 let new_snapshot_ref: SnapshotRef = Arc::new(new_snapshot.clone());
701 let manifest_list = table
702 .manifest_list_reader(&new_snapshot_ref)
703 .load()
704 .await
705 .expect("newly written encrypted manifest list should decrypt and parse");
706 assert_eq!(
707 manifest_list.entries().len(),
708 1,
709 "append should record exactly the one new data manifest"
710 );
711 }
712
713 #[tokio::test]
714 async fn test_snapshot_properties_cannot_override_computed_metrics() {
715 let table = make_v2_minimal_table();
720 let tx = Transaction::new(&table);
721
722 let mut snapshot_properties = HashMap::new();
723 snapshot_properties.insert("added-data-files".to_string(), "9999".to_string());
726 snapshot_properties.insert("added-records".to_string(), "not-a-number".to_string());
727
728 let data_file = DataFileBuilder::default()
729 .content(DataContentType::Data)
730 .file_path("test/1.parquet".to_string())
731 .file_format(DataFileFormat::Parquet)
732 .file_size_in_bytes(100)
733 .record_count(1)
734 .partition_spec_id(table.metadata().default_partition_spec_id())
735 .partition(Struct::from_iter([Some(Literal::long(300))]))
736 .build()
737 .unwrap();
738
739 let action = tx
740 .fast_append()
741 .set_snapshot_properties(snapshot_properties)
742 .add_data_files(vec![data_file]);
743 let mut action_commit = Arc::new(action).commit(&table).await.unwrap();
745 let updates = action_commit.take_updates();
746
747 let new_snapshot = if let TableUpdate::AddSnapshot { snapshot } = &updates[0] {
748 snapshot
749 } else {
750 unreachable!()
751 };
752 let props = &new_snapshot.summary().additional_properties;
753
754 assert_eq!(
756 props.get("added-data-files").unwrap(),
757 "1",
758 "computed added-data-files must override the user-supplied value"
759 );
760 assert_eq!(
761 props.get("added-records").unwrap(),
762 "1",
763 "computed added-records must override the user-supplied non-integer value"
764 );
765 }
766
767 #[tokio::test]
768 async fn test_append_snapshot_properties() {
769 let table = make_v2_minimal_table();
770 let tx = Transaction::new(&table);
771
772 let mut snapshot_properties = HashMap::new();
773 snapshot_properties.insert("key".to_string(), "val".to_string());
774
775 let action = tx
776 .fast_append()
777 .set_snapshot_properties(snapshot_properties);
778 let mut action_commit = Arc::new(action).commit(&table).await.unwrap();
779 let updates = action_commit.take_updates();
780
781 let new_snapshot = if let TableUpdate::AddSnapshot { snapshot } = &updates[0] {
783 snapshot
784 } else {
785 unreachable!()
786 };
787 assert_eq!(
788 new_snapshot
789 .summary()
790 .additional_properties
791 .get("key")
792 .unwrap(),
793 "val"
794 );
795 }
796
797 #[tokio::test]
798 async fn test_fast_append_file_with_incompatible_partition_value() {
799 let table = make_v2_minimal_table();
800 let tx = Transaction::new(&table);
801 let action = tx.fast_append();
802
803 let data_file = DataFileBuilder::default()
805 .content(DataContentType::Data)
806 .file_path("test/3.parquet".to_string())
807 .file_format(DataFileFormat::Parquet)
808 .file_size_in_bytes(100)
809 .record_count(1)
810 .partition_spec_id(table.metadata().default_partition_spec_id())
811 .partition(Struct::from_iter([Some(Literal::string("test"))]))
812 .build()
813 .unwrap();
814
815 let action = action.add_data_files(vec![data_file.clone()]);
816
817 assert!(Arc::new(action).commit(&table).await.is_err());
818 }
819
820 #[tokio::test]
821 async fn test_fast_append_dedupes_intra_batch_duplicate_paths() {
822 let table = make_v2_minimal_table();
823 let tx = Transaction::new(&table);
824
825 let make_file = |size: u64, records: u64| {
826 DataFileBuilder::default()
827 .content(DataContentType::Data)
828 .file_path("test/dup.parquet".to_string())
829 .file_format(DataFileFormat::Parquet)
830 .file_size_in_bytes(size)
831 .record_count(records)
832 .partition_spec_id(table.metadata().default_partition_spec_id())
833 .partition(Struct::from_iter([Some(Literal::long(1))]))
834 .build()
835 .unwrap()
836 };
837
838 let action = tx.fast_append().add_data_files(vec![
840 make_file(100, 10),
841 make_file(200, 20),
842 make_file(300, 30),
843 ]);
844 let mut action_commit = Arc::new(action).commit(&table).await.unwrap();
845 let files = committed_data_files(&table, &action_commit.take_updates()).await;
846 assert_eq!(1, files.len());
847 assert_eq!(100, files[0].file_size_in_bytes());
848 }
849
850 #[tokio::test]
851 async fn test_fast_append_dedupes_regardless_of_check_duplicate_flag() {
852 let table = make_v2_minimal_table();
853 let tx = Transaction::new(&table);
854
855 let make_file = || {
856 DataFileBuilder::default()
857 .content(DataContentType::Data)
858 .file_path("test/dup.parquet".to_string())
859 .file_format(DataFileFormat::Parquet)
860 .file_size_in_bytes(100)
861 .record_count(10)
862 .partition_spec_id(table.metadata().default_partition_spec_id())
863 .partition(Struct::from_iter([Some(Literal::long(1))]))
864 .build()
865 .unwrap()
866 };
867
868 let action = tx
870 .fast_append()
871 .with_check_duplicate(false)
872 .add_data_files(vec![make_file(), make_file()]);
873 let mut action_commit = Arc::new(action).commit(&table).await.unwrap();
874 let files = committed_data_files(&table, &action_commit.take_updates()).await;
875 assert_eq!(1, files.len());
876 }
877
878 #[tokio::test]
879 async fn test_fast_append() {
880 let table = make_v2_minimal_table();
881 let tx = Transaction::new(&table);
882 let action = tx.fast_append();
883
884 let data_file = DataFileBuilder::default()
885 .content(DataContentType::Data)
886 .file_path("test/3.parquet".to_string())
887 .file_format(DataFileFormat::Parquet)
888 .file_size_in_bytes(100)
889 .record_count(1)
890 .partition_spec_id(table.metadata().default_partition_spec_id())
891 .partition(Struct::from_iter([Some(Literal::long(300))]))
892 .build()
893 .unwrap();
894
895 let action = action.add_data_files(vec![data_file.clone()]);
896 let mut action_commit = Arc::new(action).commit(&table).await.unwrap();
897 let updates = action_commit.take_updates();
898 let requirements = action_commit.take_requirements();
899
900 assert!(
902 matches!((&updates[0],&updates[1]), (TableUpdate::AddSnapshot { snapshot },TableUpdate::SetSnapshotRef { reference,ref_name }) if snapshot.snapshot_id() == reference.snapshot_id && ref_name == MAIN_BRANCH)
903 );
904 assert_eq!(
905 vec![
906 TableRequirement::UuidMatch {
907 uuid: table.metadata().uuid()
908 },
909 TableRequirement::RefSnapshotIdMatch {
910 r#ref: MAIN_BRANCH.to_string(),
911 snapshot_id: table.metadata().current_snapshot_id
912 }
913 ],
914 requirements
915 );
916
917 let new_snapshot: SnapshotRef = if let TableUpdate::AddSnapshot { snapshot } = &updates[0] {
919 SnapshotRef::new(snapshot.clone())
920 } else {
921 unreachable!()
922 };
923 let manifest_list = table
924 .manifest_list_reader(&new_snapshot)
925 .load()
926 .await
927 .unwrap();
928 assert_eq!(1, manifest_list.entries().len());
929 assert_eq!(
930 manifest_list.entries()[0].sequence_number,
931 new_snapshot.sequence_number()
932 );
933
934 let manifest = table
936 .manifest_reader()
937 .read(&manifest_list.entries()[0])
938 .await
939 .unwrap();
940 assert_eq!(1, manifest.entries().len());
941 assert_eq!(
942 new_snapshot.sequence_number(),
943 manifest.entries()[0]
944 .sequence_number()
945 .expect("Inherit sequence number by load manifest")
946 );
947
948 assert_eq!(
949 new_snapshot.snapshot_id(),
950 manifest.entries()[0].snapshot_id().unwrap()
951 );
952 assert_eq!(data_file, *manifest.entries()[0].data_file());
953 }
954}