laskoviymishka commented on code in PR #3354: URL: https://github.com/apache/iceberg-rust/pull/3354#discussion_r4195561250
########## crates/iceberg/src/avro/named_types.rs: ########## @@ -0,0 +1,255 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! Reading Avro files whose schema defines a named type more than once. +//! +//! The Avro specification allows one definition per full name. iceberg-rust +//! wrote manifests that define a `fixed` or `decimal` type once for every field +//! that uses it, while `apache-avro` 0.21 accepted that. `apache-avro` 0.22 +//! rejects such a file's header. Rewriting each repeated definition as a +//! reference to the first one keeps those files readable, because a reference +//! decodes the same bytes as the definition it refers to. + +use std::collections::HashMap; +use std::io::Cursor; + +use apache_avro::Schema as AvroSchema; +use apache_avro::reader::datum::GenericDatumReader; +use apache_avro::schema::MapSchema; +use apache_avro::types::Value as AvroValue; +use apache_avro::writer::datum::GenericDatumWriter; +use serde_json::{Map, Value as JsonValue}; + +use crate::Result; +use crate::error::invalid_data; + +const MAGIC: &[u8] = b"Obj\x01"; +const SCHEMA_KEY: &str = "avro.schema"; + +/// Rewrites the header of an Avro object container file so that each named type +/// is defined once and referenced by name afterwards. Returns the rewritten +/// file and the full names of the repeated types, or `None` if a repeated +/// definition differs from the first one, which leaves the file ambiguous. +pub(crate) fn define_named_types_once(bs: &[u8]) -> Result<Option<(Vec<u8>, Vec<String>)>> { + let body = bs + .strip_prefix(MAGIC) + .ok_or_else(|| invalid_data!("Not an Avro object container file"))?; + // The header's metadata is an Avro `map<bytes>`, followed by the sync marker. + let metadata_schema = AvroSchema::Map(MapSchema { + types: Box::new(AvroSchema::Bytes), + attributes: Default::default(), + }); + let mut cursor = Cursor::new(body); + let AvroValue::Map(mut metadata) = GenericDatumReader::builder(&metadata_schema) + .build()? + .read_value(&mut cursor)? + else { + return Err(invalid_data!("Avro file metadata is not a map")); + }; + let rest = &body[cursor.position() as usize..]; Review Comment: `cursor.position() as usize` then slicing is a latent panic on untrusted header bytes. Safe in practice since a read won't advance past the end, but `body.get(usize::try_from(cursor.position()).ok()?..)` (or mapping to `invalid_data!`) removes the path. ########## Cargo.toml: ########## @@ -46,7 +46,7 @@ unused_qualifications = "deny" [workspace.dependencies] aes-gcm = "0.10" anyhow = "1.0.72" -apache-avro = { version = "0.21", features = ["snappy", "zstandard"] } +apache-avro = { version = "0.22", features = ["snappy", "zstandard"] } Review Comment: The deserializer's union handling rides on how 0.22's `deserialize_any` dispatches unions, which isn't a documented public contract — a `0.22.x` minor bump could shift it under us. Worth a comment here that the adapter semantics are pinned to this version, and keeping `test_parse_manifest_reads_union_writer_fields_as_{required,optional}` as the guard for the next bump. ########## crates/iceberg/src/avro/schema.rs: ########## @@ -254,11 +271,11 @@ impl SchemaVisitor for SchemaToAvroSchema { PrimitiveType::TimestampNs => AvroSchema::TimestampNanos, PrimitiveType::TimestamptzNs => AvroSchema::TimestampNanos, PrimitiveType::String => AvroSchema::String, - PrimitiveType::Uuid => AvroSchema::Uuid, - PrimitiveType::Fixed(len) => avro_fixed_schema((*len) as usize)?, + PrimitiveType::Uuid => AvroSchema::Uuid(UuidSchema::String), Review Comment: The new read test accepts the spec's `fixed(16)` uuid, but the writer still picks `UuidSchema::String` here — so we read the spec form and don't write it, and a Java/PyIceberg reader resolving our manifest sees a string writer schema against its fixed reader schema. The bytes match 0.21 so it isn't a regression, but the description frames the uuid change as read-only. I'd either emit `UuidSchema::Fixed` (name `uuid_fixed`, size 16, same define-once handling) or call the choice out explicitly. ########## crates/iceberg/src/avro/schema.rs: ########## @@ -975,6 +989,10 @@ mod tests { } #[test] + #[ignore = "apache-avro 0.22 drops `logicalType: map` when parsing schema JSON \ Review Comment: The `#[ignore]` reason only covers the read direction (`avro_schema_to_schema`), but the writer builds its manifest header schema from a parsed `AvroSchema` and carries `logicalType: map` through `custom_attributes`. If 0.22 drops that attribute on parse/re-serialize, the headers we write lose the map-encoding marker that Java and PyIceberg rely on when reading map-shaped `column_sizes`/bounds — a silent cross-engine interop regression, and nothing here inspects the written header. Before taking the bump I'd add a test that writes a manifest and asserts its header `avro.schema` still contains `"logicalType":"map"`. And since the ignore drops the only coverage of the map-array conversion, a replacement test exercising it on a programmatically built `AvroSchema` would keep that from going dark — plus a tracking issue in this repo. ########## crates/iceberg/src/spec/manifest/mod.rs: ########## @@ -61,43 +84,35 @@ impl Manifest { let partition_struct_type = Type::Struct(partition_type.clone()); let entries = match metadata.format_version { - FormatVersion::V1 => { - let schema = manifest_schema_v1(&partition_type)?; - let reader = AvroReader::with_schema(&schema, bs)?; - reader - .into_iter() - .map(|value| { - from_value::<_serde::ManifestEntryV1>(&value?)?.try_into( - metadata.partition_spec.spec_id(), - &partition_struct_type, - &metadata.schema, - ) - }) - .collect::<Result<Vec<_>>>()? - } + FormatVersion::V1 => reader + .into_deser_iter::<Resolved<_serde::ManifestEntryV1>>() + .map(|entry| { + entry?.0.try_into( + metadata.partition_spec.spec_id(), + &partition_struct_type, + &metadata.schema, + ) + }) + .collect::<Result<Vec<_>>>()?, // Manifest Schema & Manifest Entry did not change between V2 and V3 - FormatVersion::V2 | FormatVersion::V3 => { - let schema = manifest_schema_v2(&partition_type)?; - let reader = AvroReader::with_schema(&schema, bs)?; - reader - .into_iter() - .map(|value| { - from_value::<_serde::ManifestEntryV2>(&value?)?.try_into( - metadata.partition_spec.spec_id(), - &partition_struct_type, - &metadata.schema, - ) - }) - .collect::<Result<Vec<_>>>()? - } + FormatVersion::V2 | FormatVersion::V3 => reader + .into_deser_iter::<Resolved<_serde::ManifestEntryV2>>() Review Comment: Dropping `with_schema` resolution means these `_serde` structs are now the only compatibility contract for every on-disk manifest — Java, PyIceberg, older Rust. The reader schema used to fill defaults, promote, and reorder; now a required non-`Option` field that a foreign writer omits fails with a serde "missing field" error instead of taking the schema default, and nothing systematically audits which fields need `#[serde(default)]` (only `content` has it). I'd want this written into the module doc — any new field must be `Option` or `#[serde(default)]` — and backed by a table-driven test asserting every spec-optional/defaulted field is tolerated when absent, plus at least one real Java/PyIceberg-written manifest in the fixtures rather than only synthetic JSON through this crate's own writer. ########## crates/iceberg/src/spec/manifest/mod.rs: ########## @@ -45,9 +47,30 @@ 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, + Err(e) if matches!(e.details(), Details::AmbiguousSchemaDefinition(_)) => { + let Some((bs, repeated)) = define_named_types_once(bs)? else { + return Err(e.into()); + }; + tracing::warn!( Review Comment: This fires a `warn!` and copies the whole file into a second buffer on every read of every legacy manifest, so a large legacy table turns into a warn storm plus a per-manifest copy on each scan. It also keys off one 0.x error-details variant with no stated removal condition. I'd dedupe the warning (once per path, or `debug!` after the first) and note when the shim can be dropped. ########## crates/iceberg/src/spec/manifest/data_file.rs: ########## @@ -333,7 +333,9 @@ pub fn read_data_files_from_avro<R: Read>( FormatVersion::V3 => data_file_schema_v3(partition_type).unwrap(), }; - let reader = AvroReader::with_schema(&avro_schema, reader)?; + let reader = AvroReader::builder(reader) Review Comment: The repeated-definition fallback only runs in `try_from_avro_bytes`, but `read_data_files_from_avro` opens the same kind of bytes this crate itself wrote via `write_data_files_to_avro` — and it now hard-fails on anything a 0.21 build produced. That's a behavior change on a public, previously-readable surface. I'd centralize "open an Avro container, tolerating legacy repeated definitions" in one helper and route both readers through it, so the compat guarantee doesn't depend on which entry point you came in by. ########## crates/iceberg/src/avro/deserializer.rs: ########## @@ -0,0 +1,280 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! Deserialization straight from an Avro writer schema, with the parts of Avro +//! schema resolution that apache-avro's schema-aware deserializer leaves out. +//! +//! `Reader::into_deser_iter` decodes into serde types without building +//! `apache_avro::types::Value`s, but it requires the reader schema to equal the +//! writer schema and applies no resolution rules. [`ResolvingDeserializer`] +//! wraps it and routes every value through `deserialize_any`, which follows the +//! writer schema. As a result: +//! +//! - Avro record names don't have to match serde type names. +//! - A writer value that isn't a union reads into an `Option`. +//! - A writer union reads into a type that isn't an `Option`. A null value +//! returns an error. +//! - serde's numeric visitors convert between numeric types, so an `int` reads +//! into an `i64` and a `float` into an `f64`. An integer that doesn't fit the +//! target, such as a `long` above `i32::MAX` read into an `i32`, returns an +//! error. Conversions into `f32` or `f64` use `as` and can lose precision. +//! +//! apache-avro plans a `SchemaAwareResolvingDeserializer` that resolves against +//! a reader schema (<https://github.com/apache/avro-rs/issues/575>). Once a +//! release includes it, readers can pass their reader schema to +//! `Reader::builder` and drop this module, provided it doesn't reject writer +//! record names that differ from the reader's. The Avro spec requires record +//! names to match, and the writers this crate reads from don't all agree on +//! them. + +use std::fmt; + +use serde::de::value::{ + BorrowedBytesDeserializer, BorrowedStrDeserializer, BytesDeserializer, EnumAccessDeserializer, + MapAccessDeserializer, SeqAccessDeserializer, +}; +use serde::de::{ + DeserializeSeed, Deserializer, EnumAccess, Error, IntoDeserializer, MapAccess, SeqAccess, + Visitor, +}; +use serde::{Deserialize, forward_to_deserialize_any}; + +/// Deserializes the wrapped type through [`ResolvingDeserializer`], for use +/// with `Reader::into_deser_iter`. +pub(crate) struct Resolved<T>(pub T); Review Comment: This module is ~280 lines of visitor plumbing and the resolution contract for manifest reads, but it's only exercised indirectly through the manifest tests — `visit_byte_buf`, `visit_seq`/`visit_map`, `visit_enum` and the `visit_borrowed_*` paths aren't hit, and `OptionVisitor` has no `visit_newtype_struct`. A few table-driven unit tests here (one per writer-schema shape: union, non-union, int→i64, long→i32 overflow, null into a non-option) would lock the branches down and double as the regression guard for the next avro bump. ########## crates/iceberg/src/spec/manifest/mod.rs: ########## @@ -45,9 +47,30 @@ 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, + Err(e) if matches!(e.details(), Details::AmbiguousSchemaDefinition(_)) => { + let Some((bs, repeated)) = define_named_types_once(bs)? else { Review Comment: The `?` here surfaces any secondary failure from the rewrite — "not an Avro object container file", missing `avro.schema`, a serde_json error — in place of the real `AmbiguousSchemaDefinition` cause, so a malformed header reads as an unrelated rewrite error. I'd fall back to the original on any rewrite failure: `match define_named_types_once(bs) { Ok(Some(x)) => x, _ => return Err(e.into()) }`, or attach `e` as the source. ########## crates/iceberg/src/spec/values/tests.rs: ########## @@ -1719,3 +1722,60 @@ fn test_datum_to_decimal_rejects_scale_change() { .contains("Decimal scale conversion is not supported") ); } + +#[test] +fn raw_literal_project_by_name_reorders_and_fills_fields() { Review Comment: small thing while we're here — these two use the bare `raw_literal_project_by_name_*` names where the rest of the file prefixes with `test_` (e.g. `test_datum_to_decimal_rejects_scale_change`). ########## crates/iceberg/src/spec/values/serde.rs: ########## @@ -44,6 +46,49 @@ pub(crate) mod _serde { pub fn try_into(self, ty: &Type) -> Result<Option<Literal>, Error> { self.0.try_into(ty) } + + /// Matches the fields of a record to `struct_type` by name, the way Avro + /// schema resolution matches record fields. The result has + /// `struct_type`'s fields in its order, with null for an optional field the + /// record lacks. Record fields that `struct_type` lacks are dropped. Values + /// other than records are returned unchanged. + pub fn project_by_name(self, struct_type: &StructType) -> Result<Self, Error> { Review Comment: Matching is by field name only, where Java and PyIceberg project the partition struct by field-id, and the slow path collects into a `HashMap` so a duplicate name silently keeps the last value. After partition evolution — a dropped-then-readded field reusing a name with a new id — this can bind the wrong column. The old reader-schema resolution had the same name-based limitation so it isn't a regression, but it's the sole safeguard now. I'd at least note the limitation in the doc comment; the writer's `field-id` is in the JSON schema if we ever want a fallback. ########## crates/iceberg/testdata/manifests/README.md: ########## @@ -0,0 +1,37 @@ +<!-- + ~ Licensed to the Apache Software Foundation (ASF) under one + ~ or more contributor license agreements. See the NOTICE file + ~ distributed with this work for additional information + ~ regarding copyright ownership. The ASF licenses this file + ~ to you under the Apache License, Version 2.0 (the + ~ "License"); you may not use this file except in compliance + ~ with the License. You may obtain a copy of the License at + ~ + ~ http://www.apache.org/licenses/LICENSE-2.0 + ~ + ~ Unless required by applicable law or agreed to in writing, + ~ software distributed under the License is distributed on an + ~ "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + ~ KIND, either express or implied. See the License for the + ~ specific language governing permissions and limitations + ~ under the License. +--> + + +# Manifest test data + +`repeated-decimal-type-definitions.avro` is a V2 data manifest written by +iceberg-rust at commit `ecba0b5`, which used `apache-avro` 0.21. Its table schema +has two optional `decimal(10, 2)` columns, `d1` (id 1) and `d2` (id 2), and its +partition spec has an identity partition on each. It holds one entry for +`s3://b/t/a.parquet` with partition values `123.45` and `-6.78`. + +The Avro schema in its header defines the named fixed type `decimal_10_2` twice. +The Avro specification allows only one definition of each name, and +`apache-avro` 0.22 rejects the header, but iceberg-rust wrote manifests like +this one before it moved to 0.22. Tests use it to check that such manifests +still read. + +To regenerate it, check out a commit that uses `apache-avro` 0.21, write the Review Comment: The one legacy fixture is binary and decimal-only, and the regeneration recipe pins a commit plus a 0.21 checkout that's hard to reproduce once that dependency state is gone. Generating the legacy bytes in the test itself — write the duplicated-definition JSON schema through the `named_types` header writer — would be self-contained and reviewable, and would let us also cover the repeated-`fixed` case the fallback handles but nothing currently tests. ########## crates/iceberg/src/avro/schema.rs: ########## @@ -43,6 +43,29 @@ const LOGICAL_TYPE: &str = "logicalType"; struct SchemaToAvroSchema { schema: String, + /// Names of the `fixed` types defined so far. Avro allows one definition per + /// name, and the visitor reaches fields in the order they're serialized. + defined_names: HashSet<Name>, +} + +impl SchemaToAvroSchema { + /// Returns a reference to `schema` by name if a type with its name was + /// already defined, and `schema` itself otherwise. + fn define_once(&mut self, schema: AvroSchema) -> AvroSchema { + let name = match &schema { + AvroSchema::Fixed(FixedSchema { name, .. }) + | AvroSchema::Decimal(DecimalSchema { + inner: InnerDecimalSchema::Fixed(FixedSchema { name, .. }), + .. + }) => name.clone(), + _ => return schema, + }; + if self.defined_names.insert(name.clone()) { + schema + } else { + AvroSchema::Ref { name } Review Comment: `define_once` now emits `AvroSchema::Ref` for a repeated fixed/decimal, but `AvroSchemaToSchema` has no `Ref` arm, so an iceberg→avro→iceberg round-trip on two same-typed decimal or fixed columns hits the unsupported-primitive path. Correctness here also leans on the visitor reaching names in serialization order, which no test pins. A round-trip test with two `decimal(10,2)` fields — one nested in a map value or list element — would cover both; otherwise resolve `Ref` in the visitor or document it as write-only. ########## crates/iceberg/src/spec/manifest_list/writer.rs: ########## @@ -149,7 +149,8 @@ impl ManifestListWriter { FormatVersion::V2 => &MANIFEST_LIST_AVRO_SCHEMA_V2, FormatVersion::V3 => &MANIFEST_LIST_AVRO_SCHEMA_V3, }; - let mut avro_writer = Writer::new(avro_schema, Vec::new()); + let mut avro_writer = Writer::new(avro_schema, Vec::new()) + .expect("Manifest list Avro schemas have no named references to resolve."); Review Comment: `Writer::new` is fallible now, and this `expect` adds a panic path in library code — the repo convention here is to not panic, and `data_file.rs` / `manifest/writer.rs` already propagate the same call with `?`. The surrounding fn returns `Result`, so I'd just `?` it with an error map for consistency. -- 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]
