mbutrovich commented on code in PR #3354:
URL: https://github.com/apache/iceberg-rust/pull/3354#discussion_r4224931672
##########
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:
> have you seen a null V1 `snapshot_id` in the wild (inheritance writers
defer the id)?
I haven't.
> 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.
I added the
[note](https://github.com/apache/iceberg-rust/blob/95667f5be29b8c24a8a5155b80c1fb5986a1c0d0/crates/iceberg/src/spec/manifest/_serde.rs#L79-L81)
on `ManifestEntryV1::snapshot_id`, with a link to #3371. `main` rejects these
manifests too, and #3371 has a reproduction against `main`. Reading the field
as `Option<i64>` changes which manifests iceberg-rust accepts, so I'd like to
decide that in #3371 rather than in this PR.
--
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]