This is an automated email from the ASF dual-hosted git repository.
tustvold pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/arrow-rs.git
The following commit(s) were added to refs/heads/master by this push:
new cf1e778b8 Fix Backwards Compatible Parquet List Encodings (#1915)
(#2774)
cf1e778b8 is described below
commit cf1e778b8c34155e7b598907a829ff6c8e52a1ea
Author: Raphael Taylor-Davies <[email protected]>
AuthorDate: Sat Sep 24 19:39:13 2022 +0100
Fix Backwards Compatible Parquet List Encodings (#1915) (#2774)
* Fix schema for non-list repeated fields (#1915)
* Clippy
---
parquet/src/arrow/array_reader/builder.rs | 229 ++++++++++++++++++---------
parquet/src/arrow/array_reader/list_array.rs | 31 ++--
parquet/src/arrow/arrow_reader/mod.rs | 99 ++++++++++--
parquet/src/arrow/async_reader.rs | 28 ++--
parquet/src/arrow/schema.rs | 22 ++-
sample.parquet | Bin 0 -> 686 bytes
6 files changed, 282 insertions(+), 127 deletions(-)
diff --git a/parquet/src/arrow/array_reader/builder.rs
b/parquet/src/arrow/array_reader/builder.rs
index 5f3ce7582..c0216466d 100644
--- a/parquet/src/arrow/array_reader/builder.rs
+++ b/parquet/src/arrow/array_reader/builder.rs
@@ -17,7 +17,7 @@
use std::sync::Arc;
-use arrow::datatypes::{DataType, SchemaRef};
+use arrow::datatypes::DataType;
use crate::arrow::array_reader::empty_array::make_empty_array_reader;
use
crate::arrow::array_reader::fixed_len_byte_array::make_fixed_len_byte_array_reader;
@@ -26,40 +26,43 @@ use crate::arrow::array_reader::{
ListArrayReader, MapArrayReader, NullArrayReader, PrimitiveArrayReader,
RowGroupCollection, StructArrayReader,
};
-use crate::arrow::schema::{convert_schema, ParquetField, ParquetFieldType};
+use crate::arrow::schema::{ParquetField, ParquetFieldType};
use crate::arrow::ProjectionMask;
use crate::basic::Type as PhysicalType;
use crate::data_type::{
BoolType, DoubleType, FloatType, Int32Type, Int64Type, Int96Type,
};
-use crate::errors::Result;
+use crate::errors::{ParquetError, Result};
use crate::schema::types::{ColumnDescriptor, ColumnPath, Type};
/// Create array reader from parquet schema, projection mask, and parquet file
reader.
pub fn build_array_reader(
- arrow_schema: SchemaRef,
- mask: ProjectionMask,
+ field: Option<&ParquetField>,
+ mask: &ProjectionMask,
row_groups: &dyn RowGroupCollection,
) -> Result<Box<dyn ArrayReader>> {
- let field = convert_schema(&row_groups.schema(), mask,
Some(arrow_schema.as_ref()))?;
+ let reader = field
+ .and_then(|field| build_reader(field, mask, row_groups).transpose())
+ .transpose()?
+ .unwrap_or_else(|| make_empty_array_reader(row_groups.num_rows()));
- match &field {
- Some(field) => build_reader(field, row_groups),
- None => Ok(make_empty_array_reader(row_groups.num_rows())),
- }
+ Ok(reader)
}
fn build_reader(
field: &ParquetField,
+ mask: &ProjectionMask,
row_groups: &dyn RowGroupCollection,
-) -> Result<Box<dyn ArrayReader>> {
+) -> Result<Option<Box<dyn ArrayReader>>> {
match field.field_type {
- ParquetFieldType::Primitive { .. } => build_primitive_reader(field,
row_groups),
+ ParquetFieldType::Primitive { .. } => {
+ build_primitive_reader(field, mask, row_groups)
+ }
ParquetFieldType::Group { .. } => match &field.arrow_type {
- DataType::Map(_, _) => build_map_reader(field, row_groups),
- DataType::Struct(_) => build_struct_reader(field, row_groups),
- DataType::List(_) => build_list_reader(field, false, row_groups),
- DataType::LargeList(_) => build_list_reader(field, true,
row_groups),
+ DataType::Map(_, _) => build_map_reader(field, mask, row_groups),
+ DataType::Struct(_) => build_struct_reader(field, mask,
row_groups),
+ DataType::List(_) => build_list_reader(field, mask, false,
row_groups),
+ DataType::LargeList(_) => build_list_reader(field, mask, true,
row_groups),
d => unimplemented!("reading group type {} not implemented", d),
},
}
@@ -68,59 +71,106 @@ fn build_reader(
/// Build array reader for map type.
fn build_map_reader(
field: &ParquetField,
+ mask: &ProjectionMask,
row_groups: &dyn RowGroupCollection,
-) -> Result<Box<dyn ArrayReader>> {
+) -> Result<Option<Box<dyn ArrayReader>>> {
let children = field.children().unwrap();
assert_eq!(children.len(), 2);
- let key_reader = build_reader(&children[0], row_groups)?;
- let value_reader = build_reader(&children[1], row_groups)?;
+ let key_reader = build_reader(&children[0], mask, row_groups)?;
+ let value_reader = build_reader(&children[1], mask, row_groups)?;
- Ok(Box::new(MapArrayReader::new(
- key_reader,
- value_reader,
- field.arrow_type.clone(),
- field.def_level,
- field.rep_level,
- field.nullable,
- )))
+ match (key_reader, value_reader) {
+ (Some(key_reader), Some(value_reader)) => {
+ let key_type = key_reader.get_data_type().clone();
+ let value_type = value_reader.get_data_type().clone();
+
+ let data_type = match &field.arrow_type {
+ DataType::Map(map_field, is_sorted) => match
map_field.data_type() {
+ DataType::Struct(fields) => {
+ assert_eq!(fields.len(), 2);
+ let struct_field =
+
map_field.clone().with_data_type(DataType::Struct(vec![
+ fields[0].clone().with_data_type(key_type),
+ fields[1].clone().with_data_type(value_type),
+ ]));
+ DataType::Map(Box::new(struct_field), *is_sorted)
+ }
+ _ => unreachable!(),
+ },
+ _ => unreachable!(),
+ };
+
+ Ok(Some(Box::new(MapArrayReader::new(
+ key_reader,
+ value_reader,
+ data_type,
+ field.def_level,
+ field.rep_level,
+ field.nullable,
+ ))))
+ }
+ (None, None) => Ok(None),
+ _ => {
+ Err(general_err!(
+ "partial projection of MapArray is not supported"
+ ))
+ }
+ }
}
/// Build array reader for list type.
fn build_list_reader(
field: &ParquetField,
+ mask: &ProjectionMask,
is_large: bool,
row_groups: &dyn RowGroupCollection,
-) -> Result<Box<dyn ArrayReader>> {
+) -> Result<Option<Box<dyn ArrayReader>>> {
let children = field.children().unwrap();
assert_eq!(children.len(), 1);
- let data_type = field.arrow_type.clone();
- let item_reader = build_reader(&children[0], row_groups)?;
+ let reader = match build_reader(&children[0], mask, row_groups)? {
+ Some(item_reader) => {
+ let item_type = item_reader.get_data_type().clone();
+ let data_type = match &field.arrow_type {
+ DataType::List(f) => {
+
DataType::List(Box::new(f.clone().with_data_type(item_type)))
+ }
+ DataType::LargeList(f) => {
+
DataType::LargeList(Box::new(f.clone().with_data_type(item_type)))
+ }
+ _ => unreachable!(),
+ };
- match is_large {
- false => Ok(Box::new(ListArrayReader::<i32>::new(
- item_reader,
- data_type,
- field.def_level,
- field.rep_level,
- field.nullable,
- )) as _),
- true => Ok(Box::new(ListArrayReader::<i64>::new(
- item_reader,
- data_type,
- field.def_level,
- field.rep_level,
- field.nullable,
- )) as _),
- }
+ let reader = match is_large {
+ false => Box::new(ListArrayReader::<i32>::new(
+ item_reader,
+ data_type,
+ field.def_level,
+ field.rep_level,
+ field.nullable,
+ )) as _,
+ true => Box::new(ListArrayReader::<i64>::new(
+ item_reader,
+ data_type,
+ field.def_level,
+ field.rep_level,
+ field.nullable,
+ )) as _,
+ };
+ Some(reader)
+ }
+ None => None,
+ };
+ Ok(reader)
}
/// Creates primitive array reader for each primitive type.
fn build_primitive_reader(
field: &ParquetField,
+ mask: &ProjectionMask,
row_groups: &dyn RowGroupCollection,
-) -> Result<Box<dyn ArrayReader>> {
+) -> Result<Option<Box<dyn ArrayReader>>> {
let (col_idx, primitive_type) = match &field.field_type {
ParquetFieldType::Primitive {
col_idx,
@@ -132,6 +182,10 @@ fn build_primitive_reader(
_ => unreachable!(),
};
+ if !mask.leaf_included(col_idx) {
+ return Ok(None);
+ }
+
let physical_type = primitive_type.get_physical_type();
// We don't track the column path in ParquetField as it adds a potential
source
@@ -150,81 +204,99 @@ fn build_primitive_reader(
let page_iterator = row_groups.column_chunks(col_idx)?;
let arrow_type = Some(field.arrow_type.clone());
- match physical_type {
- PhysicalType::BOOLEAN =>
Ok(Box::new(PrimitiveArrayReader::<BoolType>::new(
+ let reader = match physical_type {
+ PhysicalType::BOOLEAN =>
Box::new(PrimitiveArrayReader::<BoolType>::new(
page_iterator,
column_desc,
arrow_type,
- )?)),
+ )?) as _,
PhysicalType::INT32 => {
if let Some(DataType::Null) = arrow_type {
- Ok(Box::new(NullArrayReader::<Int32Type>::new(
+ Box::new(NullArrayReader::<Int32Type>::new(
page_iterator,
column_desc,
- )?))
+ )?) as _
} else {
- Ok(Box::new(PrimitiveArrayReader::<Int32Type>::new(
+ Box::new(PrimitiveArrayReader::<Int32Type>::new(
page_iterator,
column_desc,
arrow_type,
- )?))
+ )?) as _
}
}
- PhysicalType::INT64 =>
Ok(Box::new(PrimitiveArrayReader::<Int64Type>::new(
+ PhysicalType::INT64 => Box::new(PrimitiveArrayReader::<Int64Type>::new(
page_iterator,
column_desc,
arrow_type,
- )?)),
- PhysicalType::INT96 =>
Ok(Box::new(PrimitiveArrayReader::<Int96Type>::new(
+ )?) as _,
+ PhysicalType::INT96 => Box::new(PrimitiveArrayReader::<Int96Type>::new(
page_iterator,
column_desc,
arrow_type,
- )?)),
- PhysicalType::FLOAT =>
Ok(Box::new(PrimitiveArrayReader::<FloatType>::new(
+ )?) as _,
+ PhysicalType::FLOAT => Box::new(PrimitiveArrayReader::<FloatType>::new(
page_iterator,
column_desc,
arrow_type,
- )?)),
- PhysicalType::DOUBLE =>
Ok(Box::new(PrimitiveArrayReader::<DoubleType>::new(
+ )?) as _,
+ PhysicalType::DOUBLE =>
Box::new(PrimitiveArrayReader::<DoubleType>::new(
page_iterator,
column_desc,
arrow_type,
- )?)),
+ )?) as _,
PhysicalType::BYTE_ARRAY => match arrow_type {
Some(DataType::Dictionary(_, _)) => {
- make_byte_array_dictionary_reader(page_iterator, column_desc,
arrow_type)
+ make_byte_array_dictionary_reader(page_iterator, column_desc,
arrow_type)?
}
- _ => make_byte_array_reader(page_iterator, column_desc,
arrow_type),
+ _ => make_byte_array_reader(page_iterator, column_desc,
arrow_type)?,
},
PhysicalType::FIXED_LEN_BYTE_ARRAY => {
- make_fixed_len_byte_array_reader(page_iterator, column_desc,
arrow_type)
+ make_fixed_len_byte_array_reader(page_iterator, column_desc,
arrow_type)?
}
- }
+ };
+ Ok(Some(reader))
}
fn build_struct_reader(
field: &ParquetField,
+ mask: &ProjectionMask,
row_groups: &dyn RowGroupCollection,
-) -> Result<Box<dyn ArrayReader>> {
+) -> Result<Option<Box<dyn ArrayReader>>> {
+ let arrow_fields = match &field.arrow_type {
+ DataType::Struct(children) => children,
+ _ => unreachable!(),
+ };
let children = field.children().unwrap();
- let children_reader = children
- .iter()
- .map(|child| build_reader(child, row_groups))
- .collect::<Result<Vec<_>>>()?;
+ assert_eq!(arrow_fields.len(), children.len());
+
+ let mut readers = Vec::with_capacity(children.len());
+ let mut projected_fields = Vec::with_capacity(children.len());
+
+ for (arrow, parquet) in arrow_fields.iter().zip(children) {
+ if let Some(reader) = build_reader(parquet, mask, row_groups)? {
+ let child_type = reader.get_data_type().clone();
+ projected_fields.push(arrow.clone().with_data_type(child_type));
+ readers.push(reader);
+ }
+ }
+
+ if readers.is_empty() {
+ return Ok(None);
+ }
- Ok(Box::new(StructArrayReader::new(
- field.arrow_type.clone(),
- children_reader,
+ Ok(Some(Box::new(StructArrayReader::new(
+ DataType::Struct(projected_fields),
+ readers,
field.def_level,
field.rep_level,
field.nullable,
- )) as _)
+ ))))
}
#[cfg(test)]
mod tests {
use super::*;
- use crate::arrow::parquet_to_arrow_schema;
+ use crate::arrow::schema::parquet_to_array_schema_and_fields;
use crate::file::reader::{FileReader, SerializedFileReader};
use crate::util::test_common::file_util::get_test_file;
use arrow::datatypes::Field;
@@ -238,14 +310,15 @@ mod tests {
let file_metadata = file_reader.metadata().file_metadata();
let mask = ProjectionMask::leaves(file_metadata.schema_descr(), [0]);
- let arrow_schema = parquet_to_arrow_schema(
+ let (_, fields) = parquet_to_array_schema_and_fields(
file_metadata.schema_descr(),
+ ProjectionMask::all(),
file_metadata.key_value_metadata(),
)
.unwrap();
let array_reader =
- build_array_reader(Arc::new(arrow_schema), mask,
&file_reader).unwrap();
+ build_array_reader(fields.as_ref(), &mask, &file_reader).unwrap();
// Create arrow types
let arrow_type = DataType::Struct(vec![Field::new(
diff --git a/parquet/src/arrow/array_reader/list_array.rs
b/parquet/src/arrow/array_reader/list_array.rs
index d2fa94611..f0b5092e1 100644
--- a/parquet/src/arrow/array_reader/list_array.rs
+++ b/parquet/src/arrow/array_reader/list_array.rs
@@ -251,6 +251,7 @@ mod tests {
use crate::arrow::array_reader::build_array_reader;
use crate::arrow::array_reader::list_array::ListArrayReader;
use crate::arrow::array_reader::test_util::InMemoryArrayReader;
+ use crate::arrow::schema::parquet_to_array_schema_and_fields;
use crate::arrow::{parquet_to_arrow_schema, ArrowWriter, ProjectionMask};
use crate::file::properties::WriterProperties;
use crate::file::reader::{FileReader, SerializedFileReader};
@@ -389,21 +390,10 @@ mod tests {
true,
);
- let l2 = ListArrayReader::<OffsetSize>::new(
- Box::new(l3),
- l2_type,
- 3,
- 2,
- false,
- );
+ let l2 = ListArrayReader::<OffsetSize>::new(Box::new(l3), l2_type, 3,
2, false);
- let mut l1 = ListArrayReader::<OffsetSize>::new(
- Box::new(l2),
- l1_type,
- 2,
- 1,
- true,
- );
+ let mut l1 =
+ ListArrayReader::<OffsetSize>::new(Box::new(l2), l1_type, 2, 1,
true);
let expected_1 = expected.slice(0, 2);
let expected_2 = expected.slice(2, 2);
@@ -573,18 +563,17 @@ mod tests {
Arc::new(SerializedFileReader::new(file).unwrap());
let file_metadata = file_reader.metadata().file_metadata();
- let arrow_schema = parquet_to_arrow_schema(
- file_metadata.schema_descr(),
+ let schema = file_metadata.schema_descr();
+ let mask = ProjectionMask::leaves(schema, vec![0]);
+ let (_, fields) = parquet_to_array_schema_and_fields(
+ schema,
+ ProjectionMask::all(),
file_metadata.key_value_metadata(),
)
.unwrap();
- let schema = file_metadata.schema_descr_ptr();
- let mask = ProjectionMask::leaves(&schema, vec![0]);
-
let mut array_reader =
- build_array_reader(Arc::new(arrow_schema), mask, &file_reader)
- .unwrap();
+ build_array_reader(fields.as_ref(), &mask, &file_reader).unwrap();
let batch = array_reader.next_batch(100).unwrap();
assert_eq!(batch.data_type(), array_reader.get_data_type());
diff --git a/parquet/src/arrow/arrow_reader/mod.rs
b/parquet/src/arrow/arrow_reader/mod.rs
index 59abf9ad8..5ee963916 100644
--- a/parquet/src/arrow/arrow_reader/mod.rs
+++ b/parquet/src/arrow/arrow_reader/mod.rs
@@ -30,8 +30,8 @@ use arrow::{array::StructArray, error::ArrowError};
use crate::arrow::array_reader::{
build_array_reader, ArrayReader, FileReaderRowGroupCollection,
RowGroupCollection,
};
-use crate::arrow::schema::parquet_to_arrow_schema;
-use crate::arrow::schema::parquet_to_arrow_schema_by_columns;
+use crate::arrow::schema::{parquet_to_array_schema_and_fields,
parquet_to_arrow_schema};
+use crate::arrow::schema::{parquet_to_arrow_schema_by_columns, ParquetField};
use crate::arrow::ProjectionMask;
use crate::errors::{ParquetError, Result};
use crate::file::metadata::{KeyValue, ParquetMetaData};
@@ -60,6 +60,8 @@ pub struct ArrowReaderBuilder<T> {
pub(crate) schema: SchemaRef,
+ pub(crate) fields: Option<ParquetField>,
+
pub(crate) batch_size: usize,
pub(crate) row_groups: Option<Vec<usize>>,
@@ -82,15 +84,17 @@ impl<T> ArrowReaderBuilder<T> {
false => metadata.file_metadata().key_value_metadata(),
};
- let schema = Arc::new(parquet_to_arrow_schema(
+ let (schema, fields) = parquet_to_array_schema_and_fields(
metadata.file_metadata().schema_descr(),
+ ProjectionMask::all(),
kv_metadata,
- )?);
+ )?;
Ok(Self {
input,
metadata,
- schema,
+ schema: Arc::new(schema),
+ fields,
batch_size: 1024,
row_groups: None,
projection: ProjectionMask::all(),
@@ -283,8 +287,16 @@ impl ArrowReader for ParquetFileArrowReader {
mask: ProjectionMask,
batch_size: usize,
) -> Result<ParquetRecordBatchReader> {
- let array_reader =
- build_array_reader(Arc::new(self.get_schema()?), mask,
&self.file_reader)?;
+ let (_, field) = parquet_to_array_schema_and_fields(
+ self.parquet_schema(),
+ mask,
+ self.get_kv_metadata(),
+ )?;
+ let array_reader = build_array_reader(
+ field.as_ref(),
+ &ProjectionMask::all(),
+ &self.file_reader,
+ )?;
// Try to avoid allocate large buffer
let batch_size = self.file_reader.num_rows().min(batch_size);
@@ -420,9 +432,11 @@ impl<T: ChunkReader + 'static>
ArrowReaderBuilder<SyncReader<T>> {
break;
}
- let projection = predicate.projection().clone();
- let array_reader =
- build_array_reader(Arc::clone(&self.schema), projection,
&reader)?;
+ let array_reader = build_array_reader(
+ self.fields.as_ref(),
+ predicate.projection(),
+ &reader,
+ )?;
selection = Some(evaluate_predicate(
batch_size,
@@ -433,7 +447,8 @@ impl<T: ChunkReader + 'static>
ArrowReaderBuilder<SyncReader<T>> {
}
}
- let array_reader = build_array_reader(self.schema, self.projection,
&reader)?;
+ let array_reader =
+ build_array_reader(self.fields.as_ref(), &self.projection,
&reader)?;
// If selection is empty, truncate
if !selects_any(selection.as_ref()) {
@@ -2313,4 +2328,66 @@ mod tests {
assert_ne!(1024, num_rows);
assert_eq!(reader.batch_size, num_rows as usize);
}
+
+ #[test]
+ fn test_raw_repetition() {
+ const MESSAGE_TYPE: &str = "
+ message Log {
+ OPTIONAL INT32 eventType;
+ REPEATED INT32 category;
+ REPEATED group filter {
+ OPTIONAL INT32 error;
+ }
+ }
+ ";
+ let schema = Arc::new(parse_message_type(MESSAGE_TYPE).unwrap());
+ let props = Arc::new(WriterProperties::builder().build());
+
+ let mut buf = Vec::with_capacity(1024);
+ let mut writer = SerializedFileWriter::new(&mut buf, schema,
props).unwrap();
+ let mut row_group_writer = writer.next_row_group().unwrap();
+
+ // column 0
+ let mut col_writer = row_group_writer.next_column().unwrap().unwrap();
+ col_writer
+ .typed::<Int32Type>()
+ .write_batch(&[1], Some(&[1]), None)
+ .unwrap();
+ col_writer.close().unwrap();
+ // column 1
+ let mut col_writer = row_group_writer.next_column().unwrap().unwrap();
+ col_writer
+ .typed::<Int32Type>()
+ .write_batch(&[1, 1], Some(&[1, 1]), Some(&[0, 1]))
+ .unwrap();
+ col_writer.close().unwrap();
+ // column 2
+ let mut col_writer = row_group_writer.next_column().unwrap().unwrap();
+ col_writer
+ .typed::<Int32Type>()
+ .write_batch(&[1], Some(&[1]), Some(&[0]))
+ .unwrap();
+ col_writer.close().unwrap();
+
+ let rg_md = row_group_writer.close().unwrap();
+ assert_eq!(rg_md.num_rows(), 1);
+ writer.close().unwrap();
+
+ let bytes = Bytes::from(buf);
+
+ let mut no_mask = ParquetRecordBatchReader::try_new(bytes.clone(),
1024).unwrap();
+ let full = no_mask.next().unwrap().unwrap();
+
+ assert_eq!(full.num_columns(), 3);
+
+ for idx in 0..3 {
+ let b =
ParquetRecordBatchReaderBuilder::try_new(bytes.clone()).unwrap();
+ let mask = ProjectionMask::leaves(b.parquet_schema(), [idx]);
+ let mut reader = b.with_projection(mask).build().unwrap();
+ let projected = reader.next().unwrap().unwrap();
+
+ assert_eq!(projected.num_columns(), 1);
+ assert_eq!(full.column(idx), projected.column(0));
+ }
+ }
}
diff --git a/parquet/src/arrow/async_reader.rs
b/parquet/src/arrow/async_reader.rs
index d444d20d5..b6b5d7ff7 100644
--- a/parquet/src/arrow/async_reader.rs
+++ b/parquet/src/arrow/async_reader.rs
@@ -101,6 +101,7 @@ use crate::arrow::arrow_reader::{
evaluate_predicate, selects_any, ArrowReaderBuilder, ArrowReaderOptions,
ParquetRecordBatchReader, RowFilter, RowSelection,
};
+use crate::arrow::schema::ParquetField;
use crate::arrow::ProjectionMask;
use crate::column::page::{PageIterator, PageReader};
@@ -337,7 +338,7 @@ impl<T: AsyncFileReader + Send + 'static>
ArrowReaderBuilder<AsyncReader<T>> {
input: self.input.0,
filter: self.filter,
metadata: self.metadata.clone(),
- schema: self.schema.clone(),
+ fields: self.fields,
};
Ok(ParquetRecordBatchStream {
@@ -360,7 +361,7 @@ type ReadResult<T> = Result<(ReaderFactory<T>,
Option<ParquetRecordBatchReader>)
struct ReaderFactory<T> {
metadata: Arc<ParquetMetaData>,
- schema: SchemaRef,
+ fields: Option<ParquetField>,
input: T,
@@ -397,13 +398,13 @@ where
return Ok((self, None));
}
- let predicate_projection = predicate.projection().clone();
+ let predicate_projection = predicate.projection();
row_group
- .fetch(&mut self.input, &predicate_projection,
selection.as_ref())
+ .fetch(&mut self.input, predicate_projection,
selection.as_ref())
.await?;
let array_reader = build_array_reader(
- self.schema.clone(),
+ self.fields.as_ref(),
predicate_projection,
&row_group,
)?;
@@ -427,7 +428,7 @@ where
let reader = ParquetRecordBatchReader::new(
batch_size,
- build_array_reader(self.schema.clone(), projection, &row_group)?,
+ build_array_reader(self.fields.as_ref(), &projection, &row_group)?,
selection,
);
@@ -792,7 +793,8 @@ mod tests {
use crate::arrow::arrow_reader::{
ArrowPredicateFn, ParquetRecordBatchReaderBuilder, RowSelector,
};
- use crate::arrow::{parquet_to_arrow_schema, ArrowWriter};
+ use crate::arrow::schema::parquet_to_array_schema_and_fields;
+ use crate::arrow::ArrowWriter;
use crate::file::footer::parse_metadata;
use crate::file::page_index::index_reader;
use arrow::array::{Array, ArrayRef, Int32Array, StringArray};
@@ -1278,10 +1280,12 @@ mod tests {
};
let requests = async_reader.requests.clone();
- let schema = Arc::new(
- parquet_to_arrow_schema(metadata.file_metadata().schema_descr(),
None)
- .expect("building arrow schema"),
- );
+ let (_, fields) = parquet_to_array_schema_and_fields(
+ metadata.file_metadata().schema_descr(),
+ ProjectionMask::all(),
+ None,
+ )
+ .unwrap();
let _schema_desc = metadata.file_metadata().schema_descr();
@@ -1290,7 +1294,7 @@ mod tests {
let reader_factory = ReaderFactory {
metadata,
- schema,
+ fields,
input: async_reader,
filter: None,
};
diff --git a/parquet/src/arrow/schema.rs b/parquet/src/arrow/schema.rs
index ad5b6b1f5..7803385e7 100644
--- a/parquet/src/arrow/schema.rs
+++ b/parquet/src/arrow/schema.rs
@@ -41,7 +41,7 @@ mod complex;
mod primitive;
use crate::arrow::ProjectionMask;
-pub(crate) use complex::{convert_schema, ParquetField, ParquetFieldType};
+pub(crate) use complex::{ParquetField, ParquetFieldType};
/// Convert Parquet schema to Arrow schema including optional metadata.
/// Attempts to decode any existing Arrow schema metadata, falling back
@@ -64,6 +64,15 @@ pub fn parquet_to_arrow_schema_by_columns(
mask: ProjectionMask,
key_value_metadata: Option<&Vec<KeyValue>>,
) -> Result<Schema> {
+ Ok(parquet_to_array_schema_and_fields(parquet_schema, mask,
key_value_metadata)?.0)
+}
+
+/// Extracts the arrow metadata
+pub(crate) fn parquet_to_array_schema_and_fields(
+ parquet_schema: &SchemaDescriptor,
+ mask: ProjectionMask,
+ key_value_metadata: Option<&Vec<KeyValue>>,
+) -> Result<(Schema, Option<ParquetField>)> {
let mut metadata =
parse_key_value_metadata(key_value_metadata).unwrap_or_default();
let maybe_schema = metadata
.remove(super::ARROW_SCHEMA_META_KEY)
@@ -77,12 +86,15 @@ pub fn parquet_to_arrow_schema_by_columns(
});
}
- match convert_schema(parquet_schema, mask, maybe_schema.as_ref())? {
- Some(field) => match field.arrow_type {
- DataType::Struct(fields) => Ok(Schema::new_with_metadata(fields,
metadata)),
+ match complex::convert_schema(parquet_schema, mask,
maybe_schema.as_ref())? {
+ Some(field) => match &field.arrow_type {
+ DataType::Struct(fields) => Ok((
+ Schema::new_with_metadata(fields.clone(), metadata),
+ Some(field),
+ )),
_ => unreachable!(),
},
- None => Ok(Schema::new_with_metadata(vec![], metadata)),
+ None => Ok((Schema::new_with_metadata(vec![], metadata), None)),
}
}
diff --git a/sample.parquet b/sample.parquet
new file mode 100644
index 000000000..093b6438a
Binary files /dev/null and b/sample.parquet differ