laskoviymishka commented on code in PR #3354:
URL: https://github.com/apache/iceberg-rust/pull/3354#discussion_r4224612703
##########
crates/iceberg/src/spec/manifest/mod.rs:
##########
@@ -45,9 +52,45 @@ pub struct Manifest {
}
impl Manifest {
- /// Parse manifest metadata and entries from bytes of avro file.
- pub(crate) fn try_from_avro_bytes(bs: &[u8]) -> Result<(ManifestMetadata,
Vec<ManifestEntry>)> {
- let reader = AvroReader::new(bs)?;
+ /// Parse manifest metadata and entries from bytes of avro file. `location`
+ /// names the manifest in warnings.
+ pub(crate) fn try_from_avro_bytes(
+ bs: &[u8],
+ location: Option<&str>,
+ ) -> Result<(ManifestMetadata, Vec<ManifestEntry>)> {
+ let rewritten;
+ let reader = match AvroReader::new(bs) {
+ Ok(reader) => reader,
+ // iceberg-rust repeated `decimal` definitions before
+ // `schema_to_avro_schema` defined each named type once, so this
+ // fallback stays while tables can contain manifests it wrote.
+ Err(e) if matches!(e.details(),
Details::AmbiguousSchemaDefinition(_)) => {
+ let Ok(Some((bs, repeated))) = define_named_types_once(bs)
else {
+ return Err(e.into());
+ };
+ let location = location.unwrap_or("<unknown location>");
+ if WARNED_REPEATED_DEFINITIONS.swap(true, Ordering::Relaxed) {
Review Comment:
The error-return is right now — the original `AmbiguousSchemaDefinition`
survives with the retry error as context. The one piece that didn't move: the
`swap(true)` and `warn!` here still fire before the retry `AvroReader::new` at
line 85, so a manifest whose rewrite also fails logs "Reading it with each
repeated definition replaced…" (and spends the warn-once slot) even though the
read then errors out. I'd move the swap/log into the success path after the
retry parse returns Ok.
While we're in here: the `let Ok(Some(..)) = define_named_types_once(bs)
else` just above swallows an `Err` from the rewriter — a header we couldn't
parse looks identical to "not a repeated-definition case". Returning the
original error is still correct, but I'd attach the rewriter err as debug
context so a failed fallback leaves a trace.
##########
crates/iceberg/src/spec/manifest/mod.rs:
##########
@@ -1361,4 +1398,861 @@ mod tests {
assert_eq!(deserialized_data_file1, original_data_file1);
assert_eq!(deserialized_data_file2, original_data_file2);
}
+
+ /// Metadata for the writer schema tests: a V2 data manifest whose spec has
+ /// identity partitions on a long, a string, and a double column.
+ fn writer_schema_test_metadata() -> ManifestMetadata {
+ let schema = Arc::new(
+ Schema::builder()
+ .with_fields(vec![
+ NestedField::optional(1, "id",
Type::Primitive(PrimitiveType::Long)).into(),
+ NestedField::optional(2, "name",
Type::Primitive(PrimitiveType::String)).into(),
+ NestedField::optional(3, "score",
Type::Primitive(PrimitiveType::Double))
+ .into(),
+ ])
+ .build()
+ .unwrap(),
+ );
+ let partition_spec = PartitionSpec::builder(schema.clone())
+ .with_spec_id(0)
+ .add_partition_field("id", "id", Transform::Identity)
+ .unwrap()
+ .add_partition_field("name", "name", Transform::Identity)
+ .unwrap()
+ .add_partition_field("score", "score", Transform::Identity)
+ .unwrap()
+ .build()
+ .unwrap();
+ ManifestMetadata {
+ schema_id: 0,
+ schema,
+ partition_spec,
+ content: ManifestContentType::Data,
+ format_version: FormatVersion::V2,
+ }
+ }
+
+ /// A V2 manifest entry writer schema with Java's record names and only the
+ /// required `data_file` fields.
+ fn v2_writer_schema(partition_fields: Value) -> Value {
+ serde_json::json!({
+ "type": "record",
+ "name": "manifest_entry",
+ "fields": [
+ {"name": "status", "type": "int", "field-id": 0},
+ {"name": "snapshot_id", "type": ["null", "long"], "default":
null, "field-id": 1},
+ {"name": "sequence_number", "type": ["null", "long"],
"default": null, "field-id": 3},
+ {"name": "file_sequence_number", "type": ["null", "long"],
"default": null, "field-id": 4},
+ {"name": "data_file", "field-id": 2, "type": {
+ "type": "record",
+ "name": "r2",
+ "fields": [
+ {"name": "content", "type": "int", "field-id": 134},
+ {"name": "file_path", "type": "string", "field-id":
100},
+ {"name": "file_format", "type": "string", "field-id":
101},
+ {"name": "partition", "field-id": 102, "type": {
+ "type": "record",
+ "name": "r102",
+ "fields": partition_fields,
+ }},
+ {"name": "record_count", "type": "long", "field-id":
103},
+ {"name": "file_size_in_bytes", "type": "long",
"field-id": 104},
+ ],
+ }},
+ ],
+ })
+ }
+
+ /// The field named by `path` in a record schema, descending through nested
+ /// record types.
+ fn writer_schema_field<'a>(record: &'a mut Value, path: &[&str]) -> &'a
mut Value {
+ let (name, rest) = path.split_first().unwrap();
+ let field = record["fields"]
+ .as_array_mut()
+ .unwrap()
+ .iter_mut()
+ .find(|field| field["name"] == *name)
+ .unwrap();
+ if rest.is_empty() {
+ field
+ } else {
+ writer_schema_field(&mut field["type"], rest)
+ }
+ }
+
+ fn v2_entry(partition: Value) -> Value {
+ serde_json::json!({
+ "status": 1,
+ "snapshot_id": 7,
+ "sequence_number": null,
+ "file_sequence_number": null,
+ "data_file": {
+ "content": 0,
+ "file_path": "s3://bucket/table/data/a.parquet",
+ "file_format": "PARQUET",
+ "partition": partition,
+ "record_count": 10,
+ "file_size_in_bytes": 100,
+ },
+ })
+ }
+
+ fn expected_entry(partition: Struct) -> ManifestEntry {
+ ManifestEntry {
+ status: ManifestStatus::Added,
+ snapshot_id: Some(7),
+ sequence_number: None,
+ file_sequence_number: None,
+ data_file: DataFile {
+ content: DataContentType::Data,
+ file_path: "s3://bucket/table/data/a.parquet".to_string(),
+ file_format: DataFileFormat::Parquet,
+ partition,
+ record_count: 10,
+ file_size_in_bytes: 100,
+ column_sizes: HashMap::new(),
+ value_counts: HashMap::new(),
+ null_value_counts: HashMap::new(),
+ nan_value_counts: HashMap::new(),
+ lower_bounds: HashMap::new(),
+ upper_bounds: HashMap::new(),
+ key_metadata: None,
+ split_offsets: None,
+ equality_ids: None,
+ sort_order_id: None,
+ partition_spec_id: 0,
+ first_row_id: None,
+ referenced_data_file: None,
+ content_offset: None,
+ content_size_in_bytes: None,
+ },
+ }
+ }
+
+ /// Writes `entries` as a manifest with the given writer schema, encoding
each
+ /// JSON value with the writer schema's types.
+ fn write_with_writer_schema(
+ metadata: &ManifestMetadata,
+ writer_schema: &Value,
+ entries: Vec<Value>,
+ ) -> Vec<u8> {
+ write_avro_values_with_writer_schema(
+ metadata,
+ writer_schema,
+ entries
+ .into_iter()
+ .map(|entry| AvroValue::try_from(entry).unwrap())
+ .collect(),
+ )
+ }
+
+ /// Like [`write_with_writer_schema`], for entries that JSON can't express.
+ fn write_avro_values_with_writer_schema(
+ metadata: &ManifestMetadata,
+ writer_schema: &Value,
+ entries: Vec<AvroValue>,
+ ) -> Vec<u8> {
+ let avro_schema = apache_avro::Schema::parse(writer_schema).unwrap();
+ let mut writer = Writer::new(&avro_schema, Vec::new()).unwrap();
+ for (key, value) in [
+ ("schema", to_vec(&metadata.schema).unwrap()),
+ ("schema-id", metadata.schema_id.to_string().into_bytes()),
+ (
+ "partition-spec",
+ to_vec(&metadata.partition_spec.fields()).unwrap(),
+ ),
+ (
+ "partition-spec-id",
+ metadata.partition_spec.spec_id().to_string().into_bytes(),
+ ),
+ (
+ "format-version",
+ (metadata.format_version as u8).to_string().into_bytes(),
+ ),
+ ("content", metadata.content.to_string().into_bytes()),
+ ] {
+ writer.add_user_metadata(key.to_string(), value).unwrap();
+ }
+ for entry in entries {
+ writer
+ .append_value(entry.resolve(&avro_schema).unwrap())
+ .unwrap();
+ }
+ writer.into_inner().unwrap()
+ }
+
+ /// Removes the field named by `path` from a record schema and from a
record
+ /// value of that schema.
+ fn remove_writer_field(schema: &mut Value, entry: &mut AvroValue, path:
&[&str]) {
+ let (name, rest) = path.split_first().unwrap();
+ let AvroValue::Record(values) = entry else {
+ unreachable!("the entry is a record");
+ };
+ let fields = schema["fields"].as_array_mut().unwrap();
+ if rest.is_empty() {
+ fields.retain(|field| field["name"] != *name);
+ values.retain(|(value_name, _)| value_name != name);
+ } else {
+ let field = fields.iter_mut().find(|field| field["name"] == *name);
+ let value = values.iter_mut().find(|(value_name, _)| value_name ==
name);
+ remove_writer_field(&mut field.unwrap()["type"], &mut
value.unwrap().1, rest);
+ }
+ }
+
+ type ResetField = fn(&mut ManifestEntry);
+
+ /// Removes each field in `optional` and `required` in turn from a manifest
+ /// that `full_schema` and `full_value` write as `full_entry`. Checks that
the
+ /// manifest reads with an optional field reset and fails naming a
required one.
+ fn assert_reads_without_each_field(
+ metadata: &ManifestMetadata,
+ full_schema: &Value,
+ full_value: &AvroValue,
+ full_entry: &ManifestEntry,
+ optional: &[(&[&str], ResetField)],
+ required: &[&[&str]],
+ ) {
+ for (path, reset) in optional {
+ let (mut schema, mut value) = (full_schema.clone(),
full_value.clone());
+ remove_writer_field(&mut schema, &mut value, path);
+ let bs = write_avro_values_with_writer_schema(metadata, &schema,
vec![value]);
+
+ let manifest = Manifest::parse_avro(&bs).unwrap();
+
+ let mut expected = full_entry.clone();
+ reset(&mut expected);
+ assert_eq!(
+ manifest,
+ Manifest::new(metadata.clone(), vec![expected]),
+ "{path:?}"
+ );
+ }
+
+ for path in required {
+ let (mut schema, mut value) = (full_schema.clone(),
full_value.clone());
+ remove_writer_field(&mut schema, &mut value, path);
+ let bs = write_avro_values_with_writer_schema(metadata, &schema,
vec![value]);
+
+ let err = Manifest::parse_avro(&bs).unwrap_err();
+
+ assert!(
+ err.to_string().contains(path.last().unwrap()),
+ "{path:?}: {err}"
+ );
+ }
+ }
+
+ #[test]
+ fn test_parse_manifest_without_each_field() {
+ // Fields that the spec doesn't require in every version must read as
+ // their default when the writer omits them.
+ let metadata = writer_schema_test_metadata();
+ let partition_type = metadata
+ .partition_spec
+ .partition_type(&metadata.schema)
+ .unwrap();
+ let full_schema =
+
serde_json::to_value(manifest_schema_v2(&partition_type).unwrap()).unwrap();
+ let mut full_entry = expected_entry(Struct::from_iter([
+ Some(Literal::long(5)),
+ Some(Literal::string("a")),
+ Some(Literal::double(2.5)),
+ ]));
+ full_entry.sequence_number = Some(3);
+ full_entry.file_sequence_number = Some(4);
+ let data_file = &mut full_entry.data_file;
+ data_file.content = DataContentType::PositionDeletes;
+ data_file.column_sizes = HashMap::from([(1, 40)]);
+ data_file.value_counts = HashMap::from([(1, 10)]);
+ data_file.null_value_counts = HashMap::from([(1, 1)]);
+ data_file.nan_value_counts = HashMap::from([(3, 2)]);
+ data_file.lower_bounds = HashMap::from([(1, Datum::long(1))]);
+ data_file.upper_bounds = HashMap::from([(1, Datum::long(9))]);
+ data_file.key_metadata = Some(vec![1, 2]);
+ data_file.split_offsets = Some(vec![4]);
+ data_file.equality_ids = Some(vec![1]);
+ data_file.sort_order_id = Some(0);
+ data_file.first_row_id = Some(100);
+ data_file.referenced_data_file =
Some("s3://bucket/table/data/b.parquet".to_string());
+ data_file.content_offset = Some(4);
+ data_file.content_size_in_bytes = Some(8);
+ let full_value = to_value(
+ _serde::ManifestEntryV2::try_from(
+ full_entry.clone(),
+ &Type::Struct(partition_type.clone()),
+ )
+ .unwrap(),
+ )
+ .unwrap();
+
+ let optional: [(&[&str], ResetField); 18] = [
+ (&["snapshot_id"], |e| e.snapshot_id = None),
+ (&["sequence_number"], |e| e.sequence_number = None),
+ (&["file_sequence_number"], |e| e.file_sequence_number = None),
+ (&["data_file", "content"], |e| {
+ e.data_file.content = DataContentType::Data
+ }),
+ (&["data_file", "column_sizes"], |e| {
+ e.data_file.column_sizes.clear()
+ }),
+ (&["data_file", "value_counts"], |e| {
+ e.data_file.value_counts.clear()
+ }),
+ (&["data_file", "null_value_counts"], |e| {
+ e.data_file.null_value_counts.clear()
+ }),
+ (&["data_file", "nan_value_counts"], |e| {
+ e.data_file.nan_value_counts.clear()
+ }),
+ (&["data_file", "lower_bounds"], |e| {
+ e.data_file.lower_bounds.clear()
+ }),
+ (&["data_file", "upper_bounds"], |e| {
+ e.data_file.upper_bounds.clear()
+ }),
+ (&["data_file", "key_metadata"], |e| {
+ e.data_file.key_metadata = None
+ }),
+ (&["data_file", "split_offsets"], |e| {
+ e.data_file.split_offsets = None
+ }),
+ (&["data_file", "equality_ids"], |e| {
+ e.data_file.equality_ids = None
+ }),
+ (&["data_file", "sort_order_id"], |e| {
+ e.data_file.sort_order_id = None
+ }),
+ (&["data_file", "first_row_id"], |e| {
+ e.data_file.first_row_id = None
+ }),
+ (&["data_file", "referenced_data_file"], |e| {
+ e.data_file.referenced_data_file = None
+ }),
+ (&["data_file", "content_offset"], |e| {
+ e.data_file.content_offset = None
+ }),
+ (&["data_file", "content_size_in_bytes"], |e| {
+ e.data_file.content_size_in_bytes = None
+ }),
+ ];
+ let required: [&[&str]; 7] = [
+ &["status"],
+ &["data_file"],
+ &["data_file", "file_path"],
+ &["data_file", "file_format"],
+ &["data_file", "partition"],
+ &["data_file", "record_count"],
+ &["data_file", "file_size_in_bytes"],
+ ];
+ assert_reads_without_each_field(
+ &metadata,
+ &full_schema,
+ &full_value,
+ &full_entry,
+ &optional,
+ &required,
+ );
+ }
+
+ #[test]
+ fn test_parse_v1_manifest_without_each_field() {
+ let mut metadata = writer_schema_test_metadata();
+ metadata.format_version = FormatVersion::V1;
+ let partition_type = metadata
+ .partition_spec
+ .partition_type(&metadata.schema)
+ .unwrap();
+ let full_schema =
+
serde_json::to_value(manifest_schema_v1(&partition_type).unwrap()).unwrap();
+ let mut full_entry = expected_entry(Struct::from_iter([
+ Some(Literal::long(5)),
+ Some(Literal::string("a")),
+ Some(Literal::double(2.5)),
+ ]));
+ // V1 has no sequence numbers, and every V1 entry reads with 0.
+ full_entry.sequence_number = Some(0);
+ full_entry.file_sequence_number = Some(0);
+ let data_file = &mut full_entry.data_file;
+ data_file.column_sizes = HashMap::from([(1, 40)]);
+ data_file.value_counts = HashMap::from([(1, 10)]);
+ data_file.null_value_counts = HashMap::from([(1, 1)]);
+ data_file.nan_value_counts = HashMap::from([(3, 2)]);
+ data_file.lower_bounds = HashMap::from([(1, Datum::long(1))]);
+ data_file.upper_bounds = HashMap::from([(1, Datum::long(9))]);
+ data_file.key_metadata = Some(vec![1, 2]);
+ data_file.split_offsets = Some(vec![4]);
+ data_file.sort_order_id = Some(0);
+ let full_value = to_value(
+ _serde::ManifestEntryV1::try_from(
+ full_entry.clone(),
+ &Type::Struct(partition_type.clone()),
+ )
+ .unwrap(),
+ )
+ .unwrap();
+
+ let optional: [(&[&str], ResetField); 10] = [
+ // Required in V1, but deprecated and not read.
+ (&["data_file", "block_size_in_bytes"], |_| {}),
+ (&["data_file", "column_sizes"], |e| {
+ e.data_file.column_sizes.clear()
+ }),
+ (&["data_file", "value_counts"], |e| {
+ e.data_file.value_counts.clear()
+ }),
+ (&["data_file", "null_value_counts"], |e| {
+ e.data_file.null_value_counts.clear()
+ }),
+ (&["data_file", "nan_value_counts"], |e| {
+ e.data_file.nan_value_counts.clear()
+ }),
+ (&["data_file", "lower_bounds"], |e| {
+ e.data_file.lower_bounds.clear()
+ }),
+ (&["data_file", "upper_bounds"], |e| {
+ e.data_file.upper_bounds.clear()
+ }),
+ (&["data_file", "key_metadata"], |e| {
+ e.data_file.key_metadata = None
+ }),
+ (&["data_file", "split_offsets"], |e| {
+ e.data_file.split_offsets = None
+ }),
+ (&["data_file", "sort_order_id"], |e| {
+ e.data_file.sort_order_id = None
+ }),
+ ];
+ let required: [&[&str]; 8] = [
+ &["status"],
+ &["snapshot_id"],
+ &["data_file"],
+ &["data_file", "file_path"],
+ &["data_file", "file_format"],
+ &["data_file", "partition"],
+ &["data_file", "record_count"],
+ &["data_file", "file_size_in_bytes"],
+ ];
+ assert_reads_without_each_field(
+ &metadata,
+ &full_schema,
+ &full_value,
+ &full_entry,
+ &optional,
+ &required,
+ );
+ }
+
+ #[test]
+ fn test_parse_manifest_matches_partition_fields_by_name() {
+ // The writer orders the partition fields differently from the spec,
omits
+ // `name`, and adds a field the spec doesn't have.
+ let metadata = writer_schema_test_metadata();
+ let writer_schema = v2_writer_schema(serde_json::json!([
+ {"name": "score", "type": ["null", "double"], "default": null,
"field-id": 1002},
+ {"name": "extra", "type": ["null", "int"], "default": null,
"field-id": 1003},
+ {"name": "id", "type": ["null", "long"], "default": null,
"field-id": 1000},
+ ]));
+ let bs = write_with_writer_schema(&metadata, &writer_schema,
vec![v2_entry(
+ serde_json::json!({"score": 2.5, "extra": 9, "id": 5}),
+ )]);
+
+ let manifest = Manifest::parse_avro(&bs).unwrap();
+
+ let partition =
+ Struct::from_iter([Some(Literal::long(5)), None,
Some(Literal::double(2.5))]);
+ assert_eq!(
+ manifest,
+ Manifest::new(metadata, vec![expected_entry(partition)])
+ );
+ }
+
+ #[test]
+ fn test_parse_manifest_promotes_writer_types() {
+ // Avro promotes int to long and float to double.
+ let metadata = writer_schema_test_metadata();
+ let mut writer_schema = v2_writer_schema(serde_json::json!([
+ {"name": "id", "type": ["null", "int"], "default": null,
"field-id": 1000},
+ {"name": "name", "type": ["null", "string"], "default": null,
"field-id": 1001},
+ {"name": "score", "type": ["null", "float"], "default": null,
"field-id": 1002},
+ ]));
+ for field in ["record_count", "file_size_in_bytes"] {
+ writer_schema_field(&mut writer_schema, &["data_file",
field])["type"] =
+ serde_json::json!("int");
+ }
+ let bs = write_with_writer_schema(&metadata, &writer_schema,
vec![v2_entry(
+ serde_json::json!({"id": 5, "name": "a", "score": 2.5}),
+ )]);
+
+ let manifest = Manifest::parse_avro(&bs).unwrap();
+
+ let partition = Struct::from_iter([
+ Some(Literal::long(5)),
+ Some(Literal::string("a")),
+ Some(Literal::double(2.5)),
+ ]);
+ assert_eq!(
+ manifest,
+ Manifest::new(metadata, vec![expected_entry(partition)])
+ );
+ }
+
+ #[test]
+ fn test_parse_manifest_reads_required_writer_fields_as_optional() {
Review Comment:
Optional follow-on, not blocking: these writer-schema variants all write
`long`/`double` matching the reader, so the one path still untested is numeric
promotion. A sibling that writes `int` for a `long` field (or `float` for
`double`) would pin the promotion the serde visitors now own — that's the piece
the old resolving reader used to handle for us. Cheap to add right alongside
these.
##########
crates/iceberg/src/spec/manifest/mod.rs:
##########
@@ -1361,4 +1398,861 @@ mod tests {
assert_eq!(deserialized_data_file1, original_data_file1);
assert_eq!(deserialized_data_file2, original_data_file2);
}
+
+ /// Metadata for the writer schema tests: a V2 data manifest whose spec has
+ /// identity partitions on a long, a string, and a double column.
+ fn writer_schema_test_metadata() -> ManifestMetadata {
+ let schema = Arc::new(
+ Schema::builder()
+ .with_fields(vec![
+ NestedField::optional(1, "id",
Type::Primitive(PrimitiveType::Long)).into(),
+ NestedField::optional(2, "name",
Type::Primitive(PrimitiveType::String)).into(),
+ NestedField::optional(3, "score",
Type::Primitive(PrimitiveType::Double))
+ .into(),
+ ])
+ .build()
+ .unwrap(),
+ );
+ let partition_spec = PartitionSpec::builder(schema.clone())
+ .with_spec_id(0)
+ .add_partition_field("id", "id", Transform::Identity)
+ .unwrap()
+ .add_partition_field("name", "name", Transform::Identity)
+ .unwrap()
+ .add_partition_field("score", "score", Transform::Identity)
+ .unwrap()
+ .build()
+ .unwrap();
+ ManifestMetadata {
+ schema_id: 0,
+ schema,
+ partition_spec,
+ content: ManifestContentType::Data,
+ format_version: FormatVersion::V2,
+ }
+ }
+
+ /// A V2 manifest entry writer schema with Java's record names and only the
+ /// required `data_file` fields.
+ fn v2_writer_schema(partition_fields: Value) -> Value {
+ serde_json::json!({
+ "type": "record",
+ "name": "manifest_entry",
+ "fields": [
+ {"name": "status", "type": "int", "field-id": 0},
+ {"name": "snapshot_id", "type": ["null", "long"], "default":
null, "field-id": 1},
+ {"name": "sequence_number", "type": ["null", "long"],
"default": null, "field-id": 3},
+ {"name": "file_sequence_number", "type": ["null", "long"],
"default": null, "field-id": 4},
+ {"name": "data_file", "field-id": 2, "type": {
+ "type": "record",
+ "name": "r2",
+ "fields": [
+ {"name": "content", "type": "int", "field-id": 134},
+ {"name": "file_path", "type": "string", "field-id":
100},
+ {"name": "file_format", "type": "string", "field-id":
101},
+ {"name": "partition", "field-id": 102, "type": {
+ "type": "record",
+ "name": "r102",
+ "fields": partition_fields,
+ }},
+ {"name": "record_count", "type": "long", "field-id":
103},
+ {"name": "file_size_in_bytes", "type": "long",
"field-id": 104},
+ ],
+ }},
+ ],
+ })
+ }
+
+ /// The field named by `path` in a record schema, descending through nested
+ /// record types.
+ fn writer_schema_field<'a>(record: &'a mut Value, path: &[&str]) -> &'a
mut Value {
+ let (name, rest) = path.split_first().unwrap();
+ let field = record["fields"]
+ .as_array_mut()
+ .unwrap()
+ .iter_mut()
+ .find(|field| field["name"] == *name)
+ .unwrap();
+ if rest.is_empty() {
+ field
+ } else {
+ writer_schema_field(&mut field["type"], rest)
+ }
+ }
+
+ fn v2_entry(partition: Value) -> Value {
+ serde_json::json!({
+ "status": 1,
+ "snapshot_id": 7,
+ "sequence_number": null,
+ "file_sequence_number": null,
+ "data_file": {
+ "content": 0,
+ "file_path": "s3://bucket/table/data/a.parquet",
+ "file_format": "PARQUET",
+ "partition": partition,
+ "record_count": 10,
+ "file_size_in_bytes": 100,
+ },
+ })
+ }
+
+ fn expected_entry(partition: Struct) -> ManifestEntry {
+ ManifestEntry {
+ status: ManifestStatus::Added,
+ snapshot_id: Some(7),
+ sequence_number: None,
+ file_sequence_number: None,
+ data_file: DataFile {
+ content: DataContentType::Data,
+ file_path: "s3://bucket/table/data/a.parquet".to_string(),
+ file_format: DataFileFormat::Parquet,
+ partition,
+ record_count: 10,
+ file_size_in_bytes: 100,
+ column_sizes: HashMap::new(),
+ value_counts: HashMap::new(),
+ null_value_counts: HashMap::new(),
+ nan_value_counts: HashMap::new(),
+ lower_bounds: HashMap::new(),
+ upper_bounds: HashMap::new(),
+ key_metadata: None,
+ split_offsets: None,
+ equality_ids: None,
+ sort_order_id: None,
+ partition_spec_id: 0,
+ first_row_id: None,
+ referenced_data_file: None,
+ content_offset: None,
+ content_size_in_bytes: None,
+ },
+ }
+ }
+
+ /// Writes `entries` as a manifest with the given writer schema, encoding
each
+ /// JSON value with the writer schema's types.
+ fn write_with_writer_schema(
+ metadata: &ManifestMetadata,
+ writer_schema: &Value,
+ entries: Vec<Value>,
+ ) -> Vec<u8> {
+ write_avro_values_with_writer_schema(
+ metadata,
+ writer_schema,
+ entries
+ .into_iter()
+ .map(|entry| AvroValue::try_from(entry).unwrap())
+ .collect(),
+ )
+ }
+
+ /// Like [`write_with_writer_schema`], for entries that JSON can't express.
+ fn write_avro_values_with_writer_schema(
+ metadata: &ManifestMetadata,
+ writer_schema: &Value,
+ entries: Vec<AvroValue>,
+ ) -> Vec<u8> {
+ let avro_schema = apache_avro::Schema::parse(writer_schema).unwrap();
+ let mut writer = Writer::new(&avro_schema, Vec::new()).unwrap();
+ for (key, value) in [
+ ("schema", to_vec(&metadata.schema).unwrap()),
+ ("schema-id", metadata.schema_id.to_string().into_bytes()),
+ (
+ "partition-spec",
+ to_vec(&metadata.partition_spec.fields()).unwrap(),
+ ),
+ (
+ "partition-spec-id",
+ metadata.partition_spec.spec_id().to_string().into_bytes(),
+ ),
+ (
+ "format-version",
+ (metadata.format_version as u8).to_string().into_bytes(),
+ ),
+ ("content", metadata.content.to_string().into_bytes()),
+ ] {
+ writer.add_user_metadata(key.to_string(), value).unwrap();
+ }
+ for entry in entries {
+ writer
+ .append_value(entry.resolve(&avro_schema).unwrap())
+ .unwrap();
+ }
+ writer.into_inner().unwrap()
+ }
+
+ /// Removes the field named by `path` from a record schema and from a
record
+ /// value of that schema.
+ fn remove_writer_field(schema: &mut Value, entry: &mut AvroValue, path:
&[&str]) {
+ let (name, rest) = path.split_first().unwrap();
+ let AvroValue::Record(values) = entry else {
+ unreachable!("the entry is a record");
+ };
+ let fields = schema["fields"].as_array_mut().unwrap();
+ if rest.is_empty() {
+ fields.retain(|field| field["name"] != *name);
+ values.retain(|(value_name, _)| value_name != name);
+ } else {
+ let field = fields.iter_mut().find(|field| field["name"] == *name);
+ let value = values.iter_mut().find(|(value_name, _)| value_name ==
name);
+ remove_writer_field(&mut field.unwrap()["type"], &mut
value.unwrap().1, rest);
+ }
+ }
+
+ type ResetField = fn(&mut ManifestEntry);
+
+ /// Removes each field in `optional` and `required` in turn from a manifest
+ /// that `full_schema` and `full_value` write as `full_entry`. Checks that
the
+ /// manifest reads with an optional field reset and fails naming a
required one.
+ fn assert_reads_without_each_field(
+ metadata: &ManifestMetadata,
+ full_schema: &Value,
+ full_value: &AvroValue,
+ full_entry: &ManifestEntry,
+ optional: &[(&[&str], ResetField)],
+ required: &[&[&str]],
+ ) {
+ for (path, reset) in optional {
+ let (mut schema, mut value) = (full_schema.clone(),
full_value.clone());
+ remove_writer_field(&mut schema, &mut value, path);
+ let bs = write_avro_values_with_writer_schema(metadata, &schema,
vec![value]);
+
+ let manifest = Manifest::parse_avro(&bs).unwrap();
+
+ let mut expected = full_entry.clone();
+ reset(&mut expected);
+ assert_eq!(
+ manifest,
+ Manifest::new(metadata.clone(), vec![expected]),
+ "{path:?}"
+ );
+ }
+
+ for path in required {
+ let (mut schema, mut value) = (full_schema.clone(),
full_value.clone());
+ remove_writer_field(&mut schema, &mut value, path);
+ let bs = write_avro_values_with_writer_schema(metadata, &schema,
vec![value]);
+
+ let err = Manifest::parse_avro(&bs).unwrap_err();
+
+ assert!(
+ err.to_string().contains(path.last().unwrap()),
+ "{path:?}: {err}"
+ );
+ }
+ }
+
+ #[test]
+ fn test_parse_manifest_without_each_field() {
+ // Fields that the spec doesn't require in every version must read as
+ // their default when the writer omits them.
+ let metadata = writer_schema_test_metadata();
+ let partition_type = metadata
+ .partition_spec
+ .partition_type(&metadata.schema)
+ .unwrap();
+ let full_schema =
+
serde_json::to_value(manifest_schema_v2(&partition_type).unwrap()).unwrap();
+ let mut full_entry = expected_entry(Struct::from_iter([
+ Some(Literal::long(5)),
+ Some(Literal::string("a")),
+ Some(Literal::double(2.5)),
+ ]));
+ full_entry.sequence_number = Some(3);
+ full_entry.file_sequence_number = Some(4);
+ let data_file = &mut full_entry.data_file;
+ data_file.content = DataContentType::PositionDeletes;
+ data_file.column_sizes = HashMap::from([(1, 40)]);
+ data_file.value_counts = HashMap::from([(1, 10)]);
+ data_file.null_value_counts = HashMap::from([(1, 1)]);
+ data_file.nan_value_counts = HashMap::from([(3, 2)]);
+ data_file.lower_bounds = HashMap::from([(1, Datum::long(1))]);
+ data_file.upper_bounds = HashMap::from([(1, Datum::long(9))]);
+ data_file.key_metadata = Some(vec![1, 2]);
+ data_file.split_offsets = Some(vec![4]);
+ data_file.equality_ids = Some(vec![1]);
+ data_file.sort_order_id = Some(0);
+ data_file.first_row_id = Some(100);
+ data_file.referenced_data_file =
Some("s3://bucket/table/data/b.parquet".to_string());
+ data_file.content_offset = Some(4);
+ data_file.content_size_in_bytes = Some(8);
+ let full_value = to_value(
+ _serde::ManifestEntryV2::try_from(
+ full_entry.clone(),
+ &Type::Struct(partition_type.clone()),
+ )
+ .unwrap(),
+ )
+ .unwrap();
+
+ let optional: [(&[&str], ResetField); 18] = [
+ (&["snapshot_id"], |e| e.snapshot_id = None),
+ (&["sequence_number"], |e| e.sequence_number = None),
+ (&["file_sequence_number"], |e| e.file_sequence_number = None),
+ (&["data_file", "content"], |e| {
+ e.data_file.content = DataContentType::Data
+ }),
+ (&["data_file", "column_sizes"], |e| {
+ e.data_file.column_sizes.clear()
+ }),
+ (&["data_file", "value_counts"], |e| {
+ e.data_file.value_counts.clear()
+ }),
+ (&["data_file", "null_value_counts"], |e| {
+ e.data_file.null_value_counts.clear()
+ }),
+ (&["data_file", "nan_value_counts"], |e| {
+ e.data_file.nan_value_counts.clear()
+ }),
+ (&["data_file", "lower_bounds"], |e| {
+ e.data_file.lower_bounds.clear()
+ }),
+ (&["data_file", "upper_bounds"], |e| {
+ e.data_file.upper_bounds.clear()
+ }),
+ (&["data_file", "key_metadata"], |e| {
+ e.data_file.key_metadata = None
+ }),
+ (&["data_file", "split_offsets"], |e| {
+ e.data_file.split_offsets = None
+ }),
+ (&["data_file", "equality_ids"], |e| {
+ e.data_file.equality_ids = None
+ }),
+ (&["data_file", "sort_order_id"], |e| {
+ e.data_file.sort_order_id = None
+ }),
+ (&["data_file", "first_row_id"], |e| {
+ e.data_file.first_row_id = None
+ }),
+ (&["data_file", "referenced_data_file"], |e| {
+ e.data_file.referenced_data_file = None
+ }),
+ (&["data_file", "content_offset"], |e| {
+ e.data_file.content_offset = None
+ }),
+ (&["data_file", "content_size_in_bytes"], |e| {
+ e.data_file.content_size_in_bytes = None
+ }),
+ ];
+ let required: [&[&str]; 7] = [
+ &["status"],
+ &["data_file"],
+ &["data_file", "file_path"],
+ &["data_file", "file_format"],
+ &["data_file", "partition"],
+ &["data_file", "record_count"],
+ &["data_file", "file_size_in_bytes"],
+ ];
+ assert_reads_without_each_field(
+ &metadata,
+ &full_schema,
+ &full_value,
+ &full_entry,
+ &optional,
+ &required,
+ );
+ }
+
+ #[test]
+ fn test_parse_v1_manifest_without_each_field() {
+ let mut metadata = writer_schema_test_metadata();
+ metadata.format_version = FormatVersion::V1;
+ let partition_type = metadata
+ .partition_spec
+ .partition_type(&metadata.schema)
+ .unwrap();
+ let full_schema =
+
serde_json::to_value(manifest_schema_v1(&partition_type).unwrap()).unwrap();
+ let mut full_entry = expected_entry(Struct::from_iter([
+ Some(Literal::long(5)),
+ Some(Literal::string("a")),
+ Some(Literal::double(2.5)),
+ ]));
+ // V1 has no sequence numbers, and every V1 entry reads with 0.
+ full_entry.sequence_number = Some(0);
+ full_entry.file_sequence_number = Some(0);
+ let data_file = &mut full_entry.data_file;
+ data_file.column_sizes = HashMap::from([(1, 40)]);
+ data_file.value_counts = HashMap::from([(1, 10)]);
+ data_file.null_value_counts = HashMap::from([(1, 1)]);
+ data_file.nan_value_counts = HashMap::from([(3, 2)]);
+ data_file.lower_bounds = HashMap::from([(1, Datum::long(1))]);
+ data_file.upper_bounds = HashMap::from([(1, Datum::long(9))]);
+ data_file.key_metadata = Some(vec![1, 2]);
+ data_file.split_offsets = Some(vec![4]);
+ data_file.sort_order_id = Some(0);
+ let full_value = to_value(
+ _serde::ManifestEntryV1::try_from(
+ full_entry.clone(),
+ &Type::Struct(partition_type.clone()),
+ )
+ .unwrap(),
+ )
+ .unwrap();
+
+ let optional: [(&[&str], ResetField); 10] = [
+ // Required in V1, but deprecated and not read.
+ (&["data_file", "block_size_in_bytes"], |_| {}),
+ (&["data_file", "column_sizes"], |e| {
+ e.data_file.column_sizes.clear()
+ }),
+ (&["data_file", "value_counts"], |e| {
+ e.data_file.value_counts.clear()
+ }),
+ (&["data_file", "null_value_counts"], |e| {
+ e.data_file.null_value_counts.clear()
+ }),
+ (&["data_file", "nan_value_counts"], |e| {
+ e.data_file.nan_value_counts.clear()
+ }),
+ (&["data_file", "lower_bounds"], |e| {
+ e.data_file.lower_bounds.clear()
+ }),
+ (&["data_file", "upper_bounds"], |e| {
+ e.data_file.upper_bounds.clear()
+ }),
+ (&["data_file", "key_metadata"], |e| {
+ e.data_file.key_metadata = None
+ }),
+ (&["data_file", "split_offsets"], |e| {
+ e.data_file.split_offsets = None
+ }),
+ (&["data_file", "sort_order_id"], |e| {
+ e.data_file.sort_order_id = None
+ }),
+ ];
+ let required: [&[&str]; 8] = [
+ &["status"],
+ &["snapshot_id"],
Review Comment:
The V1 audit landed exactly as I wanted — the required set and the 9
optional-reset cases plus the `block_size_in_bytes` no-op cover the real path.
One thing I want to confirm before we call the contract closed, not a blocker:
Java's `V1Metadata.wrapFileSchema` uses `ManifestEntry.SNAPSHOT_ID`, which is
`optional(1, …)`, so a Java-written V1 manifest can carry `["null","long"]`
here and `snapshot_id: i64` would reject the null. The spec's V1 column lists
it required, so pinning it here isn't wrong — have you seen a null V1
`snapshot_id` in the wild (inheritance writers defer the id)? If it's a shape
we need to tolerate, `Option<i64>` plus a null case covers it; if not, a
one-line note that it's required-per-spec /
Java-writes-optional-for-list-inheritance is plenty.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]