1use std::cmp::min;
19
20use apache_avro::{Writer as AvroWriter, to_value};
21use bytes::Bytes;
22use itertools::Itertools;
23use serde_json::to_vec;
24
25use super::{
26 Datum, FormatVersion, ManifestContentType, PartitionSpec, PrimitiveType,
27 UNASSIGNED_SEQUENCE_NUMBER,
28};
29use crate::encryption::EncryptedOutputFile;
30use crate::error::{Result, invalid_data};
31use crate::io::{FileMetadata, FileWrite, OutputFile};
32use crate::spec::manifest::_serde::{ManifestEntryV1, ManifestEntryV2};
33use crate::spec::manifest::{manifest_schema_v1, manifest_schema_v2};
34use crate::spec::{
35 DataContentType, DataFile, FieldSummary, ManifestEntry, ManifestFile, ManifestMetadata,
36 ManifestStatus, PrimitiveLiteral, SchemaRef, StructType, Type,
37};
38
39const UNASSIGNED_SNAPSHOT_ID: i64 = -1;
42
43pub(crate) enum ManifestOutput {
45 Plain(OutputFile),
46 Encrypted(EncryptedOutputFile),
47}
48
49impl ManifestOutput {
50 async fn writer(&self) -> Result<Box<dyn FileWrite>> {
51 match self {
52 Self::Plain(output) => output.writer().await,
53 Self::Encrypted(output) => output.writer().await,
54 }
55 }
56
57 fn encoded_key_metadata(&self, file_metadata: &FileMetadata) -> Result<Option<Vec<u8>>> {
58 match self {
59 Self::Plain(_) => Ok(None),
60 Self::Encrypted(output) => Ok(Some(
61 output
62 .key_metadata_with_saved_file_metadata(file_metadata)
63 .encode()?
64 .into_vec(),
65 )),
66 }
67 }
68}
69
70pub struct ManifestWriterBuilder {
72 output: ManifestOutput,
73 location: String,
74 snapshot_id: Option<i64>,
75 schema: SchemaRef,
76 partition_spec: PartitionSpec,
77}
78
79impl ManifestWriterBuilder {
80 pub fn new(
82 output: OutputFile,
83 snapshot_id: Option<i64>,
84 schema: SchemaRef,
85 partition_spec: PartitionSpec,
86 ) -> Self {
87 let location = output.location().to_owned();
88 Self {
89 output: ManifestOutput::Plain(output),
90 location,
91 snapshot_id,
92 schema,
93 partition_spec,
94 }
95 }
96
97 pub fn new_from_encrypted(
101 encrypted_output: EncryptedOutputFile,
102 snapshot_id: Option<i64>,
103 schema: SchemaRef,
104 partition_spec: PartitionSpec,
105 ) -> Result<Self> {
106 let location = encrypted_output.location().to_owned();
107 Ok(Self {
108 output: ManifestOutput::Encrypted(encrypted_output),
109 location,
110 snapshot_id,
111 schema,
112 partition_spec,
113 })
114 }
115
116 pub fn build_v1(self) -> ManifestWriter {
118 let metadata = ManifestMetadata::builder()
119 .schema_id(self.schema.schema_id())
120 .schema(self.schema)
121 .partition_spec(self.partition_spec)
122 .format_version(FormatVersion::V1)
123 .content(ManifestContentType::Data)
124 .build();
125 ManifestWriter::new(self.output, self.location, self.snapshot_id, metadata, None)
126 }
127
128 pub fn build_v2_data(self) -> ManifestWriter {
130 let metadata = ManifestMetadata::builder()
131 .schema_id(self.schema.schema_id())
132 .schema(self.schema)
133 .partition_spec(self.partition_spec)
134 .format_version(FormatVersion::V2)
135 .content(ManifestContentType::Data)
136 .build();
137
138 ManifestWriter::new(self.output, self.location, self.snapshot_id, metadata, None)
139 }
140
141 pub fn build_v2_deletes(self) -> ManifestWriter {
143 let metadata = ManifestMetadata::builder()
144 .schema_id(self.schema.schema_id())
145 .schema(self.schema)
146 .partition_spec(self.partition_spec)
147 .format_version(FormatVersion::V2)
148 .content(ManifestContentType::Deletes)
149 .build();
150 ManifestWriter::new(self.output, self.location, self.snapshot_id, metadata, None)
151 }
152
153 pub fn build_v3_data(self) -> ManifestWriter {
155 let metadata = ManifestMetadata::builder()
156 .schema_id(self.schema.schema_id())
157 .schema(self.schema)
158 .partition_spec(self.partition_spec)
159 .format_version(FormatVersion::V3)
160 .content(ManifestContentType::Data)
161 .build();
162 ManifestWriter::new(
163 self.output,
164 self.location,
165 self.snapshot_id,
166 metadata,
167 None,
170 )
171 }
172
173 pub fn build_v3_deletes(self) -> ManifestWriter {
175 let metadata = ManifestMetadata::builder()
176 .schema_id(self.schema.schema_id())
177 .schema(self.schema)
178 .partition_spec(self.partition_spec)
179 .format_version(FormatVersion::V3)
180 .content(ManifestContentType::Deletes)
181 .build();
182 ManifestWriter::new(self.output, self.location, self.snapshot_id, metadata, None)
183 }
184}
185
186pub struct ManifestWriter {
188 output: ManifestOutput,
189 location: String,
190
191 snapshot_id: Option<i64>,
192
193 added_files: u32,
194 added_rows: u64,
195 existing_files: u32,
196 existing_rows: u64,
197 deleted_files: u32,
198 deleted_rows: u64,
199 first_row_id: Option<u64>,
200
201 min_seq_num: Option<i64>,
202
203 manifest_entries: Vec<ManifestEntry>,
204
205 metadata: ManifestMetadata,
206}
207
208impl ManifestWriter {
209 pub(crate) fn new(
211 output: ManifestOutput,
212 location: String,
213 snapshot_id: Option<i64>,
214 metadata: ManifestMetadata,
215 first_row_id: Option<u64>,
216 ) -> Self {
217 Self {
218 output,
219 location,
220 snapshot_id,
221 added_files: 0,
222 added_rows: 0,
223 existing_files: 0,
224 existing_rows: 0,
225 deleted_files: 0,
226 deleted_rows: 0,
227 first_row_id,
228 min_seq_num: None,
229 manifest_entries: Vec::new(),
230 metadata,
231 }
232 }
233
234 fn construct_partition_summaries(
235 &mut self,
236 partition_type: &StructType,
237 ) -> Result<Vec<FieldSummary>> {
238 let mut field_stats: Vec<_> = partition_type
239 .fields()
240 .iter()
241 .map(|f| PartitionFieldStats::new(f.field_type.as_primitive_type().unwrap().clone()))
242 .collect();
243 for partition in self.manifest_entries.iter().map(|e| &e.data_file.partition) {
244 for (literal, stat) in partition.iter().zip_eq(field_stats.iter_mut()) {
245 let primitive_literal = literal.map(|v| v.as_primitive_literal().unwrap());
246 stat.update(primitive_literal)?;
247 }
248 }
249 Ok(field_stats.into_iter().map(|stat| stat.finish()).collect())
250 }
251
252 fn check_data_file(&self, data_file: &DataFile) -> Result<()> {
253 match self.metadata.content {
254 ManifestContentType::Data => {
255 if data_file.content != DataContentType::Data {
256 return Err(invalid_data!(
257 "Date file at path {} with manifest content type `data`, should have DataContentType `Data`, but has `{:?}`",
258 data_file.file_path(),
259 data_file.content
260 ));
261 }
262 }
263 ManifestContentType::Deletes => {
264 if data_file.content != DataContentType::EqualityDeletes
265 && data_file.content != DataContentType::PositionDeletes
266 {
267 return Err(invalid_data!(
268 "Date file at path {} with manifest content type `deletes`, should have DataContentType `Data`, but has `{:?}`",
269 data_file.file_path(),
270 data_file.content
271 ));
272 }
273 }
274 }
275 Ok(())
276 }
277
278 pub(crate) fn add_entry(&mut self, mut entry: ManifestEntry) -> Result<()> {
284 self.check_data_file(&entry.data_file)?;
285 if entry.sequence_number().is_some_and(|n| n >= 0) {
286 entry.status = ManifestStatus::Added;
287 entry.snapshot_id = self.snapshot_id;
288 entry.file_sequence_number = None;
289 } else {
290 entry.status = ManifestStatus::Added;
291 entry.snapshot_id = self.snapshot_id;
292 entry.sequence_number = None;
293 entry.file_sequence_number = None;
294 };
295 self.add_entry_inner(entry)?;
296 Ok(())
297 }
298
299 pub fn add_file(&mut self, data_file: DataFile, sequence_number: i64) -> Result<()> {
303 self.check_data_file(&data_file)?;
304 let entry = ManifestEntry {
305 status: ManifestStatus::Added,
306 snapshot_id: self.snapshot_id,
307 sequence_number: (sequence_number >= 0).then_some(sequence_number),
308 file_sequence_number: None,
309 data_file,
310 };
311 self.add_entry_inner(entry)?;
312 Ok(())
313 }
314
315 #[allow(dead_code)]
322 pub(crate) fn add_delete_entry(&mut self, mut entry: ManifestEntry) -> Result<()> {
323 self.check_data_file(&entry.data_file)?;
324 entry.status = ManifestStatus::Deleted;
325 entry.snapshot_id = self.snapshot_id;
326 self.add_entry_inner(entry)?;
327 Ok(())
328 }
329
330 pub fn add_delete_file(
334 &mut self,
335 data_file: DataFile,
336 sequence_number: i64,
337 file_sequence_number: Option<i64>,
338 ) -> Result<()> {
339 self.check_data_file(&data_file)?;
340 let entry = ManifestEntry {
341 status: ManifestStatus::Deleted,
342 snapshot_id: self.snapshot_id,
343 sequence_number: Some(sequence_number),
344 file_sequence_number,
345 data_file,
346 };
347 self.add_entry_inner(entry)?;
348 Ok(())
349 }
350
351 #[allow(dead_code)]
357 pub(crate) fn add_existing_entry(&mut self, mut entry: ManifestEntry) -> Result<()> {
358 self.check_data_file(&entry.data_file)?;
359 entry.status = ManifestStatus::Existing;
360 self.add_entry_inner(entry)?;
361 Ok(())
362 }
363
364 pub fn add_existing_file(
367 &mut self,
368 data_file: DataFile,
369 snapshot_id: i64,
370 sequence_number: i64,
371 file_sequence_number: Option<i64>,
372 ) -> Result<()> {
373 self.check_data_file(&data_file)?;
374 let entry = ManifestEntry {
375 status: ManifestStatus::Existing,
376 snapshot_id: Some(snapshot_id),
377 sequence_number: Some(sequence_number),
378 file_sequence_number,
379 data_file,
380 };
381 self.add_entry_inner(entry)?;
382 Ok(())
383 }
384
385 fn add_entry_inner(&mut self, entry: ManifestEntry) -> Result<()> {
386 if (entry.status == ManifestStatus::Deleted || entry.status == ManifestStatus::Existing)
388 && (entry.sequence_number.is_none() || entry.file_sequence_number.is_none())
389 {
390 return Err(invalid_data!(
391 "Manifest entry with status Existing or Deleted should have sequence number"
392 ));
393 }
394
395 match entry.status {
397 ManifestStatus::Added => {
398 self.added_files += 1;
399 self.added_rows += entry.data_file.record_count;
400 }
401 ManifestStatus::Deleted => {
402 self.deleted_files += 1;
403 self.deleted_rows += entry.data_file.record_count;
404 }
405 ManifestStatus::Existing => {
406 self.existing_files += 1;
407 self.existing_rows += entry.data_file.record_count;
408 }
409 }
410 if entry.is_alive()
411 && let Some(seq_num) = entry.sequence_number
412 {
413 self.min_seq_num = Some(self.min_seq_num.map_or(seq_num, |v| min(v, seq_num)));
414 }
415 self.manifest_entries.push(entry);
416 Ok(())
417 }
418
419 pub async fn write_manifest_file(mut self) -> Result<ManifestFile> {
421 let partition_type = self
423 .metadata
424 .partition_spec
425 .partition_type(&self.metadata.schema)?;
426 let partition_struct_type = Type::Struct(partition_type.clone());
429 let table_schema = &self.metadata.schema;
430 let avro_schema = match self.metadata.format_version {
431 FormatVersion::V1 => manifest_schema_v1(&partition_type)?,
432 FormatVersion::V2 | FormatVersion::V3 => manifest_schema_v2(&partition_type)?,
434 };
435 let mut avro_writer = AvroWriter::new(&avro_schema, Vec::new())?;
436 avro_writer.add_user_metadata(
437 "schema".to_string(),
438 to_vec(table_schema)
439 .map_err(|err| invalid_data!("Fail to serialize table schema").with_source(err))?,
440 )?;
441 avro_writer.add_user_metadata(
442 "schema-id".to_string(),
443 table_schema.schema_id().to_string(),
444 )?;
445 avro_writer.add_user_metadata(
446 "partition-spec".to_string(),
447 to_vec(&self.metadata.partition_spec.fields()).map_err(|err| {
448 invalid_data!("Fail to serialize partition spec").with_source(err)
449 })?,
450 )?;
451 avro_writer.add_user_metadata(
452 "partition-spec-id".to_string(),
453 self.metadata.partition_spec.spec_id().to_string(),
454 )?;
455 avro_writer.add_user_metadata(
456 "format-version".to_string(),
457 (self.metadata.format_version as u8).to_string(),
458 )?;
459 match self.metadata.format_version {
460 FormatVersion::V1 => {}
461 FormatVersion::V2 | FormatVersion::V3 => {
462 avro_writer
463 .add_user_metadata("content".to_string(), self.metadata.content.to_string())?;
464 }
465 }
466
467 let partition_summary = self.construct_partition_summaries(&partition_type)?;
468 for entry in std::mem::take(&mut self.manifest_entries) {
470 let value = match self.metadata.format_version {
471 FormatVersion::V1 => {
472 to_value(ManifestEntryV1::try_from(entry, &partition_struct_type)?)?
473 .resolve(&avro_schema)?
474 }
475 FormatVersion::V2 | FormatVersion::V3 => {
477 to_value(ManifestEntryV2::try_from(entry, &partition_struct_type)?)?
478 .resolve(&avro_schema)?
479 }
480 };
481
482 avro_writer.append_value(value)?;
483 }
484
485 let content = avro_writer.into_inner()?;
486 let mut writer = self.output.writer().await?;
487 writer.write(Bytes::from(content)).await?;
488 let file_metadata = writer.close().await?;
489 let key_metadata = self.output.encoded_key_metadata(&file_metadata)?;
490
491 Ok(ManifestFile {
492 manifest_path: self.location,
493 manifest_length: file_metadata.size.try_into()?,
495 partition_spec_id: self.metadata.partition_spec.spec_id(),
496 content: self.metadata.content,
497 sequence_number: UNASSIGNED_SEQUENCE_NUMBER,
500 min_sequence_number: self.min_seq_num.unwrap_or(UNASSIGNED_SEQUENCE_NUMBER),
501 added_snapshot_id: self.snapshot_id.unwrap_or(UNASSIGNED_SNAPSHOT_ID),
502 added_files_count: Some(self.added_files),
503 existing_files_count: Some(self.existing_files),
504 deleted_files_count: Some(self.deleted_files),
505 added_rows_count: Some(self.added_rows),
506 existing_rows_count: Some(self.existing_rows),
507 deleted_rows_count: Some(self.deleted_rows),
508 partitions: Some(partition_summary),
509 key_metadata,
510 first_row_id: self.first_row_id,
511 })
512 }
513}
514
515struct PartitionFieldStats {
516 partition_type: PrimitiveType,
517
518 contains_null: bool,
519 contains_nan: Option<bool>,
520 lower_bound: Option<Datum>,
521 upper_bound: Option<Datum>,
522}
523
524impl PartitionFieldStats {
525 pub(crate) fn new(partition_type: PrimitiveType) -> Self {
526 Self {
527 partition_type,
528 contains_null: false,
529 contains_nan: Some(false),
530 upper_bound: None,
531 lower_bound: None,
532 }
533 }
534
535 pub(crate) fn update(&mut self, value: Option<PrimitiveLiteral>) -> Result<()> {
536 let Some(value) = value else {
537 self.contains_null = true;
538 return Ok(());
539 };
540 if !self.partition_type.compatible(&value) {
541 return Err(invalid_data!("value is not compatible with type"));
542 }
543 let value = Datum::new(self.partition_type.clone(), value);
544
545 if value.is_nan() {
546 self.contains_nan = Some(true);
547 return Ok(());
548 }
549
550 self.lower_bound = Some(self.lower_bound.take().map_or(value.clone(), |original| {
551 if value < original {
552 value.clone()
553 } else {
554 original
555 }
556 }));
557 self.upper_bound = Some(self.upper_bound.take().map_or(value.clone(), |original| {
558 if value > original { value } else { original }
559 }));
560
561 Ok(())
562 }
563
564 pub(crate) fn finish(self) -> FieldSummary {
565 FieldSummary {
566 contains_null: self.contains_null,
567 contains_nan: self.contains_nan,
568 upper_bound: self.upper_bound.map(|v| v.to_bytes().unwrap()),
569 lower_bound: self.lower_bound.map(|v| v.to_bytes().unwrap()),
570 }
571 }
572}
573
574#[cfg(test)]
575mod tests {
576 use std::collections::HashMap;
577 use std::fs;
578 use std::sync::Arc;
579
580 use tempfile::TempDir;
581
582 use super::*;
583 use crate::avro::define_named_types_repeatedly;
584 use crate::io::FileIO;
585 use crate::spec::{
586 DataContentType, DataFileBuilder, DataFileFormat, Literal, Manifest, NestedField,
587 PrimitiveType, Schema, Struct, Transform, Type,
588 };
589
590 #[tokio::test]
591 async fn test_add_delete_existing() {
592 let schema = Arc::new(
593 Schema::builder()
594 .with_fields(vec![
595 Arc::new(NestedField::optional(
596 1,
597 "id",
598 Type::Primitive(PrimitiveType::Int),
599 )),
600 Arc::new(NestedField::optional(
601 2,
602 "name",
603 Type::Primitive(PrimitiveType::String),
604 )),
605 ])
606 .build()
607 .unwrap(),
608 );
609 let metadata = ManifestMetadata {
610 schema_id: 0,
611 schema: schema.clone(),
612 partition_spec: PartitionSpec::builder(schema)
613 .with_spec_id(0)
614 .build()
615 .unwrap(),
616 content: ManifestContentType::Data,
617 format_version: FormatVersion::V2,
618 };
619 let mut entries = vec![
620 ManifestEntry {
621 status: ManifestStatus::Added,
622 snapshot_id: None,
623 sequence_number: Some(1),
624 file_sequence_number: Some(1),
625 data_file: DataFile {
626 content: DataContentType::Data,
627 file_path: "s3a://icebergdata/demo/s1/t1/data/00000-0-ba56fbfa-f2ff-40c9-bb27-565ad6dc2be8-00000.parquet".to_string(),
628 file_format: DataFileFormat::Parquet,
629 partition: Struct::empty(),
630 record_count: 1,
631 file_size_in_bytes: 5442,
632 column_sizes: HashMap::from([(1, 61), (2, 73)]),
633 value_counts: HashMap::from([(1, 1), (2, 1)]),
634 null_value_counts: HashMap::from([(1, 0), (2, 0)]),
635 nan_value_counts: HashMap::new(),
636 lower_bounds: HashMap::new(),
637 upper_bounds: HashMap::new(),
638 key_metadata: Some(Vec::new()),
639 split_offsets: Some(vec![4]),
640 equality_ids: None,
641 sort_order_id: None,
642 partition_spec_id: 0,
643 first_row_id: None,
644 referenced_data_file: None,
645 content_offset: None,
646 content_size_in_bytes: None,
647 },
648 },
649 ManifestEntry {
650 status: ManifestStatus::Deleted,
651 snapshot_id: Some(1),
652 sequence_number: Some(1),
653 file_sequence_number: Some(1),
654 data_file: DataFile {
655 content: DataContentType::Data,
656 file_path: "s3a://icebergdata/demo/s1/t1/data/00000-0-ba56fbfa-f2ff-40c9-bb27-565ad6dc2be8-00000.parquet".to_string(),
657 file_format: DataFileFormat::Parquet,
658 partition: Struct::empty(),
659 record_count: 1,
660 file_size_in_bytes: 5442,
661 column_sizes: HashMap::from([(1, 61), (2, 73)]),
662 value_counts: HashMap::from([(1, 1), (2, 1)]),
663 null_value_counts: HashMap::from([(1, 0), (2, 0)]),
664 nan_value_counts: HashMap::new(),
665 lower_bounds: HashMap::new(),
666 upper_bounds: HashMap::new(),
667 key_metadata: Some(Vec::new()),
668 split_offsets: Some(vec![4]),
669 equality_ids: None,
670 sort_order_id: None,
671 partition_spec_id: 0,
672 first_row_id: None,
673 referenced_data_file: None,
674 content_offset: None,
675 content_size_in_bytes: None,
676 },
677 },
678 ManifestEntry {
679 status: ManifestStatus::Existing,
680 snapshot_id: Some(1),
681 sequence_number: Some(1),
682 file_sequence_number: Some(1),
683 data_file: DataFile {
684 content: DataContentType::Data,
685 file_path: "s3a://icebergdata/demo/s1/t1/data/00000-0-ba56fbfa-f2ff-40c9-bb27-565ad6dc2be8-00000.parquet".to_string(),
686 file_format: DataFileFormat::Parquet,
687 partition: Struct::empty(),
688 record_count: 1,
689 file_size_in_bytes: 5442,
690 column_sizes: HashMap::from([(1, 61), (2, 73)]),
691 value_counts: HashMap::from([(1, 1), (2, 1)]),
692 null_value_counts: HashMap::from([(1, 0), (2, 0)]),
693 nan_value_counts: HashMap::new(),
694 lower_bounds: HashMap::new(),
695 upper_bounds: HashMap::new(),
696 key_metadata: Some(Vec::new()),
697 split_offsets: Some(vec![4]),
698 equality_ids: None,
699 sort_order_id: None,
700 partition_spec_id: 0,
701 first_row_id: None,
702 referenced_data_file: None,
703 content_offset: None,
704 content_size_in_bytes: None,
705 },
706 },
707 ];
708
709 let tmp_dir = TempDir::new().unwrap();
711 let path = tmp_dir.path().join("test_manifest.avro");
712 let io = FileIO::new_with_fs();
713 let output_file = io.new_output(path.to_str().unwrap()).unwrap();
714 let mut writer = ManifestWriterBuilder::new(
715 output_file,
716 Some(3),
717 metadata.schema.clone(),
718 metadata.partition_spec.clone(),
719 )
720 .build_v2_data();
721 writer.add_entry(entries[0].clone()).unwrap();
722 writer.add_delete_entry(entries[1].clone()).unwrap();
723 writer.add_existing_entry(entries[2].clone()).unwrap();
724 writer.write_manifest_file().await.unwrap();
725
726 let actual_manifest =
728 Manifest::parse_avro(fs::read(path).expect("read_file must succeed").as_slice())
729 .unwrap();
730
731 entries[0].snapshot_id = Some(3);
733 entries[1].snapshot_id = Some(3);
734 entries[0].file_sequence_number = None;
736 assert_eq!(actual_manifest, Manifest::new(metadata, entries));
737 }
738
739 #[tokio::test]
740 async fn test_v3_delete_manifest_delete_file_roundtrip() {
741 let schema = Arc::new(
742 Schema::builder()
743 .with_fields(vec![
744 Arc::new(NestedField::optional(
745 1,
746 "id",
747 Type::Primitive(PrimitiveType::Long),
748 )),
749 Arc::new(NestedField::optional(
750 2,
751 "data",
752 Type::Primitive(PrimitiveType::String),
753 )),
754 ])
755 .build()
756 .unwrap(),
757 );
758
759 let partition_spec = PartitionSpec::builder(schema.clone())
760 .with_spec_id(0)
761 .build()
762 .unwrap();
763
764 let delete_entry = ManifestEntry {
766 status: ManifestStatus::Added,
767 snapshot_id: None,
768 sequence_number: None,
769 file_sequence_number: None,
770 data_file: DataFile {
771 content: DataContentType::PositionDeletes,
772 file_path: "s3://bucket/table/data/delete-00000.parquet".to_string(),
773 file_format: DataFileFormat::Parquet,
774 partition: Struct::empty(),
775 record_count: 10,
776 file_size_in_bytes: 1024,
777 column_sizes: HashMap::new(),
778 value_counts: HashMap::new(),
779 null_value_counts: HashMap::new(),
780 nan_value_counts: HashMap::new(),
781 lower_bounds: HashMap::new(),
782 upper_bounds: HashMap::new(),
783 key_metadata: None,
784 split_offsets: None,
785 equality_ids: None,
786 sort_order_id: None,
787 partition_spec_id: 0,
788 first_row_id: None,
789 referenced_data_file: None,
790 content_offset: None,
791 content_size_in_bytes: None,
792 },
793 };
794
795 let tmp_dir = TempDir::new().unwrap();
797 let path = tmp_dir.path().join("v3_delete_manifest.avro");
798 let io = FileIO::new_with_fs();
799 let output_file = io.new_output(path.to_str().unwrap()).unwrap();
800
801 let mut writer = ManifestWriterBuilder::new(
802 output_file,
803 Some(1),
804 schema.clone(),
805 partition_spec.clone(),
806 )
807 .build_v3_deletes();
808
809 writer.add_entry(delete_entry).unwrap();
810 let manifest_file = writer.write_manifest_file().await.unwrap();
811
812 assert_eq!(manifest_file.content, ManifestContentType::Deletes);
814
815 let bytes = fs::read(&path).expect("read_file must succeed");
817 assert_eq!(manifest_file.manifest_length, bytes.len() as i64);
818 let actual_manifest = Manifest::parse_avro(&bytes).unwrap();
819
820 assert_eq!(
822 actual_manifest.metadata().content,
823 ManifestContentType::Deletes,
824 );
825 }
826
827 async fn roundtrip_two_partition_fields_of_type(
830 field_type: Type,
831 values: [Literal; 2],
832 ) -> Vec<u8> {
833 let schema = Arc::new(
834 Schema::builder()
835 .with_fields(vec![
836 NestedField::optional(1, "a", field_type.clone()).into(),
837 NestedField::optional(2, "b", field_type).into(),
838 ])
839 .build()
840 .unwrap(),
841 );
842 let partition_spec = PartitionSpec::builder(schema.clone())
843 .with_spec_id(0)
844 .add_partition_field("a", "a", Transform::Identity)
845 .unwrap()
846 .add_partition_field("b", "b", Transform::Identity)
847 .unwrap()
848 .build()
849 .unwrap();
850 let partition = Struct::from_iter(values.map(Some));
851 let tmp_dir = TempDir::new().unwrap();
852 let path = tmp_dir.path().join("manifest.avro");
853 let output_file = FileIO::new_with_fs()
854 .new_output(path.to_str().unwrap())
855 .unwrap();
856 let mut writer = ManifestWriterBuilder::new(output_file, Some(1), schema, partition_spec)
857 .build_v2_data();
858 writer
859 .add_file(
860 DataFileBuilder::default()
861 .content(DataContentType::Data)
862 .file_path("s3://bucket/table/data/a.parquet".to_string())
863 .file_format(DataFileFormat::Parquet)
864 .partition(partition.clone())
865 .record_count(1)
866 .file_size_in_bytes(1)
867 .partition_spec_id(0)
868 .build()
869 .unwrap(),
870 1,
871 )
872 .unwrap();
873 writer.write_manifest_file().await.unwrap();
874 let bs = fs::read(&path).unwrap();
875
876 for bs in [bs.clone(), define_named_types_repeatedly(&bs)] {
879 let manifest = Manifest::parse_avro(&bs).unwrap();
880
881 assert_eq!(*manifest.entries()[0].data_file().partition(), partition);
882 }
883 bs
884 }
885
886 #[tokio::test]
887 async fn test_write_manifest_header_marks_int_keyed_maps() {
888 let bs = roundtrip_two_partition_fields_of_type(Type::Primitive(PrimitiveType::Int), [
893 Literal::int(1),
894 Literal::int(2),
895 ])
896 .await;
897
898 let file = String::from_utf8_lossy(&bs);
899 assert_eq!(file.matches(r#""logicalType":"map""#).count(), 6);
902 }
903
904 #[tokio::test]
905 async fn test_write_manifest_with_repeated_decimal_partition_type() {
906 roundtrip_two_partition_fields_of_type(
907 Type::Primitive(PrimitiveType::Decimal {
908 precision: 10,
909 scale: 2,
910 }),
911 [Literal::decimal(12345), Literal::decimal(-678)],
912 )
913 .await;
914 }
915
916 #[tokio::test]
917 async fn test_write_manifest_with_repeated_fixed_partition_type() {
918 roundtrip_two_partition_fields_of_type(Type::Primitive(PrimitiveType::Fixed(4)), [
919 Literal::fixed(vec![1, 2, 3, 4]),
920 Literal::fixed(vec![5, 6, 7, 8]),
921 ])
922 .await;
923 }
924}