Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions crates/catalog/glue/src/schema.rs
Original file line number Diff line number Diff line change
Expand Up @@ -149,6 +149,12 @@ impl SchemaVisitor for GlueSchemaBuilder {

fn primitive(&mut self, p: &PrimitiveType) -> Result<Self::T> {
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(),
Expand Down
6 changes: 6 additions & 0 deletions crates/catalog/hms/src/schema.rs
Original file line number Diff line number Diff line change
Expand Up @@ -100,6 +100,12 @@ impl SchemaVisitor for HiveSchemaBuilder {

fn primitive(&mut self, p: &PrimitiveType) -> Result<String> {
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(),
Expand Down
1 change: 1 addition & 0 deletions crates/iceberg/public-api.txt
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
93 changes: 79 additions & 14 deletions crates/iceberg/src/arrow/nan_val_cnt_visitor.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,15 +21,17 @@ use std::collections::HashMap;
use std::collections::hash_map::Entry;
use std::sync::Arc;

use arrow_array::{ArrayRef, Float32Array, Float64Array, RecordBatch, StructArray};
use arrow_array::{
ArrayRef, Float32Array, Float64Array, ListArray, MapArray, RecordBatch, StructArray,
};
use arrow_schema::DataType;

use crate::Result;
use crate::arrow::{ArrowArrayAccessor, FieldMatchMode};
use crate::arrow::FieldMatchMode;
use crate::spec::{
ListType, MapType, NestedFieldRef, PrimitiveType, Schema, SchemaRef, SchemaWithPartnerVisitor,
StructType, VariantType, visit_struct_with_partner,
StructType, Type, VariantType,
};
use crate::{Error, ErrorKind, Result};

macro_rules! cast_and_update_cnt_map {
($t:ty, $col:ident, $self:ident, $field_id:ident) => {
Expand Down Expand Up @@ -152,6 +154,78 @@ impl SchemaWithPartnerVisitor<ArrayRef> for NanValueCountVisitor {
}

impl NanValueCountVisitor {
fn visit_field(&mut self, field: &NestedFieldRef, array: &ArrayRef) -> Result<()> {
let field_id = field.id;
count_float_nans!(array, self, field_id);

match field.field_type.as_ref() {
Type::Primitive(_) | Type::Variant(_) => Ok(()),
Type::Struct(struct_type) => self.visit_struct(struct_type, array),
Type::List(list_type) => {
let list_array = array.as_any().downcast_ref::<ListArray>().ok_or_else(|| {
Error::new(
ErrorKind::DataInvalid,
format!(
"Expected list array for field {}, got {}",
field.id,
array.data_type()
),
)
})?;
self.visit_field(&list_type.element_field, list_array.values())
}
Type::Map(map_type) => {
let map_array = array.as_any().downcast_ref::<MapArray>().ok_or_else(|| {
Error::new(
ErrorKind::DataInvalid,
format!(
"Expected map array for field {}, got {}",
field.id,
array.data_type()
),
)
})?;
self.visit_field(&map_type.key_field, map_array.keys())?;
self.visit_field(&map_type.value_field, map_array.values())
}
}
}

fn visit_struct(&mut self, struct_type: &StructType, array: &ArrayRef) -> Result<()> {
let struct_array = array
.as_any()
.downcast_ref::<StructArray>()
.ok_or_else(|| {
Error::new(
ErrorKind::DataInvalid,
format!("Expected struct array, got {}", array.data_type()),
)
})?;

for field in struct_type.fields() {
if matches!(
field.field_type.as_ref(),
Type::Primitive(PrimitiveType::Unknown)
) {
continue;
}

let field_position = struct_array
.fields()
.iter()
.position(|arrow_field| self.match_mode.match_field(arrow_field, field))
.ok_or_else(|| {
Error::new(
ErrorKind::DataInvalid,
format!("Field id {} not found in struct array", field.id),
)
})?;
self.visit_field(field, struct_array.column(field_position))?;
}

Ok(())
}

/// Creates new instance of NanValueCountVisitor
pub fn new() -> Self {
Self::new_with_match_mode(FieldMatchMode::Id)
Expand All @@ -167,17 +241,8 @@ impl NanValueCountVisitor {

/// Compute nan value counts in given schema and record batch
pub fn compute(&mut self, schema: SchemaRef, batch: RecordBatch) -> Result<()> {
let arrow_arr_partner_accessor = ArrowArrayAccessor::new_with_match_mode(self.match_mode);

let struct_arr = Arc::new(StructArray::from(batch)) as ArrayRef;
visit_struct_with_partner(
schema.as_struct(),
&struct_arr,
self,
&arrow_arr_partner_accessor,
)?;

Ok(())
self.visit_struct(schema.as_struct(), &struct_arr)
}
}

Expand Down
6 changes: 4 additions & 2 deletions crates/iceberg/src/arrow/reader/pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -180,7 +180,7 @@ impl FileScanTaskReader {
} else {
// Branch 3: No name mapping - use position-based fallback IDs
// Corresponds to Java's ParquetSchemaUtil.addFallbackIds()
add_fallback_field_ids_to_arrow_schema(arrow_metadata.schema())
add_fallback_field_ids_to_arrow_schema(arrow_metadata.schema(), &task.schema)?
};

let options = ArrowReaderOptions::new().with_schema(arrow_schema);
Expand Down Expand Up @@ -345,7 +345,8 @@ impl FileScanTaskReader {
// that come back from the file, such as type promotion, default column insertion,
// column re-ordering, partition constants, and virtual field addition (like _file)
let mut record_batch_transformer_builder =
RecordBatchTransformerBuilder::new(task.schema_ref(), task.project_field_ids());
RecordBatchTransformerBuilder::new(task.schema_ref(), task.project_field_ids())
.with_position_fallback(use_position_fallback);

// Add the _file metadata column if it's in the projected fields
if task.project_field_ids().contains(&RESERVED_FIELD_ID_FILE) {
Expand Down Expand Up @@ -518,6 +519,7 @@ impl FileScanTaskReader {
record_batch_stream_builder.parquet_schema(),
record_batch_stream_builder.schema(),
&predicate,
&task.schema,
use_position_fallback,
)?;

Expand Down
Loading
Loading