1mod _serde;
19
20mod data_file;
21pub use data_file::*;
22mod entry;
23pub use entry::*;
24mod metadata;
25pub use metadata::*;
26mod reader;
27pub use reader::*;
28mod writer;
29use std::sync::Arc;
30use std::sync::atomic::{AtomicBool, Ordering};
31
32use apache_avro::Reader as AvroReader;
33use apache_avro::error::Details;
34pub use writer::*;
35
36use super::{
37 Datum, FormatVersion, ManifestContentType, PartitionSpec, PrimitiveType, Schema, Struct, Type,
38 UNASSIGNED_SEQUENCE_NUMBER,
39};
40use crate::avro::{Resolved, define_named_types_once};
41use crate::error::{Error, Result, invalid_data};
42
43static WARNED_REPEATED_DEFINITIONS: AtomicBool = AtomicBool::new(false);
46
47#[derive(Debug, PartialEq, Eq, Clone)]
49pub struct Manifest {
50 metadata: ManifestMetadata,
51 entries: Vec<ManifestEntryRef>,
52}
53
54impl Manifest {
55 pub(crate) fn try_from_avro_bytes(
58 bs: &[u8],
59 location: Option<&str>,
60 ) -> Result<(ManifestMetadata, Vec<ManifestEntry>)> {
61 let rewritten;
62 let reader = match AvroReader::new(bs) {
63 Ok(reader) => reader,
64 Err(e) if matches!(e.details(), Details::AmbiguousSchemaDefinition(_)) => {
68 let (bs, repeated) = match define_named_types_once(bs) {
69 Ok(Some(rewrite)) => rewrite,
70 Ok(None) => return Err(e.into()),
71 Err(rewrite_error) => {
72 return Err(Error::from(e).with_context(
73 "error defining each named type once",
74 rewrite_error.to_string(),
75 ));
76 }
77 };
78 rewritten = bs;
79 let reader = AvroReader::new(rewritten.as_slice()).map_err(|retry_error| {
80 Error::from(e).with_context(
81 "error after defining each named type once",
82 retry_error.to_string(),
83 )
84 })?;
85 let location = location.unwrap_or("<unknown location>");
86 if WARNED_REPEATED_DEFINITIONS.swap(true, Ordering::Relaxed) {
87 tracing::debug!(
88 "Manifest {location} defines Avro named types {repeated:?} more than once."
89 );
90 } else {
91 tracing::warn!(
92 "Manifest {location} defines Avro named types {repeated:?} more than once, \
93 which the Avro specification doesn't allow. Reading it with each repeated \
94 definition replaced by a reference to the first. Later manifests like \
95 this are logged at debug level."
96 );
97 }
98 reader
99 }
100 Err(e) => return Err(e.into()),
101 };
102
103 let meta = reader.user_metadata();
105 let metadata = ManifestMetadata::parse(meta)?;
106
107 let partition_type = metadata.partition_spec.partition_type(&metadata.schema)?;
109 let partition_struct_type = Type::Struct(partition_type.clone());
113
114 let entries = match metadata.format_version {
115 FormatVersion::V1 => reader
116 .into_deser_iter::<Resolved<_serde::ManifestEntryV1>>()
117 .map(|entry| {
118 entry?.0.try_into(
119 metadata.partition_spec.spec_id(),
120 &partition_struct_type,
121 &metadata.schema,
122 )
123 })
124 .collect::<Result<Vec<_>>>()?,
125 FormatVersion::V2 | FormatVersion::V3 => reader
127 .into_deser_iter::<Resolved<_serde::ManifestEntryV2>>()
128 .map(|entry| {
129 entry?.0.try_into(
130 metadata.partition_spec.spec_id(),
131 &partition_struct_type,
132 &metadata.schema,
133 )
134 })
135 .collect::<Result<Vec<_>>>()?,
136 };
137
138 Ok((metadata, entries))
139 }
140
141 pub fn parse_avro(bs: &[u8]) -> Result<Self> {
143 let (metadata, entries) = Self::try_from_avro_bytes(bs, None)?;
144 Ok(Self::new(metadata, entries))
145 }
146
147 pub fn entries(&self) -> &[ManifestEntryRef] {
149 &self.entries
150 }
151
152 pub fn metadata(&self) -> &ManifestMetadata {
154 &self.metadata
155 }
156
157 pub fn into_parts(self) -> (Vec<ManifestEntryRef>, ManifestMetadata) {
159 let Self { entries, metadata } = self;
160 (entries, metadata)
161 }
162
163 pub fn new(metadata: ManifestMetadata, entries: Vec<ManifestEntry>) -> Self {
165 Self {
166 metadata,
167 entries: entries.into_iter().map(Arc::new).collect(),
168 }
169 }
170}
171
172pub fn serialize_data_file_to_json(
174 data_file: DataFile,
175 partition_type: &super::StructType,
176 format_version: FormatVersion,
177) -> Result<String> {
178 let partition_struct_type = Type::Struct(partition_type.clone());
179 let serde = _serde::DataFileSerde::try_from(data_file, &partition_struct_type, format_version)?;
180 serde_json::to_string(&serde)
181 .map_err(|e| invalid_data!("Failed to serialize DataFile to JSON!").with_source(e))
182}
183
184pub fn deserialize_data_file_from_json(
186 json: &str,
187 partition_spec_id: i32,
188 partition_type: &super::StructType,
189 schema: &Schema,
190) -> Result<DataFile> {
191 let serde = serde_json::from_str::<_serde::DataFileSerde>(json)
192 .map_err(|e| invalid_data!("Failed to deserialize JSON to DataFile!").with_source(e))?;
193
194 let partition_struct_type = Type::Struct(partition_type.clone());
195 serde.try_into(partition_spec_id, &partition_struct_type, schema)
196}
197
198#[cfg(test)]
199mod tests {
200 use std::collections::HashMap;
201 use std::fs;
202 use std::sync::Arc;
203
204 use apache_avro::types::Value as AvroValue;
205 use apache_avro::{Codec, Writer, to_value};
206 use serde_json::{Value, to_vec};
207 use tempfile::TempDir;
208 use tracing::Level;
209 use tracing::subscriber::with_default;
210
211 use super::*;
212 use crate::ErrorKind;
213 use crate::io::FileIO;
214 use crate::spec::{Literal, NestedField, PrimitiveType, Struct, Transform, Type};
215
216 #[tokio::test]
217 async fn test_parse_manifest_v2_unpartition() {
218 let schema = Arc::new(
219 Schema::builder()
220 .with_fields(vec![
221 Arc::new(NestedField::optional(
223 1,
224 "id",
225 Type::Primitive(PrimitiveType::Long),
226 )),
227 Arc::new(NestedField::optional(
228 2,
229 "v_int",
230 Type::Primitive(PrimitiveType::Int),
231 )),
232 Arc::new(NestedField::optional(
233 3,
234 "v_long",
235 Type::Primitive(PrimitiveType::Long),
236 )),
237 Arc::new(NestedField::optional(
238 4,
239 "v_float",
240 Type::Primitive(PrimitiveType::Float),
241 )),
242 Arc::new(NestedField::optional(
243 5,
244 "v_double",
245 Type::Primitive(PrimitiveType::Double),
246 )),
247 Arc::new(NestedField::optional(
248 6,
249 "v_varchar",
250 Type::Primitive(PrimitiveType::String),
251 )),
252 Arc::new(NestedField::optional(
253 7,
254 "v_bool",
255 Type::Primitive(PrimitiveType::Boolean),
256 )),
257 Arc::new(NestedField::optional(
258 8,
259 "v_date",
260 Type::Primitive(PrimitiveType::Date),
261 )),
262 Arc::new(NestedField::optional(
263 9,
264 "v_timestamp",
265 Type::Primitive(PrimitiveType::Timestamptz),
266 )),
267 Arc::new(NestedField::optional(
268 10,
269 "v_decimal",
270 Type::Primitive(PrimitiveType::Decimal {
271 precision: 36,
272 scale: 10,
273 }),
274 )),
275 Arc::new(NestedField::optional(
276 11,
277 "v_ts_ntz",
278 Type::Primitive(PrimitiveType::Timestamp),
279 )),
280 Arc::new(NestedField::optional(
281 12,
282 "v_ts_ns_ntz",
283 Type::Primitive(PrimitiveType::TimestampNs),
284 )),
285 ])
286 .build()
287 .unwrap(),
288 );
289 let metadata = ManifestMetadata {
290 schema_id: 0,
291 schema: schema.clone(),
292 partition_spec: PartitionSpec::builder(schema)
293 .with_spec_id(0)
294 .build()
295 .unwrap(),
296 content: ManifestContentType::Data,
297 format_version: FormatVersion::V2,
298 };
299 let mut entries = vec![
300 ManifestEntry {
301 status: ManifestStatus::Added,
302 snapshot_id: None,
303 sequence_number: None,
304 file_sequence_number: None,
305 data_file: DataFile {content:DataContentType::Data,file_path:"s3a://icebergdata/demo/s1/t1/data/00000-0-ba56fbfa-f2ff-40c9-bb27-565ad6dc2be8-00000.parquet".to_string(),file_format:DataFileFormat::Parquet,partition:Struct::empty(),record_count:1,file_size_in_bytes:5442,column_sizes:HashMap::from([(0,73),(6,34),(2,73),(7,61),(3,61),(5,62),(9,79),(10,73),(1,61),(4,73),(8,73)]),value_counts:HashMap::from([(4,1),(5,1),(2,1),(0,1),(3,1),(6,1),(8,1),(1,1),(10,1),(7,1),(9,1)]),null_value_counts:HashMap::from([(1,0),(6,0),(2,0),(8,0),(0,0),(3,0),(5,0),(9,0),(7,0),(4,0),(10,0)]),nan_value_counts:HashMap::new(),lower_bounds:HashMap::new(),upper_bounds:HashMap::new(),key_metadata:None,split_offsets:Some(vec![4]),equality_ids:Some(Vec::new()),sort_order_id:None, partition_spec_id: 0,first_row_id: None,referenced_data_file: None,content_offset: None,content_size_in_bytes: None }
306 }
307 ];
308
309 let tmp_dir = TempDir::new().unwrap();
311 let path = tmp_dir.path().join("test_manifest.avro");
312 let io = FileIO::new_with_fs();
313 let output_file = io.new_output(path.to_str().unwrap()).unwrap();
314 let mut writer = ManifestWriterBuilder::new(
315 output_file,
316 Some(1),
317 metadata.schema.clone(),
318 metadata.partition_spec.clone(),
319 )
320 .build_v2_data();
321 for entry in &entries {
322 writer.add_entry(entry.clone()).unwrap();
323 }
324 writer.write_manifest_file().await.unwrap();
325
326 let actual_manifest =
328 Manifest::parse_avro(fs::read(path).expect("read_file must succeed").as_slice())
329 .unwrap();
330 entries[0].snapshot_id = Some(1);
332 assert_eq!(actual_manifest, Manifest::new(metadata, entries));
333 }
334
335 #[test]
336 fn test_parse_snappy_manifest_v2() {
337 let schema = Arc::new(
338 Schema::builder()
339 .with_fields(vec![Arc::new(NestedField::optional(
340 1,
341 "id",
342 Type::Primitive(PrimitiveType::Long),
343 ))])
344 .build()
345 .unwrap(),
346 );
347 let partition_spec = PartitionSpec::builder(schema.clone())
348 .with_spec_id(0)
349 .build()
350 .unwrap();
351
352 for (manifest_content, file_content, file_path) in [
353 (
354 ManifestContentType::Data,
355 DataContentType::Data,
356 "s3://bucket/table/data/data.parquet",
357 ),
358 (
359 ManifestContentType::Deletes,
360 DataContentType::PositionDeletes,
361 "s3://bucket/table/data/delete.parquet",
362 ),
363 ] {
364 let metadata = ManifestMetadata {
365 schema_id: 0,
366 schema: schema.clone(),
367 partition_spec: partition_spec.clone(),
368 content: manifest_content,
369 format_version: FormatVersion::V2,
370 };
371 let entry = ManifestEntry {
372 status: ManifestStatus::Added,
373 snapshot_id: Some(1),
374 sequence_number: None,
375 file_sequence_number: None,
376 data_file: DataFile {
377 content: file_content,
378 file_path: file_path.to_string(),
379 file_format: DataFileFormat::Parquet,
380 partition: Struct::empty(),
381 record_count: 1,
382 file_size_in_bytes: 1024,
383 column_sizes: HashMap::new(),
384 value_counts: HashMap::new(),
385 null_value_counts: HashMap::new(),
386 nan_value_counts: HashMap::new(),
387 lower_bounds: HashMap::new(),
388 upper_bounds: HashMap::new(),
389 key_metadata: None,
390 split_offsets: None,
391 equality_ids: None,
392 sort_order_id: None,
393 partition_spec_id: 0,
394 first_row_id: None,
395 referenced_data_file: None,
396 content_offset: None,
397 content_size_in_bytes: None,
398 },
399 };
400
401 let partition_type = metadata
402 .partition_spec
403 .partition_type(&metadata.schema)
404 .unwrap();
405 let avro_schema = manifest_schema_v2(&partition_type).unwrap();
406 let mut writer = Writer::with_codec(&avro_schema, Vec::new(), Codec::Snappy).unwrap();
407 writer
408 .add_user_metadata("schema".to_string(), to_vec(&metadata.schema).unwrap())
409 .unwrap();
410 writer
411 .add_user_metadata(
412 "schema-id".to_string(),
413 metadata.schema.schema_id().to_string(),
414 )
415 .unwrap();
416 writer
417 .add_user_metadata(
418 "partition-spec".to_string(),
419 to_vec(&metadata.partition_spec.fields()).unwrap(),
420 )
421 .unwrap();
422 writer
423 .add_user_metadata(
424 "partition-spec-id".to_string(),
425 metadata.partition_spec.spec_id().to_string(),
426 )
427 .unwrap();
428 writer
429 .add_user_metadata(
430 "format-version".to_string(),
431 (metadata.format_version as u8).to_string(),
432 )
433 .unwrap();
434 writer
435 .add_user_metadata("content".to_string(), metadata.content.to_string())
436 .unwrap();
437 let value = to_value(
438 _serde::ManifestEntryV2::try_from(
439 entry.clone(),
440 &Type::Struct(partition_type.clone()),
441 )
442 .unwrap(),
443 )
444 .unwrap()
445 .resolve(&avro_schema)
446 .unwrap();
447 writer.append_value(value).unwrap();
448 let bs = writer.into_inner().unwrap();
449
450 let parsed_manifest = Manifest::parse_avro(&bs).unwrap();
451
452 assert_eq!(parsed_manifest, Manifest::new(metadata, vec![entry]));
453 }
454 }
455
456 #[tokio::test]
457 async fn test_parse_manifest_v2_partition() {
458 let schema = Arc::new(
459 Schema::builder()
460 .with_fields(vec![
461 Arc::new(NestedField::optional(
462 1,
463 "id",
464 Type::Primitive(PrimitiveType::Long),
465 )),
466 Arc::new(NestedField::optional(
467 2,
468 "v_int",
469 Type::Primitive(PrimitiveType::Int),
470 )),
471 Arc::new(NestedField::optional(
472 3,
473 "v_long",
474 Type::Primitive(PrimitiveType::Long),
475 )),
476 Arc::new(NestedField::optional(
477 4,
478 "v_float",
479 Type::Primitive(PrimitiveType::Float),
480 )),
481 Arc::new(NestedField::optional(
482 5,
483 "v_double",
484 Type::Primitive(PrimitiveType::Double),
485 )),
486 Arc::new(NestedField::optional(
487 6,
488 "v_varchar",
489 Type::Primitive(PrimitiveType::String),
490 )),
491 Arc::new(NestedField::optional(
492 7,
493 "v_bool",
494 Type::Primitive(PrimitiveType::Boolean),
495 )),
496 Arc::new(NestedField::optional(
497 8,
498 "v_date",
499 Type::Primitive(PrimitiveType::Date),
500 )),
501 Arc::new(NestedField::optional(
502 9,
503 "v_timestamp",
504 Type::Primitive(PrimitiveType::Timestamptz),
505 )),
506 Arc::new(NestedField::optional(
507 10,
508 "v_decimal",
509 Type::Primitive(PrimitiveType::Decimal {
510 precision: 36,
511 scale: 10,
512 }),
513 )),
514 Arc::new(NestedField::optional(
515 11,
516 "v_ts_ntz",
517 Type::Primitive(PrimitiveType::Timestamp),
518 )),
519 Arc::new(NestedField::optional(
520 12,
521 "v_ts_ns_ntz",
522 Type::Primitive(PrimitiveType::TimestampNs),
523 )),
524 ])
525 .build()
526 .unwrap(),
527 );
528 let metadata = ManifestMetadata {
529 schema_id: 0,
530 schema: schema.clone(),
531 partition_spec: PartitionSpec::builder(schema)
532 .with_spec_id(0)
533 .add_partition_field("v_int", "v_int", Transform::Identity)
534 .unwrap()
535 .add_partition_field("v_long", "v_long", Transform::Identity)
536 .unwrap()
537 .build()
538 .unwrap(),
539 content: ManifestContentType::Data,
540 format_version: FormatVersion::V2,
541 };
542 let mut entries = vec![ManifestEntry {
543 status: ManifestStatus::Added,
544 snapshot_id: None,
545 sequence_number: None,
546 file_sequence_number: None,
547 data_file: DataFile {
548 content: DataContentType::Data,
549 file_format: DataFileFormat::Parquet,
550 file_path: "s3a://icebergdata/demo/s1/t1/data/00000-0-378b56f5-5c52-4102-a2c2-f05f8a7cbe4a-00000.parquet".to_string(),
551 partition: Struct::from_iter(
552 vec![
553 Some(Literal::int(1)),
554 Some(Literal::long(1000)),
555 ]
556 .into_iter()
557 ),
558 record_count: 1,
559 file_size_in_bytes: 5442,
560 column_sizes: HashMap::from([
561 (0, 73),
562 (6, 34),
563 (2, 73),
564 (7, 61),
565 (3, 61),
566 (5, 62),
567 (9, 79),
568 (10, 73),
569 (1, 61),
570 (4, 73),
571 (8, 73)
572 ]),
573 value_counts: HashMap::from([
574 (4, 1),
575 (5, 1),
576 (2, 1),
577 (0, 1),
578 (3, 1),
579 (6, 1),
580 (8, 1),
581 (1, 1),
582 (10, 1),
583 (7, 1),
584 (9, 1)
585 ]),
586 null_value_counts: HashMap::from([
587 (1, 0),
588 (6, 0),
589 (2, 0),
590 (8, 0),
591 (0, 0),
592 (3, 0),
593 (5, 0),
594 (9, 0),
595 (7, 0),
596 (4, 0),
597 (10, 0)
598 ]),
599 nan_value_counts: HashMap::new(),
600 lower_bounds: HashMap::new(),
601 upper_bounds: HashMap::new(),
602 key_metadata: None,
603 split_offsets: Some(vec![4]),
604 equality_ids: Some(Vec::new()),
605 sort_order_id: None,
606 partition_spec_id: 0,
607 first_row_id: None,
608 referenced_data_file: None,
609 content_offset: None,
610 content_size_in_bytes: None,
611 },
612 }];
613
614 let tmp_dir = TempDir::new().unwrap();
616 let path = tmp_dir.path().join("test_manifest.avro");
617 let io = FileIO::new_with_fs();
618 let output_file = io.new_output(path.to_str().unwrap()).unwrap();
619 let mut writer = ManifestWriterBuilder::new(
620 output_file,
621 Some(2),
622 metadata.schema.clone(),
623 metadata.partition_spec.clone(),
624 )
625 .build_v2_data();
626 for entry in &entries {
627 writer.add_entry(entry.clone()).unwrap();
628 }
629 let manifest_file = writer.write_manifest_file().await.unwrap();
630 assert_eq!(manifest_file.sequence_number, UNASSIGNED_SEQUENCE_NUMBER);
631 assert_eq!(
632 manifest_file.min_sequence_number,
633 UNASSIGNED_SEQUENCE_NUMBER
634 );
635
636 let actual_manifest =
638 Manifest::parse_avro(fs::read(path).expect("read_file must succeed").as_slice())
639 .unwrap();
640 entries[0].snapshot_id = Some(2);
642 assert_eq!(actual_manifest, Manifest::new(metadata, entries));
643 }
644
645 #[tokio::test]
646 async fn test_parse_manifest_v1_unpartition() {
647 let schema = Arc::new(
648 Schema::builder()
649 .with_schema_id(1)
650 .with_fields(vec![
651 Arc::new(NestedField::optional(
652 1,
653 "id",
654 Type::Primitive(PrimitiveType::Int),
655 )),
656 Arc::new(NestedField::optional(
657 2,
658 "data",
659 Type::Primitive(PrimitiveType::String),
660 )),
661 Arc::new(NestedField::optional(
662 3,
663 "comment",
664 Type::Primitive(PrimitiveType::String),
665 )),
666 ])
667 .build()
668 .unwrap(),
669 );
670 let metadata = ManifestMetadata {
671 schema_id: 1,
672 schema: schema.clone(),
673 partition_spec: PartitionSpec::builder(schema)
674 .with_spec_id(0)
675 .build()
676 .unwrap(),
677 content: ManifestContentType::Data,
678 format_version: FormatVersion::V1,
679 };
680 let mut entries = vec![ManifestEntry {
681 status: ManifestStatus::Added,
682 snapshot_id: Some(0),
683 sequence_number: Some(0),
684 file_sequence_number: Some(0),
685 data_file: DataFile {
686 content: DataContentType::Data,
687 file_path: "s3://testbucket/iceberg_data/iceberg_ctl/iceberg_db/iceberg_tbl/data/00000-7-45268d71-54eb-476c-b42c-942d880c04a1-00001.parquet".to_string(),
688 file_format: DataFileFormat::Parquet,
689 partition: Struct::empty(),
690 record_count: 1,
691 file_size_in_bytes: 875,
692 column_sizes: HashMap::from([(1,47),(2,48),(3,52)]),
693 value_counts: HashMap::from([(1,1),(2,1),(3,1)]),
694 null_value_counts: HashMap::from([(1,0),(2,0),(3,0)]),
695 nan_value_counts: HashMap::new(),
696 lower_bounds: HashMap::from([(1,Datum::int(1)),(2,Datum::string("a")),(3,Datum::string("AC/DC"))]),
697 upper_bounds: HashMap::from([(1,Datum::int(1)),(2,Datum::string("a")),(3,Datum::string("AC/DC"))]),
698 key_metadata: None,
699 split_offsets: Some(vec![4]),
700 equality_ids: None,
701 sort_order_id: Some(0),
702 partition_spec_id: 0,
703 first_row_id: None,
704 referenced_data_file: None,
705 content_offset: None,
706 content_size_in_bytes: None,
707 }
708 }];
709
710 let tmp_dir = TempDir::new().unwrap();
712 let path = tmp_dir.path().join("test_manifest.avro");
713 let io = FileIO::new_with_fs();
714 let output_file = io.new_output(path.to_str().unwrap()).unwrap();
715 let mut writer = ManifestWriterBuilder::new(
716 output_file,
717 Some(3),
718 metadata.schema.clone(),
719 metadata.partition_spec.clone(),
720 )
721 .build_v1();
722 for entry in &entries {
723 writer.add_entry(entry.clone()).unwrap();
724 }
725 writer.write_manifest_file().await.unwrap();
726
727 let actual_manifest =
729 Manifest::parse_avro(fs::read(path).expect("read_file must succeed").as_slice())
730 .unwrap();
731 entries[0].snapshot_id = Some(3);
733 assert_eq!(actual_manifest, Manifest::new(metadata, entries));
734 }
735
736 #[tokio::test]
737 async fn test_parse_manifest_v1_partition() {
738 let schema = Arc::new(
739 Schema::builder()
740 .with_fields(vec![
741 Arc::new(NestedField::optional(
742 1,
743 "id",
744 Type::Primitive(PrimitiveType::Long),
745 )),
746 Arc::new(NestedField::optional(
747 2,
748 "data",
749 Type::Primitive(PrimitiveType::String),
750 )),
751 Arc::new(NestedField::optional(
752 3,
753 "category",
754 Type::Primitive(PrimitiveType::String),
755 )),
756 ])
757 .build()
758 .unwrap(),
759 );
760 let metadata = ManifestMetadata {
761 schema_id: 0,
762 schema: schema.clone(),
763 partition_spec: PartitionSpec::builder(schema)
764 .add_partition_field("category", "category", Transform::Identity)
765 .unwrap()
766 .build()
767 .unwrap(),
768 content: ManifestContentType::Data,
769 format_version: FormatVersion::V1,
770 };
771 let mut entries = vec![
772 ManifestEntry {
773 status: ManifestStatus::Added,
774 snapshot_id: Some(0),
775 sequence_number: Some(0),
776 file_sequence_number: Some(0),
777 data_file: DataFile {
778 content: DataContentType::Data,
779 file_path: "s3://testbucket/prod/db/sample/data/category=x/00010-1-d5c93668-1e52-41ac-92a6-bba590cbf249-00001.parquet".to_string(),
780 file_format: DataFileFormat::Parquet,
781 partition: Struct::from_iter(
782 vec![
783 Some(
784 Literal::string("x"),
785 ),
786 ]
787 .into_iter()
788 ),
789 record_count: 1,
790 file_size_in_bytes: 874,
791 column_sizes: HashMap::from([(1, 46), (2, 48), (3, 48)]),
792 value_counts: HashMap::from([(1, 1), (2, 1), (3, 1)]),
793 null_value_counts: HashMap::from([(1, 0), (2, 0), (3, 0)]),
794 nan_value_counts: HashMap::new(),
795 lower_bounds: HashMap::from([
796 (1, Datum::long(1)),
797 (2, Datum::string("a")),
798 (3, Datum::string("x"))
799 ]),
800 upper_bounds: HashMap::from([
801 (1, Datum::long(1)),
802 (2, Datum::string("a")),
803 (3, Datum::string("x"))
804 ]),
805 key_metadata: None,
806 split_offsets: Some(vec![4]),
807 equality_ids: None,
808 sort_order_id: Some(0),
809 partition_spec_id: 0,
810 first_row_id: None,
811 referenced_data_file: None,
812 content_offset: None,
813 content_size_in_bytes: None,
814 },
815 }
816 ];
817
818 let tmp_dir = TempDir::new().unwrap();
820 let path = tmp_dir.path().join("test_manifest.avro");
821 let io = FileIO::new_with_fs();
822 let output_file = io.new_output(path.to_str().unwrap()).unwrap();
823 let mut writer = ManifestWriterBuilder::new(
824 output_file,
825 Some(2),
826 metadata.schema.clone(),
827 metadata.partition_spec.clone(),
828 )
829 .build_v1();
830 for entry in &entries {
831 writer.add_entry(entry.clone()).unwrap();
832 }
833 let manifest_file = writer.write_manifest_file().await.unwrap();
834 let partitions = manifest_file.partitions.unwrap();
835 assert_eq!(partitions.len(), 1);
836 assert_eq!(
837 partitions[0].clone().lower_bound.unwrap(),
838 Datum::string("x").to_bytes().unwrap()
839 );
840 assert_eq!(
841 partitions[0].clone().upper_bound.unwrap(),
842 Datum::string("x").to_bytes().unwrap()
843 );
844
845 let actual_manifest =
847 Manifest::parse_avro(fs::read(path).expect("read_file must succeed").as_slice())
848 .unwrap();
849 entries[0].snapshot_id = Some(2);
851 assert_eq!(actual_manifest, Manifest::new(metadata, entries));
852 }
853
854 #[tokio::test]
855 async fn test_parse_manifest_with_schema_evolution() {
856 let schema = Arc::new(
857 Schema::builder()
858 .with_fields(vec![
859 Arc::new(NestedField::optional(
860 1,
861 "id",
862 Type::Primitive(PrimitiveType::Long),
863 )),
864 Arc::new(NestedField::optional(
865 2,
866 "v_int",
867 Type::Primitive(PrimitiveType::Int),
868 )),
869 ])
870 .build()
871 .unwrap(),
872 );
873 let metadata = ManifestMetadata {
874 schema_id: 0,
875 schema: schema.clone(),
876 partition_spec: PartitionSpec::builder(schema)
877 .with_spec_id(0)
878 .build()
879 .unwrap(),
880 content: ManifestContentType::Data,
881 format_version: FormatVersion::V2,
882 };
883 let entries = vec![ManifestEntry {
884 status: ManifestStatus::Added,
885 snapshot_id: None,
886 sequence_number: None,
887 file_sequence_number: None,
888 data_file: DataFile {
889 content: DataContentType::Data,
890 file_format: DataFileFormat::Parquet,
891 file_path: "s3a://icebergdata/demo/s1/t1/data/00000-0-378b56f5-5c52-4102-a2c2-f05f8a7cbe4a-00000.parquet".to_string(),
892 partition: Struct::empty(),
893 record_count: 1,
894 file_size_in_bytes: 5442,
895 column_sizes: HashMap::from([
896 (1, 61),
897 (2, 73),
898 (3, 61),
899 ]),
900 value_counts: HashMap::default(),
901 null_value_counts: HashMap::default(),
902 nan_value_counts: HashMap::new(),
903 lower_bounds: HashMap::from([
904 (1, Datum::long(1)),
905 (2, Datum::int(2)),
906 (3, Datum::string("x"))
907 ]),
908 upper_bounds: HashMap::from([
909 (1, Datum::long(1)),
910 (2, Datum::int(2)),
911 (3, Datum::string("x"))
912 ]),
913 key_metadata: None,
914 split_offsets: Some(vec![4]),
915 equality_ids: None,
916 sort_order_id: None,
917 partition_spec_id: 0,
918 first_row_id: None,
919 referenced_data_file: None,
920 content_offset: None,
921 content_size_in_bytes: None,
922 },
923 }];
924
925 let tmp_dir = TempDir::new().unwrap();
927 let path = tmp_dir.path().join("test_manifest.avro");
928 let io = FileIO::new_with_fs();
929 let output_file = io.new_output(path.to_str().unwrap()).unwrap();
930 let mut writer = ManifestWriterBuilder::new(
931 output_file,
932 Some(2),
933 metadata.schema.clone(),
934 metadata.partition_spec.clone(),
935 )
936 .build_v2_data();
937 for entry in &entries {
938 writer.add_entry(entry.clone()).unwrap();
939 }
940 writer.write_manifest_file().await.unwrap();
941
942 let actual_manifest =
944 Manifest::parse_avro(fs::read(path).expect("read_file must succeed").as_slice())
945 .unwrap();
946
947 let schema = Arc::new(
951 Schema::builder()
952 .with_fields(vec![
953 Arc::new(NestedField::optional(
954 1,
955 "id",
956 Type::Primitive(PrimitiveType::Long),
957 )),
958 Arc::new(NestedField::optional(
959 2,
960 "v_int",
961 Type::Primitive(PrimitiveType::Int),
962 )),
963 ])
964 .build()
965 .unwrap(),
966 );
967 let expected_manifest = Manifest {
968 metadata: ManifestMetadata {
969 schema_id: 0,
970 schema: schema.clone(),
971 partition_spec: PartitionSpec::builder(schema).with_spec_id(0).build().unwrap(),
972 content: ManifestContentType::Data,
973 format_version: FormatVersion::V2,
974 },
975 entries: vec![Arc::new(ManifestEntry {
976 status: ManifestStatus::Added,
977 snapshot_id: Some(2),
978 sequence_number: None,
979 file_sequence_number: None,
980 data_file: DataFile {
981 content: DataContentType::Data,
982 file_format: DataFileFormat::Parquet,
983 file_path: "s3a://icebergdata/demo/s1/t1/data/00000-0-378b56f5-5c52-4102-a2c2-f05f8a7cbe4a-00000.parquet".to_string(),
984 partition: Struct::empty(),
985 record_count: 1,
986 file_size_in_bytes: 5442,
987 column_sizes: HashMap::from([
988 (1, 61),
989 (2, 73),
990 (3, 61),
991 ]),
992 value_counts: HashMap::default(),
993 null_value_counts: HashMap::default(),
994 nan_value_counts: HashMap::new(),
995 lower_bounds: HashMap::from([
996 (1, Datum::long(1)),
997 (2, Datum::int(2)),
998 ]),
999 upper_bounds: HashMap::from([
1000 (1, Datum::long(1)),
1001 (2, Datum::int(2)),
1002 ]),
1003 key_metadata: None,
1004 split_offsets: Some(vec![4]),
1005 equality_ids: None,
1006 sort_order_id: None,
1007 partition_spec_id: 0,
1008 first_row_id: None,
1009 referenced_data_file: None,
1010 content_offset: None,
1011 content_size_in_bytes: None,
1012 },
1013 })],
1014 };
1015
1016 assert_eq!(actual_manifest, expected_manifest);
1017 }
1018
1019 #[tokio::test]
1020 async fn test_manifest_summary() {
1021 let schema = Arc::new(
1022 Schema::builder()
1023 .with_fields(vec![
1024 Arc::new(NestedField::optional(
1025 1,
1026 "time",
1027 Type::Primitive(PrimitiveType::Date),
1028 )),
1029 Arc::new(NestedField::optional(
1030 2,
1031 "v_float",
1032 Type::Primitive(PrimitiveType::Float),
1033 )),
1034 Arc::new(NestedField::optional(
1035 3,
1036 "v_double",
1037 Type::Primitive(PrimitiveType::Double),
1038 )),
1039 ])
1040 .build()
1041 .unwrap(),
1042 );
1043 let partition_spec = PartitionSpec::builder(schema.clone())
1044 .with_spec_id(0)
1045 .add_partition_field("time", "year_of_time", Transform::Year)
1046 .unwrap()
1047 .add_partition_field("v_float", "f", Transform::Identity)
1048 .unwrap()
1049 .add_partition_field("v_double", "d", Transform::Identity)
1050 .unwrap()
1051 .build()
1052 .unwrap();
1053 let metadata = ManifestMetadata {
1054 schema_id: 0,
1055 schema,
1056 partition_spec,
1057 content: ManifestContentType::Data,
1058 format_version: FormatVersion::V2,
1059 };
1060 let entries = vec![
1061 ManifestEntry {
1062 status: ManifestStatus::Added,
1063 snapshot_id: None,
1064 sequence_number: None,
1065 file_sequence_number: None,
1066 data_file: DataFile {
1067 content: DataContentType::Data,
1068 file_path: "s3a://icebergdata/demo/s1/t1/data/00000-0-ba56fbfa-f2ff-40c9-bb27-565ad6dc2be8-00000.parquet".to_string(),
1069 file_format: DataFileFormat::Parquet,
1070 partition: Struct::from_iter(
1071 vec![
1072 Some(Literal::int(2021)),
1073 Some(Literal::float(1.0_f32)),
1074 Some(Literal::double(2.0)),
1075 ]
1076 ),
1077 record_count: 1,
1078 file_size_in_bytes: 5442,
1079 column_sizes: HashMap::from([(0,73),(6,34),(2,73),(7,61),(3,61),(5,62),(9,79),(10,73),(1,61),(4,73),(8,73)]),
1080 value_counts: HashMap::from([(4,1),(5,1),(2,1),(0,1),(3,1),(6,1),(8,1),(1,1),(10,1),(7,1),(9,1)]),
1081 null_value_counts: HashMap::from([(1,0),(6,0),(2,0),(8,0),(0,0),(3,0),(5,0),(9,0),(7,0),(4,0),(10,0)]),
1082 nan_value_counts: HashMap::new(),
1083 lower_bounds: HashMap::new(),
1084 upper_bounds: HashMap::new(),
1085 key_metadata: None,
1086 split_offsets: Some(vec![4]),
1087 equality_ids: None,
1088 sort_order_id: None,
1089 partition_spec_id: 0,
1090 first_row_id: None,
1091 referenced_data_file: None,
1092 content_offset: None,
1093 content_size_in_bytes: None,
1094 }
1095 },
1096 ManifestEntry {
1097 status: ManifestStatus::Added,
1098 snapshot_id: None,
1099 sequence_number: None,
1100 file_sequence_number: None,
1101 data_file: DataFile {
1102 content: DataContentType::Data,
1103 file_path: "s3a://icebergdata/demo/s1/t1/data/00000-0-ba56fbfa-f2ff-40c9-bb27-565ad6dc2be8-00000.parquet".to_string(),
1104 file_format: DataFileFormat::Parquet,
1105 partition: Struct::from_iter(
1106 vec![
1107 Some(Literal::int(1111)),
1108 Some(Literal::float(15.5_f32)),
1109 Some(Literal::double(25.5)),
1110 ]
1111 ),
1112 record_count: 1,
1113 file_size_in_bytes: 5442,
1114 column_sizes: HashMap::from([(0,73),(6,34),(2,73),(7,61),(3,61),(5,62),(9,79),(10,73),(1,61),(4,73),(8,73)]),
1115 value_counts: HashMap::from([(4,1),(5,1),(2,1),(0,1),(3,1),(6,1),(8,1),(1,1),(10,1),(7,1),(9,1)]),
1116 null_value_counts: HashMap::from([(1,0),(6,0),(2,0),(8,0),(0,0),(3,0),(5,0),(9,0),(7,0),(4,0),(10,0)]),
1117 nan_value_counts: HashMap::new(),
1118 lower_bounds: HashMap::new(),
1119 upper_bounds: HashMap::new(),
1120 key_metadata: None,
1121 split_offsets: Some(vec![4]),
1122 equality_ids: None,
1123 sort_order_id: None,
1124 partition_spec_id: 0,
1125 first_row_id: None,
1126 referenced_data_file: None,
1127 content_offset: None,
1128 content_size_in_bytes: None,
1129 }
1130 },
1131 ManifestEntry {
1132 status: ManifestStatus::Added,
1133 snapshot_id: None,
1134 sequence_number: None,
1135 file_sequence_number: None,
1136 data_file: DataFile {
1137 content: DataContentType::Data,
1138 file_path: "s3a://icebergdata/demo/s1/t1/data/00000-0-ba56fbfa-f2ff-40c9-bb27-565ad6dc2be8-00000.parquet".to_string(),
1139 file_format: DataFileFormat::Parquet,
1140 partition: Struct::from_iter(
1141 vec![
1142 Some(Literal::int(1211)),
1143 Some(Literal::float(f32::NAN)),
1144 Some(Literal::double(1.0)),
1145 ]
1146 ),
1147 record_count: 1,
1148 file_size_in_bytes: 5442,
1149 column_sizes: HashMap::from([(0,73),(6,34),(2,73),(7,61),(3,61),(5,62),(9,79),(10,73),(1,61),(4,73),(8,73)]),
1150 value_counts: HashMap::from([(4,1),(5,1),(2,1),(0,1),(3,1),(6,1),(8,1),(1,1),(10,1),(7,1),(9,1)]),
1151 null_value_counts: HashMap::from([(1,0),(6,0),(2,0),(8,0),(0,0),(3,0),(5,0),(9,0),(7,0),(4,0),(10,0)]),
1152 nan_value_counts: HashMap::new(),
1153 lower_bounds: HashMap::new(),
1154 upper_bounds: HashMap::new(),
1155 key_metadata: None,
1156 split_offsets: Some(vec![4]),
1157 equality_ids: None,
1158 sort_order_id: None,
1159 partition_spec_id: 0,
1160 first_row_id: None,
1161 referenced_data_file: None,
1162 content_offset: None,
1163 content_size_in_bytes: None,
1164 }
1165 },
1166 ManifestEntry {
1167 status: ManifestStatus::Added,
1168 snapshot_id: None,
1169 sequence_number: None,
1170 file_sequence_number: None,
1171 data_file: DataFile {
1172 content: DataContentType::Data,
1173 file_path: "s3a://icebergdata/demo/s1/t1/data/00000-0-ba56fbfa-f2ff-40c9-bb27-565ad6dc2be8-00000.parquet".to_string(),
1174 file_format: DataFileFormat::Parquet,
1175 partition: Struct::from_iter(
1176 vec![
1177 Some(Literal::int(1111)),
1178 None,
1179 Some(Literal::double(11.0)),
1180 ]
1181 ),
1182 record_count: 1,
1183 file_size_in_bytes: 5442,
1184 column_sizes: HashMap::from([(0,73),(6,34),(2,73),(7,61),(3,61),(5,62),(9,79),(10,73),(1,61),(4,73),(8,73)]),
1185 value_counts: HashMap::from([(4,1),(5,1),(2,1),(0,1),(3,1),(6,1),(8,1),(1,1),(10,1),(7,1),(9,1)]),
1186 null_value_counts: HashMap::from([(1,0),(6,0),(2,0),(8,0),(0,0),(3,0),(5,0),(9,0),(7,0),(4,0),(10,0)]),
1187 nan_value_counts: HashMap::new(),
1188 lower_bounds: HashMap::new(),
1189 upper_bounds: HashMap::new(),
1190 key_metadata: None,
1191 split_offsets: Some(vec![4]),
1192 equality_ids: None,
1193 sort_order_id: None,
1194 partition_spec_id: 0,
1195 first_row_id: None,
1196 referenced_data_file: None,
1197 content_offset: None,
1198 content_size_in_bytes: None,
1199 }
1200 },
1201 ];
1202
1203 let tmp_dir = TempDir::new().unwrap();
1205 let path = tmp_dir.path().join("test_manifest.avro");
1206 let io = FileIO::new_with_fs();
1207 let output_file = io.new_output(path.to_str().unwrap()).unwrap();
1208 let mut writer = ManifestWriterBuilder::new(
1209 output_file,
1210 Some(1),
1211 metadata.schema.clone(),
1212 metadata.partition_spec.clone(),
1213 )
1214 .build_v2_data();
1215 for entry in &entries {
1216 writer.add_entry(entry.clone()).unwrap();
1217 }
1218 let res = writer.write_manifest_file().await.unwrap();
1219
1220 let partitions = res.partitions.unwrap();
1221
1222 assert_eq!(partitions.len(), 3);
1223 assert_eq!(
1224 partitions[0].clone().lower_bound.unwrap(),
1225 Datum::int(1111).to_bytes().unwrap()
1226 );
1227 assert_eq!(
1228 partitions[0].clone().upper_bound.unwrap(),
1229 Datum::int(2021).to_bytes().unwrap()
1230 );
1231 assert!(!partitions[0].clone().contains_null);
1232 assert_eq!(partitions[0].clone().contains_nan, Some(false));
1233
1234 assert_eq!(
1235 partitions[1].clone().lower_bound.unwrap(),
1236 Datum::float(1.0_f32).to_bytes().unwrap()
1237 );
1238 assert_eq!(
1239 partitions[1].clone().upper_bound.unwrap(),
1240 Datum::float(15.5_f32).to_bytes().unwrap()
1241 );
1242 assert!(partitions[1].clone().contains_null);
1243 assert_eq!(partitions[1].clone().contains_nan, Some(true));
1244
1245 assert_eq!(
1246 partitions[2].clone().lower_bound.unwrap(),
1247 Datum::double(1.0).to_bytes().unwrap()
1248 );
1249 assert_eq!(
1250 partitions[2].clone().upper_bound.unwrap(),
1251 Datum::double(25.5).to_bytes().unwrap()
1252 );
1253 assert!(!partitions[2].clone().contains_null);
1254 assert_eq!(partitions[2].clone().contains_nan, Some(false));
1255 }
1256
1257 #[test]
1258 fn test_data_file_serialization() {
1259 let schema = Schema::builder()
1261 .with_schema_id(1)
1262 .with_identifier_field_ids(vec![1])
1263 .with_fields(vec![
1264 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Long)).into(),
1265 NestedField::required(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1266 ])
1267 .build()
1268 .unwrap();
1269
1270 let partition_spec = PartitionSpec::builder(schema.clone())
1272 .with_spec_id(1)
1273 .add_partition_field("id", "id_partition", Transform::Identity)
1274 .unwrap()
1275 .build()
1276 .unwrap();
1277
1278 let partition_type = partition_spec.partition_type(&schema).unwrap();
1280
1281 let data_files = vec![
1283 DataFileBuilder::default()
1284 .content(DataContentType::Data)
1285 .file_format(DataFileFormat::Parquet)
1286 .file_path("path/to/file1.parquet".to_string())
1287 .file_size_in_bytes(1024)
1288 .record_count(100)
1289 .partition_spec_id(1)
1290 .partition(Struct::empty())
1291 .column_sizes(HashMap::from([(1, 512), (2, 1024)]))
1292 .value_counts(HashMap::from([(1, 100), (2, 500)]))
1293 .null_value_counts(HashMap::from([(1, 0), (2, 1)]))
1294 .build()
1295 .unwrap(),
1296 DataFileBuilder::default()
1297 .content(DataContentType::Data)
1298 .file_format(DataFileFormat::Parquet)
1299 .file_path("path/to/file2.parquet".to_string())
1300 .file_size_in_bytes(2048)
1301 .record_count(200)
1302 .partition_spec_id(1)
1303 .partition(Struct::empty())
1304 .column_sizes(HashMap::from([(1, 1024), (2, 2048)]))
1305 .value_counts(HashMap::from([(1, 200), (2, 600)]))
1306 .null_value_counts(HashMap::from([(1, 10), (2, 999)]))
1307 .build()
1308 .unwrap(),
1309 ];
1310
1311 let serialized_files = data_files
1313 .clone()
1314 .into_iter()
1315 .map(|f| serialize_data_file_to_json(f, &partition_type, FormatVersion::V2).unwrap())
1316 .collect::<Vec<String>>();
1317
1318 assert_eq!(serialized_files.len(), 2);
1320 let pretty_json1: Value = serde_json::from_str(serialized_files.first().unwrap()).unwrap();
1321 let pretty_json2: Value = serde_json::from_str(serialized_files.get(1).unwrap()).unwrap();
1322 let expected_serialized_file1 = serde_json::json!({
1323 "content": 0,
1324 "file_path": "path/to/file1.parquet",
1325 "file_format": "PARQUET",
1326 "partition": {},
1327 "record_count": 100,
1328 "file_size_in_bytes": 1024,
1329 "column_sizes": [
1330 { "key": 1, "value": 512 },
1331 { "key": 2, "value": 1024 }
1332 ],
1333 "value_counts": [
1334 { "key": 1, "value": 100 },
1335 { "key": 2, "value": 500 }
1336 ],
1337 "null_value_counts": [
1338 { "key": 1, "value": 0 },
1339 { "key": 2, "value": 1 }
1340 ],
1341 "nan_value_counts": [],
1342 "lower_bounds": [],
1343 "upper_bounds": [],
1344 "key_metadata": null,
1345 "split_offsets": null,
1346 "equality_ids": null,
1347 "sort_order_id": null,
1348 "first_row_id": null,
1349 "referenced_data_file": null,
1350 "content_offset": null,
1351 "content_size_in_bytes": null
1352 });
1353 let expected_serialized_file2 = serde_json::json!({
1354 "content": 0,
1355 "file_path": "path/to/file2.parquet",
1356 "file_format": "PARQUET",
1357 "partition": {},
1358 "record_count": 200,
1359 "file_size_in_bytes": 2048,
1360 "column_sizes": [
1361 { "key": 1, "value": 1024 },
1362 { "key": 2, "value": 2048 }
1363 ],
1364 "value_counts": [
1365 { "key": 1, "value": 200 },
1366 { "key": 2, "value": 600 }
1367 ],
1368 "null_value_counts": [
1369 { "key": 1, "value": 10 },
1370 { "key": 2, "value": 999 }
1371 ],
1372 "nan_value_counts": [],
1373 "lower_bounds": [],
1374 "upper_bounds": [],
1375 "key_metadata": null,
1376 "split_offsets": null,
1377 "equality_ids": null,
1378 "sort_order_id": null,
1379 "first_row_id": null,
1380 "referenced_data_file": null,
1381 "content_offset": null,
1382 "content_size_in_bytes": null
1383 });
1384 assert_eq!(pretty_json1, expected_serialized_file1);
1385 assert_eq!(pretty_json2, expected_serialized_file2);
1386
1387 let deserialized_files: Vec<DataFile> = serialized_files
1389 .into_iter()
1390 .map(|json| {
1391 deserialize_data_file_from_json(
1392 &json,
1393 partition_spec.spec_id(),
1394 &partition_type,
1395 &schema,
1396 )
1397 .unwrap()
1398 })
1399 .collect();
1400
1401 assert_eq!(deserialized_files.len(), 2);
1403 let deserialized_data_file1 = deserialized_files.first().unwrap();
1404 let deserialized_data_file2 = deserialized_files.get(1).unwrap();
1405 let original_data_file1 = data_files.first().unwrap();
1406 let original_data_file2 = data_files.get(1).unwrap();
1407
1408 assert_eq!(deserialized_data_file1, original_data_file1);
1409 assert_eq!(deserialized_data_file2, original_data_file2);
1410 }
1411
1412 fn writer_schema_test_metadata() -> ManifestMetadata {
1415 let schema = Arc::new(
1416 Schema::builder()
1417 .with_fields(vec![
1418 NestedField::optional(1, "id", Type::Primitive(PrimitiveType::Long)).into(),
1419 NestedField::optional(2, "name", Type::Primitive(PrimitiveType::String)).into(),
1420 NestedField::optional(3, "score", Type::Primitive(PrimitiveType::Double))
1421 .into(),
1422 ])
1423 .build()
1424 .unwrap(),
1425 );
1426 let partition_spec = PartitionSpec::builder(schema.clone())
1427 .with_spec_id(0)
1428 .add_partition_field("id", "id", Transform::Identity)
1429 .unwrap()
1430 .add_partition_field("name", "name", Transform::Identity)
1431 .unwrap()
1432 .add_partition_field("score", "score", Transform::Identity)
1433 .unwrap()
1434 .build()
1435 .unwrap();
1436 ManifestMetadata {
1437 schema_id: 0,
1438 schema,
1439 partition_spec,
1440 content: ManifestContentType::Data,
1441 format_version: FormatVersion::V2,
1442 }
1443 }
1444
1445 fn v2_writer_schema(partition_fields: Value) -> Value {
1448 serde_json::json!({
1449 "type": "record",
1450 "name": "manifest_entry",
1451 "fields": [
1452 {"name": "status", "type": "int", "field-id": 0},
1453 {"name": "snapshot_id", "type": ["null", "long"], "default": null, "field-id": 1},
1454 {"name": "sequence_number", "type": ["null", "long"], "default": null, "field-id": 3},
1455 {"name": "file_sequence_number", "type": ["null", "long"], "default": null, "field-id": 4},
1456 {"name": "data_file", "field-id": 2, "type": {
1457 "type": "record",
1458 "name": "r2",
1459 "fields": [
1460 {"name": "content", "type": "int", "field-id": 134},
1461 {"name": "file_path", "type": "string", "field-id": 100},
1462 {"name": "file_format", "type": "string", "field-id": 101},
1463 {"name": "partition", "field-id": 102, "type": {
1464 "type": "record",
1465 "name": "r102",
1466 "fields": partition_fields,
1467 }},
1468 {"name": "record_count", "type": "long", "field-id": 103},
1469 {"name": "file_size_in_bytes", "type": "long", "field-id": 104},
1470 ],
1471 }},
1472 ],
1473 })
1474 }
1475
1476 fn writer_schema_field<'a>(record: &'a mut Value, path: &[&str]) -> &'a mut Value {
1479 let (name, rest) = path.split_first().unwrap();
1480 let field = record["fields"]
1481 .as_array_mut()
1482 .unwrap()
1483 .iter_mut()
1484 .find(|field| field["name"] == *name)
1485 .unwrap();
1486 if rest.is_empty() {
1487 field
1488 } else {
1489 writer_schema_field(&mut field["type"], rest)
1490 }
1491 }
1492
1493 fn v2_entry(partition: Value) -> Value {
1494 serde_json::json!({
1495 "status": 1,
1496 "snapshot_id": 7,
1497 "sequence_number": null,
1498 "file_sequence_number": null,
1499 "data_file": {
1500 "content": 0,
1501 "file_path": "s3://bucket/table/data/a.parquet",
1502 "file_format": "PARQUET",
1503 "partition": partition,
1504 "record_count": 10,
1505 "file_size_in_bytes": 100,
1506 },
1507 })
1508 }
1509
1510 fn expected_entry(partition: Struct) -> ManifestEntry {
1511 ManifestEntry {
1512 status: ManifestStatus::Added,
1513 snapshot_id: Some(7),
1514 sequence_number: None,
1515 file_sequence_number: None,
1516 data_file: DataFile {
1517 content: DataContentType::Data,
1518 file_path: "s3://bucket/table/data/a.parquet".to_string(),
1519 file_format: DataFileFormat::Parquet,
1520 partition,
1521 record_count: 10,
1522 file_size_in_bytes: 100,
1523 column_sizes: HashMap::new(),
1524 value_counts: HashMap::new(),
1525 null_value_counts: HashMap::new(),
1526 nan_value_counts: HashMap::new(),
1527 lower_bounds: HashMap::new(),
1528 upper_bounds: HashMap::new(),
1529 key_metadata: None,
1530 split_offsets: None,
1531 equality_ids: None,
1532 sort_order_id: None,
1533 partition_spec_id: 0,
1534 first_row_id: None,
1535 referenced_data_file: None,
1536 content_offset: None,
1537 content_size_in_bytes: None,
1538 },
1539 }
1540 }
1541
1542 fn write_with_writer_schema(
1545 metadata: &ManifestMetadata,
1546 writer_schema: &Value,
1547 entries: Vec<Value>,
1548 ) -> Vec<u8> {
1549 write_avro_values_with_writer_schema(
1550 metadata,
1551 writer_schema,
1552 entries
1553 .into_iter()
1554 .map(|entry| AvroValue::try_from(entry).unwrap())
1555 .collect(),
1556 )
1557 }
1558
1559 fn write_avro_values_with_writer_schema(
1561 metadata: &ManifestMetadata,
1562 writer_schema: &Value,
1563 entries: Vec<AvroValue>,
1564 ) -> Vec<u8> {
1565 let avro_schema = apache_avro::Schema::parse(writer_schema).unwrap();
1566 let mut writer = Writer::new(&avro_schema, Vec::new()).unwrap();
1567 for (key, value) in [
1568 ("schema", to_vec(&metadata.schema).unwrap()),
1569 ("schema-id", metadata.schema_id.to_string().into_bytes()),
1570 (
1571 "partition-spec",
1572 to_vec(&metadata.partition_spec.fields()).unwrap(),
1573 ),
1574 (
1575 "partition-spec-id",
1576 metadata.partition_spec.spec_id().to_string().into_bytes(),
1577 ),
1578 (
1579 "format-version",
1580 (metadata.format_version as u8).to_string().into_bytes(),
1581 ),
1582 ("content", metadata.content.to_string().into_bytes()),
1583 ] {
1584 writer.add_user_metadata(key.to_string(), value).unwrap();
1585 }
1586 for entry in entries {
1587 writer
1588 .append_value(entry.resolve(&avro_schema).unwrap())
1589 .unwrap();
1590 }
1591 writer.into_inner().unwrap()
1592 }
1593
1594 fn remove_writer_field(schema: &mut Value, entry: &mut AvroValue, path: &[&str]) {
1597 let (name, rest) = path.split_first().unwrap();
1598 let AvroValue::Record(values) = entry else {
1599 unreachable!("the entry is a record");
1600 };
1601 let fields = schema["fields"].as_array_mut().unwrap();
1602 if rest.is_empty() {
1603 fields.retain(|field| field["name"] != *name);
1604 values.retain(|(value_name, _)| value_name != name);
1605 } else {
1606 let field = fields.iter_mut().find(|field| field["name"] == *name);
1607 let value = values.iter_mut().find(|(value_name, _)| value_name == name);
1608 remove_writer_field(&mut field.unwrap()["type"], &mut value.unwrap().1, rest);
1609 }
1610 }
1611
1612 type ResetField = fn(&mut ManifestEntry);
1613
1614 fn assert_reads_without_each_field(
1618 metadata: &ManifestMetadata,
1619 full_schema: &Value,
1620 full_value: &AvroValue,
1621 full_entry: &ManifestEntry,
1622 optional: &[(&[&str], ResetField)],
1623 required: &[&[&str]],
1624 ) {
1625 for (path, reset) in optional {
1626 let (mut schema, mut value) = (full_schema.clone(), full_value.clone());
1627 remove_writer_field(&mut schema, &mut value, path);
1628 let bs = write_avro_values_with_writer_schema(metadata, &schema, vec![value]);
1629
1630 let manifest = Manifest::parse_avro(&bs).unwrap();
1631
1632 let mut expected = full_entry.clone();
1633 reset(&mut expected);
1634 assert_eq!(
1635 manifest,
1636 Manifest::new(metadata.clone(), vec![expected]),
1637 "{path:?}"
1638 );
1639 }
1640
1641 for path in required {
1642 let (mut schema, mut value) = (full_schema.clone(), full_value.clone());
1643 remove_writer_field(&mut schema, &mut value, path);
1644 let bs = write_avro_values_with_writer_schema(metadata, &schema, vec![value]);
1645
1646 let err = Manifest::parse_avro(&bs).unwrap_err();
1647
1648 assert!(
1649 err.to_string().contains(path.last().unwrap()),
1650 "{path:?}: {err}"
1651 );
1652 }
1653 }
1654
1655 #[test]
1656 fn test_parse_manifest_without_each_field() {
1657 let metadata = writer_schema_test_metadata();
1660 let partition_type = metadata
1661 .partition_spec
1662 .partition_type(&metadata.schema)
1663 .unwrap();
1664 let full_schema =
1665 serde_json::to_value(manifest_schema_v2(&partition_type).unwrap()).unwrap();
1666 let mut full_entry = expected_entry(Struct::from_iter([
1667 Some(Literal::long(5)),
1668 Some(Literal::string("a")),
1669 Some(Literal::double(2.5)),
1670 ]));
1671 full_entry.sequence_number = Some(3);
1672 full_entry.file_sequence_number = Some(4);
1673 let data_file = &mut full_entry.data_file;
1674 data_file.content = DataContentType::PositionDeletes;
1675 data_file.column_sizes = HashMap::from([(1, 40)]);
1676 data_file.value_counts = HashMap::from([(1, 10)]);
1677 data_file.null_value_counts = HashMap::from([(1, 1)]);
1678 data_file.nan_value_counts = HashMap::from([(3, 2)]);
1679 data_file.lower_bounds = HashMap::from([(1, Datum::long(1))]);
1680 data_file.upper_bounds = HashMap::from([(1, Datum::long(9))]);
1681 data_file.key_metadata = Some(vec![1, 2]);
1682 data_file.split_offsets = Some(vec![4]);
1683 data_file.equality_ids = Some(vec![1]);
1684 data_file.sort_order_id = Some(0);
1685 data_file.first_row_id = Some(100);
1686 data_file.referenced_data_file = Some("s3://bucket/table/data/b.parquet".to_string());
1687 data_file.content_offset = Some(4);
1688 data_file.content_size_in_bytes = Some(8);
1689 let full_value = to_value(
1690 _serde::ManifestEntryV2::try_from(
1691 full_entry.clone(),
1692 &Type::Struct(partition_type.clone()),
1693 )
1694 .unwrap(),
1695 )
1696 .unwrap();
1697
1698 let optional: [(&[&str], ResetField); 18] = [
1699 (&["snapshot_id"], |e| e.snapshot_id = None),
1700 (&["sequence_number"], |e| e.sequence_number = None),
1701 (&["file_sequence_number"], |e| e.file_sequence_number = None),
1702 (&["data_file", "content"], |e| {
1703 e.data_file.content = DataContentType::Data
1704 }),
1705 (&["data_file", "column_sizes"], |e| {
1706 e.data_file.column_sizes.clear()
1707 }),
1708 (&["data_file", "value_counts"], |e| {
1709 e.data_file.value_counts.clear()
1710 }),
1711 (&["data_file", "null_value_counts"], |e| {
1712 e.data_file.null_value_counts.clear()
1713 }),
1714 (&["data_file", "nan_value_counts"], |e| {
1715 e.data_file.nan_value_counts.clear()
1716 }),
1717 (&["data_file", "lower_bounds"], |e| {
1718 e.data_file.lower_bounds.clear()
1719 }),
1720 (&["data_file", "upper_bounds"], |e| {
1721 e.data_file.upper_bounds.clear()
1722 }),
1723 (&["data_file", "key_metadata"], |e| {
1724 e.data_file.key_metadata = None
1725 }),
1726 (&["data_file", "split_offsets"], |e| {
1727 e.data_file.split_offsets = None
1728 }),
1729 (&["data_file", "equality_ids"], |e| {
1730 e.data_file.equality_ids = None
1731 }),
1732 (&["data_file", "sort_order_id"], |e| {
1733 e.data_file.sort_order_id = None
1734 }),
1735 (&["data_file", "first_row_id"], |e| {
1736 e.data_file.first_row_id = None
1737 }),
1738 (&["data_file", "referenced_data_file"], |e| {
1739 e.data_file.referenced_data_file = None
1740 }),
1741 (&["data_file", "content_offset"], |e| {
1742 e.data_file.content_offset = None
1743 }),
1744 (&["data_file", "content_size_in_bytes"], |e| {
1745 e.data_file.content_size_in_bytes = None
1746 }),
1747 ];
1748 let required: [&[&str]; 7] = [
1749 &["status"],
1750 &["data_file"],
1751 &["data_file", "file_path"],
1752 &["data_file", "file_format"],
1753 &["data_file", "partition"],
1754 &["data_file", "record_count"],
1755 &["data_file", "file_size_in_bytes"],
1756 ];
1757 assert_reads_without_each_field(
1758 &metadata,
1759 &full_schema,
1760 &full_value,
1761 &full_entry,
1762 &optional,
1763 &required,
1764 );
1765 }
1766
1767 #[test]
1768 fn test_parse_v1_manifest_without_each_field() {
1769 let mut metadata = writer_schema_test_metadata();
1770 metadata.format_version = FormatVersion::V1;
1771 let partition_type = metadata
1772 .partition_spec
1773 .partition_type(&metadata.schema)
1774 .unwrap();
1775 let full_schema =
1776 serde_json::to_value(manifest_schema_v1(&partition_type).unwrap()).unwrap();
1777 let mut full_entry = expected_entry(Struct::from_iter([
1778 Some(Literal::long(5)),
1779 Some(Literal::string("a")),
1780 Some(Literal::double(2.5)),
1781 ]));
1782 full_entry.sequence_number = Some(0);
1784 full_entry.file_sequence_number = Some(0);
1785 let data_file = &mut full_entry.data_file;
1786 data_file.column_sizes = HashMap::from([(1, 40)]);
1787 data_file.value_counts = HashMap::from([(1, 10)]);
1788 data_file.null_value_counts = HashMap::from([(1, 1)]);
1789 data_file.nan_value_counts = HashMap::from([(3, 2)]);
1790 data_file.lower_bounds = HashMap::from([(1, Datum::long(1))]);
1791 data_file.upper_bounds = HashMap::from([(1, Datum::long(9))]);
1792 data_file.key_metadata = Some(vec![1, 2]);
1793 data_file.split_offsets = Some(vec![4]);
1794 data_file.sort_order_id = Some(0);
1795 let full_value = to_value(
1796 _serde::ManifestEntryV1::try_from(
1797 full_entry.clone(),
1798 &Type::Struct(partition_type.clone()),
1799 )
1800 .unwrap(),
1801 )
1802 .unwrap();
1803
1804 let optional: [(&[&str], ResetField); 10] = [
1805 (&["data_file", "block_size_in_bytes"], |_| {}),
1807 (&["data_file", "column_sizes"], |e| {
1808 e.data_file.column_sizes.clear()
1809 }),
1810 (&["data_file", "value_counts"], |e| {
1811 e.data_file.value_counts.clear()
1812 }),
1813 (&["data_file", "null_value_counts"], |e| {
1814 e.data_file.null_value_counts.clear()
1815 }),
1816 (&["data_file", "nan_value_counts"], |e| {
1817 e.data_file.nan_value_counts.clear()
1818 }),
1819 (&["data_file", "lower_bounds"], |e| {
1820 e.data_file.lower_bounds.clear()
1821 }),
1822 (&["data_file", "upper_bounds"], |e| {
1823 e.data_file.upper_bounds.clear()
1824 }),
1825 (&["data_file", "key_metadata"], |e| {
1826 e.data_file.key_metadata = None
1827 }),
1828 (&["data_file", "split_offsets"], |e| {
1829 e.data_file.split_offsets = None
1830 }),
1831 (&["data_file", "sort_order_id"], |e| {
1832 e.data_file.sort_order_id = None
1833 }),
1834 ];
1835 let required: [&[&str]; 8] = [
1836 &["status"],
1837 &["snapshot_id"],
1838 &["data_file"],
1839 &["data_file", "file_path"],
1840 &["data_file", "file_format"],
1841 &["data_file", "partition"],
1842 &["data_file", "record_count"],
1843 &["data_file", "file_size_in_bytes"],
1844 ];
1845 assert_reads_without_each_field(
1846 &metadata,
1847 &full_schema,
1848 &full_value,
1849 &full_entry,
1850 &optional,
1851 &required,
1852 );
1853 }
1854
1855 #[test]
1856 fn test_parse_manifest_matches_partition_fields_by_name() {
1857 let metadata = writer_schema_test_metadata();
1860 let writer_schema = v2_writer_schema(serde_json::json!([
1861 {"name": "score", "type": ["null", "double"], "default": null, "field-id": 1002},
1862 {"name": "extra", "type": ["null", "int"], "default": null, "field-id": 1003},
1863 {"name": "id", "type": ["null", "long"], "default": null, "field-id": 1000},
1864 ]));
1865 let bs = write_with_writer_schema(&metadata, &writer_schema, vec![v2_entry(
1866 serde_json::json!({"score": 2.5, "extra": 9, "id": 5}),
1867 )]);
1868
1869 let manifest = Manifest::parse_avro(&bs).unwrap();
1870
1871 let partition =
1872 Struct::from_iter([Some(Literal::long(5)), None, Some(Literal::double(2.5))]);
1873 assert_eq!(
1874 manifest,
1875 Manifest::new(metadata, vec![expected_entry(partition)])
1876 );
1877 }
1878
1879 #[test]
1880 fn test_parse_manifest_promotes_writer_types() {
1881 let metadata = writer_schema_test_metadata();
1883 let mut writer_schema = v2_writer_schema(serde_json::json!([
1884 {"name": "id", "type": ["null", "int"], "default": null, "field-id": 1000},
1885 {"name": "name", "type": ["null", "string"], "default": null, "field-id": 1001},
1886 {"name": "score", "type": ["null", "float"], "default": null, "field-id": 1002},
1887 ]));
1888 for field in ["record_count", "file_size_in_bytes"] {
1889 writer_schema_field(&mut writer_schema, &["data_file", field])["type"] =
1890 serde_json::json!("int");
1891 }
1892 let bs = write_with_writer_schema(&metadata, &writer_schema, vec![v2_entry(
1893 serde_json::json!({"id": 5, "name": "a", "score": 2.5}),
1894 )]);
1895
1896 let manifest = Manifest::parse_avro(&bs).unwrap();
1897
1898 let partition = Struct::from_iter([
1899 Some(Literal::long(5)),
1900 Some(Literal::string("a")),
1901 Some(Literal::double(2.5)),
1902 ]);
1903 assert_eq!(
1904 manifest,
1905 Manifest::new(metadata, vec![expected_entry(partition)])
1906 );
1907 }
1908
1909 #[test]
1910 fn test_parse_manifest_reads_required_writer_fields_as_optional() {
1911 let metadata = writer_schema_test_metadata();
1914 let mut writer_schema = v2_writer_schema(serde_json::json!([
1915 {"name": "id", "type": "long", "field-id": 1000},
1916 {"name": "name", "type": "string", "field-id": 1001},
1917 {"name": "score", "type": "double", "field-id": 1002},
1918 ]));
1919 *writer_schema_field(&mut writer_schema, &["snapshot_id"]) =
1920 serde_json::json!({"name": "snapshot_id", "type": "long", "field-id": 1});
1921 writer_schema_field(&mut writer_schema, &["data_file"])["type"]["fields"]
1922 .as_array_mut()
1923 .unwrap()
1924 .push(serde_json::json!({"name": "sort_order_id", "type": "int", "field-id": 140}));
1925 let mut entry = v2_entry(serde_json::json!({"id": 5, "name": "a", "score": 2.5}));
1926 entry["data_file"]["sort_order_id"] = serde_json::json!(3);
1927 let bs = write_with_writer_schema(&metadata, &writer_schema, vec![entry]);
1928
1929 let manifest = Manifest::parse_avro(&bs).unwrap();
1930
1931 let mut expected = expected_entry(Struct::from_iter([
1932 Some(Literal::long(5)),
1933 Some(Literal::string("a")),
1934 Some(Literal::double(2.5)),
1935 ]));
1936 expected.data_file.sort_order_id = Some(3);
1937 assert_eq!(manifest, Manifest::new(metadata, vec![expected]));
1938 }
1939
1940 #[test]
1941 fn test_parse_manifest_reads_union_writer_fields_as_required() {
1942 let metadata = writer_schema_test_metadata();
1945 let mut writer_schema = v2_writer_schema(serde_json::json!([
1946 {"name": "id", "type": ["null", "long"], "default": null, "field-id": 1000},
1947 {"name": "name", "type": ["null", "string"], "default": null, "field-id": 1001},
1948 {"name": "score", "type": ["null", "double"], "default": null, "field-id": 1002},
1949 ]));
1950 writer_schema_field(&mut writer_schema, &["data_file", "record_count"])["type"] =
1951 serde_json::json!(["null", "long"]);
1952 let partition = serde_json::json!({"id": 5, "name": "a", "score": 2.5});
1953 let bs =
1954 write_with_writer_schema(&metadata, &writer_schema, vec![v2_entry(partition.clone())]);
1955
1956 let manifest = Manifest::parse_avro(&bs).unwrap();
1957
1958 let expected = expected_entry(Struct::from_iter([
1959 Some(Literal::long(5)),
1960 Some(Literal::string("a")),
1961 Some(Literal::double(2.5)),
1962 ]));
1963 assert_eq!(manifest, Manifest::new(metadata.clone(), vec![expected]));
1964
1965 let mut entry = v2_entry(partition);
1966 entry["data_file"]["record_count"] = Value::Null;
1967 let bs = write_with_writer_schema(&metadata, &writer_schema, vec![entry]);
1968
1969 let err = Manifest::parse_avro(&bs).unwrap_err();
1970
1971 assert_eq!(err.kind(), ErrorKind::DataInvalid);
1972 }
1973
1974 #[test]
1975 fn test_parse_manifest_with_fixed_uuid_partition() {
1976 let schema = Arc::new(
1978 Schema::builder()
1979 .with_fields(vec![
1980 NestedField::optional(1, "u", Type::Primitive(PrimitiveType::Uuid)).into(),
1981 ])
1982 .build()
1983 .unwrap(),
1984 );
1985 let partition_spec = PartitionSpec::builder(schema.clone())
1986 .with_spec_id(0)
1987 .add_partition_field("u", "u", Transform::Identity)
1988 .unwrap()
1989 .build()
1990 .unwrap();
1991 let metadata = ManifestMetadata {
1992 schema_id: 0,
1993 schema,
1994 partition_spec,
1995 content: ManifestContentType::Data,
1996 format_version: FormatVersion::V2,
1997 };
1998 let writer_schema = v2_writer_schema(serde_json::json!([
1999 {"name": "u", "default": null, "field-id": 1000, "type": ["null", {
2000 "type": "fixed", "name": "uuid_fixed", "size": 16, "logicalType": "uuid",
2001 }]},
2002 ]));
2003 let uuid = uuid::Uuid::from_u128(0xf79c3e09_677c_4bbd_a479_3f349cb785e7);
2004 let mut entry = AvroValue::try_from(v2_entry(serde_json::json!({}))).unwrap();
2007 let AvroValue::Map(fields) = &mut entry else {
2008 unreachable!("a JSON object converts to an Avro map");
2009 };
2010 let Some(AvroValue::Map(data_file)) = fields.get_mut("data_file") else {
2011 unreachable!("a JSON object converts to an Avro map");
2012 };
2013 data_file.insert(
2014 "partition".to_string(),
2015 AvroValue::Map(HashMap::from([("u".to_string(), AvroValue::Uuid(uuid))])),
2016 );
2017 let bs = write_avro_values_with_writer_schema(&metadata, &writer_schema, vec![entry]);
2018
2019 let manifest = Manifest::parse_avro(&bs).unwrap();
2020
2021 let expected = expected_entry(Struct::from_iter([Some(Literal::uuid(uuid))]));
2022 assert_eq!(manifest, Manifest::new(metadata, vec![expected]));
2023 }
2024
2025 #[test]
2026 fn test_parse_manifest_ignores_record_names_and_unknown_fields() {
2027 let metadata = writer_schema_test_metadata();
2030 let mut writer_schema = v2_writer_schema(serde_json::json!([
2031 {"name": "id", "type": ["null", "long"], "default": null, "field-id": 1000},
2032 {"name": "name", "type": ["null", "string"], "default": null, "field-id": 1001},
2033 {"name": "score", "type": ["null", "double"], "default": null, "field-id": 1002},
2034 ]));
2035 writer_schema["name"] = serde_json::json!("entry");
2036 writer_schema["fields"]
2037 .as_array_mut()
2038 .unwrap()
2039 .push(serde_json::json!({"name": "unknown", "type": "string"}));
2040 let data_file = &mut writer_schema_field(&mut writer_schema, &["data_file"])["type"];
2041 data_file["name"] = serde_json::json!("file");
2042 data_file["fields"][3]["type"]["name"] = serde_json::json!("part");
2043 let data_file_fields = data_file["fields"].as_array_mut().unwrap();
2044 data_file_fields.push(serde_json::json!({
2045 "name": "column_sizes",
2046 "field-id": 108,
2047 "default": null,
2048 "type": ["null", {
2049 "type": "array",
2050 "logicalType": "map",
2051 "items": {
2052 "type": "record",
2053 "name": "sizes",
2054 "fields": [
2055 {"name": "key", "type": "int", "field-id": 117},
2056 {"name": "value", "type": "long", "field-id": 118},
2057 ],
2058 },
2059 }],
2060 }));
2061 data_file_fields.push(
2062 serde_json::json!({"name": "block_size_in_bytes", "type": "long", "field-id": 105}),
2063 );
2064 let mut entry = v2_entry(serde_json::json!({"id": 5, "name": "a", "score": 2.5}));
2065 entry["unknown"] = serde_json::json!("ignored");
2066 entry["data_file"]["column_sizes"] = serde_json::json!([{"key": 1, "value": 40}]);
2067 entry["data_file"]["block_size_in_bytes"] = serde_json::json!(64);
2068 let bs = write_with_writer_schema(&metadata, &writer_schema, vec![entry]);
2069
2070 let manifest = Manifest::parse_avro(&bs).unwrap();
2071
2072 let mut expected = expected_entry(Struct::from_iter([
2073 Some(Literal::long(5)),
2074 Some(Literal::string("a")),
2075 Some(Literal::double(2.5)),
2076 ]));
2077 expected.data_file.column_sizes = HashMap::from([(1, 40)]);
2078 assert_eq!(manifest, Manifest::new(metadata, vec![expected]));
2079 }
2080
2081 #[test]
2082 fn test_parse_manifest_equality_ids_written_as_long() {
2083 let metadata = writer_schema_test_metadata();
2086 let mut writer_schema = v2_writer_schema(serde_json::json!([
2087 {"name": "id", "type": ["null", "long"], "default": null, "field-id": 1000},
2088 {"name": "name", "type": ["null", "string"], "default": null, "field-id": 1001},
2089 {"name": "score", "type": ["null", "double"], "default": null, "field-id": 1002},
2090 ]));
2091 writer_schema_field(&mut writer_schema, &["data_file"])["type"]["fields"]
2092 .as_array_mut()
2093 .unwrap()
2094 .push(serde_json::json!({
2095 "name": "equality_ids",
2096 "field-id": 135,
2097 "default": null,
2098 "type": ["null", {"type": "array", "items": "long", "element-id": 136}],
2099 }));
2100 let mut entry = v2_entry(serde_json::json!({"id": 5, "name": "a", "score": 2.5}));
2101 entry["data_file"]["equality_ids"] = serde_json::json!([1, 2]);
2102 let bs = write_with_writer_schema(&metadata, &writer_schema, vec![entry]);
2103
2104 let manifest = Manifest::parse_avro(&bs).unwrap();
2105
2106 let mut expected = expected_entry(Struct::from_iter([
2107 Some(Literal::long(5)),
2108 Some(Literal::string("a")),
2109 Some(Literal::double(2.5)),
2110 ]));
2111 expected.data_file.equality_ids = Some(vec![1, 2]);
2112 assert_eq!(manifest, Manifest::new(metadata, vec![expected]));
2113 }
2114
2115 #[test]
2116 fn test_parse_manifest_rejects_equality_id_above_int_range() {
2117 let metadata = writer_schema_test_metadata();
2118 let mut writer_schema = v2_writer_schema(serde_json::json!([
2119 {"name": "id", "type": ["null", "long"], "default": null, "field-id": 1000},
2120 {"name": "name", "type": ["null", "string"], "default": null, "field-id": 1001},
2121 {"name": "score", "type": ["null", "double"], "default": null, "field-id": 1002},
2122 ]));
2123 writer_schema_field(&mut writer_schema, &["data_file"])["type"]["fields"]
2124 .as_array_mut()
2125 .unwrap()
2126 .push(serde_json::json!({
2127 "name": "equality_ids",
2128 "field-id": 135,
2129 "default": null,
2130 "type": ["null", {"type": "array", "items": "long", "element-id": 136}],
2131 }));
2132 let mut entry = v2_entry(serde_json::json!({"id": 5, "name": "a", "score": 2.5}));
2133 entry["data_file"]["equality_ids"] = serde_json::json!([i64::from(i32::MAX) + 1]);
2134 let bs = write_with_writer_schema(&metadata, &writer_schema, vec![entry]);
2135
2136 let err = Manifest::parse_avro(&bs).unwrap_err();
2137
2138 assert_eq!(err.kind(), ErrorKind::DataInvalid);
2139 }
2140
2141 #[test]
2142 fn test_parse_manifest_written_by_pyiceberg() {
2143 let bs = fs::read(format!(
2144 "{}/testdata/manifests/pyiceberg-v2-data.avro",
2145 env!("CARGO_MANIFEST_DIR")
2146 ))
2147 .unwrap();
2148
2149 let manifest = Manifest::parse_avro(&bs).unwrap();
2150
2151 let schema = Arc::new(
2152 Schema::builder()
2153 .with_fields(vec![
2154 NestedField::optional(1, "id", Type::Primitive(PrimitiveType::Long)).into(),
2155 NestedField::optional(2, "category", Type::Primitive(PrimitiveType::String))
2156 .into(),
2157 NestedField::optional(3, "score", Type::Primitive(PrimitiveType::Double))
2158 .into(),
2159 ])
2160 .build()
2161 .unwrap(),
2162 );
2163 let partition_spec = PartitionSpec::builder(schema.clone())
2164 .with_spec_id(0)
2165 .add_partition_field("category", "category", Transform::Identity)
2166 .unwrap()
2167 .build()
2168 .unwrap();
2169 let metadata = ManifestMetadata {
2170 schema_id: 0,
2171 schema,
2172 partition_spec,
2173 content: ManifestContentType::Data,
2174 format_version: FormatVersion::V2,
2175 };
2176 let entries = (0..2)
2177 .map(|i: i64| {
2178 let category = format!("c{i}");
2179 ManifestEntry {
2180 status: ManifestStatus::Added,
2181 snapshot_id: Some(7),
2182 sequence_number: Some(1),
2183 file_sequence_number: Some(1),
2184 data_file: DataFile {
2185 content: DataContentType::Data,
2186 file_path: format!(
2187 "s3://bucket/table/data/category={category}/0000{i}.parquet"
2188 ),
2189 file_format: DataFileFormat::Parquet,
2190 partition: Struct::from_iter([Some(Literal::string(&category))]),
2191 record_count: 100 + i as u64,
2192 file_size_in_bytes: 1000 + i as u64,
2193 column_sizes: HashMap::from([
2194 (1, 10 + i as u64),
2195 (2, 20 + i as u64),
2196 (3, 30 + i as u64),
2197 ]),
2198 value_counts: HashMap::from([
2199 (1, 100 + i as u64),
2200 (2, 100 + i as u64),
2201 (3, 100 + i as u64),
2202 ]),
2203 null_value_counts: HashMap::from([(1, 0), (2, i as u64), (3, 1)]),
2204 nan_value_counts: HashMap::from([(3, i as u64)]),
2205 lower_bounds: HashMap::from([
2206 (1, Datum::long(i)),
2207 (2, Datum::string(&category)),
2208 ]),
2209 upper_bounds: HashMap::from([
2210 (1, Datum::long(i + 50)),
2211 (2, Datum::string(&category)),
2212 ]),
2213 key_metadata: None,
2214 split_offsets: Some(vec![4]),
2215 equality_ids: None,
2216 sort_order_id: Some(0),
2217 partition_spec_id: 0,
2218 first_row_id: None,
2219 referenced_data_file: None,
2220 content_offset: None,
2221 content_size_in_bytes: None,
2222 },
2223 }
2224 })
2225 .collect();
2226 assert_eq!(manifest, Manifest::new(metadata, entries));
2227 }
2228
2229 #[test]
2230 fn test_parse_manifest_with_repeated_named_type_definitions() {
2231 let bs = fs::read(format!(
2232 "{}/testdata/manifests/repeated-decimal-type-definitions.avro",
2233 env!("CARGO_MANIFEST_DIR")
2234 ))
2235 .unwrap();
2236
2237 let manifest = Manifest::parse_avro(&bs).unwrap();
2238
2239 assert_eq!(
2240 *manifest.entries()[0].data_file().partition(),
2241 Struct::from_iter([Some(Literal::decimal(12345)), Some(Literal::decimal(-678))])
2242 );
2243 }
2244
2245 #[test]
2246 fn test_parse_manifest_with_repeated_named_type_definitions_reports_original_error() {
2247 let bs = fs::read(format!(
2248 "{}/testdata/manifests/repeated-decimal-type-definitions.avro",
2249 env!("CARGO_MANIFEST_DIR")
2250 ))
2251 .unwrap();
2252 let marker = &bs[bs.len() - 16..];
2255 let header_marker = bs.windows(16).position(|w| w == marker).unwrap();
2256 let truncated = &bs[..header_marker + 8];
2257 let log_dir = TempDir::new().unwrap();
2258 let log_path = log_dir.path().join("log");
2259 let subscriber = tracing_subscriber::fmt()
2260 .with_max_level(Level::DEBUG)
2261 .with_writer(Arc::new(fs::File::create(&log_path).unwrap()))
2262 .finish();
2263
2264 let err = with_default(subscriber, || Manifest::parse_avro(truncated)).unwrap_err();
2265
2266 assert_eq!(err.kind(), ErrorKind::DataInvalid);
2267 let message = err.to_string();
2268 assert!(
2269 message.contains("Two named schema defined for same fullname"),
2270 "{message}"
2271 );
2272 assert!(message.contains("Failed to read marker bytes"), "{message}");
2273 let log = fs::read_to_string(&log_path).unwrap();
2275 assert!(!log.contains("more than once"), "{log}");
2276 }
2277}