From 718e831aa6eb52a360df3232a943815072aebdbb Mon Sep 17 00:00:00 2001 From: Ed Seidl Date: Wed, 29 Apr 2026 17:24:46 -0700 Subject: [PATCH 1/4] allow UNKNOWN logical type annotation on any physical type --- parquet/src/arrow/array_reader/builder.rs | 36 +++++++++++---------- parquet/src/arrow/arrow_reader/mod.rs | 39 +++++++++++++++++++++++ parquet/src/arrow/schema/primitive.rs | 29 +++++++++-------- parquet/src/schema/types.rs | 2 +- 4 files changed, 76 insertions(+), 30 deletions(-) diff --git a/parquet/src/arrow/array_reader/builder.rs b/parquet/src/arrow/array_reader/builder.rs index d806b2147a09..0c575bf2e62f 100644 --- a/parquet/src/arrow/array_reader/builder.rs +++ b/parquet/src/arrow/array_reader/builder.rs @@ -422,6 +422,20 @@ impl<'a> ArrayReaderBuilder<'a> { let page_iterator = self.row_groups.column_chunks(col_idx)?; let arrow_type = Some(field.arrow_type.clone()); + // LogicalType::Unknown maps to DataType::Null. In the past it has been assumed + // that only INT32 can have this annotation, but this is not required by the Parquet + // specification. Since this can only annotate an entirely null column, the data type + // used for the NullArrayReader should be irrelevant. It's just needed to read the + // repetition and definition level data. + if matches!(arrow_type, Some(DataType::Null)) { + let reader = Box::new(NullArrayReader::::new( + page_iterator, + column_desc, + self.batch_size, + )?) as _; + return Ok(Some(reader)); + } + let reader = match physical_type { PhysicalType::BOOLEAN => Box::new(PrimitiveArrayReader::::new( page_iterator, @@ -429,22 +443,12 @@ impl<'a> ArrayReaderBuilder<'a> { arrow_type, self.batch_size, )?) as _, - PhysicalType::INT32 => { - if let Some(DataType::Null) = arrow_type { - Box::new(NullArrayReader::::new( - page_iterator, - column_desc, - self.batch_size, - )?) as _ - } else { - Box::new(PrimitiveArrayReader::::new( - page_iterator, - column_desc, - arrow_type, - self.batch_size, - )?) as _ - } - } + PhysicalType::INT32 => Box::new(PrimitiveArrayReader::::new( + page_iterator, + column_desc, + arrow_type, + self.batch_size, + )?) as _, PhysicalType::INT64 => Box::new(PrimitiveArrayReader::::new( page_iterator, column_desc, diff --git a/parquet/src/arrow/arrow_reader/mod.rs b/parquet/src/arrow/arrow_reader/mod.rs index f9ef4eac651f..fe64a0a989b5 100644 --- a/parquet/src/arrow/arrow_reader/mod.rs +++ b/parquet/src/arrow/arrow_reader/mod.rs @@ -3686,6 +3686,45 @@ pub(crate) mod tests { } } + // test that we can handle the UNKNOWN logical type annotation on any physical type + #[test] + fn test_unknown_logical_type() { + let message_type = "message uk { + OPTIONAL INT32 uki32 (UNKNOWN); + OPTIONAL INT64 uki64 (UNKNOWN); + OPTIONAL INT96 uki96 (UNKNOWN); + OPTIONAL BOOLEAN ukbool (UNKNOWN); + OPTIONAL FLOAT ukfloat (UNKNOWN); + OPTIONAL DOUBLE ukdbl (UNKNOWN); + OPTIONAL BYTE_ARRAY ukbytes (UNKNOWN); + OPTIONAL FIXED_LEN_BYTE_ARRAY(10) ukflba (UNKNOWN); + }"; + + let schema = Arc::new(parse_message_type(message_type).unwrap()); + let file = tempfile::tempfile().unwrap(); + + { + let mut writer = + SerializedFileWriter::new(file.try_clone().unwrap(), schema, Default::default()) + .unwrap(); + let schema_desc = writer.schema_descr_ptr().clone(); + + { + let mut row_group_writer = writer.next_row_group().unwrap(); + for _ in schema_desc.columns() { + let column_writer = row_group_writer.next_column().unwrap().unwrap(); + column_writer.close().unwrap(); + } + row_group_writer.close().unwrap(); + } + + writer.close().unwrap(); + } + + let _ = ParquetRecordBatchReader::try_new(file, 2).unwrap(); + // if we get this far we've succeeded + } + #[test] fn test_nested_nullability() { let message_type = "message nested { diff --git a/parquet/src/arrow/schema/primitive.rs b/parquet/src/arrow/schema/primitive.rs index 8959081bcb41..b440753cc83b 100644 --- a/parquet/src/arrow/schema/primitive.rs +++ b/parquet/src/arrow/schema/primitive.rs @@ -115,17 +115,22 @@ fn from_parquet(parquet_type: &Type) -> Result { scale, precision, .. - } => match physical_type { - PhysicalType::BOOLEAN => Ok(DataType::Boolean), - PhysicalType::INT32 => from_int32(basic_info, *scale, *precision), - PhysicalType::INT64 => from_int64(basic_info, *scale, *precision), - PhysicalType::INT96 => Ok(DataType::Timestamp(TimeUnit::Nanosecond, None)), - PhysicalType::FLOAT => Ok(DataType::Float32), - PhysicalType::DOUBLE => Ok(DataType::Float64), - PhysicalType::BYTE_ARRAY => from_byte_array(basic_info, *precision, *scale), - PhysicalType::FIXED_LEN_BYTE_ARRAY => { - from_fixed_len_byte_array(basic_info, *scale, *precision, *type_length) - } + } => match basic_info.logical_type_ref() { + // Any physical type can have the UNKNOWN logical type annotation. Check for that first. + // https://github.com/apache/parquet-format/blob/master/LogicalTypes.md#unknown-always-null + Some(&LogicalType::Unknown) => Ok(DataType::Null), + _ => match physical_type { + PhysicalType::BOOLEAN => Ok(DataType::Boolean), + PhysicalType::INT32 => from_int32(basic_info, *scale, *precision), + PhysicalType::INT64 => from_int64(basic_info, *scale, *precision), + PhysicalType::INT96 => Ok(DataType::Timestamp(TimeUnit::Nanosecond, None)), + PhysicalType::FLOAT => Ok(DataType::Float32), + PhysicalType::DOUBLE => Ok(DataType::Float64), + PhysicalType::BYTE_ARRAY => from_byte_array(basic_info, *precision, *scale), + PhysicalType::FIXED_LEN_BYTE_ARRAY => { + from_fixed_len_byte_array(basic_info, *scale, *precision, *type_length) + } + }, }, Type::GroupType { .. } => unreachable!(), } @@ -194,8 +199,6 @@ fn from_int32(info: &BasicTypeInfo, scale: i32, precision: i32) -> Result Ok(DataType::Null), (None, ConvertedType::UINT_8) => Ok(DataType::UInt8), (None, ConvertedType::UINT_16) => Ok(DataType::UInt16), (None, ConvertedType::UINT_32) => Ok(DataType::UInt32), diff --git a/parquet/src/schema/types.rs b/parquet/src/schema/types.rs index 2925557e7b86..2c63da74df47 100644 --- a/parquet/src/schema/types.rs +++ b/parquet/src/schema/types.rs @@ -398,7 +398,7 @@ impl<'a> PrimitiveTypeBuilder<'a> { (LogicalType::Integer { bit_width, .. }, PhysicalType::INT64) if *bit_width == 64 => {} // Null type - (LogicalType::Unknown, PhysicalType::INT32) => {} + (LogicalType::Unknown, _) => {} (LogicalType::String, PhysicalType::BYTE_ARRAY) => {} (LogicalType::Json, PhysicalType::BYTE_ARRAY) => {} (LogicalType::Bson, PhysicalType::BYTE_ARRAY) => {} From 9cda257a70e59aae35e365a35a6f1b72f747b08e Mon Sep 17 00:00:00 2001 From: Ed Seidl Date: Wed, 29 Apr 2026 21:27:45 -0700 Subject: [PATCH 2/4] write some data in test --- parquet/src/arrow/arrow_reader/mod.rs | 100 +++++++++++++++++++++----- 1 file changed, 81 insertions(+), 19 deletions(-) diff --git a/parquet/src/arrow/arrow_reader/mod.rs b/parquet/src/arrow/arrow_reader/mod.rs index fe64a0a989b5..2839b7c12bbc 100644 --- a/parquet/src/arrow/arrow_reader/mod.rs +++ b/parquet/src/arrow/arrow_reader/mod.rs @@ -1600,8 +1600,8 @@ pub(crate) mod tests { use crate::basic::{ConvertedType, Encoding, LogicalType, Repetition, Type as PhysicalType}; use crate::column::reader::decoder::REPETITION_LEVELS_BATCH_SIZE; use crate::data_type::{ - BoolType, ByteArray, ByteArrayType, DataType, FixedLenByteArray, FixedLenByteArrayType, - FloatType, Int32Type, Int64Type, Int96, Int96Type, + BoolType, ByteArray, ByteArrayType, DataType, DoubleType, FixedLenByteArray, + FixedLenByteArrayType, FloatType, Int32Type, Int64Type, Int96, Int96Type, }; use crate::errors::Result; use crate::file::metadata::{PageIndexPolicy, ParquetMetaData, ParquetStatisticsPolicy}; @@ -3703,26 +3703,88 @@ pub(crate) mod tests { let schema = Arc::new(parse_message_type(message_type).unwrap()); let file = tempfile::tempfile().unwrap(); - { - let mut writer = - SerializedFileWriter::new(file.try_clone().unwrap(), schema, Default::default()) - .unwrap(); - let schema_desc = writer.schema_descr_ptr().clone(); + let mut writer = + SerializedFileWriter::new(file.try_clone().unwrap(), schema, Default::default()) + .unwrap(); - { - let mut row_group_writer = writer.next_row_group().unwrap(); - for _ in schema_desc.columns() { - let column_writer = row_group_writer.next_column().unwrap().unwrap(); - column_writer.close().unwrap(); - } - row_group_writer.close().unwrap(); - } + let mut row_group_writer = writer.next_row_group().unwrap(); - writer.close().unwrap(); - } + // INT32 + let mut column_writer = row_group_writer.next_column().unwrap().unwrap(); + // write out a bunch of nulls + column_writer + .typed::() + .write_batch(&[], Some(&[0, 0, 0, 0]), None) + .unwrap(); + column_writer.close().unwrap(); + + // INT64 + let mut column_writer = row_group_writer.next_column().unwrap().unwrap(); + column_writer + .typed::() + .write_batch(&[], Some(&[0, 0, 0, 0]), None) + .unwrap(); + column_writer.close().unwrap(); - let _ = ParquetRecordBatchReader::try_new(file, 2).unwrap(); - // if we get this far we've succeeded + // INT96 + let mut column_writer = row_group_writer.next_column().unwrap().unwrap(); + column_writer + .typed::() + .write_batch(&[], Some(&[0, 0, 0, 0]), None) + .unwrap(); + column_writer.close().unwrap(); + + // BOOLEAN + let mut column_writer = row_group_writer.next_column().unwrap().unwrap(); + column_writer + .typed::() + .write_batch(&[], Some(&[0, 0, 0, 0]), None) + .unwrap(); + column_writer.close().unwrap(); + + // FLOAT + let mut column_writer = row_group_writer.next_column().unwrap().unwrap(); + column_writer + .typed::() + .write_batch(&[], Some(&[0, 0, 0, 0]), None) + .unwrap(); + column_writer.close().unwrap(); + + // DOUBLE + let mut column_writer = row_group_writer.next_column().unwrap().unwrap(); + column_writer + .typed::() + .write_batch(&[], Some(&[0, 0, 0, 0]), None) + .unwrap(); + column_writer.close().unwrap(); + + // BYTE_ARRAY + let mut column_writer = row_group_writer.next_column().unwrap().unwrap(); + column_writer + .typed::() + .write_batch(&[], Some(&[0, 0, 0, 0]), None) + .unwrap(); + column_writer.close().unwrap(); + + // FIXED_LEN_BYTE_ARRAY + let mut column_writer = row_group_writer.next_column().unwrap().unwrap(); + column_writer + .typed::() + .write_batch(&[], Some(&[0, 0, 0, 0]), None) + .unwrap(); + column_writer.close().unwrap(); + + row_group_writer.close().unwrap(); + + writer.close().unwrap(); + + let mut reader = ParquetRecordBatchReader::try_new(file, 4).unwrap(); + let batch = reader.next().unwrap().unwrap(); + + for col in batch.columns() { + assert_eq!(col.len(), 4); + assert_eq!(col.logical_null_count(), 4); + } } #[test] From fd36afb940ba20c7be20fc4a51390063b82eced9 Mon Sep 17 00:00:00 2001 From: Ed Seidl Date: Wed, 29 Apr 2026 21:39:39 -0700 Subject: [PATCH 3/4] reduce clutter --- parquet/src/arrow/arrow_reader/mod.rs | 69 ++++++++------------------- 1 file changed, 19 insertions(+), 50 deletions(-) diff --git a/parquet/src/arrow/arrow_reader/mod.rs b/parquet/src/arrow/arrow_reader/mod.rs index 2839b7c12bbc..16bce2fbc5e8 100644 --- a/parquet/src/arrow/arrow_reader/mod.rs +++ b/parquet/src/arrow/arrow_reader/mod.rs @@ -1606,7 +1606,7 @@ pub(crate) mod tests { use crate::errors::Result; use crate::file::metadata::{PageIndexPolicy, ParquetMetaData, ParquetStatisticsPolicy}; use crate::file::properties::{EnabledStatistics, WriterProperties, WriterVersion}; - use crate::file::writer::SerializedFileWriter; + use crate::file::writer::{SerializedFileWriter, SerializedRowGroupWriter}; use crate::schema::parser::parse_message_type; use crate::schema::types::{Type, TypePtr}; use crate::util::test_common::rand_gen::RandGen; @@ -3709,70 +3709,39 @@ pub(crate) mod tests { let mut row_group_writer = writer.next_row_group().unwrap(); + fn write_nulls(row_group_writer: &mut SerializedRowGroupWriter<'_, File>) { + let mut column_writer = row_group_writer.next_column().unwrap().unwrap(); + // write out a bunch of nulls + column_writer + .typed::() + .write_batch(&[], Some(&[0, 0, 0, 0]), None) + .unwrap(); + column_writer.close().unwrap(); + } + // INT32 - let mut column_writer = row_group_writer.next_column().unwrap().unwrap(); - // write out a bunch of nulls - column_writer - .typed::() - .write_batch(&[], Some(&[0, 0, 0, 0]), None) - .unwrap(); - column_writer.close().unwrap(); + write_nulls::(&mut row_group_writer); // INT64 - let mut column_writer = row_group_writer.next_column().unwrap().unwrap(); - column_writer - .typed::() - .write_batch(&[], Some(&[0, 0, 0, 0]), None) - .unwrap(); - column_writer.close().unwrap(); + write_nulls::(&mut row_group_writer); // INT96 - let mut column_writer = row_group_writer.next_column().unwrap().unwrap(); - column_writer - .typed::() - .write_batch(&[], Some(&[0, 0, 0, 0]), None) - .unwrap(); - column_writer.close().unwrap(); + write_nulls::(&mut row_group_writer); // BOOLEAN - let mut column_writer = row_group_writer.next_column().unwrap().unwrap(); - column_writer - .typed::() - .write_batch(&[], Some(&[0, 0, 0, 0]), None) - .unwrap(); - column_writer.close().unwrap(); + write_nulls::(&mut row_group_writer); // FLOAT - let mut column_writer = row_group_writer.next_column().unwrap().unwrap(); - column_writer - .typed::() - .write_batch(&[], Some(&[0, 0, 0, 0]), None) - .unwrap(); - column_writer.close().unwrap(); + write_nulls::(&mut row_group_writer); // DOUBLE - let mut column_writer = row_group_writer.next_column().unwrap().unwrap(); - column_writer - .typed::() - .write_batch(&[], Some(&[0, 0, 0, 0]), None) - .unwrap(); - column_writer.close().unwrap(); + write_nulls::(&mut row_group_writer); // BYTE_ARRAY - let mut column_writer = row_group_writer.next_column().unwrap().unwrap(); - column_writer - .typed::() - .write_batch(&[], Some(&[0, 0, 0, 0]), None) - .unwrap(); - column_writer.close().unwrap(); + write_nulls::(&mut row_group_writer); // FIXED_LEN_BYTE_ARRAY - let mut column_writer = row_group_writer.next_column().unwrap().unwrap(); - column_writer - .typed::() - .write_batch(&[], Some(&[0, 0, 0, 0]), None) - .unwrap(); - column_writer.close().unwrap(); + write_nulls::(&mut row_group_writer); row_group_writer.close().unwrap(); From ecbeaf141d4b9b4865fdae734d55fa8a151efbf6 Mon Sep 17 00:00:00 2001 From: Ed Seidl Date: Fri, 1 May 2026 07:05:48 -0700 Subject: [PATCH 4/4] add check suggested in review --- parquet/src/arrow/arrow_reader/mod.rs | 1 + 1 file changed, 1 insertion(+) diff --git a/parquet/src/arrow/arrow_reader/mod.rs b/parquet/src/arrow/arrow_reader/mod.rs index 16bce2fbc5e8..70d3ce7cf9a9 100644 --- a/parquet/src/arrow/arrow_reader/mod.rs +++ b/parquet/src/arrow/arrow_reader/mod.rs @@ -3753,6 +3753,7 @@ pub(crate) mod tests { for col in batch.columns() { assert_eq!(col.len(), 4); assert_eq!(col.logical_null_count(), 4); + assert_eq!(*col.data_type(), ArrowDataType::Null); } }