diff --git a/crates/catalog/glue/src/schema.rs b/crates/catalog/glue/src/schema.rs index 59d51f3dc1..747435611b 100644 --- a/crates/catalog/glue/src/schema.rs +++ b/crates/catalog/glue/src/schema.rs @@ -149,6 +149,12 @@ impl SchemaVisitor for GlueSchemaBuilder { fn primitive(&mut self, p: &PrimitiveType) -> Result { let glue_type = match p { + PrimitiveType::Unknown => { + return Err(Error::new( + ErrorKind::FeatureUnsupported, + format!("Conversion from {p:?} is not supported"), + )); + } PrimitiveType::Boolean => "boolean".to_string(), PrimitiveType::Int => "int".to_string(), PrimitiveType::Long => "bigint".to_string(), diff --git a/crates/catalog/hms/src/schema.rs b/crates/catalog/hms/src/schema.rs index e21a75f976..f04ee11c17 100644 --- a/crates/catalog/hms/src/schema.rs +++ b/crates/catalog/hms/src/schema.rs @@ -100,6 +100,12 @@ impl SchemaVisitor for HiveSchemaBuilder { fn primitive(&mut self, p: &PrimitiveType) -> Result { let hive_type = match p { + PrimitiveType::Unknown => { + return Err(Error::new( + ErrorKind::FeatureUnsupported, + format!("Conversion from {p:?} is not supported"), + )); + } PrimitiveType::Boolean => "boolean".to_string(), PrimitiveType::Int => "int".to_string(), PrimitiveType::Long => "bigint".to_string(), diff --git a/crates/iceberg/public-api.txt b/crates/iceberg/public-api.txt index 3d9abf2275..18e03cb257 100644 --- a/crates/iceberg/public-api.txt +++ b/crates/iceberg/public-api.txt @@ -1655,6 +1655,7 @@ pub iceberg::spec::PrimitiveType::Timestamp pub iceberg::spec::PrimitiveType::TimestampNs pub iceberg::spec::PrimitiveType::Timestamptz pub iceberg::spec::PrimitiveType::TimestamptzNs +pub iceberg::spec::PrimitiveType::Unknown pub iceberg::spec::PrimitiveType::Uuid impl iceberg::spec::PrimitiveType pub fn iceberg::spec::PrimitiveType::compatible(&self, literal: &iceberg::spec::PrimitiveLiteral) -> bool diff --git a/crates/iceberg/src/arrow/schema.rs b/crates/iceberg/src/arrow/schema.rs index dd26824c3a..d1912b159f 100644 --- a/crates/iceberg/src/arrow/schema.rs +++ b/crates/iceberg/src/arrow/schema.rs @@ -191,6 +191,7 @@ fn visit_type(r#type: &DataType, visitor: &mut V) -> Resu | DataType::Utf8 | DataType::LargeUtf8 | DataType::Utf8View + | DataType::Null | DataType::Binary | DataType::LargeBinary | DataType::BinaryView @@ -509,6 +510,7 @@ impl ArrowSchemaVisitor for ArrowSchemaConverter { fn primitive(&mut self, p: &DataType) -> Result { match p { + DataType::Null => Ok(Type::Primitive(PrimitiveType::Unknown)), DataType::Boolean => Ok(Type::Primitive(PrimitiveType::Boolean)), DataType::Int8 | DataType::Int16 | DataType::Int32 => { Ok(Type::Primitive(PrimitiveType::Int)) @@ -697,6 +699,7 @@ impl SchemaVisitor for ToArrowSchemaConverter { fn primitive(&mut self, p: &PrimitiveType) -> Result { match p { + PrimitiveType::Unknown => Ok(ArrowSchemaOrFieldOrType::Type(DataType::Null)), PrimitiveType::Boolean => Ok(ArrowSchemaOrFieldOrType::Type(DataType::Boolean)), PrimitiveType::Int => Ok(ArrowSchemaOrFieldOrType::Type(DataType::Int32)), PrimitiveType::Long => Ok(ArrowSchemaOrFieldOrType::Type(DataType::Int64)), @@ -1209,6 +1212,7 @@ pub(crate) fn primitive_type_to_arrow_type_with_ree(primitive_type: &PrimitiveTy }; match primitive_type { + PrimitiveType::Unknown => make_ree(DataType::Null), PrimitiveType::Boolean => make_ree(DataType::Boolean), PrimitiveType::Int => make_ree(DataType::Int32), PrimitiveType::Long => make_ree(DataType::Int64), @@ -2227,6 +2231,13 @@ mod tests { assert_eq!(iceberg_type, arrow_type_to_type(&arrow_type).unwrap()); } + { + let arrow_type = DataType::Null; + let iceberg_type = Type::Primitive(PrimitiveType::Unknown); + assert_eq!(arrow_type, type_to_arrow_type(&iceberg_type).unwrap()); + assert_eq!(iceberg_type, arrow_type_to_type(&arrow_type).unwrap()); + } + // test struct type { // no metadata will cause error diff --git a/crates/iceberg/src/arrow/value.rs b/crates/iceberg/src/arrow/value.rs index d22e565a4a..7abfe94b1e 100644 --- a/crates/iceberg/src/arrow/value.rs +++ b/crates/iceberg/src/arrow/value.rs @@ -206,6 +206,7 @@ impl SchemaWithPartnerVisitor for ArrowArrayToIcebergStructConverter { fn primitive(&mut self, p: &PrimitiveType, partner: &ArrayRef) -> Result>> { match p { + PrimitiveType::Unknown => Ok(vec![None; partner.len()]), PrimitiveType::Boolean => { let array = partner .as_any() @@ -636,6 +637,7 @@ pub(crate) fn create_primitive_array_single_element( prim_lit: &Option, ) -> Result { match (data_type, prim_lit) { + (DataType::Null, None) => Ok(Arc::new(arrow_array::NullArray::new(1))), (DataType::Boolean, Some(PrimitiveLiteral::Boolean(v))) => { Ok(Arc::new(BooleanArray::from(vec![*v]))) } @@ -943,7 +945,7 @@ pub(crate) fn create_primitive_array_repeated( Some(NullBuffer::new_null(num_rows)), )) } - (DataType::Null, _) => Arc::new(arrow_array::NullArray::new(num_rows)), + (DataType::Null, None) => Arc::new(arrow_array::NullArray::new(num_rows)), // --- Catch-all null arm: use arrow-rs new_null_array for any remaining DataType --- (dt, None) => new_null_array(dt, num_rows), @@ -1866,6 +1868,26 @@ mod test { } } + #[test] + fn test_create_null_array_rejects_non_null_literal() { + let literal = Some(PrimitiveLiteral::Int(1)); + + assert!(create_primitive_array_single_element(&DataType::Null, &literal).is_err()); + assert!(create_primitive_array_repeated(&DataType::Null, &literal, 2).is_err()); + assert_eq!( + create_primitive_array_single_element(&DataType::Null, &None) + .unwrap() + .len(), + 1 + ); + assert_eq!( + create_primitive_array_repeated(&DataType::Null, &None, 2) + .unwrap() + .len(), + 2 + ); + } + #[test] fn test_create_decimal_array_repeated_respects_precision() { // Ensure repeated arrays also respect target precision, not Arrow's default. diff --git a/crates/iceberg/src/avro/schema.rs b/crates/iceberg/src/avro/schema.rs index 528b89b7c1..107cf00414 100644 --- a/crates/iceberg/src/avro/schema.rs +++ b/crates/iceberg/src/avro/schema.rs @@ -74,7 +74,7 @@ impl SchemaVisitor for SchemaToAvroSchema { record.name = Name::from(format!("r{}", field.id).as_str()); } - if !field.required { + if !field.required && !is_avro_null(&field_schema) { field_schema = avro_optional(field_schema)?; } @@ -126,7 +126,7 @@ impl SchemaVisitor for SchemaToAvroSchema { record.name = Name::from(format!("r{}", list.element_field.id).as_str()); } - if !list.element_field.required { + if !list.element_field.required && !is_avro_null(&field_schema) { field_schema = avro_optional(field_schema)?; } @@ -147,7 +147,7 @@ impl SchemaVisitor for SchemaToAvroSchema { ) -> Result { let key_field_schema = key_value.unwrap_left(); let mut value_field_schema = value.unwrap_left(); - if !map.value_field.required { + if !map.value_field.required && !is_avro_null(&value_field_schema) { value_field_schema = avro_optional(value_field_schema)?; } @@ -222,6 +222,7 @@ impl SchemaVisitor for SchemaToAvroSchema { fn primitive(&mut self, p: &PrimitiveType) -> Result { let avro_schema = match p { + PrimitiveType::Unknown => AvroSchema::Null, PrimitiveType::Boolean => AvroSchema::Boolean, PrimitiveType::Int => AvroSchema::Int, PrimitiveType::Long => AvroSchema::Long, @@ -311,6 +312,10 @@ pub(crate) fn avro_decimal_schema(precision: usize, scale: usize) -> Result Result { + if is_avro_null(&avro_schema) { + return Ok(AvroSchema::Null); + } + Ok(AvroSchema::Union(UnionSchema::new(vec![ AvroSchema::Null, avro_schema, @@ -324,6 +329,14 @@ fn is_avro_optional(avro_schema: &AvroSchema) -> bool { } } +fn is_avro_null(avro_schema: &AvroSchema) -> bool { + matches!(avro_schema, AvroSchema::Null) +} + +fn unwrap_or_unknown(field_type: Option) -> Type { + field_type.unwrap_or(Type::Primitive(PrimitiveType::Unknown)) +} + /// Post order avro schema visitor. pub(crate) trait AvroSchemaVisitor { type T; @@ -447,10 +460,10 @@ impl AvroSchemaVisitor for AvroSchemaToSchema { let field_id = Self::get_element_id_from_attributes(&avro_field.custom_attributes, FIELD_ID_PROP)?; - let optional = is_avro_optional(&avro_field.schema); + let optional = is_avro_optional(&avro_field.schema) || is_avro_null(&avro_field.schema); - let mut field = - NestedField::new(field_id, &avro_field.name, field_type.unwrap(), !optional); + let field_type = unwrap_or_unknown(field_type); + let mut field = NestedField::new(field_id, &avro_field.name, field_type, !optional); if let Some(doc) = &avro_field.doc { field = field.with_doc(doc); @@ -482,7 +495,7 @@ impl AvroSchemaVisitor for AvroSchemaToSchema { } if options.len() == 1 { - Ok(Some(options.remove(0).unwrap())) + Ok(Some(unwrap_or_unknown(options.remove(0)))) } else { Ok(Some(options.remove(1).unwrap())) } @@ -490,10 +503,11 @@ impl AvroSchemaVisitor for AvroSchemaToSchema { fn array(&mut self, array: &ArraySchema, item: Option) -> Result { let element_field_id = Self::get_element_id_from_attributes(&array.attributes, ELEMENT_ID)?; + let item = unwrap_or_unknown(item); let element_field = NestedField::list_element( element_field_id, - item.unwrap(), - !is_avro_optional(&array.items), + item, + !is_avro_optional(&array.items) && !is_avro_null(&array.items), ) .into(); Ok(Some(Type::List(ListType { element_field }))) @@ -504,10 +518,11 @@ impl AvroSchemaVisitor for AvroSchemaToSchema { let key_field = NestedField::map_key_element(key_field_id, Type::Primitive(PrimitiveType::String)); let value_field_id = Self::get_element_id_from_attributes(&map.attributes, VALUE_ID)?; + let value = unwrap_or_unknown(value); let value_field = NestedField::map_value_element( value_field_id, - value.unwrap(), - !is_avro_optional(&map.types), + value, + !is_avro_optional(&map.types) && !is_avro_null(&map.types), ); Ok(Some(Type::Map(MapType { key_field: key_field.into(), @@ -557,12 +572,7 @@ impl AvroSchemaVisitor for AvroSchemaToSchema { "Can't convert avro map schema, missing key schema.", ) })?; - let value = value.ok_or_else(|| { - Error::new( - ErrorKind::DataInvalid, - "Can't convert avro map schema, missing value schema.", - ) - })?; + let value = unwrap_or_unknown(value); let key_id = Self::get_element_id_from_attributes( &array.fields[0].custom_attributes, FIELD_ID_PROP, @@ -575,7 +585,7 @@ impl AvroSchemaVisitor for AvroSchemaToSchema { let value_field = NestedField::map_value_element( value_id, value, - !is_avro_optional(&array.fields[1].schema), + !is_avro_optional(&array.fields[1].schema) && !is_avro_null(&array.fields[1].schema), ); Ok(Some(Type::Map(MapType { key_field: key_field.into(), @@ -659,6 +669,25 @@ mod tests { assert_eq!(iceberg_schema, converted_avro_converted_iceberg_schema); } + #[test] + fn test_unknown_type_schema_conversion() { + let schema = Schema::builder() + .with_fields(vec![ + NestedField::optional(1, "empty", PrimitiveType::Unknown.into()).into(), + ]) + .build() + .unwrap(); + + let avro_schema = schema_to_avro_schema("table", &schema).unwrap(); + let AvroSchema::Record(record) = &avro_schema else { + panic!("expected avro record schema"); + }; + assert!(is_avro_null(&record.fields[0].schema)); + assert_eq!(record.fields[0].default, Some(Value::Null)); + + assert_eq!(schema, avro_schema_to_schema(&avro_schema).unwrap()); + } + #[test] fn test_manifest_file_v1_schema() { let fields = vec![ diff --git a/crates/iceberg/src/spec/datatypes.rs b/crates/iceberg/src/spec/datatypes.rs index 79c48c1318..1f3f9a403a 100644 --- a/crates/iceberg/src/spec/datatypes.rs +++ b/crates/iceberg/src/spec/datatypes.rs @@ -19,7 +19,6 @@ * Data Types */ use std::collections::HashMap; -use std::convert::identity; use std::fmt; use std::ops::Index; use std::sync::{Arc, OnceLock}; @@ -134,8 +133,8 @@ impl Type { /// Minimum [`FormatVersion`] required to support this type, **without** taking /// nested field types into account. /// - /// `TimestampNs` / `TimestamptzNs` / `Variant` require [`FormatVersion::V3`]; every - /// other type is valid from [`FormatVersion::V1`]. Mirrors Java's + /// `Unknown` / `TimestampNs` / `TimestamptzNs` / `Variant` require + /// [`FormatVersion::V3`]; every other type is valid from [`FormatVersion::V1`]. Mirrors Java's /// `Schema.MIN_FORMAT_VERSIONS` (a shallow lookup keyed by type id), so it /// intentionally does not recurse: callers needing the floor for a whole schema /// iterate its flattened fields (see [`Schema::calc_min_compatible_format`]). @@ -143,7 +142,9 @@ impl Type { /// [`Schema::calc_min_compatible_format`]: crate::spec::Schema::calc_min_compatible_format pub(crate) fn min_format_version(&self) -> FormatVersion { match self { - Type::Primitive(PrimitiveType::TimestampNs | PrimitiveType::TimestamptzNs) + Type::Primitive( + PrimitiveType::Unknown | PrimitiveType::TimestampNs | PrimitiveType::TimestamptzNs, + ) | Type::Variant(_) => FormatVersion::V3, _ => FormatVersion::V1, } @@ -274,6 +275,8 @@ pub enum PrimitiveType { Fixed(u64), /// Arbitrary-length byte array. Binary, + /// Default / null column type used when a more specific type is not known. + Unknown, } impl PrimitiveType { @@ -391,6 +394,7 @@ where S: Serializer { impl fmt::Display for PrimitiveType { fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result { match self { + PrimitiveType::Unknown => write!(f, "unknown"), PrimitiveType::Boolean => write!(f, "boolean"), PrimitiveType::Int => write!(f, "int"), PrimitiveType::Long => write!(f, "long"), @@ -554,7 +558,7 @@ impl fmt::Display for StructType { } #[derive(Debug, PartialEq, Serialize, Deserialize, Eq, Clone)] -#[serde(from = "SerdeNestedField", into = "SerdeNestedField")] +#[serde(try_from = "SerdeNestedField", into = "SerdeNestedField")] /// A struct is a tuple of typed values. Each field in the tuple is named and has an integer id that is unique in the table schema. /// Each field can be either optional or required, meaning that values can (or cannot) be null. Fields may be any type. /// Fields may have an optional comment or doc string. Fields can have default values. @@ -591,25 +595,115 @@ struct SerdeNestedField { pub write_default: Option, } -impl From for NestedField { - fn from(value: SerdeNestedField) -> Self { - NestedField { +impl TryFrom for NestedField { + type Error = crate::Error; + + fn try_from(value: SerdeNestedField) -> Result { + fn validate_unknown_values(default: &JsonValue, field_type: &Type) -> Result<()> { + if default.is_null() { + return Ok(()); + } + + match (field_type, default) { + (Type::Primitive(PrimitiveType::Unknown), _) => { + ensure_data_valid!(false, "Unknown type only supports null default values",); + } + (Type::Struct(struct_type), JsonValue::Object(object)) => { + for field in struct_type.fields() { + if let Some(value) = object.get(&field.id.to_string()) { + validate_unknown_values(value, &field.field_type)?; + } + } + } + (Type::List(list_type), JsonValue::Array(values)) => { + for value in values { + validate_unknown_values(value, &list_type.element_field.field_type)?; + } + } + (Type::Map(map_type), JsonValue::Object(object)) => { + if let Some(JsonValue::Array(keys)) = object.get("keys") { + for key in keys { + validate_unknown_values(key, &map_type.key_field.field_type)?; + } + } + if let Some(JsonValue::Array(values)) = object.get("values") { + for value in values { + validate_unknown_values(value, &map_type.value_field.field_type)?; + } + } + } + _ => {} + } + + Ok(()) + } + + fn validate_unknown_is_optional( + required: bool, + field_type: &Type, + field_name: &str, + ) -> Result<()> { + match field_type { + Type::Primitive(PrimitiveType::Unknown) => { + ensure_data_valid!( + !required, + "Field {} cannot be required because unknown type must be optional", + field_name + ); + } + Type::Struct(struct_type) => { + for field in struct_type.fields() { + validate_unknown_is_optional( + field.required, + &field.field_type, + &field.name, + )?; + } + } + Type::List(list_type) => { + let field = &list_type.element_field; + validate_unknown_is_optional(field.required, &field.field_type, &field.name)?; + } + Type::Map(map_type) => { + for field in [&map_type.key_field, &map_type.value_field] { + validate_unknown_is_optional( + field.required, + &field.field_type, + &field.name, + )?; + } + } + Type::Primitive(_) | Type::Variant(_) => {} + } + + Ok(()) + } + + fn parse_default(default: Option, field_type: &Type) -> Result> { + let Some(default) = default else { + return Ok(None); + }; + validate_unknown_values(&default, field_type)?; + match Literal::try_from_json(default, field_type) { + Ok(default) => Ok(default), + Err(_) => Ok(None), + } + } + + validate_unknown_is_optional(value.required, &value.field_type, &value.name)?; + + let initial_default = parse_default(value.initial_default, &value.field_type)?; + let write_default = parse_default(value.write_default, &value.field_type)?; + + Ok(NestedField { id: value.id, name: value.name, required: value.required, - initial_default: value.initial_default.and_then(|x| { - Literal::try_from_json(x, &value.field_type) - .ok() - .and_then(identity) - }), - write_default: value.write_default.and_then(|x| { - Literal::try_from_json(x, &value.field_type) - .ok() - .and_then(identity) - }), + initial_default, + write_default, field_type: value.field_type, doc: value.doc, - } + }) } } @@ -936,6 +1030,7 @@ mod tests { { "type": "struct", "fields": [ + {"id": 17, "name": "unknown_field", "required": false, "type": "unknown"}, {"id": 1, "name": "bool_field", "required": true, "type": "boolean"}, {"id": 2, "name": "int_field", "required": true, "type": "int"}, {"id": 3, "name": "long_field", "required": true, "type": "long"}, @@ -960,6 +1055,12 @@ mod tests { record, Type::Struct(StructType { fields: vec![ + NestedField::optional( + 17, + "unknown_field", + Type::Primitive(PrimitiveType::Unknown), + ) + .into(), NestedField::required(1, "bool_field", Type::Primitive(PrimitiveType::Boolean)) .into(), NestedField::required(2, "int_field", Type::Primitive(PrimitiveType::Int)) @@ -1341,6 +1442,71 @@ mod tests { for (ty, literal) in pairs { assert!(ty.compatible(&literal)); } + + assert!(!PrimitiveType::Unknown.compatible(&PrimitiveLiteral::Int(1))); + } + + #[test] + fn unknown_nested_field_deserialization_rejects_required() { + let json = r#"{"id":1,"name":"empty","required":true,"type":"unknown"}"#; + + let error = serde_json::from_str::(json).unwrap_err(); + assert!( + error.to_string().contains("unknown type must be optional"), + "unexpected error: {error}" + ); + } + + #[test] + fn unknown_nested_field_deserialization_rejects_required_in_containers() { + let cases = [ + ( + "list element", + serde_json::json!({ + "type": "list", + "element-id": 2, + "element-required": true, + "element": "unknown" + }), + ), + ( + "map key", + serde_json::json!({ + "type": "map", + "key-id": 2, + "key": "unknown", + "value-id": 3, + "value-required": false, + "value": "string" + }), + ), + ( + "map value", + serde_json::json!({ + "type": "map", + "key-id": 2, + "key": "string", + "value-id": 3, + "value-required": true, + "value": "unknown" + }), + ), + ]; + + for (name, field_type) in cases { + let field = serde_json::json!({ + "id": 1, + "name": name, + "required": false, + "type": field_type + }); + + let error = serde_json::from_value::(field).unwrap_err(); + assert!( + error.to_string().contains("unknown type must be optional"), + "unexpected error for {name}: {error}" + ); + } } #[test] diff --git a/crates/iceberg/src/spec/schema/mod.rs b/crates/iceberg/src/spec/schema/mod.rs index 652f98b649..ca51f0e6c3 100644 --- a/crates/iceberg/src/spec/schema/mod.rs +++ b/crates/iceberg/src/spec/schema/mod.rs @@ -39,11 +39,11 @@ pub use self::prune_columns::prune_columns; use super::NestedField; use crate::error::Result; use crate::expr::accessor::StructAccessor; -use crate::spec::FormatVersion; use crate::spec::datatypes::{ LIST_FIELD_NAME, ListType, MAP_KEY_FIELD_NAME, MAP_VALUE_FIELD_NAME, MapType, NestedFieldRef, PrimitiveType, StructType, Type, }; +use crate::spec::{FormatVersion, Literal}; use crate::{Error, ErrorKind, ensure_data_valid}; /// Type alias for schema id. @@ -133,6 +133,10 @@ impl SchemaBuilder { /// Builds the schema. pub fn build(self) -> Result { + for field in &self.fields { + Self::validate_unknown_type_field(field)?; + } + let field_id_to_accessor = self.build_accessors(); let r#struct = StructType::new(self.fields); @@ -190,6 +194,77 @@ impl SchemaBuilder { Ok(schema) } + fn validate_unknown_type_field(field: &NestedFieldRef) -> Result<()> { + ensure_data_valid!( + !field + .initial_default + .iter() + .chain(field.write_default.iter()) + .any(|default| { + Self::default_contains_non_null_unknown(default, &field.field_type) + }), + "Field {} cannot have non-null defaults because unknown type requires null defaults", + field.name + ); + + match field.field_type.as_ref() { + Type::Primitive(PrimitiveType::Unknown) => { + ensure_data_valid!( + !field.required, + "Field {} cannot be required because unknown type must be optional", + field.name + ); + } + Type::Struct(struct_type) => { + for nested_field in struct_type.fields() { + Self::validate_unknown_type_field(nested_field)?; + } + } + Type::List(list_type) => { + Self::validate_unknown_type_field(&list_type.element_field)?; + } + Type::Map(map_type) => { + Self::validate_unknown_type_field(&map_type.key_field)?; + Self::validate_unknown_type_field(&map_type.value_field)?; + } + Type::Primitive(_) | Type::Variant(_) => {} + } + + Ok(()) + } + + fn default_contains_non_null_unknown(default: &Literal, field_type: &Type) -> bool { + match (default, field_type) { + (_, Type::Primitive(PrimitiveType::Unknown)) => true, + (Literal::Struct(value), Type::Struct(struct_type)) => value + .iter() + .zip(struct_type.fields()) + .any(|(value, field)| { + value.is_some_and(|value| { + Self::default_contains_non_null_unknown(value, &field.field_type) + }) + }), + (Literal::List(values), Type::List(list_type)) => values.iter().any(|value| { + value.as_ref().is_some_and(|value| { + Self::default_contains_non_null_unknown( + value, + &list_type.element_field.field_type, + ) + }) + }), + (Literal::Map(map), Type::Map(map_type)) => map.iter().any(|(key, value)| { + Self::default_contains_non_null_unknown(key, &map_type.key_field.field_type) + || value.as_ref().is_some_and(|value| { + Self::default_contains_non_null_unknown( + value, + &map_type.value_field.field_type, + ) + }) + }), + _ => false, + } + } + fn build_accessors(&self) -> HashMap> { let mut map = HashMap::new(); @@ -647,6 +722,17 @@ mod tests { ]); assert_eq!(variant.calc_min_compatible_format(), FormatVersion::V3); + // Unknown is a v3-only primitive type. + let unknown = schema_with(vec![ + NestedField::optional(1, "u", Primitive(PrimitiveType::Unknown)).into(), + ]); + assert_eq!(unknown.calc_min_compatible_format(), FormatVersion::V3); + assert!( + unknown + .check_format_compatibility(FormatVersion::V2) + .is_err() + ); + // A v3-only type nested inside a list inside a struct → V3 (flattened fields). let nested = schema_with(vec![ NestedField::required( @@ -1475,4 +1561,245 @@ table { .is_err() ); } + + #[test] + fn test_unknown_type_deserialization_rejects_non_null_default() { + let field_json = serde_json::json!({ + "id": 1, + "name": "empty", + "required": false, + "type": "unknown", + "initial-default": 1 + }); + + let error = serde_json::from_value::(field_json.clone()).unwrap_err(); + assert!( + error + .to_string() + .contains("Unknown type only supports null default values"), + "unexpected error: {error}" + ); + + let schema_json = serde_json::json!({ + "type": "struct", + "schema-id": 1, + "fields": [field_json] + }); + assert!(serde_json::from_value::(schema_json).is_err()); + } + + #[test] + fn test_unknown_type_deserialization_accepts_null_defaults() { + let schema_json = serde_json::json!({ + "type": "struct", + "schema-id": 1, + "fields": [ + { + "id": 1, + "name": "empty", + "required": false, + "type": "unknown", + "initial-default": null, + "write-default": null + } + ] + }); + + serde_json::from_value::(schema_json).unwrap(); + } + + #[test] + fn test_unknown_type_must_be_optional_with_null_defaults() { + assert!( + Schema::builder() + .with_schema_id(1) + .with_fields(vec![ + NestedField::optional(1, "empty", Primitive(PrimitiveType::Unknown)).into() + ]) + .build() + .is_ok() + ); + + let required_error = Schema::builder() + .with_schema_id(1) + .with_fields(vec![ + NestedField::required(1, "empty", Primitive(PrimitiveType::Unknown)).into(), + ]) + .build() + .unwrap_err(); + assert!( + required_error + .message() + .contains("unknown type must be optional") + ); + + let default_error = Schema::builder() + .with_schema_id(1) + .with_fields(vec![ + NestedField::optional(1, "empty", Primitive(PrimitiveType::Unknown)) + .with_initial_default(Literal::int(1)) + .into(), + ]) + .build() + .unwrap_err(); + assert!( + default_error + .message() + .contains("unknown type requires null defaults") + ); + } + + #[test] + fn test_unknown_type_rejects_non_null_container_defaults() { + let cases = [ + ( + "struct", + Struct(StructType::new(vec![ + NestedField::optional(2, "empty", Primitive(PrimitiveType::Unknown)).into(), + ])), + Literal::Struct(crate::spec::Struct::from_iter([Some(Literal::int(1))])), + ), + ( + "list", + List(ListType::new( + NestedField::list_element(2, Primitive(PrimitiveType::Unknown), false).into(), + )), + Literal::List(vec![Some(Literal::int(1))]), + ), + ( + "map", + Map(MapType::optional( + 2, + Primitive(PrimitiveType::String), + 3, + Primitive(PrimitiveType::Unknown), + )), + Literal::Map(MapValue::from([( + Literal::string("key"), + Some(Literal::int(1)), + )])), + ), + ]; + + for (name, field_type, default) in cases { + let error = Schema::builder() + .with_schema_id(1) + .with_fields(vec![ + NestedField::optional(1, name, field_type) + .with_initial_default(default) + .into(), + ]) + .build() + .unwrap_err(); + assert!( + error + .message() + .contains("unknown type requires null defaults"), + "unexpected error for {name}: {error}" + ); + } + } + + #[test] + fn test_unknown_type_accepts_null_container_defaults() { + let cases = [ + ( + "struct", + Struct(StructType::new(vec![ + NestedField::optional(2, "empty", Primitive(PrimitiveType::Unknown)).into(), + ])), + Literal::Struct(crate::spec::Struct::from_iter([None])), + ), + ( + "list", + List(ListType::new( + NestedField::list_element(2, Primitive(PrimitiveType::Unknown), false).into(), + )), + Literal::List(vec![None]), + ), + ( + "map", + Map(MapType::optional( + 2, + Primitive(PrimitiveType::String), + 3, + Primitive(PrimitiveType::Unknown), + )), + Literal::Map(MapValue::from([(Literal::string("key"), None)])), + ), + ]; + + for (name, field_type, default) in cases { + let schema = Schema::builder() + .with_schema_id(1) + .with_fields(vec![ + NestedField::optional(1, name, field_type) + .with_write_default(default) + .into(), + ]) + .build() + .unwrap(); + serde_json::to_value(schema).unwrap(); + } + } + + #[test] + fn test_unknown_type_deserialization_rejects_non_null_container_defaults() { + let cases = [ + ( + "struct", + serde_json::json!({ + "type": "struct", + "fields": [{ + "id": 2, + "name": "empty", + "required": false, + "type": "unknown" + }] + }), + serde_json::json!({"2": 1}), + ), + ( + "list", + serde_json::json!({ + "type": "list", + "element-id": 2, + "element-required": false, + "element": "unknown" + }), + serde_json::json!([1]), + ), + ( + "map", + serde_json::json!({ + "type": "map", + "key-id": 2, + "key": "string", + "value-id": 3, + "value-required": false, + "value": "unknown" + }), + serde_json::json!({"keys": ["key"], "values": [1]}), + ), + ]; + + for (name, field_type, default) in cases { + let schema_json = serde_json::json!({ + "type": "struct", + "schema-id": 1, + "fields": [{ + "id": 1, + "name": name, + "required": false, + "type": field_type, + "initial-default": default + }] + }); + + assert!( + serde_json::from_value::(schema_json).is_err(), + "non-null unknown default in {name} should be rejected" + ); + } + } } diff --git a/crates/iceberg/src/spec/values/datum.rs b/crates/iceberg/src/spec/values/datum.rs index f170a09df5..593544b94b 100644 --- a/crates/iceberg/src/spec/values/datum.rs +++ b/crates/iceberg/src/spec/values/datum.rs @@ -368,6 +368,12 @@ impl Datum { /// See [this spec](https://iceberg.apache.org/spec/#binary-single-value-serialization) for reference. pub fn try_from_bytes(bytes: &[u8], data_type: PrimitiveType) -> Result { let literal = match data_type { + PrimitiveType::Unknown => { + return Err(Error::new( + ErrorKind::FeatureUnsupported, + "Cannot create datum for unknown type from bytes", + )); + } PrimitiveType::Boolean => { if bytes.len() == 1 && bytes[0] == 0u8 { PrimitiveLiteral::Boolean(false) diff --git a/crates/iceberg/src/spec/values/map.rs b/crates/iceberg/src/spec/values/map.rs index e0f75205f0..81adb8da23 100644 --- a/crates/iceberg/src/spec/values/map.rs +++ b/crates/iceberg/src/spec/values/map.rs @@ -87,6 +87,10 @@ impl Map { self.index.get(key).map(|index| &self.pair[*index].1) } + pub(crate) fn iter(&self) -> impl Iterator)> { + self.pair.iter().map(|(key, value)| (key, value)) + } + /// The order of map is matter, so this method used to compare two maps has same key-value pairs without considering the order. pub fn has_same_content(&self, other: &Map) -> bool { if self.len() != other.len() {