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]

Reply via email to