1use arrow_array::cast::AsArray;
21use arrow_array::{BooleanArray, RecordBatch, RecordBatchReader};
22use arrow_buffer::{BooleanBuffer, BooleanBufferBuilder};
23use arrow_schema::{ArrowError, DataType as ArrowType, FieldRef, Schema, SchemaRef};
24use arrow_select::filter::filter_record_batch;
25pub use filter::{ArrowPredicate, ArrowPredicateFn, RowFilter};
26use selection::MaskCursor;
27pub use selection::{
28 MaskRunIter, RowSelection, RowSelectionCursor, RowSelectionPolicy, RowSelector,
29};
30use std::fmt::{Debug, Formatter};
31use std::sync::Arc;
32
33pub use crate::arrow::array_reader::RowGroups;
34use crate::arrow::array_reader::{ArrayReader, ArrayReaderBuilder};
35use crate::arrow::schema::{
36 ParquetField, parquet_to_arrow_schema_and_fields, virtual_type::is_virtual_column,
37};
38use crate::arrow::{FieldLevels, ProjectionMask, parquet_to_arrow_field_levels_with_virtual};
39use crate::basic::{BloomFilterAlgorithm, BloomFilterCompression, BloomFilterHash};
40use crate::bloom_filter::{
41 SBBF_HEADER_SIZE_ESTIMATE, Sbbf, chunk_read_bloom_filter_header_and_offset,
42};
43use crate::column::page::{PageIterator, PageReader};
44#[cfg(feature = "encryption")]
45use crate::encryption::decrypt::FileDecryptionProperties;
46use crate::errors::{ParquetError, Result};
47use crate::file::metadata::{
48 PageIndexPolicy, ParquetMetaData, ParquetMetaDataOptions, ParquetMetaDataReader,
49 ParquetStatisticsPolicy, RowGroupMetaData,
50};
51use crate::file::reader::{ChunkReader, SerializedPageReader};
52use crate::schema::types::SchemaDescriptor;
53
54use crate::arrow::arrow_reader::metrics::ArrowReaderMetrics;
55pub use read_plan::{PredicateOptions, ReadPlan, ReadPlanBuilder};
57
58mod filter;
59pub mod metrics;
60mod read_plan;
61pub(crate) mod selection;
62pub mod statistics;
63
64pub const DEFAULT_BATCH_SIZE: usize = 1024;
66
67#[derive(Debug, Clone, PartialEq, Eq)]
76pub struct RowGroupSelection {
77 pub(crate) row_group_index: usize,
78 pub(crate) selection: Option<RowSelection>,
79}
80
81impl RowGroupSelection {
82 pub fn new(row_group_index: usize, selection: Option<RowSelection>) -> Self {
84 Self {
85 row_group_index,
86 selection,
87 }
88 }
89
90 pub fn row_group_index(&self) -> usize {
92 self.row_group_index
93 }
94
95 pub fn selection(&self) -> Option<&RowSelection> {
98 self.selection.as_ref()
99 }
100}
101
102#[derive(Debug)]
104pub(crate) enum RowGroupPlan {
105 Global {
111 row_groups: Option<Vec<usize>>,
112 selection: Option<RowSelection>,
113 },
114 PerRowGroup(Vec<RowGroupSelection>),
120 Conflicting,
123}
124
125impl RowGroupPlan {
126 fn set_row_groups(&mut self, new_row_groups: Vec<usize>) {
127 match self {
128 Self::Global { row_groups, .. } => *row_groups = Some(new_row_groups),
129 Self::PerRowGroup(_) => *self = Self::Conflicting,
130 Self::Conflicting => {}
131 }
132 }
133
134 fn set_row_selection(&mut self, new_selection: RowSelection) {
135 match self {
136 Self::Global { selection, .. } => *selection = Some(new_selection),
137 Self::PerRowGroup(_) => *self = Self::Conflicting,
138 Self::Conflicting => {}
139 }
140 }
141
142 pub(crate) fn set_row_group_selections(
143 &mut self,
144 row_group_selections: Vec<RowGroupSelection>,
145 ) {
146 match self {
147 Self::Global {
148 row_groups: None,
149 selection: None,
150 }
151 | Self::PerRowGroup(_) => {
152 *self = Self::PerRowGroup(row_group_selections);
153 }
154 Self::Global { .. } => *self = Self::Conflicting,
155 Self::Conflicting => {}
156 }
157 }
158
159 pub(crate) fn conflict_error() -> ParquetError {
160 ParquetError::General(
161 "with_row_group_selections cannot be combined with with_row_groups or with_row_selection"
162 .to_string(),
163 )
164 }
165
166 fn into_global(self) -> Result<(Option<Vec<usize>>, Option<RowSelection>)> {
167 match self {
168 Self::Global {
169 row_groups,
170 selection,
171 } => Ok((row_groups, selection)),
172 Self::PerRowGroup(_) => Err(ParquetError::General(
173 "Row-group-local selections are not supported by the synchronous reader"
174 .to_string(),
175 )),
176 Self::Conflicting => Err(Self::conflict_error()),
177 }
178 }
179}
180
181pub struct ArrowReaderBuilder<T> {
229 pub(crate) input: T,
238
239 pub(crate) metadata: Arc<ParquetMetaData>,
240
241 pub(crate) schema: SchemaRef,
242
243 pub(crate) fields: Option<Arc<ParquetField>>,
244
245 pub(crate) batch_size: usize,
246
247 pub(crate) row_group_plan: RowGroupPlan,
248
249 pub(crate) projection: ProjectionMask,
250
251 pub(crate) filter: Option<RowFilter>,
252
253 pub(crate) row_selection_policy: RowSelectionPolicy,
254
255 pub(crate) limit: Option<usize>,
256
257 pub(crate) offset: Option<usize>,
258
259 pub(crate) metrics: ArrowReaderMetrics,
260
261 pub(crate) max_predicate_cache_size: usize,
262}
263
264impl<T: Debug> Debug for ArrowReaderBuilder<T> {
265 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
266 f.debug_struct("ArrowReaderBuilder<T>")
267 .field("input", &self.input)
268 .field("metadata", &self.metadata)
269 .field("schema", &self.schema)
270 .field("fields", &self.fields)
271 .field("batch_size", &self.batch_size)
272 .field("row_group_plan", &self.row_group_plan)
273 .field("projection", &self.projection)
274 .field("filter", &self.filter)
275 .field("row_selection_policy", &self.row_selection_policy)
276 .field("limit", &self.limit)
277 .field("offset", &self.offset)
278 .field("metrics", &self.metrics)
279 .finish()
280 }
281}
282
283impl<T> ArrowReaderBuilder<T> {
284 pub(crate) fn new_builder(input: T, metadata: ArrowReaderMetadata) -> Self {
285 Self {
286 input,
287 metadata: metadata.metadata,
288 schema: metadata.schema,
289 fields: metadata.fields,
290 batch_size: DEFAULT_BATCH_SIZE,
291 row_group_plan: RowGroupPlan::Global {
292 row_groups: None,
293 selection: None,
294 },
295 projection: ProjectionMask::all(),
296 filter: None,
297 row_selection_policy: RowSelectionPolicy::default(),
298 limit: None,
299 offset: None,
300 metrics: ArrowReaderMetrics::Disabled,
301 max_predicate_cache_size: 100 * 1024 * 1024, }
303 }
304
305 pub fn metadata(&self) -> &Arc<ParquetMetaData> {
307 &self.metadata
308 }
309
310 pub fn parquet_schema(&self) -> &SchemaDescriptor {
312 self.metadata.file_metadata().schema_descr()
313 }
314
315 pub fn schema(&self) -> &SchemaRef {
317 &self.schema
318 }
319
320 pub fn with_batch_size(self, batch_size: usize) -> Self {
327 let batch_size = batch_size.min(self.metadata.file_metadata().num_rows() as usize);
329 Self { batch_size, ..self }
330 }
331
332 pub fn with_row_groups(mut self, row_groups: Vec<usize>) -> Self {
344 self.row_group_plan.set_row_groups(row_groups);
345 self
346 }
347
348 pub fn with_projection(self, mask: ProjectionMask) -> Self {
350 Self {
351 projection: mask,
352 ..self
353 }
354 }
355
356 pub fn with_row_selection_policy(self, policy: RowSelectionPolicy) -> Self {
360 Self {
361 row_selection_policy: policy,
362 ..self
363 }
364 }
365
366 pub fn with_row_selection(mut self, selection: RowSelection) -> Self {
434 self.row_group_plan.set_row_selection(selection);
435 self
436 }
437
438 pub fn with_row_filter(self, filter: RowFilter) -> Self {
478 Self {
479 filter: Some(filter),
480 ..self
481 }
482 }
483
484 pub fn with_limit(self, limit: usize) -> Self {
492 Self {
493 limit: Some(limit),
494 ..self
495 }
496 }
497
498 pub fn with_offset(self, offset: usize) -> Self {
506 Self {
507 offset: Some(offset),
508 ..self
509 }
510 }
511
512 pub fn with_metrics(self, metrics: ArrowReaderMetrics) -> Self {
547 Self { metrics, ..self }
548 }
549
550 pub fn with_max_predicate_cache_size(self, max_predicate_cache_size: usize) -> Self {
565 Self {
566 max_predicate_cache_size,
567 ..self
568 }
569 }
570}
571
572#[derive(Debug, Clone, Default)]
587pub struct ArrowReaderOptions {
588 skip_arrow_metadata: bool,
590 supplied_schema: Option<SchemaRef>,
595
596 pub(crate) column_index: PageIndexPolicy,
597 pub(crate) offset_index: PageIndexPolicy,
598
599 metadata_options: ParquetMetaDataOptions,
601 #[cfg(feature = "encryption")]
603 pub(crate) file_decryption_properties: Option<Arc<FileDecryptionProperties>>,
604
605 virtual_columns: Vec<FieldRef>,
606}
607
608impl ArrowReaderOptions {
609 pub fn new() -> Self {
611 Self::default()
612 }
613
614 pub fn with_skip_arrow_metadata(self, skip_arrow_metadata: bool) -> Self {
621 Self {
622 skip_arrow_metadata,
623 ..self
624 }
625 }
626
627 pub fn with_schema(self, schema: SchemaRef) -> Self {
736 Self {
737 supplied_schema: Some(schema),
738 skip_arrow_metadata: true,
739 ..self
740 }
741 }
742
743 pub fn with_page_index_policy(self, policy: PageIndexPolicy) -> Self {
751 self.with_column_index_policy(policy)
752 .with_offset_index_policy(policy)
753 }
754
755 pub fn with_column_index_policy(mut self, policy: PageIndexPolicy) -> Self {
762 self.column_index = policy;
763 self
764 }
765
766 pub fn with_offset_index_policy(mut self, policy: PageIndexPolicy) -> Self {
773 self.offset_index = policy;
774 self
775 }
776
777 pub fn with_parquet_schema(mut self, schema: Arc<SchemaDescriptor>) -> Self {
783 self.metadata_options.set_schema(schema);
784 self
785 }
786
787 pub fn with_encoding_stats_as_mask(mut self, val: bool) -> Self {
798 self.metadata_options.set_encoding_stats_as_mask(val);
799 self
800 }
801
802 pub fn with_encoding_stats_policy(mut self, policy: ParquetStatisticsPolicy) -> Self {
807 self.metadata_options.set_encoding_stats_policy(policy);
808 self
809 }
810
811 pub fn with_column_stats_policy(mut self, policy: ParquetStatisticsPolicy) -> Self {
816 self.metadata_options.set_column_stats_policy(policy);
817 self
818 }
819
820 pub fn with_size_stats_policy(mut self, policy: ParquetStatisticsPolicy) -> Self {
825 self.metadata_options.set_size_stats_policy(policy);
826 self
827 }
828
829 #[cfg(feature = "encryption")]
833 pub fn with_file_decryption_properties(
834 self,
835 file_decryption_properties: Arc<FileDecryptionProperties>,
836 ) -> Self {
837 Self {
838 file_decryption_properties: Some(file_decryption_properties),
839 ..self
840 }
841 }
842
843 pub fn with_virtual_columns(self, virtual_columns: Vec<FieldRef>) -> Result<Self> {
896 for field in &virtual_columns {
898 if !is_virtual_column(field) {
899 return Err(ParquetError::General(format!(
900 "Field '{}' is not a virtual column. Virtual columns must have extension type names starting with 'arrow.virtual.'",
901 field.name()
902 )));
903 }
904 }
905 Ok(Self {
906 virtual_columns,
907 ..self
908 })
909 }
910
911 pub fn offset_index_policy(&self) -> PageIndexPolicy {
916 self.offset_index
917 }
918
919 pub fn column_index_policy(&self) -> PageIndexPolicy {
924 self.column_index
925 }
926
927 pub fn metadata_options(&self) -> &ParquetMetaDataOptions {
929 &self.metadata_options
930 }
931
932 #[cfg(feature = "encryption")]
937 pub fn file_decryption_properties(&self) -> Option<&Arc<FileDecryptionProperties>> {
938 self.file_decryption_properties.as_ref()
939 }
940}
941
942impl ParquetMetaDataReader {
943 pub fn with_arrow_reader_options(mut self, options: Option<&ArrowReaderOptions>) -> Self {
957 let Some(options) = options else { return self };
958
959 self = self.with_metadata_options(Some(options.metadata_options().clone()));
960
961 #[cfg(feature = "encryption")]
962 {
963 self = self.with_decryption_properties(
964 options.file_decryption_properties.as_ref().map(Arc::clone),
965 );
966 }
967
968 if options.column_index_policy() != PageIndexPolicy::Skip
969 || options.offset_index_policy() != PageIndexPolicy::Skip
970 {
971 self = self
972 .with_column_index_policy(options.column_index_policy())
973 .with_offset_index_policy(options.offset_index_policy());
974 }
975
976 self
977 }
978}
979
980#[derive(Debug, Clone)]
995pub struct ArrowReaderMetadata {
996 pub(crate) metadata: Arc<ParquetMetaData>,
998 pub(crate) schema: SchemaRef,
1000 pub(crate) fields: Option<Arc<ParquetField>>,
1002}
1003
1004impl ArrowReaderMetadata {
1005 pub fn load<T: ChunkReader>(reader: &T, options: ArrowReaderOptions) -> Result<Self> {
1019 let metadata = ParquetMetaDataReader::new()
1020 .with_column_index_policy(options.column_index)
1021 .with_offset_index_policy(options.offset_index)
1022 .with_metadata_options(Some(options.metadata_options.clone()));
1023 #[cfg(feature = "encryption")]
1024 let metadata = metadata.with_decryption_properties(
1025 options.file_decryption_properties.as_ref().map(Arc::clone),
1026 );
1027 let metadata = metadata.parse_and_finish(reader)?;
1028 Self::try_new(Arc::new(metadata), options)
1029 }
1030
1031 pub fn try_new(metadata: Arc<ParquetMetaData>, options: ArrowReaderOptions) -> Result<Self> {
1039 match options.supplied_schema {
1040 Some(supplied_schema) => Self::with_supplied_schema(
1041 metadata,
1042 supplied_schema.clone(),
1043 &options.virtual_columns,
1044 ),
1045 None => {
1046 let kv_metadata = match options.skip_arrow_metadata {
1047 true => None,
1048 false => metadata.file_metadata().key_value_metadata(),
1049 };
1050
1051 let (schema, fields) = parquet_to_arrow_schema_and_fields(
1052 metadata.file_metadata().schema_descr(),
1053 ProjectionMask::all(),
1054 kv_metadata,
1055 &options.virtual_columns,
1056 )?;
1057
1058 Ok(Self {
1059 metadata,
1060 schema: Arc::new(schema),
1061 fields: fields.map(Arc::new),
1062 })
1063 }
1064 }
1065 }
1066
1067 fn with_supplied_schema(
1068 metadata: Arc<ParquetMetaData>,
1069 supplied_schema: SchemaRef,
1070 virtual_columns: &[FieldRef],
1071 ) -> Result<Self> {
1072 let parquet_schema = metadata.file_metadata().schema_descr();
1073 let field_levels = parquet_to_arrow_field_levels_with_virtual(
1074 parquet_schema,
1075 ProjectionMask::all(),
1076 Some(supplied_schema.fields()),
1077 virtual_columns,
1078 )?;
1079 let fields = field_levels.fields;
1080 let inferred_len = fields.len();
1081 let supplied_len = supplied_schema.fields().len() + virtual_columns.len();
1082 if inferred_len != supplied_len {
1086 return Err(arrow_err!(format!(
1087 "Incompatible supplied Arrow schema: expected {} columns received {}",
1088 inferred_len, supplied_len
1089 )));
1090 }
1091
1092 let mut errors = Vec::new();
1093
1094 let field_iter = supplied_schema.fields().iter().zip(fields.iter());
1095
1096 for (field1, field2) in field_iter {
1097 if field1.data_type() != field2.data_type() {
1098 errors.push(format!(
1099 "data type mismatch for field {}: requested {} but found {}",
1100 field1.name(),
1101 field1.data_type(),
1102 field2.data_type()
1103 ));
1104 }
1105 if field1.is_nullable() != field2.is_nullable() {
1106 errors.push(format!(
1107 "nullability mismatch for field {}: expected {:?} but found {:?}",
1108 field1.name(),
1109 field1.is_nullable(),
1110 field2.is_nullable()
1111 ));
1112 }
1113 if field1.metadata() != field2.metadata() {
1114 errors.push(format!(
1115 "metadata mismatch for field {}: expected {:?} but found {:?}",
1116 field1.name(),
1117 field1.metadata(),
1118 field2.metadata()
1119 ));
1120 }
1121 }
1122
1123 if !errors.is_empty() {
1124 let message = errors.join(", ");
1125 return Err(ParquetError::ArrowError(format!(
1126 "Incompatible supplied Arrow schema: {message}",
1127 )));
1128 }
1129
1130 Ok(Self {
1131 metadata,
1132 schema: supplied_schema,
1133 fields: field_levels.levels.map(Arc::new),
1134 })
1135 }
1136
1137 pub fn metadata(&self) -> &Arc<ParquetMetaData> {
1139 &self.metadata
1140 }
1141
1142 pub fn parquet_schema(&self) -> &SchemaDescriptor {
1144 self.metadata.file_metadata().schema_descr()
1145 }
1146
1147 pub fn schema(&self) -> &SchemaRef {
1149 &self.schema
1150 }
1151}
1152
1153#[doc(hidden)]
1154pub struct SyncReader<T: ChunkReader>(T);
1156
1157impl<T: Debug + ChunkReader> Debug for SyncReader<T> {
1158 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
1159 f.debug_tuple("SyncReader").field(&self.0).finish()
1160 }
1161}
1162
1163pub type ParquetRecordBatchReaderBuilder<T> = ArrowReaderBuilder<SyncReader<T>>;
1170
1171impl<T: ChunkReader + 'static> ParquetRecordBatchReaderBuilder<T> {
1172 pub fn try_new(reader: T) -> Result<Self> {
1203 Self::try_new_with_options(reader, Default::default())
1204 }
1205
1206 pub fn try_new_with_options(reader: T, options: ArrowReaderOptions) -> Result<Self> {
1211 let metadata = ArrowReaderMetadata::load(&reader, options)?;
1212 Ok(Self::new_with_metadata(reader, metadata))
1213 }
1214
1215 pub fn new_with_metadata(input: T, metadata: ArrowReaderMetadata) -> Self {
1254 Self::new_builder(SyncReader(input), metadata)
1255 }
1256
1257 pub fn get_row_group_column_bloom_filter(
1263 &self,
1264 row_group_idx: usize,
1265 column_idx: usize,
1266 ) -> Result<Option<Sbbf>> {
1267 let metadata = self.metadata.row_group(row_group_idx);
1268 let column_metadata = metadata.column(column_idx);
1269
1270 let offset: u64 = if let Some(offset) = column_metadata.bloom_filter_offset() {
1271 offset
1272 .try_into()
1273 .map_err(|_| ParquetError::General("Bloom filter offset is invalid".to_string()))?
1274 } else {
1275 return Ok(None);
1276 };
1277
1278 let buffer = match column_metadata.bloom_filter_length() {
1279 Some(length) => self.input.0.get_bytes(offset, length as usize),
1280 None => self.input.0.get_bytes(offset, SBBF_HEADER_SIZE_ESTIMATE),
1281 }?;
1282
1283 let (header, bitset_offset) =
1284 chunk_read_bloom_filter_header_and_offset(offset, buffer.clone())?;
1285
1286 match header.algorithm {
1287 BloomFilterAlgorithm::BLOCK => {
1288 }
1290 }
1291 match header.compression {
1292 BloomFilterCompression::UNCOMPRESSED => {
1293 }
1295 }
1296 match header.hash {
1297 BloomFilterHash::XXHASH => {
1298 }
1300 }
1301
1302 let bitset = match column_metadata.bloom_filter_length() {
1303 Some(_) => {
1304 let bitset_start = bitset_offset
1305 .checked_sub(offset)
1306 .and_then(|start| usize::try_from(start).ok())
1307 .ok_or_else(|| {
1308 ParquetError::General("Bloom filter offset is invalid".to_string())
1309 })?;
1310 buffer.slice(bitset_start..)
1311 }
1312 None => {
1313 let bitset_length: usize = header.num_bytes.try_into().map_err(|_| {
1314 ParquetError::General("Bloom filter length is invalid".to_string())
1315 })?;
1316 self.input.0.get_bytes(bitset_offset, bitset_length)?
1317 }
1318 };
1319 Ok(Some(Sbbf::new(&bitset)))
1320 }
1321
1322 pub fn build(self) -> Result<ParquetRecordBatchReader> {
1326 let Self {
1327 input,
1328 metadata,
1329 schema: _,
1330 fields,
1331 batch_size,
1332 row_group_plan,
1333 projection,
1334 mut filter,
1335 row_selection_policy,
1336 limit,
1337 offset,
1338 metrics,
1339 max_predicate_cache_size: _,
1341 } = self;
1342
1343 let batch_size = batch_size.min(metadata.file_metadata().num_rows() as usize);
1345
1346 let (row_groups, selection) = row_group_plan.into_global()?;
1347
1348 let row_groups = row_groups.unwrap_or_else(|| (0..metadata.num_row_groups()).collect());
1349
1350 let reader = ReaderRowGroups {
1351 reader: Arc::new(input.0),
1352 metadata,
1353 row_groups,
1354 };
1355
1356 let mut plan_builder = ReadPlanBuilder::new(batch_size)
1357 .with_selection(selection)
1358 .with_row_selection_policy(row_selection_policy);
1359
1360 if let Some(filter) = filter.as_mut() {
1362 for predicate in &mut filter.predicates {
1363 if !plan_builder.selects_any() {
1365 break;
1366 }
1367
1368 let mut cache_projection = predicate.projection().clone();
1369 cache_projection.intersect(&projection);
1370
1371 let array_reader = ArrayReaderBuilder::new(&reader, &metrics)
1372 .with_batch_size(batch_size)
1373 .with_parquet_metadata(&reader.metadata)
1374 .build_array_reader(fields.as_deref(), predicate.projection())?;
1375
1376 plan_builder = plan_builder.with_predicate(array_reader, predicate.as_mut())?;
1377 }
1378 }
1379
1380 let array_reader = ArrayReaderBuilder::new(&reader, &metrics)
1381 .with_batch_size(batch_size)
1382 .with_parquet_metadata(&reader.metadata)
1383 .build_array_reader(fields.as_deref(), &projection)?;
1384
1385 let read_plan = plan_builder
1386 .limited(reader.num_rows())
1387 .with_offset(offset)
1388 .with_limit(limit)
1389 .build_limited()
1390 .build();
1391
1392 Ok(ParquetRecordBatchReader::new(array_reader, read_plan))
1393 }
1394}
1395
1396struct ReaderRowGroups<T: ChunkReader> {
1397 reader: Arc<T>,
1398
1399 metadata: Arc<ParquetMetaData>,
1400 row_groups: Vec<usize>,
1402}
1403
1404impl<T: ChunkReader + 'static> RowGroups for ReaderRowGroups<T> {
1405 fn num_rows(&self) -> usize {
1406 let meta = self.metadata.row_groups();
1407 self.row_groups
1408 .iter()
1409 .map(|x| meta[*x].num_rows() as usize)
1410 .sum()
1411 }
1412
1413 fn column_chunks(&self, i: usize) -> Result<Box<dyn PageIterator>> {
1414 Ok(Box::new(ReaderPageIterator {
1415 column_idx: i,
1416 reader: self.reader.clone(),
1417 metadata: self.metadata.clone(),
1418 row_groups: self.row_groups.clone().into_iter(),
1419 }))
1420 }
1421
1422 fn row_groups(&self) -> Box<dyn Iterator<Item = &RowGroupMetaData> + '_> {
1423 Box::new(
1424 self.row_groups
1425 .iter()
1426 .map(move |i| self.metadata.row_group(*i)),
1427 )
1428 }
1429
1430 fn metadata(&self) -> &ParquetMetaData {
1431 self.metadata.as_ref()
1432 }
1433}
1434
1435struct ReaderPageIterator<T: ChunkReader> {
1436 reader: Arc<T>,
1437 column_idx: usize,
1438 row_groups: std::vec::IntoIter<usize>,
1439 metadata: Arc<ParquetMetaData>,
1440}
1441
1442impl<T: ChunkReader + 'static> ReaderPageIterator<T> {
1443 fn next_page_reader(&self, rg_idx: usize) -> Result<SerializedPageReader<T>> {
1445 let rg = self.metadata.row_group(rg_idx);
1446 let column_chunk_metadata = rg.column(self.column_idx);
1447 let page_locations = self
1448 .metadata
1449 .page_index()
1450 .map(|i| i.page_locations(rg_idx, self.column_idx).cloned())
1451 .unwrap_or(None);
1452 let total_rows = rg.num_rows() as usize;
1453 let reader = self.reader.clone();
1454
1455 SerializedPageReader::new(reader, column_chunk_metadata, total_rows, page_locations)?
1456 .add_crypto_context(
1457 rg_idx,
1458 self.column_idx,
1459 self.metadata.as_ref(),
1460 column_chunk_metadata,
1461 )
1462 }
1463}
1464
1465impl<T: ChunkReader + 'static> Iterator for ReaderPageIterator<T> {
1466 type Item = Result<Box<dyn PageReader>>;
1467
1468 fn next(&mut self) -> Option<Self::Item> {
1469 let rg_idx = self.row_groups.next()?;
1470 let page_reader = self
1471 .next_page_reader(rg_idx)
1472 .map(|page_reader| Box::new(page_reader) as _);
1473 Some(page_reader)
1474 }
1475}
1476
1477impl<T: ChunkReader + 'static> PageIterator for ReaderPageIterator<T> {}
1478
1479pub struct ParquetRecordBatchReader {
1490 array_reader: Box<dyn ArrayReader>,
1491 schema: SchemaRef,
1492 read_plan: ReadPlan,
1493}
1494
1495#[derive(Default)]
1518enum FilterMaskAccumulator {
1519 #[default]
1520 Empty,
1521 Single(BooleanBuffer),
1522 Combined(BooleanBufferBuilder),
1523}
1524
1525impl FilterMaskAccumulator {
1526 fn append(&mut self, mask: BooleanBuffer) {
1527 *self = match std::mem::take(self) {
1528 Self::Empty => Self::Single(mask),
1529 Self::Single(first) => {
1530 let mut combined = BooleanBufferBuilder::new(first.len() + mask.len());
1531 combined.append_buffer(&first);
1532 combined.append_buffer(&mask);
1533 Self::Combined(combined)
1534 }
1535 Self::Combined(mut combined) => {
1536 combined.append_buffer(&mask);
1537 Self::Combined(combined)
1538 }
1539 };
1540 }
1541
1542 fn finish(self) -> Option<BooleanBuffer> {
1543 match self {
1544 Self::Empty => None,
1545 Self::Single(mask) => Some(mask),
1546 Self::Combined(combined) => Some(combined.build()),
1547 }
1548 }
1549}
1550
1551fn consume_record_batch(array_reader: &mut dyn ArrayReader) -> Result<RecordBatch> {
1553 let array = array_reader.consume_batch()?;
1554 let struct_array = array.as_struct_opt().ok_or_else(|| {
1555 ArrowError::ParquetError("Struct array reader should return struct array".to_string())
1556 })?;
1557 Ok(RecordBatch::from(struct_array))
1558}
1559
1560fn read_mask_batch(
1568 array_reader: &mut dyn ArrayReader,
1569 mask_cursor: &mut MaskCursor,
1570 batch_size: usize,
1571) -> Result<Option<RecordBatch>> {
1572 let mut selected_rows = 0;
1573 let mut filter_mask = FilterMaskAccumulator::default();
1574
1575 while selected_rows < batch_size && !mask_cursor.is_empty() {
1576 let mask_chunk = mask_cursor.next_chunk(batch_size - selected_rows)?;
1577
1578 if mask_chunk.initial_skip > 0 {
1579 let skipped = array_reader.skip_records(mask_chunk.initial_skip)?;
1580 if skipped != mask_chunk.initial_skip {
1581 return Err(general_err!(
1582 "failed to skip rows, expected {}, got {}",
1583 mask_chunk.initial_skip,
1584 skipped
1585 ));
1586 }
1587 }
1588
1589 let mask = mask_cursor.mask_values_for(&mask_chunk)?;
1590 let read = array_reader.read_records(mask_chunk.chunk_rows)?;
1591 if read == 0 {
1592 return Err(general_err!(
1593 "reached end of column while expecting {} rows",
1594 mask_chunk.chunk_rows
1595 ));
1596 }
1597 if read != mask_chunk.chunk_rows {
1598 return Err(general_err!(
1599 "insufficient rows read from array reader - expected {}, got {}",
1600 mask_chunk.chunk_rows,
1601 read
1602 ));
1603 }
1604
1605 filter_mask.append(mask.values().clone());
1606 selected_rows += mask_chunk.selected_rows;
1607 }
1608
1609 if selected_rows == 0 {
1610 return Ok(None);
1611 }
1612
1613 let filter_mask = filter_mask
1614 .finish()
1615 .ok_or_else(|| general_err!("Internal Error: decoded Mask batch has no filter values"))?;
1616 let batch = consume_record_batch(array_reader)?;
1617 let filtered_batch = filter_record_batch(&batch, &BooleanArray::from(filter_mask))?;
1618 if filtered_batch.num_rows() != selected_rows {
1619 return Err(general_err!(
1620 "filtered rows mismatch selection - expected {}, got {}",
1621 selected_rows,
1622 filtered_batch.num_rows()
1623 ));
1624 }
1625
1626 Ok(Some(filtered_batch))
1627}
1628
1629impl Debug for ParquetRecordBatchReader {
1630 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
1631 f.debug_struct("ParquetRecordBatchReader")
1632 .field("array_reader", &"...")
1633 .field("schema", &self.schema)
1634 .field("read_plan", &self.read_plan)
1635 .finish()
1636 }
1637}
1638
1639impl Iterator for ParquetRecordBatchReader {
1640 type Item = Result<RecordBatch, ArrowError>;
1641
1642 fn next(&mut self) -> Option<Self::Item> {
1643 self.next_inner()
1644 .map_err(|arrow_err| arrow_err.into())
1645 .transpose()
1646 }
1647}
1648
1649impl ParquetRecordBatchReader {
1650 fn next_inner(&mut self) -> Result<Option<RecordBatch>> {
1656 let mut read_records = 0;
1657 let batch_size = self.batch_size();
1658 if batch_size == 0 {
1659 return Ok(None);
1660 }
1661 match self.read_plan.row_selection_cursor_mut() {
1662 RowSelectionCursor::Mask(mask_cursor) => {
1663 return read_mask_batch(self.array_reader.as_mut(), mask_cursor, batch_size);
1664 }
1665 RowSelectionCursor::Selectors(selectors_cursor) => {
1666 while read_records < batch_size && !selectors_cursor.is_empty() {
1667 let front = selectors_cursor.next_selector();
1668 if front.skip {
1669 let skipped = self.array_reader.skip_records(front.row_count)?;
1670
1671 if skipped != front.row_count {
1672 return Err(general_err!(
1673 "failed to skip rows, expected {}, got {}",
1674 front.row_count,
1675 skipped
1676 ));
1677 }
1678 continue;
1679 }
1680
1681 if front.row_count == 0 {
1684 continue;
1685 }
1686
1687 let need_read = batch_size - read_records;
1689 let to_read = match front.row_count.checked_sub(need_read) {
1690 Some(remaining) if remaining != 0 => {
1691 selectors_cursor.return_selector(RowSelector::select(remaining));
1694 need_read
1695 }
1696 _ => front.row_count,
1697 };
1698 match self.array_reader.read_records(to_read)? {
1699 0 => break,
1700 rec => read_records += rec,
1701 }
1702 }
1703 }
1704 RowSelectionCursor::All => {
1705 self.array_reader.read_records(batch_size)?;
1706 }
1707 }
1708
1709 let batch = consume_record_batch(self.array_reader.as_mut())?;
1710 Ok(if batch.num_rows() > 0 {
1711 Some(batch)
1712 } else {
1713 None
1714 })
1715 }
1716}
1717
1718impl RecordBatchReader for ParquetRecordBatchReader {
1719 fn schema(&self) -> SchemaRef {
1724 self.schema.clone()
1725 }
1726}
1727
1728impl ParquetRecordBatchReader {
1729 pub fn try_new<T: ChunkReader + 'static>(reader: T, batch_size: usize) -> Result<Self> {
1733 ParquetRecordBatchReaderBuilder::try_new(reader)?
1734 .with_batch_size(batch_size)
1735 .build()
1736 }
1737
1738 pub fn try_new_with_row_groups(
1743 levels: &FieldLevels,
1744 row_groups: &dyn RowGroups,
1745 batch_size: usize,
1746 selection: Option<RowSelection>,
1747 ) -> Result<Self> {
1748 let metrics = ArrowReaderMetrics::disabled();
1750 let array_reader = ArrayReaderBuilder::new(row_groups, &metrics)
1751 .with_batch_size(batch_size)
1752 .with_parquet_metadata(row_groups.metadata())
1753 .build_array_reader(levels.levels.as_ref(), &ProjectionMask::all())?;
1754
1755 let read_plan = ReadPlanBuilder::new(batch_size)
1756 .with_selection(selection)
1757 .build();
1758
1759 Ok(Self {
1760 array_reader,
1761 schema: Arc::new(Schema::new(levels.fields.clone())),
1762 read_plan,
1763 })
1764 }
1765
1766 pub(crate) fn new(array_reader: Box<dyn ArrayReader>, read_plan: ReadPlan) -> Self {
1770 let schema = match array_reader.get_data_type() {
1771 ArrowType::Struct(fields) => Schema::new(fields.clone()),
1772 _ => unreachable!("Struct array reader's data type is not struct!"),
1773 };
1774
1775 Self {
1776 array_reader,
1777 schema: Arc::new(schema),
1778 read_plan,
1779 }
1780 }
1781
1782 #[inline(always)]
1783 pub(crate) fn batch_size(&self) -> usize {
1784 self.read_plan.batch_size()
1785 }
1786}
1787
1788#[cfg(test)]
1789pub(crate) mod tests {
1790 use std::cmp::min;
1791 use std::collections::{HashMap, VecDeque};
1792 use std::fmt::Formatter;
1793 use std::fs::File;
1794 use std::io::Seek;
1795 use std::path::PathBuf;
1796 use std::sync::Arc;
1797
1798 use rand::rngs::StdRng;
1799 use rand::{Rng, RngExt, SeedableRng, random, rng};
1800 use tempfile::tempfile;
1801
1802 use crate::arrow::arrow_reader::{
1803 ArrowPredicateFn, ArrowReaderMetadata, ArrowReaderOptions, ParquetRecordBatchReader,
1804 ParquetRecordBatchReaderBuilder, RowFilter, RowGroupPlan, RowGroupSelection, RowSelection,
1805 RowSelector,
1806 };
1807 use crate::arrow::schema::{
1808 add_encoded_arrow_schema_to_metadata,
1809 virtual_type::{RowGroupIndex, RowNumber},
1810 };
1811 use crate::arrow::{ArrowWriter, ProjectionMask};
1812 use crate::basic::{ConvertedType, Encoding, Repetition, Type as PhysicalType};
1813 use crate::column::reader::decoder::REPETITION_LEVELS_BATCH_SIZE;
1814 use crate::data_type::{
1815 BoolType, ByteArray, ByteArrayType, DataType, DoubleType, FixedLenByteArray,
1816 FixedLenByteArrayType, FloatType, Int32Type, Int64Type, Int96, Int96Type,
1817 };
1818 use crate::errors::Result;
1819 use crate::file::metadata::{PageIndexPolicy, ParquetMetaData, ParquetStatisticsPolicy};
1820 use crate::file::properties::{EnabledStatistics, WriterProperties, WriterVersion};
1821 use crate::file::writer::{SerializedFileWriter, SerializedRowGroupWriter};
1822 use crate::schema::parser::parse_message_type;
1823 use crate::schema::types::{Type, TypePtr};
1824 use crate::util::test_common::rand_gen::RandGen;
1825 use arrow_array::builder::*;
1826 use arrow_array::cast::AsArray;
1827 use arrow_array::types::{
1828 Date32Type, Date64Type, Decimal32Type, Decimal64Type, Decimal128Type, Decimal256Type,
1829 DecimalType, Float16Type, Float32Type, Float64Type, Time32MillisecondType,
1830 Time64MicrosecondType,
1831 };
1832 use arrow_array::*;
1833 use arrow_buffer::{ArrowNativeType, BooleanBuffer, Buffer, IntervalDayTime, NullBuffer, i256};
1834 use arrow_data::{ArrayData, ArrayDataBuilder};
1835 use arrow_schema::{DataType as ArrowDataType, Field, Fields, Schema, SchemaRef, TimeUnit};
1836 use arrow_select::concat::concat_batches;
1837 use bytes::Bytes;
1838 use half::f16;
1839 use num_traits::PrimInt;
1840
1841 fn row_selection(rows: usize) -> RowSelection {
1842 RowSelection::from(vec![RowSelector::select(rows)])
1843 }
1844
1845 #[test]
1846 fn row_group_selection_accessors() {
1847 let row_group = RowGroupSelection::new(3, Some(row_selection(5)));
1848 assert_eq!(row_group.row_group_index(), 3);
1849 assert_eq!(row_group.selection().unwrap().row_count(), 5);
1850
1851 let row_group = RowGroupSelection::new(4, None);
1852 assert_eq!(row_group.row_group_index(), 4);
1853 assert!(row_group.selection().is_none());
1854 }
1855
1856 #[test]
1857 fn row_group_plan_tracks_global_configuration() {
1858 let mut plan = RowGroupPlan::Global {
1859 row_groups: None,
1860 selection: None,
1861 };
1862 plan.set_row_groups(vec![0]);
1863 plan.set_row_groups(vec![1, 2]);
1864 plan.set_row_selection(row_selection(3));
1865 plan.set_row_selection(row_selection(4));
1866
1867 let (row_groups, selection) = plan.into_global().unwrap();
1868 assert_eq!(row_groups, Some(vec![1, 2]));
1869 assert_eq!(selection.unwrap().row_count(), 4);
1870 }
1871
1872 #[test]
1873 fn row_group_plan_replaces_local_configuration() {
1874 let mut plan = RowGroupPlan::Global {
1875 row_groups: None,
1876 selection: None,
1877 };
1878 plan.set_row_group_selections(vec![RowGroupSelection::new(0, None)]);
1879 plan.set_row_group_selections(vec![RowGroupSelection::new(1, None)]);
1880
1881 let RowGroupPlan::PerRowGroup(row_groups) = plan else {
1882 panic!("expected per-row-group plan");
1883 };
1884 assert_eq!(row_groups, vec![RowGroupSelection::new(1, None)]);
1885 }
1886
1887 #[test]
1888 fn row_group_plan_rejects_mixed_configuration() {
1889 let mut row_groups_then_local = RowGroupPlan::Global {
1890 row_groups: None,
1891 selection: None,
1892 };
1893 row_groups_then_local.set_row_groups(vec![0]);
1894 row_groups_then_local.set_row_group_selections(vec![RowGroupSelection::new(0, None)]);
1895 assert!(matches!(row_groups_then_local, RowGroupPlan::Conflicting));
1896 row_groups_then_local.set_row_groups(vec![1]);
1897 row_groups_then_local.set_row_selection(row_selection(1));
1898 row_groups_then_local.set_row_group_selections(vec![RowGroupSelection::new(1, None)]);
1899 assert!(row_groups_then_local.into_global().is_err());
1900
1901 let mut local_then_row_groups =
1902 RowGroupPlan::PerRowGroup(vec![RowGroupSelection::new(0, None)]);
1903 local_then_row_groups.set_row_groups(vec![0]);
1904 assert!(matches!(local_then_row_groups, RowGroupPlan::Conflicting));
1905
1906 let mut local_then_selection =
1907 RowGroupPlan::PerRowGroup(vec![RowGroupSelection::new(0, None)]);
1908 local_then_selection.set_row_selection(row_selection(1));
1909 assert!(matches!(local_then_selection, RowGroupPlan::Conflicting));
1910
1911 let local = RowGroupPlan::PerRowGroup(vec![RowGroupSelection::new(0, None)]);
1912 assert!(local.into_global().is_err());
1913 }
1914
1915 #[test]
1916 fn filter_mask_accumulator_handles_empty_single_and_multiple_chunks() {
1917 let first = BooleanBuffer::from(vec![true, false, false, false]);
1918 let second = BooleanBuffer::from(vec![true]);
1919 let third = BooleanBuffer::from(vec![false, true]);
1920
1921 assert!(super::FilterMaskAccumulator::default().finish().is_none());
1922
1923 let mut single = super::FilterMaskAccumulator::default();
1924 single.append(first.clone());
1925 assert_eq!(single.finish().unwrap(), first);
1926
1927 let mut combined = super::FilterMaskAccumulator::default();
1928 combined.append(BooleanBuffer::from(vec![true, false, false, false]));
1929 combined.append(second);
1930 combined.append(third);
1931 assert_eq!(
1932 combined.finish().unwrap(),
1933 BooleanBuffer::from(vec![true, false, false, false, true, false, true])
1934 );
1935 }
1936
1937 #[test]
1938 fn test_arrow_reader_all_columns() {
1939 let file = get_test_file("parquet/generated_simple_numerics/blogs.parquet");
1940
1941 let builder = ParquetRecordBatchReaderBuilder::try_new(file).unwrap();
1942 let original_schema = Arc::clone(builder.schema());
1943 let reader = builder.build().unwrap();
1944
1945 assert_eq!(original_schema.fields(), reader.schema().fields());
1947 }
1948
1949 #[test]
1950 fn test_reuse_schema() {
1951 let file = get_test_file("parquet/alltypes-java.parquet");
1952
1953 let builder = ParquetRecordBatchReaderBuilder::try_new(file.try_clone().unwrap()).unwrap();
1954 let expected = builder.metadata;
1955 let schema = expected.file_metadata().schema_descr_ptr();
1956
1957 let arrow_options = ArrowReaderOptions::new().with_parquet_schema(schema.clone());
1958 let builder =
1959 ParquetRecordBatchReaderBuilder::try_new_with_options(file, arrow_options).unwrap();
1960
1961 assert_eq!(expected.as_ref(), builder.metadata.as_ref());
1963 }
1964
1965 #[test]
1966 fn test_page_encoding_stats_mask() {
1967 let testdata = arrow::util::test_util::parquet_test_data();
1968 let path = format!("{testdata}/alltypes_tiny_pages.parquet");
1969 let file = File::open(path).unwrap();
1970
1971 let arrow_options = ArrowReaderOptions::new().with_encoding_stats_as_mask(true);
1972 let builder =
1973 ParquetRecordBatchReaderBuilder::try_new_with_options(file, arrow_options).unwrap();
1974
1975 let row_group_metadata = builder.metadata.row_group(0);
1976
1977 let page_encoding_stats = row_group_metadata
1979 .column(0)
1980 .page_encoding_stats_mask()
1981 .unwrap();
1982 assert!(page_encoding_stats.is_only(Encoding::PLAIN));
1983 let page_encoding_stats = row_group_metadata
1984 .column(2)
1985 .page_encoding_stats_mask()
1986 .unwrap();
1987 assert!(page_encoding_stats.is_only(Encoding::PLAIN_DICTIONARY));
1988 }
1989
1990 #[test]
1991 fn test_stats_stats_skipped() {
1992 let testdata = arrow::util::test_util::parquet_test_data();
1993 let path = format!("{testdata}/alltypes_tiny_pages.parquet");
1994 let file = File::open(path).unwrap();
1995
1996 let arrow_options = ArrowReaderOptions::new()
1998 .with_encoding_stats_policy(ParquetStatisticsPolicy::SkipAll)
1999 .with_column_stats_policy(ParquetStatisticsPolicy::SkipAll);
2000 let builder = ParquetRecordBatchReaderBuilder::try_new_with_options(
2001 file.try_clone().unwrap(),
2002 arrow_options,
2003 )
2004 .unwrap();
2005
2006 let row_group_metadata = builder.metadata.row_group(0);
2007 for column in row_group_metadata.columns() {
2008 assert!(column.page_encoding_stats().is_none());
2009 assert!(column.page_encoding_stats_mask().is_none());
2010 assert!(column.statistics().is_none());
2011 }
2012
2013 let arrow_options = ArrowReaderOptions::new()
2015 .with_encoding_stats_as_mask(true)
2016 .with_encoding_stats_policy(ParquetStatisticsPolicy::skip_except(&[0]))
2017 .with_column_stats_policy(ParquetStatisticsPolicy::skip_except(&[0]));
2018 let builder = ParquetRecordBatchReaderBuilder::try_new_with_options(
2019 file.try_clone().unwrap(),
2020 arrow_options,
2021 )
2022 .unwrap();
2023
2024 let row_group_metadata = builder.metadata.row_group(0);
2025 for (idx, column) in row_group_metadata.columns().iter().enumerate() {
2026 assert!(column.page_encoding_stats().is_none());
2027 assert_eq!(column.page_encoding_stats_mask().is_some(), idx == 0);
2028 assert_eq!(column.statistics().is_some(), idx == 0);
2029 }
2030 }
2031
2032 #[test]
2033 fn test_size_stats_stats_skipped() {
2034 let testdata = arrow::util::test_util::parquet_test_data();
2035 let path = format!("{testdata}/repeated_primitive_no_list.parquet");
2036 let file = File::open(path).unwrap();
2037
2038 let arrow_options =
2040 ArrowReaderOptions::new().with_size_stats_policy(ParquetStatisticsPolicy::SkipAll);
2041 let builder = ParquetRecordBatchReaderBuilder::try_new_with_options(
2042 file.try_clone().unwrap(),
2043 arrow_options,
2044 )
2045 .unwrap();
2046
2047 let row_group_metadata = builder.metadata.row_group(0);
2048 for column in row_group_metadata.columns() {
2049 assert!(column.repetition_level_histogram().is_none());
2050 assert!(column.definition_level_histogram().is_none());
2051 assert!(column.unencoded_byte_array_data_bytes().is_none());
2052 }
2053
2054 let arrow_options = ArrowReaderOptions::new()
2056 .with_encoding_stats_as_mask(true)
2057 .with_size_stats_policy(ParquetStatisticsPolicy::skip_except(&[1]));
2058 let builder = ParquetRecordBatchReaderBuilder::try_new_with_options(
2059 file.try_clone().unwrap(),
2060 arrow_options,
2061 )
2062 .unwrap();
2063
2064 let row_group_metadata = builder.metadata.row_group(0);
2065 for (idx, column) in row_group_metadata.columns().iter().enumerate() {
2066 assert_eq!(column.repetition_level_histogram().is_some(), idx == 1);
2067 assert_eq!(column.definition_level_histogram().is_some(), idx == 1);
2068 assert_eq!(column.unencoded_byte_array_data_bytes().is_some(), idx == 1);
2069 }
2070 }
2071
2072 #[test]
2073 fn test_arrow_reader_single_column() {
2074 let file = get_test_file("parquet/generated_simple_numerics/blogs.parquet");
2075
2076 let builder = ParquetRecordBatchReaderBuilder::try_new(file).unwrap();
2077 let original_schema = Arc::clone(builder.schema());
2078
2079 let mask = ProjectionMask::leaves(builder.parquet_schema(), [2]);
2080 let reader = builder.with_projection(mask).build().unwrap();
2081
2082 assert_eq!(1, reader.schema().fields().len());
2084 assert_eq!(original_schema.fields()[1], reader.schema().fields()[0]);
2085 }
2086
2087 #[test]
2088 fn test_arrow_reader_single_column_by_name() {
2089 let file = get_test_file("parquet/generated_simple_numerics/blogs.parquet");
2090
2091 let builder = ParquetRecordBatchReaderBuilder::try_new(file).unwrap();
2092 let original_schema = Arc::clone(builder.schema());
2093
2094 let mask = ProjectionMask::columns(builder.parquet_schema(), ["blog_id"]);
2095 let reader = builder.with_projection(mask).build().unwrap();
2096
2097 assert_eq!(1, reader.schema().fields().len());
2099 assert_eq!(original_schema.fields()[1], reader.schema().fields()[0]);
2100 }
2101
2102 #[test]
2103 fn test_null_column_reader_test() {
2104 let mut file = tempfile::tempfile().unwrap();
2105
2106 let schema = "
2107 message message {
2108 OPTIONAL INT32 int32;
2109 }
2110 ";
2111 let schema = Arc::new(parse_message_type(schema).unwrap());
2112
2113 let def_levels = vec![vec![0, 0, 0], vec![0, 0, 0, 0]];
2114 generate_single_column_file_with_data::<Int32Type>(
2115 &[vec![], vec![]],
2116 Some(&def_levels),
2117 file.try_clone().unwrap(), schema,
2119 Some(Field::new("int32", ArrowDataType::Null, true)),
2120 &Default::default(),
2121 )
2122 .unwrap();
2123
2124 file.rewind().unwrap();
2125
2126 let record_reader = ParquetRecordBatchReader::try_new(file, 2).unwrap();
2127 let batches = record_reader.collect::<Result<Vec<_>, _>>().unwrap();
2128
2129 assert_eq!(batches.len(), 4);
2130 for batch in &batches[0..3] {
2131 assert_eq!(batch.num_rows(), 2);
2132 assert_eq!(batch.num_columns(), 1);
2133 assert_eq!(batch.column(0).null_count(), 2);
2134 }
2135
2136 assert_eq!(batches[3].num_rows(), 1);
2137 assert_eq!(batches[3].num_columns(), 1);
2138 assert_eq!(batches[3].column(0).null_count(), 1);
2139 }
2140
2141 #[test]
2142 #[cfg_attr(miri, ignore)] fn test_primitive_single_column_reader_test() {
2144 run_single_column_reader_tests::<BoolType, _, BoolType>(
2145 2,
2146 ConvertedType::NONE,
2147 None,
2148 |vals| Arc::new(BooleanArray::from_iter(vals.iter().copied())),
2149 &[Encoding::PLAIN, Encoding::RLE, Encoding::RLE_DICTIONARY],
2150 );
2151 run_single_column_reader_tests::<Int32Type, _, Int32Type>(
2152 2,
2153 ConvertedType::NONE,
2154 None,
2155 |vals| Arc::new(Int32Array::from_iter(vals.iter().copied())),
2156 &[
2157 Encoding::PLAIN,
2158 Encoding::RLE_DICTIONARY,
2159 Encoding::DELTA_BINARY_PACKED,
2160 Encoding::BYTE_STREAM_SPLIT,
2161 ],
2162 );
2163 run_single_column_reader_tests::<Int64Type, _, Int64Type>(
2164 2,
2165 ConvertedType::NONE,
2166 None,
2167 |vals| Arc::new(Int64Array::from_iter(vals.iter().copied())),
2168 &[
2169 Encoding::PLAIN,
2170 Encoding::RLE_DICTIONARY,
2171 Encoding::DELTA_BINARY_PACKED,
2172 Encoding::BYTE_STREAM_SPLIT,
2173 ],
2174 );
2175 run_single_column_reader_tests::<FloatType, _, FloatType>(
2176 2,
2177 ConvertedType::NONE,
2178 None,
2179 |vals| Arc::new(Float32Array::from_iter(vals.iter().copied())),
2180 &[Encoding::PLAIN, Encoding::BYTE_STREAM_SPLIT],
2181 );
2182 }
2183
2184 #[test]
2185 #[cfg_attr(miri, ignore)] fn test_unsigned_primitive_single_column_reader_test() {
2187 run_single_column_reader_tests::<Int32Type, _, Int32Type>(
2188 2,
2189 ConvertedType::UINT_32,
2190 Some(ArrowDataType::UInt32),
2191 |vals| {
2192 Arc::new(UInt32Array::from_iter(
2193 vals.iter().map(|x| x.map(|x| x as u32)),
2194 ))
2195 },
2196 &[
2197 Encoding::PLAIN,
2198 Encoding::RLE_DICTIONARY,
2199 Encoding::DELTA_BINARY_PACKED,
2200 ],
2201 );
2202 run_single_column_reader_tests::<Int64Type, _, Int64Type>(
2203 2,
2204 ConvertedType::UINT_64,
2205 Some(ArrowDataType::UInt64),
2206 |vals| {
2207 Arc::new(UInt64Array::from_iter(
2208 vals.iter().map(|x| x.map(|x| x as u64)),
2209 ))
2210 },
2211 &[
2212 Encoding::PLAIN,
2213 Encoding::RLE_DICTIONARY,
2214 Encoding::DELTA_BINARY_PACKED,
2215 ],
2216 );
2217 }
2218
2219 #[test]
2220 fn test_unsigned_roundtrip() {
2221 let schema = Arc::new(Schema::new(vec![
2222 Field::new("uint32", ArrowDataType::UInt32, true),
2223 Field::new("uint64", ArrowDataType::UInt64, true),
2224 ]));
2225
2226 let mut buf = Vec::with_capacity(1024);
2227 let mut writer = ArrowWriter::try_new(&mut buf, schema.clone(), None).unwrap();
2228
2229 let original = RecordBatch::try_new(
2230 schema,
2231 vec![
2232 Arc::new(UInt32Array::from_iter_values([
2233 0,
2234 i32::MAX as u32,
2235 u32::MAX,
2236 ])),
2237 Arc::new(UInt64Array::from_iter_values([
2238 0,
2239 i64::MAX as u64,
2240 u64::MAX,
2241 ])),
2242 ],
2243 )
2244 .unwrap();
2245
2246 writer.write(&original).unwrap();
2247 writer.close().unwrap();
2248
2249 let mut reader = ParquetRecordBatchReader::try_new(Bytes::from(buf), 1024).unwrap();
2250 let ret = reader.next().unwrap().unwrap();
2251 assert_eq!(ret, original);
2252
2253 ret.column(0)
2255 .as_any()
2256 .downcast_ref::<UInt32Array>()
2257 .unwrap();
2258
2259 ret.column(1)
2260 .as_any()
2261 .downcast_ref::<UInt64Array>()
2262 .unwrap();
2263 }
2264
2265 #[test]
2266 fn test_float16_roundtrip() -> Result<()> {
2267 let schema = Arc::new(Schema::new(vec![
2268 Field::new("float16", ArrowDataType::Float16, false),
2269 Field::new("float16-nullable", ArrowDataType::Float16, true),
2270 ]));
2271
2272 let mut buf = Vec::with_capacity(1024);
2273 let mut writer = ArrowWriter::try_new(&mut buf, schema.clone(), None)?;
2274
2275 let original = RecordBatch::try_new(
2276 schema,
2277 vec![
2278 Arc::new(Float16Array::from_iter_values([
2279 f16::EPSILON,
2280 f16::MIN,
2281 f16::MAX,
2282 f16::NAN,
2283 f16::INFINITY,
2284 f16::NEG_INFINITY,
2285 f16::ONE,
2286 f16::NEG_ONE,
2287 f16::ZERO,
2288 f16::NEG_ZERO,
2289 f16::E,
2290 f16::PI,
2291 f16::FRAC_1_PI,
2292 ])),
2293 Arc::new(Float16Array::from(vec![
2294 None,
2295 None,
2296 None,
2297 Some(f16::NAN),
2298 Some(f16::INFINITY),
2299 Some(f16::NEG_INFINITY),
2300 None,
2301 None,
2302 None,
2303 None,
2304 None,
2305 None,
2306 Some(f16::FRAC_1_PI),
2307 ])),
2308 ],
2309 )?;
2310
2311 writer.write(&original)?;
2312 writer.close()?;
2313
2314 let mut reader = ParquetRecordBatchReader::try_new(Bytes::from(buf), 1024)?;
2315 let ret = reader.next().unwrap()?;
2316 assert_eq!(ret, original);
2317
2318 ret.column(0).as_primitive::<Float16Type>();
2320 ret.column(1).as_primitive::<Float16Type>();
2321
2322 Ok(())
2323 }
2324
2325 #[test]
2326 fn test_time_utc_roundtrip() -> Result<()> {
2327 let schema = Arc::new(Schema::new(vec![
2328 Field::new(
2329 "time_millis",
2330 ArrowDataType::Time32(TimeUnit::Millisecond),
2331 true,
2332 )
2333 .with_metadata(HashMap::from_iter(vec![(
2334 "adjusted_to_utc".to_string(),
2335 String::new(),
2336 )])),
2337 Field::new(
2338 "time_micros",
2339 ArrowDataType::Time64(TimeUnit::Microsecond),
2340 true,
2341 )
2342 .with_metadata(HashMap::from_iter(vec![(
2343 "adjusted_to_utc".to_string(),
2344 String::new(),
2345 )])),
2346 ]));
2347
2348 let mut buf = Vec::with_capacity(1024);
2349 let mut writer = ArrowWriter::try_new(&mut buf, schema.clone(), None)?;
2350
2351 let original = RecordBatch::try_new(
2352 schema,
2353 vec![
2354 Arc::new(Time32MillisecondArray::from(vec![
2355 Some(-1),
2356 Some(0),
2357 Some(86_399_000),
2358 Some(86_400_000),
2359 Some(86_401_000),
2360 None,
2361 ])),
2362 Arc::new(Time64MicrosecondArray::from(vec![
2363 Some(-1),
2364 Some(0),
2365 Some(86_399 * 1_000_000),
2366 Some(86_400 * 1_000_000),
2367 Some(86_401 * 1_000_000),
2368 None,
2369 ])),
2370 ],
2371 )?;
2372
2373 writer.write(&original)?;
2374 writer.close()?;
2375
2376 let mut reader = ParquetRecordBatchReader::try_new(Bytes::from(buf), 1024)?;
2377 let ret = reader.next().unwrap()?;
2378 assert_eq!(ret, original);
2379
2380 ret.column(0).as_primitive::<Time32MillisecondType>();
2382 ret.column(1).as_primitive::<Time64MicrosecondType>();
2383
2384 Ok(())
2385 }
2386
2387 #[test]
2388 fn test_date32_roundtrip() -> Result<()> {
2389 use arrow_array::Date32Array;
2390
2391 let schema = Arc::new(Schema::new(vec![Field::new(
2392 "date32",
2393 ArrowDataType::Date32,
2394 false,
2395 )]));
2396
2397 let mut buf = Vec::with_capacity(1024);
2398
2399 let mut writer = ArrowWriter::try_new(&mut buf, schema.clone(), None)?;
2400
2401 let original = RecordBatch::try_new(
2402 schema,
2403 vec![Arc::new(Date32Array::from(vec![
2404 -1_000_000, -100_000, -10_000, -1_000, 0, 1_000, 10_000, 100_000, 1_000_000,
2405 ]))],
2406 )?;
2407
2408 writer.write(&original)?;
2409 writer.close()?;
2410
2411 let mut reader = ParquetRecordBatchReader::try_new(Bytes::from(buf), 1024)?;
2412 let ret = reader.next().unwrap()?;
2413 assert_eq!(ret, original);
2414
2415 ret.column(0).as_primitive::<Date32Type>();
2417
2418 Ok(())
2419 }
2420
2421 #[test]
2422 fn test_date64_roundtrip() -> Result<()> {
2423 use arrow_array::Date64Array;
2424
2425 let schema = Arc::new(Schema::new(vec![
2426 Field::new("small-date64", ArrowDataType::Date64, false),
2427 Field::new("big-date64", ArrowDataType::Date64, false),
2428 Field::new("invalid-date64", ArrowDataType::Date64, false),
2429 ]));
2430
2431 let mut default_buf = Vec::with_capacity(1024);
2432 let mut coerce_buf = Vec::with_capacity(1024);
2433
2434 let coerce_props = WriterProperties::builder().set_coerce_types(true).build();
2435
2436 let mut default_writer = ArrowWriter::try_new(&mut default_buf, schema.clone(), None)?;
2437 let mut coerce_writer =
2438 ArrowWriter::try_new(&mut coerce_buf, schema.clone(), Some(coerce_props))?;
2439
2440 static NUM_MILLISECONDS_IN_DAY: i64 = 1000 * 60 * 60 * 24;
2441
2442 let original = RecordBatch::try_new(
2443 schema,
2444 vec![
2445 Arc::new(Date64Array::from(vec![
2447 -1_000_000 * NUM_MILLISECONDS_IN_DAY,
2448 -1_000 * NUM_MILLISECONDS_IN_DAY,
2449 0,
2450 1_000 * NUM_MILLISECONDS_IN_DAY,
2451 1_000_000 * NUM_MILLISECONDS_IN_DAY,
2452 ])),
2453 Arc::new(Date64Array::from(vec![
2455 -10_000_000_000 * NUM_MILLISECONDS_IN_DAY,
2456 -1_000_000_000 * NUM_MILLISECONDS_IN_DAY,
2457 0,
2458 1_000_000_000 * NUM_MILLISECONDS_IN_DAY,
2459 10_000_000_000 * NUM_MILLISECONDS_IN_DAY,
2460 ])),
2461 Arc::new(Date64Array::from(vec![
2463 -1_000_000 * NUM_MILLISECONDS_IN_DAY + 1,
2464 -1_000 * NUM_MILLISECONDS_IN_DAY + 1,
2465 1,
2466 1_000 * NUM_MILLISECONDS_IN_DAY + 1,
2467 1_000_000 * NUM_MILLISECONDS_IN_DAY + 1,
2468 ])),
2469 ],
2470 )?;
2471
2472 default_writer.write(&original)?;
2473 coerce_writer.write(&original)?;
2474
2475 default_writer.close()?;
2476 coerce_writer.close()?;
2477
2478 let mut default_reader = ParquetRecordBatchReader::try_new(Bytes::from(default_buf), 1024)?;
2479 let mut coerce_reader = ParquetRecordBatchReader::try_new(Bytes::from(coerce_buf), 1024)?;
2480
2481 let default_ret = default_reader.next().unwrap()?;
2482 let coerce_ret = coerce_reader.next().unwrap()?;
2483
2484 assert_eq!(default_ret, original);
2486
2487 assert_eq!(coerce_ret.column(0), original.column(0));
2489 assert_ne!(coerce_ret.column(1), original.column(1));
2490 assert_ne!(coerce_ret.column(2), original.column(2));
2491
2492 default_ret.column(0).as_primitive::<Date64Type>();
2494 coerce_ret.column(0).as_primitive::<Date64Type>();
2495
2496 Ok(())
2497 }
2498 struct RandFixedLenGen {}
2499
2500 impl RandGen<FixedLenByteArrayType> for RandFixedLenGen {
2501 fn r#gen(len: i32) -> FixedLenByteArray {
2502 let mut v = vec![0u8; len as usize];
2503 rng().fill_bytes(&mut v);
2504 ByteArray::from(v).into()
2505 }
2506 }
2507
2508 #[test]
2509 #[cfg_attr(miri, ignore)] fn test_fixed_length_binary_column_reader() {
2511 run_single_column_reader_tests::<FixedLenByteArrayType, _, RandFixedLenGen>(
2512 20,
2513 ConvertedType::NONE,
2514 None,
2515 |vals| {
2516 let mut builder = FixedSizeBinaryBuilder::with_capacity(vals.len(), 20);
2517 for val in vals {
2518 match val {
2519 Some(b) => builder.append_value(b).unwrap(),
2520 None => builder.append_null(),
2521 }
2522 }
2523 Arc::new(builder.finish())
2524 },
2525 &[Encoding::PLAIN, Encoding::RLE_DICTIONARY],
2526 );
2527 }
2528
2529 #[test]
2530 #[cfg_attr(miri, ignore)] fn test_interval_day_time_column_reader() {
2532 run_single_column_reader_tests::<FixedLenByteArrayType, _, RandFixedLenGen>(
2533 12,
2534 ConvertedType::INTERVAL,
2535 None,
2536 |vals| {
2537 Arc::new(
2538 vals.iter()
2539 .map(|x| {
2540 x.as_ref().map(|b| IntervalDayTime {
2541 days: i32::from_le_bytes(b.as_ref()[4..8].try_into().unwrap()),
2542 milliseconds: i32::from_le_bytes(
2543 b.as_ref()[8..12].try_into().unwrap(),
2544 ),
2545 })
2546 })
2547 .collect::<IntervalDayTimeArray>(),
2548 )
2549 },
2550 &[Encoding::PLAIN, Encoding::RLE_DICTIONARY],
2551 );
2552 }
2553
2554 #[test]
2555 #[cfg_attr(miri, ignore)] fn test_int96_single_column_reader_test() {
2557 let encodings = &[Encoding::PLAIN, Encoding::RLE_DICTIONARY];
2558
2559 type TypeHintAndConversionFunction =
2560 (Option<ArrowDataType>, fn(&[Option<Int96>]) -> ArrayRef);
2561
2562 let resolutions: Vec<TypeHintAndConversionFunction> = vec![
2563 (None, |vals: &[Option<Int96>]| {
2565 Arc::new(TimestampNanosecondArray::from_iter(
2566 vals.iter().map(|x| x.map(|x| x.to_nanos())),
2567 )) as ArrayRef
2568 }),
2569 (
2571 Some(ArrowDataType::Timestamp(TimeUnit::Second, None)),
2572 |vals: &[Option<Int96>]| {
2573 Arc::new(TimestampSecondArray::from_iter(
2574 vals.iter().map(|x| x.map(|x| x.to_seconds())),
2575 )) as ArrayRef
2576 },
2577 ),
2578 (
2579 Some(ArrowDataType::Timestamp(TimeUnit::Millisecond, None)),
2580 |vals: &[Option<Int96>]| {
2581 Arc::new(TimestampMillisecondArray::from_iter(
2582 vals.iter().map(|x| x.map(|x| x.to_millis())),
2583 )) as ArrayRef
2584 },
2585 ),
2586 (
2587 Some(ArrowDataType::Timestamp(TimeUnit::Microsecond, None)),
2588 |vals: &[Option<Int96>]| {
2589 Arc::new(TimestampMicrosecondArray::from_iter(
2590 vals.iter().map(|x| x.map(|x| x.to_micros())),
2591 )) as ArrayRef
2592 },
2593 ),
2594 (
2595 Some(ArrowDataType::Timestamp(TimeUnit::Nanosecond, None)),
2596 |vals: &[Option<Int96>]| {
2597 Arc::new(TimestampNanosecondArray::from_iter(
2598 vals.iter().map(|x| x.map(|x| x.to_nanos())),
2599 )) as ArrayRef
2600 },
2601 ),
2602 (
2604 Some(ArrowDataType::Timestamp(
2605 TimeUnit::Second,
2606 Some(Arc::from("-05:00")),
2607 )),
2608 |vals: &[Option<Int96>]| {
2609 Arc::new(
2610 TimestampSecondArray::from_iter(
2611 vals.iter().map(|x| x.map(|x| x.to_seconds())),
2612 )
2613 .with_timezone("-05:00"),
2614 ) as ArrayRef
2615 },
2616 ),
2617 ];
2618
2619 resolutions.iter().for_each(|(arrow_type, converter)| {
2620 run_single_column_reader_tests::<Int96Type, _, Int96Type>(
2621 2,
2622 ConvertedType::NONE,
2623 arrow_type.clone(),
2624 converter,
2625 encodings,
2626 );
2627 })
2628 }
2629
2630 struct RandUtf8Gen {}
2631
2632 impl RandGen<ByteArrayType> for RandUtf8Gen {
2633 fn r#gen(len: i32) -> ByteArray {
2634 Int32Type::r#gen(len).to_string().as_str().into()
2635 }
2636 }
2637
2638 #[test]
2639 #[cfg_attr(miri, ignore)] fn test_utf8_single_column_reader_test() {
2641 fn string_converter<O: OffsetSizeTrait>(vals: &[Option<ByteArray>]) -> ArrayRef {
2642 Arc::new(GenericStringArray::<O>::from_iter(vals.iter().map(|x| {
2643 x.as_ref().map(|b| std::str::from_utf8(b.data()).unwrap())
2644 })))
2645 }
2646
2647 let encodings = &[
2648 Encoding::PLAIN,
2649 Encoding::RLE_DICTIONARY,
2650 Encoding::DELTA_LENGTH_BYTE_ARRAY,
2651 Encoding::DELTA_BYTE_ARRAY,
2652 ];
2653
2654 run_single_column_reader_tests::<ByteArrayType, _, RandUtf8Gen>(
2655 2,
2656 ConvertedType::NONE,
2657 None,
2658 |vals| {
2659 Arc::new(BinaryArray::from_iter(
2660 vals.iter().map(|x| x.as_ref().map(|x| x.data())),
2661 ))
2662 },
2663 encodings,
2664 );
2665
2666 run_single_column_reader_tests::<ByteArrayType, _, RandUtf8Gen>(
2667 2,
2668 ConvertedType::UTF8,
2669 None,
2670 string_converter::<i32>,
2671 encodings,
2672 );
2673
2674 run_single_column_reader_tests::<ByteArrayType, _, RandUtf8Gen>(
2675 2,
2676 ConvertedType::UTF8,
2677 Some(ArrowDataType::Utf8),
2678 string_converter::<i32>,
2679 encodings,
2680 );
2681
2682 run_single_column_reader_tests::<ByteArrayType, _, RandUtf8Gen>(
2683 2,
2684 ConvertedType::UTF8,
2685 Some(ArrowDataType::LargeUtf8),
2686 string_converter::<i64>,
2687 encodings,
2688 );
2689
2690 let small_key_types = [ArrowDataType::Int8, ArrowDataType::UInt8];
2691 for key in &small_key_types {
2692 for encoding in encodings {
2693 let mut opts = TestOptions::new(2, 20, 15).with_null_percent(50);
2694 opts.encoding = *encoding;
2695
2696 let data_type =
2697 ArrowDataType::Dictionary(Box::new(key.clone()), Box::new(ArrowDataType::Utf8));
2698
2699 single_column_reader_test::<ByteArrayType, _, RandUtf8Gen>(
2701 opts,
2702 2,
2703 ConvertedType::UTF8,
2704 Some(data_type.clone()),
2705 move |vals| {
2706 let vals = string_converter::<i32>(vals);
2707 arrow::compute::cast(&vals, &data_type).unwrap()
2708 },
2709 );
2710 }
2711 }
2712
2713 let key_types = [
2714 ArrowDataType::Int16,
2715 ArrowDataType::UInt16,
2716 ArrowDataType::Int32,
2717 ArrowDataType::UInt32,
2718 ArrowDataType::Int64,
2719 ArrowDataType::UInt64,
2720 ];
2721
2722 for key in &key_types {
2723 let data_type =
2724 ArrowDataType::Dictionary(Box::new(key.clone()), Box::new(ArrowDataType::Utf8));
2725
2726 run_single_column_reader_tests::<ByteArrayType, _, RandUtf8Gen>(
2727 2,
2728 ConvertedType::UTF8,
2729 Some(data_type.clone()),
2730 move |vals| {
2731 let vals = string_converter::<i32>(vals);
2732 arrow::compute::cast(&vals, &data_type).unwrap()
2733 },
2734 encodings,
2735 );
2736
2737 let data_type = ArrowDataType::Dictionary(
2738 Box::new(key.clone()),
2739 Box::new(ArrowDataType::LargeUtf8),
2740 );
2741
2742 run_single_column_reader_tests::<ByteArrayType, _, RandUtf8Gen>(
2743 2,
2744 ConvertedType::UTF8,
2745 Some(data_type.clone()),
2746 move |vals| {
2747 let vals = string_converter::<i64>(vals);
2748 arrow::compute::cast(&vals, &data_type).unwrap()
2749 },
2750 encodings,
2751 );
2752 }
2753 }
2754
2755 #[test]
2756 fn test_decimal_nullable_struct() {
2757 let decimals = Decimal256Array::from_iter_values(
2758 [1, 2, 3, 4, 5, 6, 7, 8].into_iter().map(i256::from_i128),
2759 );
2760
2761 let data = ArrayDataBuilder::new(ArrowDataType::Struct(Fields::from(vec![Field::new(
2762 "decimals",
2763 decimals.data_type().clone(),
2764 false,
2765 )])))
2766 .len(8)
2767 .null_bit_buffer(Some(Buffer::from(&[0b11101111])))
2768 .child_data(vec![decimals.into_data()])
2769 .build()
2770 .unwrap();
2771
2772 let written =
2773 RecordBatch::try_from_iter([("struct", Arc::new(StructArray::from(data)) as ArrayRef)])
2774 .unwrap();
2775
2776 let mut buffer = Vec::with_capacity(1024);
2777 let mut writer = ArrowWriter::try_new(&mut buffer, written.schema(), None).unwrap();
2778 writer.write(&written).unwrap();
2779 writer.close().unwrap();
2780
2781 let read = ParquetRecordBatchReader::try_new(Bytes::from(buffer), 3)
2782 .unwrap()
2783 .collect::<Result<Vec<_>, _>>()
2784 .unwrap();
2785
2786 assert_eq!(&written.slice(0, 3), &read[0]);
2787 assert_eq!(&written.slice(3, 3), &read[1]);
2788 assert_eq!(&written.slice(6, 2), &read[2]);
2789 }
2790
2791 #[test]
2792 fn test_int32_nullable_struct() {
2793 let int32 = Int32Array::from_iter_values([1, 2, 3, 4, 5, 6, 7, 8]);
2794 let data = ArrayDataBuilder::new(ArrowDataType::Struct(Fields::from(vec![Field::new(
2795 "int32",
2796 int32.data_type().clone(),
2797 false,
2798 )])))
2799 .len(8)
2800 .null_bit_buffer(Some(Buffer::from(&[0b11101111])))
2801 .child_data(vec![int32.into_data()])
2802 .build()
2803 .unwrap();
2804
2805 let written =
2806 RecordBatch::try_from_iter([("struct", Arc::new(StructArray::from(data)) as ArrayRef)])
2807 .unwrap();
2808
2809 let mut buffer = Vec::with_capacity(1024);
2810 let mut writer = ArrowWriter::try_new(&mut buffer, written.schema(), None).unwrap();
2811 writer.write(&written).unwrap();
2812 writer.close().unwrap();
2813
2814 let read = ParquetRecordBatchReader::try_new(Bytes::from(buffer), 3)
2815 .unwrap()
2816 .collect::<Result<Vec<_>, _>>()
2817 .unwrap();
2818
2819 assert_eq!(&written.slice(0, 3), &read[0]);
2820 assert_eq!(&written.slice(3, 3), &read[1]);
2821 assert_eq!(&written.slice(6, 2), &read[2]);
2822 }
2823
2824 #[test]
2825 fn test_decimal_list() {
2826 let decimals = Decimal128Array::from_iter_values([1, 2, 3, 4, 5, 6, 7, 8]);
2827
2828 let data = ArrayDataBuilder::new(ArrowDataType::List(Arc::new(Field::new_list_field(
2830 decimals.data_type().clone(),
2831 false,
2832 ))))
2833 .len(7)
2834 .add_buffer(Buffer::from_iter([0_i32, 0, 1, 3, 3, 4, 5, 8]))
2835 .null_bit_buffer(Some(Buffer::from(&[0b01010111])))
2836 .child_data(vec![decimals.into_data()])
2837 .build()
2838 .unwrap();
2839
2840 let written =
2841 RecordBatch::try_from_iter([("list", Arc::new(ListArray::from(data)) as ArrayRef)])
2842 .unwrap();
2843
2844 let mut buffer = Vec::with_capacity(1024);
2845 let mut writer = ArrowWriter::try_new(&mut buffer, written.schema(), None).unwrap();
2846 writer.write(&written).unwrap();
2847 writer.close().unwrap();
2848
2849 let read = ParquetRecordBatchReader::try_new(Bytes::from(buffer), 3)
2850 .unwrap()
2851 .collect::<Result<Vec<_>, _>>()
2852 .unwrap();
2853
2854 assert_eq!(&written.slice(0, 3), &read[0]);
2855 assert_eq!(&written.slice(3, 3), &read[1]);
2856 assert_eq!(&written.slice(6, 1), &read[2]);
2857 }
2858
2859 #[test]
2860 fn test_read_decimal_file() {
2861 use arrow_array::Decimal128Array;
2862 let testdata = arrow::util::test_util::parquet_test_data();
2863 let file_variants = vec![
2864 ("byte_array", 4),
2865 ("fixed_length", 25),
2866 ("int32", 4),
2867 ("int64", 10),
2868 ];
2869 for (prefix, target_precision) in file_variants {
2870 let path = format!("{testdata}/{prefix}_decimal.parquet");
2871 let file = File::open(path).unwrap();
2872 let mut record_reader = ParquetRecordBatchReader::try_new(file, 32).unwrap();
2873
2874 let batch = record_reader.next().unwrap().unwrap();
2875 assert_eq!(batch.num_rows(), 24);
2876 let col = batch
2877 .column(0)
2878 .as_any()
2879 .downcast_ref::<Decimal128Array>()
2880 .unwrap();
2881
2882 let expected = 1..25;
2883
2884 assert_eq!(col.precision(), target_precision);
2885 assert_eq!(col.scale(), 2);
2886
2887 for (i, v) in expected.enumerate() {
2888 assert_eq!(col.value(i), v * 100_i128);
2889 }
2890 }
2891 }
2892
2893 #[test]
2894 #[cfg_attr(miri, ignore)] fn test_read_float16_nonzeros_file() {
2896 use arrow_array::Float16Array;
2897 let testdata = arrow::util::test_util::parquet_test_data();
2898 let path = format!("{testdata}/float16_nonzeros_and_nans.parquet");
2900 let file = File::open(path).unwrap();
2901 let mut record_reader = ParquetRecordBatchReader::try_new(file, 32).unwrap();
2902
2903 let batch = record_reader.next().unwrap().unwrap();
2904 assert_eq!(batch.num_rows(), 8);
2905 let col = batch
2906 .column(0)
2907 .as_any()
2908 .downcast_ref::<Float16Array>()
2909 .unwrap();
2910
2911 let f16_two = f16::ONE + f16::ONE;
2912
2913 assert_eq!(col.null_count(), 1);
2914 assert!(col.is_null(0));
2915 assert_eq!(col.value(1), f16::ONE);
2916 assert_eq!(col.value(2), -f16_two);
2917 assert!(col.value(3).is_nan());
2918 assert_eq!(col.value(4), f16::ZERO);
2919 assert!(col.value(4).is_sign_positive());
2920 assert_eq!(col.value(5), f16::NEG_ONE);
2921 assert_eq!(col.value(6), f16::NEG_ZERO);
2922 assert!(col.value(6).is_sign_negative());
2923 assert_eq!(col.value(7), f16_two);
2924 }
2925
2926 #[test]
2927 fn test_read_float16_zeros_file() {
2928 use arrow_array::Float16Array;
2929 let testdata = arrow::util::test_util::parquet_test_data();
2930 let path = format!("{testdata}/float16_zeros_and_nans.parquet");
2932 let file = File::open(path).unwrap();
2933 let mut record_reader = ParquetRecordBatchReader::try_new(file, 32).unwrap();
2934
2935 let batch = record_reader.next().unwrap().unwrap();
2936 assert_eq!(batch.num_rows(), 3);
2937 let col = batch
2938 .column(0)
2939 .as_any()
2940 .downcast_ref::<Float16Array>()
2941 .unwrap();
2942
2943 assert_eq!(col.null_count(), 1);
2944 assert!(col.is_null(0));
2945 assert_eq!(col.value(1), f16::ZERO);
2946 assert!(col.value(1).is_sign_positive());
2947 assert!(col.value(2).is_nan());
2948 }
2949
2950 #[test]
2951 #[cfg_attr(miri, ignore)] fn test_read_float32_float64_byte_stream_split() {
2953 let path = format!(
2954 "{}/byte_stream_split.zstd.parquet",
2955 arrow::util::test_util::parquet_test_data(),
2956 );
2957 let file = File::open(path).unwrap();
2958 let record_reader = ParquetRecordBatchReader::try_new(file, 128).unwrap();
2959
2960 let mut row_count = 0;
2961 for batch in record_reader {
2962 let batch = batch.unwrap();
2963 row_count += batch.num_rows();
2964 let f32_col = batch.column(0).as_primitive::<Float32Type>();
2965 let f64_col = batch.column(1).as_primitive::<Float64Type>();
2966
2967 for &x in f32_col.values() {
2969 assert!(x > -10.0);
2970 assert!(x < 10.0);
2971 }
2972 for &x in f64_col.values() {
2973 assert!(x > -10.0);
2974 assert!(x < 10.0);
2975 }
2976 }
2977 assert_eq!(row_count, 300);
2978 }
2979
2980 #[test]
2981 #[cfg_attr(miri, ignore)] fn test_read_extended_byte_stream_split() {
2983 let path = format!(
2984 "{}/byte_stream_split_extended.gzip.parquet",
2985 arrow::util::test_util::parquet_test_data(),
2986 );
2987 let file = File::open(path).unwrap();
2988 let record_reader = ParquetRecordBatchReader::try_new(file, 128).unwrap();
2989
2990 let mut row_count = 0;
2991 for batch in record_reader {
2992 let batch = batch.unwrap();
2993 row_count += batch.num_rows();
2994
2995 let f16_col = batch.column(0).as_primitive::<Float16Type>();
2997 let f16_bss = batch.column(1).as_primitive::<Float16Type>();
2998 assert_eq!(f16_col.len(), f16_bss.len());
2999 f16_col
3000 .iter()
3001 .zip(f16_bss.iter())
3002 .for_each(|(l, r)| assert_eq!(l.unwrap(), r.unwrap()));
3003
3004 let f32_col = batch.column(2).as_primitive::<Float32Type>();
3006 let f32_bss = batch.column(3).as_primitive::<Float32Type>();
3007 assert_eq!(f32_col.len(), f32_bss.len());
3008 f32_col
3009 .iter()
3010 .zip(f32_bss.iter())
3011 .for_each(|(l, r)| assert_eq!(l.unwrap(), r.unwrap()));
3012
3013 let f64_col = batch.column(4).as_primitive::<Float64Type>();
3015 let f64_bss = batch.column(5).as_primitive::<Float64Type>();
3016 assert_eq!(f64_col.len(), f64_bss.len());
3017 f64_col
3018 .iter()
3019 .zip(f64_bss.iter())
3020 .for_each(|(l, r)| assert_eq!(l.unwrap(), r.unwrap()));
3021
3022 let i32_col = batch.column(6).as_primitive::<types::Int32Type>();
3024 let i32_bss = batch.column(7).as_primitive::<types::Int32Type>();
3025 assert_eq!(i32_col.len(), i32_bss.len());
3026 i32_col
3027 .iter()
3028 .zip(i32_bss.iter())
3029 .for_each(|(l, r)| assert_eq!(l.unwrap(), r.unwrap()));
3030
3031 let i64_col = batch.column(8).as_primitive::<types::Int64Type>();
3033 let i64_bss = batch.column(9).as_primitive::<types::Int64Type>();
3034 assert_eq!(i64_col.len(), i64_bss.len());
3035 i64_col
3036 .iter()
3037 .zip(i64_bss.iter())
3038 .for_each(|(l, r)| assert_eq!(l.unwrap(), r.unwrap()));
3039
3040 let flba_col = batch.column(10).as_fixed_size_binary();
3042 let flba_bss = batch.column(11).as_fixed_size_binary();
3043 assert_eq!(flba_col.len(), flba_bss.len());
3044 flba_col
3045 .iter()
3046 .zip(flba_bss.iter())
3047 .for_each(|(l, r)| assert_eq!(l.unwrap(), r.unwrap()));
3048
3049 let dec_col = batch.column(12).as_primitive::<Decimal128Type>();
3051 let dec_bss = batch.column(13).as_primitive::<Decimal128Type>();
3052 assert_eq!(dec_col.len(), dec_bss.len());
3053 dec_col
3054 .iter()
3055 .zip(dec_bss.iter())
3056 .for_each(|(l, r)| assert_eq!(l.unwrap(), r.unwrap()));
3057 }
3058 assert_eq!(row_count, 200);
3059 }
3060
3061 #[test]
3062 fn test_read_incorrect_map_schema_file() {
3063 let testdata = arrow::util::test_util::parquet_test_data();
3064 let path = format!("{testdata}/incorrect_map_schema.parquet");
3066 let file = File::open(path).unwrap();
3067 let mut record_reader = ParquetRecordBatchReader::try_new(file, 32).unwrap();
3068
3069 let batch = record_reader.next().unwrap().unwrap();
3070 assert_eq!(batch.num_rows(), 1);
3071
3072 let expected_schema = Schema::new(vec![Field::new(
3073 "my_map",
3074 ArrowDataType::Map(
3075 Arc::new(Field::new(
3076 "key_value",
3077 ArrowDataType::Struct(Fields::from(vec![
3078 Field::new("key", ArrowDataType::Utf8, false),
3079 Field::new("value", ArrowDataType::Utf8, true),
3080 ])),
3081 false,
3082 )),
3083 false,
3084 ),
3085 true,
3086 )]);
3087 assert_eq!(batch.schema().as_ref(), &expected_schema);
3088
3089 assert_eq!(batch.num_rows(), 1);
3090 assert_eq!(batch.column(0).null_count(), 0);
3091 assert_eq!(
3092 batch.column(0).as_map().keys().as_ref(),
3093 &StringArray::from(vec!["parent", "name"])
3094 );
3095 assert_eq!(
3096 batch.column(0).as_map().values().as_ref(),
3097 &StringArray::from(vec!["another", "report"])
3098 );
3099 }
3100
3101 #[test]
3102 fn test_read_dict_fixed_size_binary() {
3103 let schema = Arc::new(Schema::new(vec![Field::new(
3104 "a",
3105 ArrowDataType::Dictionary(
3106 Box::new(ArrowDataType::UInt8),
3107 Box::new(ArrowDataType::FixedSizeBinary(8)),
3108 ),
3109 true,
3110 )]));
3111 let keys = UInt8Array::from_iter_values(vec![0, 0, 1]);
3112 let values = FixedSizeBinaryArray::try_from_iter(
3113 vec![
3114 (0u8..8u8).collect::<Vec<u8>>(),
3115 (24u8..32u8).collect::<Vec<u8>>(),
3116 ]
3117 .into_iter(),
3118 )
3119 .unwrap();
3120 let arr = UInt8DictionaryArray::new(keys, Arc::new(values));
3121 let batch = RecordBatch::try_new(schema, vec![Arc::new(arr)]).unwrap();
3122
3123 let mut buffer = Vec::with_capacity(1024);
3124 let mut writer = ArrowWriter::try_new(&mut buffer, batch.schema(), None).unwrap();
3125 writer.write(&batch).unwrap();
3126 writer.close().unwrap();
3127 let read = ParquetRecordBatchReader::try_new(Bytes::from(buffer), 3)
3128 .unwrap()
3129 .collect::<Result<Vec<_>, _>>()
3130 .unwrap();
3131
3132 assert_eq!(read.len(), 1);
3133 assert_eq!(&batch, &read[0])
3134 }
3135
3136 #[test]
3137 fn test_read_nullable_structs_with_binary_dict_as_first_child_column() {
3138 let struct_fields = Fields::from(vec![
3145 Field::new(
3146 "city",
3147 ArrowDataType::Dictionary(
3148 Box::new(ArrowDataType::UInt8),
3149 Box::new(ArrowDataType::Utf8),
3150 ),
3151 true,
3152 ),
3153 Field::new("name", ArrowDataType::Utf8, true),
3154 ]);
3155 let schema = Arc::new(Schema::new(vec![Field::new(
3156 "items",
3157 ArrowDataType::Struct(struct_fields.clone()),
3158 true,
3159 )]));
3160
3161 let items_arr = StructArray::new(
3162 struct_fields,
3163 vec![
3164 Arc::new(DictionaryArray::new(
3165 UInt8Array::from_iter_values(vec![0, 1, 1, 0, 2]),
3166 Arc::new(StringArray::from_iter_values(vec![
3167 "quebec",
3168 "fredericton",
3169 "halifax",
3170 ])),
3171 )),
3172 Arc::new(StringArray::from_iter_values(vec![
3173 "albert", "terry", "lance", "", "tim",
3174 ])),
3175 ],
3176 Some(NullBuffer::from_iter(vec![true, true, true, false, true])),
3177 );
3178
3179 let batch = RecordBatch::try_new(schema, vec![Arc::new(items_arr)]).unwrap();
3180 let mut buffer = Vec::with_capacity(1024);
3181 let mut writer = ArrowWriter::try_new(&mut buffer, batch.schema(), None).unwrap();
3182 writer.write(&batch).unwrap();
3183 writer.close().unwrap();
3184 let read = ParquetRecordBatchReader::try_new(Bytes::from(buffer), 8)
3185 .unwrap()
3186 .collect::<Result<Vec<_>, _>>()
3187 .unwrap();
3188
3189 assert_eq!(read.len(), 1);
3190 assert_eq!(&batch, &read[0])
3191 }
3192
3193 #[derive(Clone)]
3195 struct TestOptions {
3196 num_row_groups: usize,
3199 num_rows: usize,
3201 record_batch_size: usize,
3203 null_percent: Option<usize>,
3205 write_batch_size: usize,
3210 max_data_page_size: usize,
3212 max_dict_page_size: usize,
3214 writer_version: WriterVersion,
3216 enabled_statistics: EnabledStatistics,
3218 encoding: Encoding,
3220 row_selections: Option<(RowSelection, usize)>,
3222 row_filter: Option<Vec<bool>>,
3224 limit: Option<usize>,
3226 offset: Option<usize>,
3228 }
3229
3230 impl std::fmt::Debug for TestOptions {
3232 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
3233 f.debug_struct("TestOptions")
3234 .field("num_row_groups", &self.num_row_groups)
3235 .field("num_rows", &self.num_rows)
3236 .field("record_batch_size", &self.record_batch_size)
3237 .field("null_percent", &self.null_percent)
3238 .field("write_batch_size", &self.write_batch_size)
3239 .field("max_data_page_size", &self.max_data_page_size)
3240 .field("max_dict_page_size", &self.max_dict_page_size)
3241 .field("writer_version", &self.writer_version)
3242 .field("enabled_statistics", &self.enabled_statistics)
3243 .field("encoding", &self.encoding)
3244 .field("row_selections", &self.row_selections.is_some())
3245 .field("row_filter", &self.row_filter.is_some())
3246 .field("limit", &self.limit)
3247 .field("offset", &self.offset)
3248 .finish()
3249 }
3250 }
3251
3252 impl Default for TestOptions {
3253 fn default() -> Self {
3254 Self {
3255 num_row_groups: 2,
3256 num_rows: 100,
3257 record_batch_size: 15,
3258 null_percent: None,
3259 write_batch_size: 64,
3260 max_data_page_size: 1024 * 1024,
3261 max_dict_page_size: 1024 * 1024,
3262 writer_version: WriterVersion::PARQUET_1_0,
3263 enabled_statistics: EnabledStatistics::Page,
3264 encoding: Encoding::PLAIN,
3265 row_selections: None,
3266 row_filter: None,
3267 limit: None,
3268 offset: None,
3269 }
3270 }
3271 }
3272
3273 impl TestOptions {
3274 fn new(num_row_groups: usize, num_rows: usize, record_batch_size: usize) -> Self {
3275 Self {
3276 num_row_groups,
3277 num_rows,
3278 record_batch_size,
3279 ..Default::default()
3280 }
3281 }
3282
3283 fn with_null_percent(self, null_percent: usize) -> Self {
3284 Self {
3285 null_percent: Some(null_percent),
3286 ..self
3287 }
3288 }
3289
3290 fn with_max_data_page_size(self, max_data_page_size: usize) -> Self {
3291 Self {
3292 max_data_page_size,
3293 ..self
3294 }
3295 }
3296
3297 fn with_max_dict_page_size(self, max_dict_page_size: usize) -> Self {
3298 Self {
3299 max_dict_page_size,
3300 ..self
3301 }
3302 }
3303
3304 fn with_enabled_statistics(self, enabled_statistics: EnabledStatistics) -> Self {
3305 Self {
3306 enabled_statistics,
3307 ..self
3308 }
3309 }
3310
3311 fn with_row_selections(self) -> Self {
3312 assert!(self.row_filter.is_none(), "Must set row selection first");
3313
3314 let mut rng = rng();
3315 let step = rng.random_range(self.record_batch_size..self.num_rows);
3316 let row_selections = create_test_selection(
3317 step,
3318 self.num_row_groups * self.num_rows,
3319 rng.random::<bool>(),
3320 );
3321 Self {
3322 row_selections: Some(row_selections),
3323 ..self
3324 }
3325 }
3326
3327 fn with_row_filter(self) -> Self {
3328 let row_count = match &self.row_selections {
3329 Some((_, count)) => *count,
3330 None => self.num_row_groups * self.num_rows,
3331 };
3332
3333 let mut rng = rng();
3334 Self {
3335 row_filter: Some((0..row_count).map(|_| rng.random_bool(0.9)).collect()),
3336 ..self
3337 }
3338 }
3339
3340 fn with_limit(self, limit: usize) -> Self {
3341 Self {
3342 limit: Some(limit),
3343 ..self
3344 }
3345 }
3346
3347 fn with_offset(self, offset: usize) -> Self {
3348 Self {
3349 offset: Some(offset),
3350 ..self
3351 }
3352 }
3353
3354 fn writer_props(&self) -> WriterProperties {
3355 let builder = WriterProperties::builder()
3356 .set_data_page_size_limit(self.max_data_page_size)
3357 .set_write_batch_size(self.write_batch_size)
3358 .set_writer_version(self.writer_version)
3359 .set_statistics_enabled(self.enabled_statistics);
3360
3361 let builder = match self.encoding {
3362 Encoding::RLE_DICTIONARY | Encoding::PLAIN_DICTIONARY => builder
3363 .set_dictionary_enabled(true)
3364 .set_dictionary_page_size_limit(self.max_dict_page_size),
3365 _ => builder
3366 .set_dictionary_enabled(false)
3367 .set_encoding(self.encoding),
3368 };
3369
3370 builder.build()
3371 }
3372 }
3373
3374 fn run_single_column_reader_tests<T, F, G>(
3381 rand_max: i32,
3382 converted_type: ConvertedType,
3383 arrow_type: Option<ArrowDataType>,
3384 converter: F,
3385 encodings: &[Encoding],
3386 ) where
3387 T: DataType,
3388 G: RandGen<T>,
3389 F: Fn(&[Option<T::T>]) -> ArrayRef,
3390 {
3391 let all_options = vec![
3392 TestOptions::new(2, 100, 15),
3395 TestOptions::new(3, 25, 5),
3400 TestOptions::new(4, 100, 25),
3404 TestOptions::new(3, 256, 73).with_max_data_page_size(128),
3406 TestOptions::new(3, 256, 57).with_max_dict_page_size(128),
3408 TestOptions::new(2, 256, 127).with_null_percent(0),
3410 TestOptions::new(2, 256, 93).with_null_percent(25),
3412 TestOptions::new(4, 100, 25).with_limit(0),
3414 TestOptions::new(4, 100, 25).with_limit(50),
3416 TestOptions::new(4, 100, 25).with_limit(10),
3418 TestOptions::new(4, 100, 25).with_limit(101),
3420 TestOptions::new(4, 100, 25).with_offset(30).with_limit(20),
3422 TestOptions::new(4, 100, 25).with_offset(20).with_limit(80),
3424 TestOptions::new(4, 100, 25).with_offset(20).with_limit(81),
3426 TestOptions::new(2, 256, 91)
3428 .with_null_percent(25)
3429 .with_enabled_statistics(EnabledStatistics::Chunk),
3430 TestOptions::new(2, 256, 91)
3432 .with_null_percent(25)
3433 .with_enabled_statistics(EnabledStatistics::None),
3434 TestOptions::new(2, 128, 91)
3436 .with_null_percent(100)
3437 .with_enabled_statistics(EnabledStatistics::None),
3438 TestOptions::new(2, 100, 15).with_row_selections(),
3443 TestOptions::new(3, 25, 5).with_row_selections(),
3448 TestOptions::new(4, 100, 25).with_row_selections(),
3452 TestOptions::new(3, 256, 73)
3454 .with_max_data_page_size(128)
3455 .with_row_selections(),
3456 TestOptions::new(3, 256, 57)
3458 .with_max_dict_page_size(128)
3459 .with_row_selections(),
3460 TestOptions::new(2, 256, 127)
3462 .with_null_percent(0)
3463 .with_row_selections(),
3464 TestOptions::new(2, 256, 93)
3466 .with_null_percent(25)
3467 .with_row_selections(),
3468 TestOptions::new(2, 256, 93)
3470 .with_null_percent(25)
3471 .with_row_selections()
3472 .with_limit(10),
3473 TestOptions::new(2, 256, 93)
3475 .with_null_percent(25)
3476 .with_row_selections()
3477 .with_offset(20)
3478 .with_limit(10),
3479 TestOptions::new(4, 100, 25).with_row_filter(),
3483 TestOptions::new(4, 100, 25)
3485 .with_row_selections()
3486 .with_row_filter(),
3487 TestOptions::new(2, 256, 93)
3489 .with_null_percent(25)
3490 .with_max_data_page_size(10)
3491 .with_row_filter(),
3492 TestOptions::new(2, 256, 93)
3494 .with_null_percent(25)
3495 .with_max_data_page_size(10)
3496 .with_row_selections()
3497 .with_row_filter(),
3498 TestOptions::new(2, 256, 93)
3500 .with_enabled_statistics(EnabledStatistics::None)
3501 .with_max_data_page_size(10)
3502 .with_row_selections(),
3503 ];
3504
3505 all_options.into_iter().for_each(|opts| {
3506 for writer_version in [WriterVersion::PARQUET_1_0, WriterVersion::PARQUET_2_0] {
3507 for encoding in encodings {
3508 let opts = TestOptions {
3509 writer_version,
3510 encoding: *encoding,
3511 ..opts.clone()
3512 };
3513
3514 single_column_reader_test::<T, _, G>(
3515 opts,
3516 rand_max,
3517 converted_type,
3518 arrow_type.clone(),
3519 &converter,
3520 )
3521 }
3522 }
3523 });
3524 }
3525
3526 fn single_column_reader_test<T, F, G>(
3530 opts: TestOptions,
3531 rand_max: i32,
3532 converted_type: ConvertedType,
3533 arrow_type: Option<ArrowDataType>,
3534 converter: F,
3535 ) where
3536 T: DataType,
3537 G: RandGen<T>,
3538 F: Fn(&[Option<T::T>]) -> ArrayRef,
3539 {
3540 println!(
3542 "Running type {:?} single_column_reader_test ConvertedType::{}/ArrowType::{:?} with Options: {:?}",
3543 T::get_physical_type(),
3544 converted_type,
3545 arrow_type,
3546 opts
3547 );
3548
3549 let (repetition, def_levels) = match opts.null_percent.as_ref() {
3551 Some(null_percent) => {
3552 let mut rng = rng();
3553
3554 let def_levels: Vec<Vec<i16>> = (0..opts.num_row_groups)
3555 .map(|_| {
3556 std::iter::from_fn(|| {
3557 Some((rng.next_u32() as usize % 100 >= *null_percent) as i16)
3558 })
3559 .take(opts.num_rows)
3560 .collect()
3561 })
3562 .collect();
3563 (Repetition::OPTIONAL, Some(def_levels))
3564 }
3565 None => (Repetition::REQUIRED, None),
3566 };
3567
3568 let values: Vec<Vec<T::T>> = (0..opts.num_row_groups)
3570 .map(|idx| {
3571 let null_count = match def_levels.as_ref() {
3572 Some(d) => d[idx].iter().filter(|x| **x == 0).count(),
3573 None => 0,
3574 };
3575 G::gen_vec(rand_max, opts.num_rows - null_count)
3576 })
3577 .collect();
3578
3579 let len = match T::get_physical_type() {
3580 crate::basic::Type::FIXED_LEN_BYTE_ARRAY => rand_max,
3581 crate::basic::Type::INT96 => 12,
3582 _ => -1,
3583 };
3584
3585 let fields = vec![Arc::new(
3586 Type::primitive_type_builder("leaf", T::get_physical_type())
3587 .with_repetition(repetition)
3588 .with_converted_type(converted_type)
3589 .with_length(len)
3590 .build()
3591 .unwrap(),
3592 )];
3593
3594 let schema = Arc::new(
3595 Type::group_type_builder("test_schema")
3596 .with_fields(fields)
3597 .build()
3598 .unwrap(),
3599 );
3600
3601 let arrow_field = arrow_type.map(|t| Field::new("leaf", t, false));
3602
3603 let mut file = tempfile::tempfile().unwrap();
3604
3605 generate_single_column_file_with_data::<T>(
3606 &values,
3607 def_levels.as_ref(),
3608 file.try_clone().unwrap(), schema,
3610 arrow_field,
3611 &opts,
3612 )
3613 .unwrap();
3614
3615 file.rewind().unwrap();
3616
3617 let options = ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::from(
3618 opts.enabled_statistics == EnabledStatistics::Page,
3619 ));
3620
3621 let mut builder =
3622 ParquetRecordBatchReaderBuilder::try_new_with_options(file, options).unwrap();
3623
3624 let expected_data = match opts.row_selections {
3625 Some((selections, row_count)) => {
3626 let mut without_skip_data = gen_expected_data::<T>(def_levels.as_ref(), &values);
3627
3628 let mut skip_data: Vec<Option<T::T>> = vec![];
3629 let dequeue: VecDeque<RowSelector> = selections.clone().into();
3630 for select in dequeue {
3631 if select.skip {
3632 without_skip_data.drain(0..select.row_count);
3633 } else {
3634 skip_data.extend(without_skip_data.drain(0..select.row_count));
3635 }
3636 }
3637 builder = builder.with_row_selection(selections);
3638
3639 assert_eq!(skip_data.len(), row_count);
3640 skip_data
3641 }
3642 None => {
3643 let expected_data = gen_expected_data::<T>(def_levels.as_ref(), &values);
3645 assert_eq!(expected_data.len(), opts.num_rows * opts.num_row_groups);
3646 expected_data
3647 }
3648 };
3649
3650 let mut expected_data = match opts.row_filter {
3651 Some(filter) => {
3652 let expected_data = expected_data
3653 .into_iter()
3654 .zip(filter.iter())
3655 .filter_map(|(d, f)| f.then(|| d))
3656 .collect();
3657
3658 let mut filter_offset = 0;
3659 let filter = RowFilter::new(vec![Box::new(ArrowPredicateFn::new(
3660 ProjectionMask::all(),
3661 move |b| {
3662 let array = BooleanArray::from_iter(
3663 filter
3664 .iter()
3665 .skip(filter_offset)
3666 .take(b.num_rows())
3667 .map(|x| Some(*x)),
3668 );
3669 filter_offset += b.num_rows();
3670 Ok(array)
3671 },
3672 ))]);
3673
3674 builder = builder.with_row_filter(filter);
3675 expected_data
3676 }
3677 None => expected_data,
3678 };
3679
3680 if let Some(offset) = opts.offset {
3681 builder = builder.with_offset(offset);
3682 expected_data = expected_data.into_iter().skip(offset).collect();
3683 }
3684
3685 if let Some(limit) = opts.limit {
3686 builder = builder.with_limit(limit);
3687 expected_data = expected_data.into_iter().take(limit).collect();
3688 }
3689
3690 let mut record_reader = builder
3691 .with_batch_size(opts.record_batch_size)
3692 .build()
3693 .unwrap();
3694
3695 let mut total_read = 0;
3696 loop {
3697 let maybe_batch = record_reader.next();
3698 if total_read < expected_data.len() {
3699 let end = min(total_read + opts.record_batch_size, expected_data.len());
3700 let batch = maybe_batch.unwrap().unwrap();
3701 assert_eq!(end - total_read, batch.num_rows());
3702
3703 let a = converter(&expected_data[total_read..end]);
3704 let b = batch.column(0);
3705
3706 assert_eq!(a.data_type(), b.data_type());
3707 assert_eq!(a.to_data(), b.to_data());
3708 assert_eq!(
3709 a.as_any().type_id(),
3710 b.as_any().type_id(),
3711 "incorrect type ids"
3712 );
3713
3714 total_read = end;
3715 } else {
3716 assert!(maybe_batch.is_none());
3717 break;
3718 }
3719 }
3720 }
3721
3722 fn gen_expected_data<T: DataType>(
3723 def_levels: Option<&Vec<Vec<i16>>>,
3724 values: &[Vec<T::T>],
3725 ) -> Vec<Option<T::T>> {
3726 let data: Vec<Option<T::T>> = match def_levels {
3727 Some(levels) => {
3728 let mut values_iter = values.iter().flatten();
3729 levels
3730 .iter()
3731 .flatten()
3732 .map(|d| match d {
3733 1 => Some(values_iter.next().cloned().unwrap()),
3734 0 => None,
3735 _ => unreachable!(),
3736 })
3737 .collect()
3738 }
3739 None => values.iter().flatten().cloned().map(Some).collect(),
3740 };
3741 data
3742 }
3743
3744 fn generate_single_column_file_with_data<T: DataType>(
3745 values: &[Vec<T::T>],
3746 def_levels: Option<&Vec<Vec<i16>>>,
3747 file: File,
3748 schema: TypePtr,
3749 field: Option<Field>,
3750 opts: &TestOptions,
3751 ) -> Result<ParquetMetaData> {
3752 let mut writer_props = opts.writer_props();
3753 if let Some(field) = field {
3754 let arrow_schema = Schema::new(vec![field]);
3755 add_encoded_arrow_schema_to_metadata(&arrow_schema, &mut writer_props);
3756 }
3757
3758 let mut writer = SerializedFileWriter::new(file, schema, Arc::new(writer_props))?;
3759
3760 for (idx, v) in values.iter().enumerate() {
3761 let def_levels = def_levels.map(|d| d[idx].as_slice());
3762 let mut row_group_writer = writer.next_row_group()?;
3763 {
3764 let mut column_writer = row_group_writer
3765 .next_column()?
3766 .expect("Column writer is none!");
3767
3768 column_writer
3769 .typed::<T>()
3770 .write_batch(v, def_levels, None)?;
3771
3772 column_writer.close()?;
3773 }
3774 row_group_writer.close()?;
3775 }
3776
3777 writer.close()
3778 }
3779
3780 fn get_test_file(file_name: &str) -> File {
3781 let path = PathBuf::from(arrow::util::test_util::arrow_test_data()).join(file_name);
3782
3783 File::open(path.as_path()).expect("File not found!")
3784 }
3785
3786 #[cfg_attr(miri, ignore)] #[test]
3788 fn test_read_structs() {
3789 let testdata = arrow::util::test_util::parquet_test_data();
3793 let path = format!("{testdata}/nested_structs.rust.parquet");
3794 let file = File::open(&path).unwrap();
3795 let record_batch_reader = ParquetRecordBatchReader::try_new(file, 60).unwrap();
3796
3797 for batch in record_batch_reader {
3798 batch.unwrap();
3799 }
3800
3801 let file = File::open(&path).unwrap();
3802 let builder = ParquetRecordBatchReaderBuilder::try_new(file).unwrap();
3803
3804 let mask = ProjectionMask::leaves(builder.parquet_schema(), [3, 8, 10]);
3805 let projected_reader = builder
3806 .with_projection(mask)
3807 .with_batch_size(60)
3808 .build()
3809 .unwrap();
3810
3811 let expected_schema = Schema::new(vec![
3812 Field::new(
3813 "roll_num",
3814 ArrowDataType::Struct(Fields::from(vec![Field::new(
3815 "count",
3816 ArrowDataType::UInt64,
3817 false,
3818 )])),
3819 false,
3820 ),
3821 Field::new(
3822 "PC_CUR",
3823 ArrowDataType::Struct(Fields::from(vec![
3824 Field::new("mean", ArrowDataType::Int64, false),
3825 Field::new("sum", ArrowDataType::Int64, false),
3826 ])),
3827 false,
3828 ),
3829 ]);
3830
3831 assert_eq!(&expected_schema, projected_reader.schema().as_ref());
3833
3834 for batch in projected_reader {
3835 let batch = batch.unwrap();
3836 assert_eq!(batch.schema().as_ref(), &expected_schema);
3837 }
3838 }
3839
3840 #[cfg_attr(miri, ignore)] #[test]
3842 fn test_read_structs_by_name() {
3844 let testdata = arrow::util::test_util::parquet_test_data();
3845 let path = format!("{testdata}/nested_structs.rust.parquet");
3846 let file = File::open(&path).unwrap();
3847 let record_batch_reader = ParquetRecordBatchReader::try_new(file, 60).unwrap();
3848
3849 for batch in record_batch_reader {
3850 batch.unwrap();
3851 }
3852
3853 let file = File::open(&path).unwrap();
3854 let builder = ParquetRecordBatchReaderBuilder::try_new(file).unwrap();
3855
3856 let mask = ProjectionMask::columns(
3857 builder.parquet_schema(),
3858 ["roll_num.count", "PC_CUR.mean", "PC_CUR.sum"],
3859 );
3860 let projected_reader = builder
3861 .with_projection(mask)
3862 .with_batch_size(60)
3863 .build()
3864 .unwrap();
3865
3866 let expected_schema = Schema::new(vec![
3867 Field::new(
3868 "roll_num",
3869 ArrowDataType::Struct(Fields::from(vec![Field::new(
3870 "count",
3871 ArrowDataType::UInt64,
3872 false,
3873 )])),
3874 false,
3875 ),
3876 Field::new(
3877 "PC_CUR",
3878 ArrowDataType::Struct(Fields::from(vec![
3879 Field::new("mean", ArrowDataType::Int64, false),
3880 Field::new("sum", ArrowDataType::Int64, false),
3881 ])),
3882 false,
3883 ),
3884 ]);
3885
3886 assert_eq!(&expected_schema, projected_reader.schema().as_ref());
3887
3888 for batch in projected_reader {
3889 let batch = batch.unwrap();
3890 assert_eq!(batch.schema().as_ref(), &expected_schema);
3891 }
3892 }
3893
3894 #[test]
3895 fn test_read_maps() {
3896 let testdata = arrow::util::test_util::parquet_test_data();
3897 let path = format!("{testdata}/nested_maps.snappy.parquet");
3898 let file = File::open(path).unwrap();
3899 let record_batch_reader = ParquetRecordBatchReader::try_new(file, 60).unwrap();
3900
3901 for batch in record_batch_reader {
3902 batch.unwrap();
3903 }
3904 }
3905
3906 #[test]
3908 fn test_unknown_logical_type() {
3909 let message_type = "message uk {
3910 OPTIONAL INT32 uki32 (UNKNOWN);
3911 OPTIONAL INT64 uki64 (UNKNOWN);
3912 OPTIONAL INT96 uki96 (UNKNOWN);
3913 OPTIONAL BOOLEAN ukbool (UNKNOWN);
3914 OPTIONAL FLOAT ukfloat (UNKNOWN);
3915 OPTIONAL DOUBLE ukdbl (UNKNOWN);
3916 OPTIONAL BYTE_ARRAY ukbytes (UNKNOWN);
3917 OPTIONAL FIXED_LEN_BYTE_ARRAY(10) ukflba (UNKNOWN);
3918 }";
3919
3920 let schema = Arc::new(parse_message_type(message_type).unwrap());
3921 let file = tempfile::tempfile().unwrap();
3922
3923 let mut writer =
3924 SerializedFileWriter::new(file.try_clone().unwrap(), schema, Default::default())
3925 .unwrap();
3926
3927 let mut row_group_writer = writer.next_row_group().unwrap();
3928
3929 fn write_nulls<T: DataType>(row_group_writer: &mut SerializedRowGroupWriter<'_, File>) {
3930 let mut column_writer = row_group_writer.next_column().unwrap().unwrap();
3931 column_writer
3933 .typed::<T>()
3934 .write_batch(&[], Some(&[0, 0, 0, 0]), None)
3935 .unwrap();
3936 column_writer.close().unwrap();
3937 }
3938
3939 write_nulls::<Int32Type>(&mut row_group_writer);
3941
3942 write_nulls::<Int64Type>(&mut row_group_writer);
3944
3945 write_nulls::<Int96Type>(&mut row_group_writer);
3947
3948 write_nulls::<BoolType>(&mut row_group_writer);
3950
3951 write_nulls::<FloatType>(&mut row_group_writer);
3953
3954 write_nulls::<DoubleType>(&mut row_group_writer);
3956
3957 write_nulls::<ByteArrayType>(&mut row_group_writer);
3959
3960 write_nulls::<FixedLenByteArrayType>(&mut row_group_writer);
3962
3963 row_group_writer.close().unwrap();
3964
3965 writer.close().unwrap();
3966
3967 let mut reader = ParquetRecordBatchReader::try_new(file, 4).unwrap();
3968 let batch = reader.next().unwrap().unwrap();
3969
3970 for col in batch.columns() {
3971 assert_eq!(col.len(), 4);
3972 assert_eq!(col.logical_null_count(), 4);
3973 assert_eq!(*col.data_type(), ArrowDataType::Null);
3974 }
3975 }
3976
3977 #[test]
3978 fn test_nested_nullability() {
3979 let message_type = "message nested {
3980 OPTIONAL Group group {
3981 REQUIRED INT32 leaf;
3982 }
3983 }";
3984
3985 let file = tempfile::tempfile().unwrap();
3986 let schema = Arc::new(parse_message_type(message_type).unwrap());
3987
3988 {
3989 let mut writer =
3991 SerializedFileWriter::new(file.try_clone().unwrap(), schema, Default::default())
3992 .unwrap();
3993
3994 {
3995 let mut row_group_writer = writer.next_row_group().unwrap();
3996 let mut column_writer = row_group_writer.next_column().unwrap().unwrap();
3997
3998 column_writer
3999 .typed::<Int32Type>()
4000 .write_batch(&[34, 76], Some(&[0, 1, 0, 1]), None)
4001 .unwrap();
4002
4003 column_writer.close().unwrap();
4004 row_group_writer.close().unwrap();
4005 }
4006
4007 writer.close().unwrap();
4008 }
4009
4010 let builder = ParquetRecordBatchReaderBuilder::try_new(file).unwrap();
4011 let mask = ProjectionMask::leaves(builder.parquet_schema(), [0]);
4012
4013 let reader = builder.with_projection(mask).build().unwrap();
4014
4015 let expected_schema = Schema::new(vec![Field::new(
4016 "group",
4017 ArrowDataType::Struct(vec![Field::new("leaf", ArrowDataType::Int32, false)].into()),
4018 true,
4019 )]);
4020
4021 let batch = reader.into_iter().next().unwrap().unwrap();
4022 assert_eq!(batch.schema().as_ref(), &expected_schema);
4023 assert_eq!(batch.num_rows(), 4);
4024 assert_eq!(batch.column(0).null_count(), 2);
4025 }
4026
4027 #[test]
4028 fn test_dictionary_preservation() {
4029 let fields = vec![Arc::new(
4030 Type::primitive_type_builder("leaf", PhysicalType::BYTE_ARRAY)
4031 .with_repetition(Repetition::OPTIONAL)
4032 .with_converted_type(ConvertedType::UTF8)
4033 .build()
4034 .unwrap(),
4035 )];
4036
4037 let schema = Arc::new(
4038 Type::group_type_builder("test_schema")
4039 .with_fields(fields)
4040 .build()
4041 .unwrap(),
4042 );
4043
4044 let dict_type = ArrowDataType::Dictionary(
4045 Box::new(ArrowDataType::Int32),
4046 Box::new(ArrowDataType::Utf8),
4047 );
4048
4049 let arrow_field = Field::new("leaf", dict_type, true);
4050
4051 let mut file = tempfile::tempfile().unwrap();
4052
4053 let values = vec![
4054 vec![
4055 ByteArray::from("hello"),
4056 ByteArray::from("a"),
4057 ByteArray::from("b"),
4058 ByteArray::from("d"),
4059 ],
4060 vec![
4061 ByteArray::from("c"),
4062 ByteArray::from("a"),
4063 ByteArray::from("b"),
4064 ],
4065 ];
4066
4067 let def_levels = vec![
4068 vec![1, 0, 0, 1, 0, 0, 1, 1],
4069 vec![0, 0, 1, 1, 0, 0, 1, 0, 0],
4070 ];
4071
4072 let opts = TestOptions {
4073 encoding: Encoding::RLE_DICTIONARY,
4074 ..Default::default()
4075 };
4076
4077 generate_single_column_file_with_data::<ByteArrayType>(
4078 &values,
4079 Some(&def_levels),
4080 file.try_clone().unwrap(), schema,
4082 Some(arrow_field),
4083 &opts,
4084 )
4085 .unwrap();
4086
4087 file.rewind().unwrap();
4088
4089 let record_reader = ParquetRecordBatchReader::try_new(file, 3).unwrap();
4090
4091 let batches = record_reader
4092 .collect::<Result<Vec<RecordBatch>, _>>()
4093 .unwrap();
4094
4095 assert_eq!(batches.len(), 6);
4096 assert!(batches.iter().all(|x| x.num_columns() == 1));
4097
4098 let row_counts = batches
4099 .iter()
4100 .map(|x| (x.num_rows(), x.column(0).null_count()))
4101 .collect::<Vec<_>>();
4102
4103 assert_eq!(
4104 row_counts,
4105 vec![(3, 2), (3, 2), (3, 1), (3, 1), (3, 2), (2, 2)]
4106 );
4107
4108 let get_dict = |batch: &RecordBatch| batch.column(0).to_data().child_data()[0].clone();
4109
4110 assert_eq!(get_dict(&batches[0]), get_dict(&batches[1]));
4112 assert_ne!(get_dict(&batches[1]), get_dict(&batches[2]));
4114 assert_ne!(get_dict(&batches[2]), get_dict(&batches[3]));
4115 assert_eq!(get_dict(&batches[3]), get_dict(&batches[4]));
4117 assert_eq!(get_dict(&batches[4]), get_dict(&batches[5]));
4118 }
4119
4120 #[test]
4121 fn test_read_null_list() {
4122 let testdata = arrow::util::test_util::parquet_test_data();
4123 let path = format!("{testdata}/null_list.parquet");
4124 let file = File::open(path).unwrap();
4125 let mut record_batch_reader = ParquetRecordBatchReader::try_new(file, 60).unwrap();
4126
4127 let batch = record_batch_reader.next().unwrap().unwrap();
4128 assert_eq!(batch.num_rows(), 1);
4129 assert_eq!(batch.num_columns(), 1);
4130 assert_eq!(batch.column(0).len(), 1);
4131
4132 let list = batch
4133 .column(0)
4134 .as_any()
4135 .downcast_ref::<ListArray>()
4136 .unwrap();
4137 assert_eq!(list.len(), 1);
4138 assert!(list.is_valid(0));
4139
4140 let val = list.value(0);
4141 assert_eq!(val.len(), 0);
4142 }
4143
4144 #[test]
4145 fn test_null_schema_inference() {
4146 let testdata = arrow::util::test_util::parquet_test_data();
4147 let path = format!("{testdata}/null_list.parquet");
4148 let file = File::open(path).unwrap();
4149
4150 let arrow_field = Field::new(
4151 "emptylist",
4152 ArrowDataType::List(Arc::new(Field::new_list_field(ArrowDataType::Null, true))),
4153 true,
4154 );
4155
4156 let options = ArrowReaderOptions::new().with_skip_arrow_metadata(true);
4157 let builder = ParquetRecordBatchReaderBuilder::try_new_with_options(file, options).unwrap();
4158 let schema = builder.schema();
4159 assert_eq!(schema.fields().len(), 1);
4160 assert_eq!(schema.field(0), &arrow_field);
4161 }
4162
4163 #[test]
4164 fn test_skip_metadata() {
4165 let col = Arc::new(TimestampNanosecondArray::from_iter_values(vec![0, 1, 2]));
4166 let field = Field::new("col", col.data_type().clone(), true);
4167
4168 let schema_without_metadata = Arc::new(Schema::new(vec![field.clone()]));
4169
4170 let metadata = arrow_schema::Metadata::from([("key".to_string(), "value".to_string())]);
4171
4172 let schema_with_metadata = Arc::new(Schema::new(vec![field.with_metadata(metadata)]));
4173
4174 assert_ne!(schema_with_metadata, schema_without_metadata);
4175
4176 let batch =
4177 RecordBatch::try_new(schema_with_metadata.clone(), vec![col as ArrayRef]).unwrap();
4178
4179 let file = |version: WriterVersion| {
4180 let props = WriterProperties::builder()
4181 .set_writer_version(version)
4182 .build();
4183
4184 let file = tempfile().unwrap();
4185 let mut writer =
4186 ArrowWriter::try_new(file.try_clone().unwrap(), batch.schema(), Some(props))
4187 .unwrap();
4188 writer.write(&batch).unwrap();
4189 writer.close().unwrap();
4190 file
4191 };
4192
4193 let skip_options = ArrowReaderOptions::new().with_skip_arrow_metadata(true);
4194
4195 let v1_reader = file(WriterVersion::PARQUET_1_0);
4196 let v2_reader = file(WriterVersion::PARQUET_2_0);
4197
4198 let arrow_reader =
4199 ParquetRecordBatchReader::try_new(v1_reader.try_clone().unwrap(), 1024).unwrap();
4200 assert_eq!(arrow_reader.schema(), schema_with_metadata);
4201
4202 let reader =
4203 ParquetRecordBatchReaderBuilder::try_new_with_options(v1_reader, skip_options.clone())
4204 .unwrap()
4205 .build()
4206 .unwrap();
4207 assert_eq!(reader.schema(), schema_without_metadata);
4208
4209 let arrow_reader =
4210 ParquetRecordBatchReader::try_new(v2_reader.try_clone().unwrap(), 1024).unwrap();
4211 assert_eq!(arrow_reader.schema(), schema_with_metadata);
4212
4213 let reader = ParquetRecordBatchReaderBuilder::try_new_with_options(v2_reader, skip_options)
4214 .unwrap()
4215 .build()
4216 .unwrap();
4217 assert_eq!(reader.schema(), schema_without_metadata);
4218 }
4219
4220 fn write_parquet_from_iter<I, F>(value: I) -> File
4221 where
4222 I: IntoIterator<Item = (F, ArrayRef)>,
4223 F: AsRef<str>,
4224 {
4225 let batch = RecordBatch::try_from_iter(value).unwrap();
4226 let file = tempfile().unwrap();
4227 let mut writer =
4228 ArrowWriter::try_new(file.try_clone().unwrap(), batch.schema().clone(), None).unwrap();
4229 writer.write(&batch).unwrap();
4230 writer.close().unwrap();
4231 file
4232 }
4233
4234 fn run_schema_test_with_error<I, F>(value: I, schema: SchemaRef, expected_error: &str)
4235 where
4236 I: IntoIterator<Item = (F, ArrayRef)>,
4237 F: AsRef<str>,
4238 {
4239 let file = write_parquet_from_iter(value);
4240 let options_with_schema = ArrowReaderOptions::new().with_schema(schema.clone());
4241 let builder = ParquetRecordBatchReaderBuilder::try_new_with_options(
4242 file.try_clone().unwrap(),
4243 options_with_schema,
4244 );
4245 assert_eq!(builder.err().unwrap().to_string(), expected_error);
4246 }
4247
4248 #[test]
4249 fn test_schema_too_few_columns() {
4250 run_schema_test_with_error(
4251 vec![
4252 ("int64", Arc::new(Int64Array::from(vec![0])) as ArrayRef),
4253 ("int32", Arc::new(Int32Array::from(vec![0])) as ArrayRef),
4254 ],
4255 Arc::new(Schema::new(vec![Field::new(
4256 "int64",
4257 ArrowDataType::Int64,
4258 false,
4259 )])),
4260 "Arrow: incompatible arrow schema, expected 2 struct fields got 1",
4261 );
4262 }
4263
4264 #[test]
4265 fn test_schema_too_many_columns() {
4266 run_schema_test_with_error(
4267 vec![("int64", Arc::new(Int64Array::from(vec![0])) as ArrayRef)],
4268 Arc::new(Schema::new(vec![
4269 Field::new("int64", ArrowDataType::Int64, false),
4270 Field::new("int32", ArrowDataType::Int32, false),
4271 ])),
4272 "Arrow: incompatible arrow schema, expected 1 struct fields got 2",
4273 );
4274 }
4275
4276 #[test]
4277 fn test_schema_mismatched_column_names() {
4278 run_schema_test_with_error(
4279 vec![("int64", Arc::new(Int64Array::from(vec![0])) as ArrayRef)],
4280 Arc::new(Schema::new(vec![Field::new(
4281 "other",
4282 ArrowDataType::Int64,
4283 false,
4284 )])),
4285 "Arrow: incompatible arrow schema, expected field named int64 got other",
4286 );
4287 }
4288
4289 #[test]
4290 fn test_schema_incompatible_columns() {
4291 run_schema_test_with_error(
4292 vec![
4293 (
4294 "col1_invalid",
4295 Arc::new(Int64Array::from(vec![0])) as ArrayRef,
4296 ),
4297 (
4298 "col2_valid",
4299 Arc::new(Int32Array::from(vec![0])) as ArrayRef,
4300 ),
4301 (
4302 "col3_invalid",
4303 Arc::new(Date64Array::from(vec![0])) as ArrayRef,
4304 ),
4305 ],
4306 Arc::new(Schema::new(vec![
4307 Field::new("col1_invalid", ArrowDataType::Int32, false),
4308 Field::new("col2_valid", ArrowDataType::Int32, false),
4309 Field::new("col3_invalid", ArrowDataType::Int32, false),
4310 ])),
4311 "Arrow: Incompatible supplied Arrow schema: data type mismatch for field col1_invalid: requested Int32 but found Int64, data type mismatch for field col3_invalid: requested Int32 but found Int64",
4312 );
4313 }
4314
4315 #[test]
4316 fn test_one_incompatible_nested_column() {
4317 let nested_fields = Fields::from(vec![
4318 Field::new("nested1_valid", ArrowDataType::Utf8, false),
4319 Field::new("nested1_invalid", ArrowDataType::Int64, false),
4320 ]);
4321 let nested = StructArray::try_new(
4322 nested_fields,
4323 vec![
4324 Arc::new(StringArray::from(vec!["a"])) as ArrayRef,
4325 Arc::new(Int64Array::from(vec![0])) as ArrayRef,
4326 ],
4327 None,
4328 )
4329 .expect("struct array");
4330 let supplied_nested_fields = Fields::from(vec![
4331 Field::new("nested1_valid", ArrowDataType::Utf8, false),
4332 Field::new("nested1_invalid", ArrowDataType::Int32, false),
4333 ]);
4334 run_schema_test_with_error(
4335 vec![
4336 ("col1", Arc::new(Int64Array::from(vec![0])) as ArrayRef),
4337 ("col2", Arc::new(Int32Array::from(vec![0])) as ArrayRef),
4338 ("nested", Arc::new(nested) as ArrayRef),
4339 ],
4340 Arc::new(Schema::new(vec![
4341 Field::new("col1", ArrowDataType::Int64, false),
4342 Field::new("col2", ArrowDataType::Int32, false),
4343 Field::new(
4344 "nested",
4345 ArrowDataType::Struct(supplied_nested_fields),
4346 false,
4347 ),
4348 ])),
4349 "Arrow: Incompatible supplied Arrow schema: data type mismatch for field nested: \
4350 requested Struct(\"nested1_valid\": non-null Utf8, \"nested1_invalid\": non-null Int32) \
4351 but found Struct(\"nested1_valid\": non-null Utf8, \"nested1_invalid\": non-null Int64)",
4352 );
4353 }
4354
4355 fn utf8_parquet() -> Bytes {
4357 let input = StringArray::from_iter_values(vec!["foo", "bar", "baz"]);
4358 let batch = RecordBatch::try_from_iter(vec![("column1", Arc::new(input) as _)]).unwrap();
4359 let props = None;
4360 let mut parquet_data = vec![];
4362 let mut writer = ArrowWriter::try_new(&mut parquet_data, batch.schema(), props).unwrap();
4363 writer.write(&batch).unwrap();
4364 writer.close().unwrap();
4365 Bytes::from(parquet_data)
4366 }
4367
4368 #[test]
4369 fn test_schema_error_bad_types() {
4370 let parquet_data = utf8_parquet();
4372
4373 let input_schema: SchemaRef = Arc::new(Schema::new(vec![Field::new(
4375 "column1",
4376 arrow::datatypes::DataType::Int32,
4377 false,
4378 )]));
4379
4380 let reader_options = ArrowReaderOptions::new().with_schema(input_schema.clone());
4382 let err =
4383 ParquetRecordBatchReaderBuilder::try_new_with_options(parquet_data, reader_options)
4384 .unwrap_err();
4385 assert_eq!(
4386 err.to_string(),
4387 "Arrow: Incompatible supplied Arrow schema: data type mismatch for field column1: requested Int32 but found Utf8"
4388 )
4389 }
4390
4391 #[test]
4392 fn test_schema_error_bad_nullability() {
4393 let parquet_data = utf8_parquet();
4395
4396 let input_schema: SchemaRef = Arc::new(Schema::new(vec![Field::new(
4398 "column1",
4399 arrow::datatypes::DataType::Utf8,
4400 true,
4401 )]));
4402
4403 let reader_options = ArrowReaderOptions::new().with_schema(input_schema.clone());
4405 let err =
4406 ParquetRecordBatchReaderBuilder::try_new_with_options(parquet_data, reader_options)
4407 .unwrap_err();
4408 assert_eq!(
4409 err.to_string(),
4410 "Arrow: Incompatible supplied Arrow schema: nullability mismatch for field column1: expected true but found false"
4411 )
4412 }
4413
4414 #[test]
4415 fn test_read_binary_as_utf8() {
4416 let file = write_parquet_from_iter(vec![
4417 (
4418 "binary_to_utf8",
4419 Arc::new(BinaryArray::from(vec![
4420 b"one".as_ref(),
4421 b"two".as_ref(),
4422 b"three".as_ref(),
4423 ])) as ArrayRef,
4424 ),
4425 (
4426 "large_binary_to_large_utf8",
4427 Arc::new(LargeBinaryArray::from(vec![
4428 b"one".as_ref(),
4429 b"two".as_ref(),
4430 b"three".as_ref(),
4431 ])) as ArrayRef,
4432 ),
4433 (
4434 "binary_view_to_utf8_view",
4435 Arc::new(BinaryViewArray::from(vec![
4436 b"one".as_ref(),
4437 b"two".as_ref(),
4438 b"three".as_ref(),
4439 ])) as ArrayRef,
4440 ),
4441 ]);
4442 let supplied_fields = Fields::from(vec![
4443 Field::new("binary_to_utf8", ArrowDataType::Utf8, false),
4444 Field::new(
4445 "large_binary_to_large_utf8",
4446 ArrowDataType::LargeUtf8,
4447 false,
4448 ),
4449 Field::new("binary_view_to_utf8_view", ArrowDataType::Utf8View, false),
4450 ]);
4451
4452 let options = ArrowReaderOptions::new().with_schema(Arc::new(Schema::new(supplied_fields)));
4453 let mut arrow_reader = ParquetRecordBatchReaderBuilder::try_new_with_options(
4454 file.try_clone().unwrap(),
4455 options,
4456 )
4457 .expect("reader builder with schema")
4458 .build()
4459 .expect("reader with schema");
4460
4461 let batch = arrow_reader.next().unwrap().unwrap();
4462 assert_eq!(batch.num_columns(), 3);
4463 assert_eq!(batch.num_rows(), 3);
4464 assert_eq!(
4465 batch
4466 .column(0)
4467 .as_string::<i32>()
4468 .iter()
4469 .collect::<Vec<_>>(),
4470 vec![Some("one"), Some("two"), Some("three")]
4471 );
4472
4473 assert_eq!(
4474 batch
4475 .column(1)
4476 .as_string::<i64>()
4477 .iter()
4478 .collect::<Vec<_>>(),
4479 vec![Some("one"), Some("two"), Some("three")]
4480 );
4481
4482 assert_eq!(
4483 batch.column(2).as_string_view().iter().collect::<Vec<_>>(),
4484 vec![Some("one"), Some("two"), Some("three")]
4485 );
4486 }
4487
4488 #[test]
4489 #[should_panic(expected = "Invalid UTF8 sequence at")]
4490 fn test_read_non_utf8_binary_as_utf8() {
4491 let file = write_parquet_from_iter(vec![(
4492 "non_utf8_binary",
4493 Arc::new(BinaryArray::from(vec![
4494 b"\xDE\x00\xFF".as_ref(),
4495 b"\xDE\x01\xAA".as_ref(),
4496 b"\xDE\x02\xFF".as_ref(),
4497 ])) as ArrayRef,
4498 )]);
4499 let supplied_fields = Fields::from(vec![Field::new(
4500 "non_utf8_binary",
4501 ArrowDataType::Utf8,
4502 false,
4503 )]);
4504
4505 let options = ArrowReaderOptions::new().with_schema(Arc::new(Schema::new(supplied_fields)));
4506 let mut arrow_reader = ParquetRecordBatchReaderBuilder::try_new_with_options(
4507 file.try_clone().unwrap(),
4508 options,
4509 )
4510 .expect("reader builder with schema")
4511 .build()
4512 .expect("reader with schema");
4513 arrow_reader.next().unwrap().unwrap_err();
4514 }
4515
4516 #[test]
4517 fn test_with_schema() {
4518 let nested_fields = Fields::from(vec![
4519 Field::new("utf8_to_dict", ArrowDataType::Utf8, false),
4520 Field::new("int64_to_ts_nano", ArrowDataType::Int64, false),
4521 ]);
4522
4523 let nested_arrays: Vec<ArrayRef> = vec![
4524 Arc::new(StringArray::from(vec!["a", "a", "a", "b"])) as ArrayRef,
4525 Arc::new(Int64Array::from(vec![1, 2, 3, 4])) as ArrayRef,
4526 ];
4527
4528 let nested = StructArray::try_new(nested_fields, nested_arrays, None).unwrap();
4529
4530 let file = write_parquet_from_iter(vec![
4531 (
4532 "int32_to_ts_second",
4533 Arc::new(Int32Array::from(vec![0, 1, 2, 3])) as ArrayRef,
4534 ),
4535 (
4536 "date32_to_date64",
4537 Arc::new(Date32Array::from(vec![0, 1, 2, 3])) as ArrayRef,
4538 ),
4539 ("nested", Arc::new(nested) as ArrayRef),
4540 ]);
4541
4542 let supplied_nested_fields = Fields::from(vec![
4543 Field::new(
4544 "utf8_to_dict",
4545 ArrowDataType::Dictionary(
4546 Box::new(ArrowDataType::Int32),
4547 Box::new(ArrowDataType::Utf8),
4548 ),
4549 false,
4550 ),
4551 Field::new(
4552 "int64_to_ts_nano",
4553 ArrowDataType::Timestamp(
4554 arrow::datatypes::TimeUnit::Nanosecond,
4555 Some("+10:00".into()),
4556 ),
4557 false,
4558 ),
4559 ]);
4560
4561 let supplied_schema = Arc::new(Schema::new(vec![
4562 Field::new(
4563 "int32_to_ts_second",
4564 ArrowDataType::Timestamp(arrow::datatypes::TimeUnit::Second, Some("+01:00".into())),
4565 false,
4566 ),
4567 Field::new("date32_to_date64", ArrowDataType::Date64, false),
4568 Field::new(
4569 "nested",
4570 ArrowDataType::Struct(supplied_nested_fields),
4571 false,
4572 ),
4573 ]));
4574
4575 let options = ArrowReaderOptions::new().with_schema(supplied_schema.clone());
4576 let mut arrow_reader = ParquetRecordBatchReaderBuilder::try_new_with_options(
4577 file.try_clone().unwrap(),
4578 options,
4579 )
4580 .expect("reader builder with schema")
4581 .build()
4582 .expect("reader with schema");
4583
4584 assert_eq!(arrow_reader.schema(), supplied_schema);
4585 let batch = arrow_reader.next().unwrap().unwrap();
4586 assert_eq!(batch.num_columns(), 3);
4587 assert_eq!(batch.num_rows(), 4);
4588 assert_eq!(
4589 batch
4590 .column(0)
4591 .as_any()
4592 .downcast_ref::<TimestampSecondArray>()
4593 .expect("downcast to timestamp second")
4594 .value_as_datetime_with_tz(0, "+01:00".parse().unwrap())
4595 .map(|v| v.to_string())
4596 .expect("value as datetime"),
4597 "1970-01-01 01:00:00 +01:00"
4598 );
4599 assert_eq!(
4600 batch
4601 .column(1)
4602 .as_any()
4603 .downcast_ref::<Date64Array>()
4604 .expect("downcast to date64")
4605 .value_as_date(0)
4606 .map(|v| v.to_string())
4607 .expect("value as date"),
4608 "1970-01-01"
4609 );
4610
4611 let nested = batch
4612 .column(2)
4613 .as_any()
4614 .downcast_ref::<StructArray>()
4615 .expect("downcast to struct");
4616
4617 let nested_dict = nested
4618 .column(0)
4619 .as_any()
4620 .downcast_ref::<Int32DictionaryArray>()
4621 .expect("downcast to dictionary");
4622
4623 assert_eq!(
4624 nested_dict
4625 .values()
4626 .as_any()
4627 .downcast_ref::<StringArray>()
4628 .expect("downcast to string")
4629 .iter()
4630 .collect::<Vec<_>>(),
4631 vec![Some("a"), Some("b")]
4632 );
4633
4634 assert_eq!(
4635 nested_dict.keys().iter().collect::<Vec<_>>(),
4636 vec![Some(0), Some(0), Some(0), Some(1)]
4637 );
4638
4639 assert_eq!(
4640 nested
4641 .column(1)
4642 .as_any()
4643 .downcast_ref::<TimestampNanosecondArray>()
4644 .expect("downcast to timestamp nanosecond")
4645 .value_as_datetime_with_tz(0, "+10:00".parse().unwrap())
4646 .map(|v| v.to_string())
4647 .expect("value as datetime"),
4648 "1970-01-01 10:00:00.000000001 +10:00"
4649 );
4650 }
4651
4652 #[test]
4653 fn test_empty_projection() {
4654 let testdata = arrow::util::test_util::parquet_test_data();
4655 let path = format!("{testdata}/alltypes_plain.parquet");
4656 let file = File::open(path).unwrap();
4657
4658 let builder = ParquetRecordBatchReaderBuilder::try_new(file).unwrap();
4659 let file_metadata = builder.metadata().file_metadata();
4660 let expected_rows = file_metadata.num_rows() as usize;
4661
4662 let mask = ProjectionMask::leaves(builder.parquet_schema(), []);
4663 let batch_reader = builder
4664 .with_projection(mask)
4665 .with_batch_size(2)
4666 .build()
4667 .unwrap();
4668
4669 let mut total_rows = 0;
4670 for maybe_batch in batch_reader {
4671 let batch = maybe_batch.unwrap();
4672 total_rows += batch.num_rows();
4673 assert_eq!(batch.num_columns(), 0);
4674 assert!(batch.num_rows() <= 2);
4675 }
4676
4677 assert_eq!(total_rows, expected_rows);
4678 }
4679
4680 fn test_row_group_batch(row_group_size: usize, batch_size: usize) {
4681 let schema = Arc::new(Schema::new(vec![Field::new(
4682 "list",
4683 ArrowDataType::List(Arc::new(Field::new_list_field(ArrowDataType::Int32, true))),
4684 true,
4685 )]));
4686
4687 let mut buf = Vec::with_capacity(1024);
4688
4689 let mut writer = ArrowWriter::try_new(
4690 &mut buf,
4691 schema.clone(),
4692 Some(
4693 WriterProperties::builder()
4694 .set_max_row_group_row_count(Some(row_group_size))
4695 .build(),
4696 ),
4697 )
4698 .unwrap();
4699 for _ in 0..2 {
4700 let mut list_builder = ListBuilder::new(Int32Builder::with_capacity(batch_size));
4701 for _ in 0..(batch_size) {
4702 list_builder.append(true);
4703 }
4704 let batch = RecordBatch::try_new(schema.clone(), vec![Arc::new(list_builder.finish())])
4705 .unwrap();
4706 writer.write(&batch).unwrap();
4707 }
4708 writer.close().unwrap();
4709
4710 let mut record_reader =
4711 ParquetRecordBatchReader::try_new(Bytes::from(buf), batch_size).unwrap();
4712 assert_eq!(
4713 batch_size,
4714 record_reader.next().unwrap().unwrap().num_rows()
4715 );
4716 assert_eq!(
4717 batch_size,
4718 record_reader.next().unwrap().unwrap().num_rows()
4719 );
4720 }
4721
4722 #[test]
4723 #[cfg_attr(miri, ignore)] fn test_row_group_exact_multiple() {
4725 const BATCH_SIZE: usize = REPETITION_LEVELS_BATCH_SIZE;
4726 test_row_group_batch(8, 8);
4727 test_row_group_batch(10, 8);
4728 test_row_group_batch(8, 10);
4729 test_row_group_batch(BATCH_SIZE, BATCH_SIZE);
4730 test_row_group_batch(BATCH_SIZE + 1, BATCH_SIZE);
4731 test_row_group_batch(BATCH_SIZE, BATCH_SIZE + 1);
4732 test_row_group_batch(BATCH_SIZE, BATCH_SIZE - 1);
4733 test_row_group_batch(BATCH_SIZE - 1, BATCH_SIZE);
4734 }
4735
4736 fn get_expected_batches(
4739 column: &RecordBatch,
4740 selection: &RowSelection,
4741 batch_size: usize,
4742 ) -> Vec<RecordBatch> {
4743 let mut expected_batches = vec![];
4744
4745 let mut selection: VecDeque<_> = selection.clone().into();
4746 let mut row_offset = 0;
4747 let mut last_start = None;
4748 while row_offset < column.num_rows() && !selection.is_empty() {
4749 let mut batch_remaining = batch_size.min(column.num_rows() - row_offset);
4750 while batch_remaining > 0 && !selection.is_empty() {
4751 let (to_read, skip) = match selection.front_mut() {
4752 Some(selection) if selection.row_count > batch_remaining => {
4753 selection.row_count -= batch_remaining;
4754 (batch_remaining, selection.skip)
4755 }
4756 Some(_) => {
4757 let select = selection.pop_front().unwrap();
4758 (select.row_count, select.skip)
4759 }
4760 None => break,
4761 };
4762
4763 batch_remaining -= to_read;
4764
4765 match skip {
4766 true => {
4767 if let Some(last_start) = last_start.take() {
4768 expected_batches.push(column.slice(last_start, row_offset - last_start))
4769 }
4770 row_offset += to_read
4771 }
4772 false => {
4773 last_start.get_or_insert(row_offset);
4774 row_offset += to_read
4775 }
4776 }
4777 }
4778 }
4779
4780 if let Some(last_start) = last_start.take() {
4781 expected_batches.push(column.slice(last_start, row_offset - last_start))
4782 }
4783
4784 for batch in &expected_batches[..expected_batches.len() - 1] {
4786 assert_eq!(batch.num_rows(), batch_size);
4787 }
4788
4789 expected_batches
4790 }
4791
4792 fn create_test_selection(
4793 step_len: usize,
4794 total_len: usize,
4795 skip_first: bool,
4796 ) -> (RowSelection, usize) {
4797 let mut remaining = total_len;
4798 let mut skip = skip_first;
4799 let mut vec = vec![];
4800 let mut selected_count = 0;
4801 while remaining != 0 {
4802 let step = if remaining > step_len {
4803 step_len
4804 } else {
4805 remaining
4806 };
4807 vec.push(RowSelector {
4808 row_count: step,
4809 skip,
4810 });
4811 remaining -= step;
4812 if !skip {
4813 selected_count += step;
4814 }
4815 skip = !skip;
4816 }
4817 (vec.into(), selected_count)
4818 }
4819
4820 #[test]
4821 #[cfg_attr(miri, ignore)] fn test_scan_row_with_selection() {
4823 let testdata = arrow::util::test_util::parquet_test_data();
4824 let path = format!("{testdata}/alltypes_tiny_pages_plain.parquet");
4825 let test_file = File::open(&path).unwrap();
4826
4827 let mut serial_reader =
4828 ParquetRecordBatchReader::try_new(File::open(&path).unwrap(), 7300).unwrap();
4829 let data = serial_reader.next().unwrap().unwrap();
4830
4831 let do_test = |batch_size: usize, selection_len: usize| {
4832 for skip_first in [false, true] {
4833 let selections = create_test_selection(batch_size, data.num_rows(), skip_first).0;
4834
4835 let expected = get_expected_batches(&data, &selections, batch_size);
4836 let skip_reader = create_skip_reader(&test_file, batch_size, selections);
4837 assert_eq!(
4838 skip_reader.collect::<Result<Vec<_>, _>>().unwrap(),
4839 expected,
4840 "batch_size: {batch_size}, selection_len: {selection_len}, skip_first: {skip_first}"
4841 );
4842 }
4843 };
4844
4845 do_test(1000, 1000);
4848
4849 do_test(20, 20);
4851
4852 do_test(20, 5);
4854
4855 do_test(20, 5);
4858
4859 fn create_skip_reader(
4860 test_file: &File,
4861 batch_size: usize,
4862 selections: RowSelection,
4863 ) -> ParquetRecordBatchReader {
4864 let options =
4865 ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required);
4866 let file = test_file.try_clone().unwrap();
4867 ParquetRecordBatchReaderBuilder::try_new_with_options(file, options)
4868 .unwrap()
4869 .with_batch_size(batch_size)
4870 .with_row_selection(selections)
4871 .build()
4872 .unwrap()
4873 }
4874 }
4875
4876 #[test]
4877 fn test_batch_size_overallocate() {
4878 let testdata = arrow::util::test_util::parquet_test_data();
4879 let path = format!("{testdata}/alltypes_plain.parquet");
4881 let test_file = File::open(path).unwrap();
4882
4883 let builder = ParquetRecordBatchReaderBuilder::try_new(test_file).unwrap();
4884 let num_rows = builder.metadata.file_metadata().num_rows();
4885 let reader = builder
4886 .with_batch_size(1024)
4887 .with_projection(ProjectionMask::all())
4888 .build()
4889 .unwrap();
4890 assert_ne!(1024, num_rows);
4891 assert_eq!(reader.read_plan.batch_size(), num_rows as usize);
4892 }
4893
4894 #[test]
4895 #[cfg_attr(miri, ignore)] fn test_read_with_page_index_enabled() {
4897 let testdata = arrow::util::test_util::parquet_test_data();
4898
4899 {
4900 let path = format!("{testdata}/alltypes_tiny_pages.parquet");
4902 let test_file = File::open(path).unwrap();
4903 let builder = ParquetRecordBatchReaderBuilder::try_new_with_options(
4904 test_file,
4905 ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required),
4906 )
4907 .unwrap();
4908 let page_index = builder
4909 .metadata()
4910 .page_index()
4911 .expect("page index should be present");
4912 let num_columns = builder.metadata().row_group(0).num_columns();
4913 let offset_indexes = page_index.offset_indexes_for_rowgroup(0);
4914 assert!(offset_indexes.is_some_and(|ois| ois.len() == num_columns));
4915 let column_indexes = page_index.offset_indexes_for_rowgroup(0);
4916 assert!(column_indexes.is_some_and(|cis| cis.len() == num_columns));
4917 assert!(page_index.offset_index(0, 0).is_some());
4918 assert!(page_index.column_index(0, 0).is_some());
4919 assert!(page_index.page_locations(0, 0).is_some());
4920 assert_eq!(page_index.num_data_pages(0, 0), Some(325));
4921 let reader = builder.build().unwrap();
4922 let batches = reader.collect::<Result<Vec<_>, _>>().unwrap();
4923 assert_eq!(batches.len(), 8);
4924 }
4925
4926 {
4927 let path = format!("{testdata}/alltypes_plain.parquet");
4929 let test_file = File::open(path).unwrap();
4930 let builder = ParquetRecordBatchReaderBuilder::try_new_with_options(
4931 test_file,
4932 ArrowReaderOptions::new().with_page_index_policy(PageIndexPolicy::Required),
4933 )
4934 .unwrap();
4935 assert!(builder.metadata().page_index().is_none());
4938 let reader = builder.build().unwrap();
4939 let batches = reader.collect::<Result<Vec<_>, _>>().unwrap();
4940 assert_eq!(batches.len(), 1);
4941 }
4942 }
4943
4944 #[test]
4945 fn test_raw_repetition() {
4946 const MESSAGE_TYPE: &str = "
4947 message Log {
4948 OPTIONAL INT32 eventType;
4949 REPEATED INT32 category;
4950 REPEATED group filter {
4951 OPTIONAL INT32 error;
4952 }
4953 }
4954 ";
4955 let schema = Arc::new(parse_message_type(MESSAGE_TYPE).unwrap());
4956 let props = Default::default();
4957
4958 let mut buf = Vec::with_capacity(1024);
4959 let mut writer = SerializedFileWriter::new(&mut buf, schema, props).unwrap();
4960 let mut row_group_writer = writer.next_row_group().unwrap();
4961
4962 let mut col_writer = row_group_writer.next_column().unwrap().unwrap();
4964 col_writer
4965 .typed::<Int32Type>()
4966 .write_batch(&[1], Some(&[1]), None)
4967 .unwrap();
4968 col_writer.close().unwrap();
4969 let mut col_writer = row_group_writer.next_column().unwrap().unwrap();
4971 col_writer
4972 .typed::<Int32Type>()
4973 .write_batch(&[1, 1], Some(&[1, 1]), Some(&[0, 1]))
4974 .unwrap();
4975 col_writer.close().unwrap();
4976 let mut col_writer = row_group_writer.next_column().unwrap().unwrap();
4978 col_writer
4979 .typed::<Int32Type>()
4980 .write_batch(&[1], Some(&[1]), Some(&[0]))
4981 .unwrap();
4982 col_writer.close().unwrap();
4983
4984 let rg_md = row_group_writer.close().unwrap();
4985 assert_eq!(rg_md.num_rows(), 1);
4986 writer.close().unwrap();
4987
4988 let bytes = Bytes::from(buf);
4989
4990 let mut no_mask = ParquetRecordBatchReader::try_new(bytes.clone(), 1024).unwrap();
4991 let full = no_mask.next().unwrap().unwrap();
4992
4993 assert_eq!(full.num_columns(), 3);
4994
4995 for idx in 0..3 {
4996 let b = ParquetRecordBatchReaderBuilder::try_new(bytes.clone()).unwrap();
4997 let mask = ProjectionMask::leaves(b.parquet_schema(), [idx]);
4998 let mut reader = b.with_projection(mask).build().unwrap();
4999 let projected = reader.next().unwrap().unwrap();
5000
5001 assert_eq!(projected.num_columns(), 1);
5002 assert_eq!(full.column(idx), projected.column(0));
5003 }
5004 }
5005
5006 #[test]
5007 fn test_read_lz4_raw() {
5008 let testdata = arrow::util::test_util::parquet_test_data();
5009 let path = format!("{testdata}/lz4_raw_compressed.parquet");
5010 let file = File::open(path).unwrap();
5011
5012 let batches = ParquetRecordBatchReader::try_new(file, 1024)
5013 .unwrap()
5014 .collect::<Result<Vec<_>, _>>()
5015 .unwrap();
5016 assert_eq!(batches.len(), 1);
5017 let batch = &batches[0];
5018
5019 assert_eq!(batch.num_columns(), 3);
5020 assert_eq!(batch.num_rows(), 4);
5021
5022 let a: &Int64Array = batch.column(0).as_any().downcast_ref().unwrap();
5024 assert_eq!(
5025 a.values(),
5026 &[1593604800, 1593604800, 1593604801, 1593604801]
5027 );
5028
5029 let a: &BinaryArray = batch.column(1).as_any().downcast_ref().unwrap();
5030 let a: Vec<_> = a.iter().flatten().collect();
5031 assert_eq!(a, &[b"abc", b"def", b"abc", b"def"]);
5032
5033 let a: &Float64Array = batch.column(2).as_any().downcast_ref().unwrap();
5034 assert_eq!(a.values(), &[42.000000, 7.700000, 42.125000, 7.700000]);
5035 }
5036
5037 #[test]
5047 #[cfg_attr(miri, ignore)] fn test_read_lz4_hadoop_fallback() {
5049 for file in [
5050 "hadoop_lz4_compressed.parquet",
5051 "non_hadoop_lz4_compressed.parquet",
5052 ] {
5053 let testdata = arrow::util::test_util::parquet_test_data();
5054 let path = format!("{testdata}/{file}");
5055 let file = File::open(path).unwrap();
5056 let expected_rows = 4;
5057
5058 let batches = ParquetRecordBatchReader::try_new(file, expected_rows)
5059 .unwrap()
5060 .collect::<Result<Vec<_>, _>>()
5061 .unwrap();
5062 assert_eq!(batches.len(), 1);
5063 let batch = &batches[0];
5064
5065 assert_eq!(batch.num_columns(), 3);
5066 assert_eq!(batch.num_rows(), expected_rows);
5067
5068 let a: &Int64Array = batch.column(0).as_any().downcast_ref().unwrap();
5069 assert_eq!(
5070 a.values(),
5071 &[1593604800, 1593604800, 1593604801, 1593604801]
5072 );
5073
5074 let b: &BinaryArray = batch.column(1).as_any().downcast_ref().unwrap();
5075 let b: Vec<_> = b.iter().flatten().collect();
5076 assert_eq!(b, &[b"abc", b"def", b"abc", b"def"]);
5077
5078 let c: &Float64Array = batch.column(2).as_any().downcast_ref().unwrap();
5079 assert_eq!(c.values(), &[42.0, 7.7, 42.125, 7.7]);
5080 }
5081 }
5082
5083 #[test]
5084 #[cfg_attr(miri, ignore)] fn test_read_lz4_hadoop_large() {
5086 let testdata = arrow::util::test_util::parquet_test_data();
5087 let path = format!("{testdata}/hadoop_lz4_compressed_larger.parquet");
5088 let file = File::open(path).unwrap();
5089 let expected_rows = 10000;
5090
5091 let batches = ParquetRecordBatchReader::try_new(file, expected_rows)
5092 .unwrap()
5093 .collect::<Result<Vec<_>, _>>()
5094 .unwrap();
5095 assert_eq!(batches.len(), 1);
5096 let batch = &batches[0];
5097
5098 assert_eq!(batch.num_columns(), 1);
5099 assert_eq!(batch.num_rows(), expected_rows);
5100
5101 let a: &StringArray = batch.column(0).as_any().downcast_ref().unwrap();
5102 let a: Vec<_> = a.iter().flatten().collect();
5103 assert_eq!(a[0], "c7ce6bef-d5b0-4863-b199-8ea8c7fb117b");
5104 assert_eq!(a[1], "e8fb9197-cb9f-4118-b67f-fbfa65f61843");
5105 assert_eq!(a[expected_rows - 2], "ab52a0cc-c6bb-4d61-8a8f-166dc4b8b13c");
5106 assert_eq!(a[expected_rows - 1], "85440778-460a-41ac-aa2e-ac3ee41696bf");
5107 }
5108
5109 #[test]
5110 #[cfg(feature = "snap")]
5111 fn test_read_nested_lists() {
5112 let testdata = arrow::util::test_util::parquet_test_data();
5113 let path = format!("{testdata}/nested_lists.snappy.parquet");
5114 let file = File::open(path).unwrap();
5115
5116 let f = file.try_clone().unwrap();
5117 let mut reader = ParquetRecordBatchReader::try_new(f, 60).unwrap();
5118 let expected = reader.next().unwrap().unwrap();
5119 assert_eq!(expected.num_rows(), 3);
5120
5121 let selection = RowSelection::from(vec![
5122 RowSelector::skip(1),
5123 RowSelector::select(1),
5124 RowSelector::skip(1),
5125 ]);
5126 let mut reader = ParquetRecordBatchReaderBuilder::try_new(file)
5127 .unwrap()
5128 .with_row_selection(selection)
5129 .build()
5130 .unwrap();
5131
5132 let actual = reader.next().unwrap().unwrap();
5133 assert_eq!(actual.num_rows(), 1);
5134 assert_eq!(actual.column(0), &expected.column(0).slice(1, 1));
5135 }
5136
5137 #[test]
5138 fn test_arbitrary_decimal() {
5139 let values = [1, 2, 3, 4, 5, 6, 7, 8];
5140 let decimals_19_0 = Decimal128Array::from_iter_values(values)
5141 .with_precision_and_scale(19, 0)
5142 .unwrap();
5143 let decimals_12_0 = Decimal128Array::from_iter_values(values)
5144 .with_precision_and_scale(12, 0)
5145 .unwrap();
5146 let decimals_17_10 = Decimal128Array::from_iter_values(values)
5147 .with_precision_and_scale(17, 10)
5148 .unwrap();
5149
5150 let written = RecordBatch::try_from_iter([
5151 ("decimal_values_19_0", Arc::new(decimals_19_0) as ArrayRef),
5152 ("decimal_values_12_0", Arc::new(decimals_12_0) as ArrayRef),
5153 ("decimal_values_17_10", Arc::new(decimals_17_10) as ArrayRef),
5154 ])
5155 .unwrap();
5156
5157 let mut buffer = Vec::with_capacity(1024);
5158 let mut writer = ArrowWriter::try_new(&mut buffer, written.schema(), None).unwrap();
5159 writer.write(&written).unwrap();
5160 writer.close().unwrap();
5161
5162 let read = ParquetRecordBatchReader::try_new(Bytes::from(buffer), 8)
5163 .unwrap()
5164 .collect::<Result<Vec<_>, _>>()
5165 .unwrap();
5166
5167 assert_eq!(&written.slice(0, 8), &read[0]);
5168 }
5169
5170 #[test]
5171 fn test_list_skip() {
5172 let mut list = ListBuilder::new(Int32Builder::new());
5173 list.append_value([Some(1), Some(2)]);
5174 list.append_value([Some(3)]);
5175 list.append_value([Some(4)]);
5176 let list = list.finish();
5177 let batch = RecordBatch::try_from_iter([("l", Arc::new(list) as _)]).unwrap();
5178
5179 let props = WriterProperties::builder()
5181 .set_data_page_row_count_limit(1)
5182 .set_write_batch_size(2)
5183 .build();
5184
5185 let mut buffer = Vec::with_capacity(1024);
5186 let mut writer = ArrowWriter::try_new(&mut buffer, batch.schema(), Some(props)).unwrap();
5187 writer.write(&batch).unwrap();
5188 writer.close().unwrap();
5189
5190 let selection = vec![RowSelector::skip(2), RowSelector::select(1)];
5191 let mut reader = ParquetRecordBatchReaderBuilder::try_new(Bytes::from(buffer))
5192 .unwrap()
5193 .with_row_selection(selection.into())
5194 .build()
5195 .unwrap();
5196 let out = reader.next().unwrap().unwrap();
5197 assert_eq!(out.num_rows(), 1);
5198 assert_eq!(out, batch.slice(2, 1));
5199 }
5200
5201 fn test_decimal32_roundtrip() {
5202 let d = |values: Vec<i32>, p: u8| {
5203 let iter = values.into_iter();
5204 PrimitiveArray::<Decimal32Type>::from_iter_values(iter)
5205 .with_precision_and_scale(p, 2)
5206 .unwrap()
5207 };
5208
5209 let d1 = d(vec![1, 2, 3, 4, 5], 9);
5210 let batch = RecordBatch::try_from_iter([("d1", Arc::new(d1) as ArrayRef)]).unwrap();
5211
5212 let mut buffer = Vec::with_capacity(1024);
5213 let mut writer = ArrowWriter::try_new(&mut buffer, batch.schema(), None).unwrap();
5214 writer.write(&batch).unwrap();
5215 writer.close().unwrap();
5216
5217 let builder = ParquetRecordBatchReaderBuilder::try_new(Bytes::from(buffer)).unwrap();
5218 let t1 = builder.parquet_schema().columns()[0].physical_type();
5219 assert_eq!(t1, PhysicalType::INT32);
5220
5221 let mut reader = builder.build().unwrap();
5222 assert_eq!(batch.schema(), reader.schema());
5223
5224 let out = reader.next().unwrap().unwrap();
5225 assert_eq!(batch, out);
5226 }
5227
5228 fn test_decimal64_roundtrip() {
5229 let d = |values: Vec<i64>, p: u8| {
5233 let iter = values.into_iter();
5234 PrimitiveArray::<Decimal64Type>::from_iter_values(iter)
5235 .with_precision_and_scale(p, 2)
5236 .unwrap()
5237 };
5238
5239 let d1 = d(vec![1, 2, 3, 4, 5], 9);
5240 let d2 = d(vec![1, 2, 3, 4, 10.pow(10) - 1], 10);
5241 let d3 = d(vec![1, 2, 3, 4, 10.pow(18) - 1], 18);
5242
5243 let batch = RecordBatch::try_from_iter([
5244 ("d1", Arc::new(d1) as ArrayRef),
5245 ("d2", Arc::new(d2) as ArrayRef),
5246 ("d3", Arc::new(d3) as ArrayRef),
5247 ])
5248 .unwrap();
5249
5250 let mut buffer = Vec::with_capacity(1024);
5251 let mut writer = ArrowWriter::try_new(&mut buffer, batch.schema(), None).unwrap();
5252 writer.write(&batch).unwrap();
5253 writer.close().unwrap();
5254
5255 let builder = ParquetRecordBatchReaderBuilder::try_new(Bytes::from(buffer)).unwrap();
5256 let t1 = builder.parquet_schema().columns()[0].physical_type();
5257 assert_eq!(t1, PhysicalType::INT32);
5258 let t2 = builder.parquet_schema().columns()[1].physical_type();
5259 assert_eq!(t2, PhysicalType::INT64);
5260 let t3 = builder.parquet_schema().columns()[2].physical_type();
5261 assert_eq!(t3, PhysicalType::INT64);
5262
5263 let mut reader = builder.build().unwrap();
5264 assert_eq!(batch.schema(), reader.schema());
5265
5266 let out = reader.next().unwrap().unwrap();
5267 assert_eq!(batch, out);
5268 }
5269
5270 fn test_decimal_roundtrip<T: DecimalType>() {
5271 let d = |values: Vec<usize>, p: u8| {
5276 let iter = values.into_iter().map(T::Native::usize_as);
5277 PrimitiveArray::<T>::from_iter_values(iter)
5278 .with_precision_and_scale(p, 2)
5279 .unwrap()
5280 };
5281
5282 let d1 = d(vec![1, 2, 3, 4, 5], 9);
5283 let d2 = d(vec![1, 2, 3, 4, 10.pow(10) - 1], 10);
5284 let d3 = d(vec![1, 2, 3, 4, 10.pow(18) - 1], 18);
5285 let d4 = d(vec![1, 2, 3, 4, 10.pow(19) - 1], 19);
5286
5287 let batch = RecordBatch::try_from_iter([
5288 ("d1", Arc::new(d1) as ArrayRef),
5289 ("d2", Arc::new(d2) as ArrayRef),
5290 ("d3", Arc::new(d3) as ArrayRef),
5291 ("d4", Arc::new(d4) as ArrayRef),
5292 ])
5293 .unwrap();
5294
5295 let mut buffer = Vec::with_capacity(1024);
5296 let mut writer = ArrowWriter::try_new(&mut buffer, batch.schema(), None).unwrap();
5297 writer.write(&batch).unwrap();
5298 writer.close().unwrap();
5299
5300 let builder = ParquetRecordBatchReaderBuilder::try_new(Bytes::from(buffer)).unwrap();
5301 let t1 = builder.parquet_schema().columns()[0].physical_type();
5302 assert_eq!(t1, PhysicalType::INT32);
5303 let t2 = builder.parquet_schema().columns()[1].physical_type();
5304 assert_eq!(t2, PhysicalType::INT64);
5305 let t3 = builder.parquet_schema().columns()[2].physical_type();
5306 assert_eq!(t3, PhysicalType::INT64);
5307 let t4 = builder.parquet_schema().columns()[3].physical_type();
5308 assert_eq!(t4, PhysicalType::FIXED_LEN_BYTE_ARRAY);
5309
5310 let mut reader = builder.build().unwrap();
5311 assert_eq!(batch.schema(), reader.schema());
5312
5313 let out = reader.next().unwrap().unwrap();
5314 assert_eq!(batch, out);
5315 }
5316
5317 #[test]
5318 fn test_decimal() {
5319 test_decimal32_roundtrip();
5320 test_decimal64_roundtrip();
5321 test_decimal_roundtrip::<Decimal128Type>();
5322 test_decimal_roundtrip::<Decimal256Type>();
5323 }
5324
5325 #[test]
5326 #[cfg_attr(miri, ignore)] fn test_list_selection() {
5328 let schema = Arc::new(Schema::new(vec![Field::new_list(
5329 "list",
5330 Field::new_list_field(ArrowDataType::Utf8, true),
5331 false,
5332 )]));
5333 let mut buf = Vec::with_capacity(1024);
5334
5335 let mut writer = ArrowWriter::try_new(&mut buf, schema.clone(), None).unwrap();
5336
5337 for i in 0..2 {
5338 let mut list_a_builder = ListBuilder::new(StringBuilder::new());
5339 for j in 0..1024 {
5340 list_a_builder.values().append_value(format!("{i} {j}"));
5341 list_a_builder.append(true);
5342 }
5343 let batch =
5344 RecordBatch::try_new(schema.clone(), vec![Arc::new(list_a_builder.finish())])
5345 .unwrap();
5346 writer.write(&batch).unwrap();
5347 }
5348 let _metadata = writer.close().unwrap();
5349
5350 let buf = Bytes::from(buf);
5351 let reader = ParquetRecordBatchReaderBuilder::try_new(buf)
5352 .unwrap()
5353 .with_row_selection(RowSelection::from(vec![
5354 RowSelector::skip(100),
5355 RowSelector::select(924),
5356 RowSelector::skip(100),
5357 RowSelector::select(924),
5358 ]))
5359 .build()
5360 .unwrap();
5361
5362 let batches = reader.collect::<Result<Vec<_>, _>>().unwrap();
5363 let batch = concat_batches(&schema, &batches).unwrap();
5364
5365 assert_eq!(batch.num_rows(), 924 * 2);
5366 let list = batch.column(0).as_list::<i32>();
5367
5368 for w in list.value_offsets().windows(2) {
5369 assert_eq!(w[0] + 1, w[1])
5370 }
5371 let mut values = list.values().as_string::<i32>().iter();
5372
5373 for i in 0..2 {
5374 for j in 100..1024 {
5375 let expected = format!("{i} {j}");
5376 assert_eq!(values.next().unwrap().unwrap(), &expected);
5377 }
5378 }
5379 }
5380
5381 #[test]
5382 #[cfg_attr(miri, ignore)] fn test_list_selection_fuzz() {
5384 let mut rng = rng();
5385 let schema = Arc::new(Schema::new(vec![Field::new_list(
5386 "list",
5387 Field::new_list(
5388 Field::LIST_FIELD_DEFAULT_NAME,
5389 Field::new_list_field(ArrowDataType::Int32, true),
5390 true,
5391 ),
5392 true,
5393 )]));
5394 let mut buf = Vec::with_capacity(1024);
5395 let mut writer = ArrowWriter::try_new(&mut buf, schema.clone(), None).unwrap();
5396
5397 let mut list_a_builder = ListBuilder::new(ListBuilder::new(Int32Builder::new()));
5398
5399 for _ in 0..2048 {
5400 if rng.random_bool(0.2) {
5401 list_a_builder.append(false);
5402 continue;
5403 }
5404
5405 let list_a_len = rng.random_range(0..10);
5406 let list_b_builder = list_a_builder.values();
5407
5408 for _ in 0..list_a_len {
5409 if rng.random_bool(0.2) {
5410 list_b_builder.append(false);
5411 continue;
5412 }
5413
5414 let list_b_len = rng.random_range(0..10);
5415 let int_builder = list_b_builder.values();
5416 for _ in 0..list_b_len {
5417 match rng.random_bool(0.2) {
5418 true => int_builder.append_null(),
5419 false => int_builder.append_value(rng.random()),
5420 }
5421 }
5422 list_b_builder.append(true)
5423 }
5424 list_a_builder.append(true);
5425 }
5426
5427 let array = Arc::new(list_a_builder.finish());
5428 let batch = RecordBatch::try_new(schema, vec![array]).unwrap();
5429
5430 writer.write(&batch).unwrap();
5431 let _metadata = writer.close().unwrap();
5432
5433 let buf = Bytes::from(buf);
5434
5435 let cases = [
5436 vec![
5437 RowSelector::skip(100),
5438 RowSelector::select(924),
5439 RowSelector::skip(100),
5440 RowSelector::select(924),
5441 ],
5442 vec![
5443 RowSelector::select(924),
5444 RowSelector::skip(100),
5445 RowSelector::select(924),
5446 RowSelector::skip(100),
5447 ],
5448 vec![
5449 RowSelector::skip(1023),
5450 RowSelector::select(1),
5451 RowSelector::skip(1023),
5452 RowSelector::select(1),
5453 ],
5454 vec![
5455 RowSelector::select(1),
5456 RowSelector::skip(1023),
5457 RowSelector::select(1),
5458 RowSelector::skip(1023),
5459 ],
5460 ];
5461
5462 for batch_size in [100, 1024, 2048] {
5463 for selection in &cases {
5464 let selection = RowSelection::from(selection.clone());
5465 let reader = ParquetRecordBatchReaderBuilder::try_new(buf.clone())
5466 .unwrap()
5467 .with_row_selection(selection.clone())
5468 .with_batch_size(batch_size)
5469 .build()
5470 .unwrap();
5471
5472 let batches = reader.collect::<Result<Vec<_>, _>>().unwrap();
5473 let actual = concat_batches(batch.schema_ref(), &batches).unwrap();
5474 assert_eq!(actual.num_rows(), selection.row_count());
5475
5476 let mut batch_offset = 0;
5477 let mut actual_offset = 0;
5478 for selector in selection.iter() {
5479 if selector.skip {
5480 batch_offset += selector.row_count;
5481 continue;
5482 }
5483
5484 assert_eq!(
5485 batch.slice(batch_offset, selector.row_count),
5486 actual.slice(actual_offset, selector.row_count)
5487 );
5488
5489 batch_offset += selector.row_count;
5490 actual_offset += selector.row_count;
5491 }
5492 }
5493 }
5494 }
5495
5496 #[test]
5497 fn test_read_old_nested_list() {
5498 use arrow::datatypes::DataType;
5499 use arrow::datatypes::ToByteSlice;
5500
5501 let testdata = arrow::util::test_util::parquet_test_data();
5502 let path = format!("{testdata}/old_list_structure.parquet");
5511 let test_file = File::open(path).unwrap();
5512
5513 let a_values = Int32Array::from(vec![1, 2, 3, 4]);
5515
5516 let a_value_offsets = arrow::buffer::Buffer::from([0, 2, 4].to_byte_slice());
5518
5519 let a_list_data = ArrayData::builder(DataType::List(Arc::new(Field::new(
5521 "array",
5522 DataType::Int32,
5523 false,
5524 ))))
5525 .len(2)
5526 .add_buffer(a_value_offsets)
5527 .add_child_data(a_values.into_data())
5528 .build()
5529 .unwrap();
5530 let a = ListArray::from(a_list_data);
5531
5532 let builder = ParquetRecordBatchReaderBuilder::try_new(test_file).unwrap();
5533 let mut reader = builder.build().unwrap();
5534 let out = reader.next().unwrap().unwrap();
5535 assert_eq!(out.num_rows(), 1);
5536 assert_eq!(out.num_columns(), 1);
5537 let c0 = out.column(0);
5539 let c0arr = c0.as_any().downcast_ref::<ListArray>().unwrap();
5540 let r0 = c0arr.value(0);
5542 let r0arr = r0.as_any().downcast_ref::<ListArray>().unwrap();
5543 assert_eq!(r0arr, &a);
5544 }
5545
5546 #[test]
5547 fn test_read_row_numbers() {
5548 let file = write_parquet_from_iter(vec![(
5549 "value",
5550 Arc::new(Int64Array::from(vec![1, 2, 3])) as ArrayRef,
5551 )]);
5552 let supplied_fields = Fields::from(vec![Field::new("value", ArrowDataType::Int64, false)]);
5553
5554 let row_number_field = Arc::new(
5555 Field::new("row_number", ArrowDataType::Int64, false).with_extension_type(RowNumber),
5556 );
5557
5558 let options = ArrowReaderOptions::new()
5559 .with_schema(Arc::new(Schema::new(supplied_fields)))
5560 .with_virtual_columns(vec![row_number_field.clone()])
5561 .unwrap();
5562 let mut arrow_reader = ParquetRecordBatchReaderBuilder::try_new_with_options(
5563 file.try_clone().unwrap(),
5564 options,
5565 )
5566 .expect("reader builder with schema")
5567 .build()
5568 .expect("reader with schema");
5569
5570 let batch = arrow_reader.next().unwrap().unwrap();
5571 let schema = Arc::new(Schema::new(vec![
5572 Field::new("value", ArrowDataType::Int64, false),
5573 (*row_number_field).clone(),
5574 ]));
5575
5576 assert_eq!(batch.schema(), schema);
5577 assert_eq!(batch.num_columns(), 2);
5578 assert_eq!(batch.num_rows(), 3);
5579 assert_eq!(
5580 batch
5581 .column(0)
5582 .as_primitive::<types::Int64Type>()
5583 .iter()
5584 .collect::<Vec<_>>(),
5585 vec![Some(1), Some(2), Some(3)]
5586 );
5587 assert_eq!(
5588 batch
5589 .column(1)
5590 .as_primitive::<types::Int64Type>()
5591 .iter()
5592 .collect::<Vec<_>>(),
5593 vec![Some(0), Some(1), Some(2)]
5594 );
5595 }
5596
5597 #[test]
5598 fn test_read_only_row_numbers() {
5599 let file = write_parquet_from_iter(vec![(
5600 "value",
5601 Arc::new(Int64Array::from(vec![1, 2, 3])) as ArrayRef,
5602 )]);
5603 let row_number_field = Arc::new(
5604 Field::new("row_number", ArrowDataType::Int64, false).with_extension_type(RowNumber),
5605 );
5606 let options = ArrowReaderOptions::new()
5607 .with_virtual_columns(vec![row_number_field.clone()])
5608 .unwrap();
5609 let metadata = ArrowReaderMetadata::load(&file, options).unwrap();
5610 let num_columns = metadata
5611 .metadata
5612 .file_metadata()
5613 .schema_descr()
5614 .num_columns();
5615
5616 let mut arrow_reader = ParquetRecordBatchReaderBuilder::new_with_metadata(file, metadata)
5617 .with_projection(ProjectionMask::none(num_columns))
5618 .build()
5619 .expect("reader with schema");
5620
5621 let batch = arrow_reader.next().unwrap().unwrap();
5622 let schema = Arc::new(Schema::new(vec![row_number_field]));
5623
5624 assert_eq!(batch.schema(), schema);
5625 assert_eq!(batch.num_columns(), 1);
5626 assert_eq!(batch.num_rows(), 3);
5627 assert_eq!(
5628 batch
5629 .column(0)
5630 .as_primitive::<types::Int64Type>()
5631 .iter()
5632 .collect::<Vec<_>>(),
5633 vec![Some(0), Some(1), Some(2)]
5634 );
5635 }
5636
5637 #[test]
5638 fn test_read_row_numbers_row_group_order() -> Result<()> {
5639 let array = Int64Array::from_iter_values(5000..5100);
5641 let batch = RecordBatch::try_from_iter([("col", Arc::new(array) as ArrayRef)])?;
5642 let mut buffer = Vec::new();
5643 let options = WriterProperties::builder()
5644 .set_max_row_group_row_count(Some(50))
5645 .build();
5646 let mut writer = ArrowWriter::try_new(&mut buffer, batch.schema().clone(), Some(options))?;
5647 for batch_chunk in (0..10).map(|i| batch.slice(i * 10, 10)) {
5649 writer.write(&batch_chunk)?;
5650 }
5651 writer.close()?;
5652
5653 let row_number_field = Arc::new(
5654 Field::new("row_number", ArrowDataType::Int64, false).with_extension_type(RowNumber),
5655 );
5656
5657 let buffer = Bytes::from(buffer);
5658
5659 let options =
5660 ArrowReaderOptions::new().with_virtual_columns(vec![row_number_field.clone()])?;
5661
5662 let arrow_reader =
5664 ParquetRecordBatchReaderBuilder::try_new_with_options(buffer.clone(), options.clone())?
5665 .build()?;
5666
5667 assert_eq!(
5668 ValuesAndRowNumbers {
5669 values: (5000..5100).collect(),
5670 row_numbers: (0..100).collect()
5671 },
5672 ValuesAndRowNumbers::new_from_reader(arrow_reader)
5673 );
5674
5675 let arrow_reader = ParquetRecordBatchReaderBuilder::try_new_with_options(buffer, options)?
5677 .with_row_groups(vec![1, 0])
5678 .build()?;
5679
5680 assert_eq!(
5681 ValuesAndRowNumbers {
5682 values: (5050..5100).chain(5000..5050).collect(),
5683 row_numbers: (50..100).chain(0..50).collect(),
5684 },
5685 ValuesAndRowNumbers::new_from_reader(arrow_reader)
5686 );
5687
5688 Ok(())
5689 }
5690
5691 #[test]
5697 fn test_mixed_row_group_ordinals() -> Result<()> {
5698 use crate::file::metadata::{ParquetMetaDataReader, RowGroupMetaData};
5699
5700 let array = Int64Array::from_iter_values(5000..5100);
5702 let batch = RecordBatch::try_from_iter([("col", Arc::new(array) as ArrayRef)])?;
5703 let mut buffer = Vec::new();
5704 let props = WriterProperties::builder()
5705 .set_max_row_group_row_count(Some(25))
5706 .build();
5707 let mut writer = ArrowWriter::try_new(&mut buffer, batch.schema().clone(), Some(props))?;
5708 for batch_chunk in (0..10).map(|i| batch.slice(i * 10, 10)) {
5709 writer.write(&batch_chunk)?;
5710 }
5711 writer.close()?;
5712 let buffer = Bytes::from(buffer);
5713
5714 let metadata = ParquetMetaDataReader::new().parse_and_finish(&buffer)?;
5717 let schema_descr = metadata.file_metadata().schema_descr_ptr();
5718 let mut row_groups = metadata.row_groups().to_vec();
5719 let stripped = row_groups[1].clone();
5720 let mut builder = RowGroupMetaData::builder(schema_descr)
5721 .set_num_rows(stripped.num_rows())
5722 .set_total_byte_size(stripped.total_byte_size())
5723 .set_sorting_columns(stripped.sorting_columns().cloned())
5724 .set_column_metadata(stripped.columns().to_vec());
5725 if let Some(offset) = stripped.file_offset() {
5726 builder = builder.set_file_offset(offset);
5727 }
5728 row_groups[1] = builder.build()?;
5729 assert_eq!(row_groups[1].ordinal(), None);
5730 let metadata = Arc::new(metadata.into_builder().set_row_groups(row_groups).build());
5731
5732 let arrow_metadata =
5734 ArrowReaderMetadata::try_new(Arc::clone(&metadata), ArrowReaderOptions::new())?;
5735 let reader =
5736 ParquetRecordBatchReaderBuilder::new_with_metadata(buffer.clone(), arrow_metadata)
5737 .build()?;
5738 let values: Vec<i64> = reader
5739 .flat_map(|batch| {
5740 let batch = batch.expect("could not read batch");
5741 batch
5742 .column(0)
5743 .as_primitive::<types::Int64Type>()
5744 .values()
5745 .to_vec()
5746 })
5747 .collect();
5748 assert_eq!(values, (5000..5100).collect::<Vec<_>>());
5749
5750 let row_number_field = Arc::new(
5753 Field::new("row_number", ArrowDataType::Int64, false).with_extension_type(RowNumber),
5754 );
5755 let options = ArrowReaderOptions::new().with_virtual_columns(vec![row_number_field])?;
5756 let arrow_metadata = ArrowReaderMetadata::try_new(Arc::clone(&metadata), options)?;
5757 let result = ParquetRecordBatchReaderBuilder::new_with_metadata(buffer, arrow_metadata)
5758 .with_row_groups(vec![0]) .build()
5760 .and_then(|mut reader| reader.next().transpose().map_err(|e| e.into()));
5761 let err = result.expect_err("row numbers over mixed ordinals must fail");
5762 assert!(
5763 err.to_string().contains("inconsistent row-group ordinals"),
5764 "unexpected error: {err}"
5765 );
5766
5767 Ok(())
5768 }
5769
5770 #[derive(Debug, PartialEq)]
5771 struct ValuesAndRowNumbers {
5772 values: Vec<i64>,
5773 row_numbers: Vec<i64>,
5774 }
5775 impl ValuesAndRowNumbers {
5776 fn new_from_reader(reader: ParquetRecordBatchReader) -> Self {
5777 let mut values = vec![];
5778 let mut row_numbers = vec![];
5779 for batch in reader {
5780 let batch = batch.expect("Could not read batch");
5781 values.extend(
5782 batch
5783 .column_by_name("col")
5784 .expect("Could not get col column")
5785 .as_primitive::<arrow::datatypes::Int64Type>()
5786 .iter()
5787 .map(|v| v.expect("Could not get value")),
5788 );
5789
5790 row_numbers.extend(
5791 batch
5792 .column_by_name("row_number")
5793 .expect("Could not get row_number column")
5794 .as_primitive::<arrow::datatypes::Int64Type>()
5795 .iter()
5796 .map(|v| v.expect("Could not get row number"))
5797 .collect::<Vec<_>>(),
5798 );
5799 }
5800 Self {
5801 values,
5802 row_numbers,
5803 }
5804 }
5805 }
5806
5807 #[test]
5808 fn test_with_virtual_columns_rejects_non_virtual_fields() {
5809 let regular_field = Arc::new(Field::new("regular_column", ArrowDataType::Int64, false));
5811 assert_eq!(
5812 ArrowReaderOptions::new()
5813 .with_virtual_columns(vec![regular_field])
5814 .unwrap_err()
5815 .to_string(),
5816 "Parquet error: Field 'regular_column' is not a virtual column. Virtual columns must have extension type names starting with 'arrow.virtual.'"
5817 );
5818 }
5819
5820 #[test]
5821 #[cfg_attr(miri, ignore)] fn test_row_numbers_with_multiple_row_groups() {
5823 test_row_numbers_with_multiple_row_groups_helper(
5824 false,
5825 |path, selection, _row_filter, batch_size| {
5826 let file = File::open(path).unwrap();
5827 let row_number_field = Arc::new(
5828 Field::new("row_number", ArrowDataType::Int64, false)
5829 .with_extension_type(RowNumber),
5830 );
5831 let options = ArrowReaderOptions::new()
5832 .with_virtual_columns(vec![row_number_field])
5833 .unwrap();
5834 let reader = ParquetRecordBatchReaderBuilder::try_new_with_options(file, options)
5835 .unwrap()
5836 .with_row_selection(selection)
5837 .with_batch_size(batch_size)
5838 .build()
5839 .expect("Could not create reader");
5840 reader
5841 .collect::<Result<Vec<_>, _>>()
5842 .expect("Could not read")
5843 },
5844 );
5845 }
5846
5847 #[test]
5848 #[cfg_attr(miri, ignore)] fn test_row_numbers_with_multiple_row_groups_and_filter() {
5850 test_row_numbers_with_multiple_row_groups_helper(
5851 true,
5852 |path, selection, row_filter, batch_size| {
5853 let file = File::open(path).unwrap();
5854 let row_number_field = Arc::new(
5855 Field::new("row_number", ArrowDataType::Int64, false)
5856 .with_extension_type(RowNumber),
5857 );
5858 let options = ArrowReaderOptions::new()
5859 .with_virtual_columns(vec![row_number_field])
5860 .unwrap();
5861 let reader = ParquetRecordBatchReaderBuilder::try_new_with_options(file, options)
5862 .unwrap()
5863 .with_row_selection(selection)
5864 .with_batch_size(batch_size)
5865 .with_row_filter(row_filter.expect("No filter"))
5866 .build()
5867 .expect("Could not create reader");
5868 reader
5869 .collect::<Result<Vec<_>, _>>()
5870 .expect("Could not read")
5871 },
5872 );
5873 }
5874
5875 #[test]
5876 fn test_read_row_group_indices() {
5877 let array1 = Int64Array::from(vec![1, 2]);
5879 let array2 = Int64Array::from(vec![3, 4]);
5880 let array3 = Int64Array::from(vec![5, 6]);
5881
5882 let batch1 =
5883 RecordBatch::try_from_iter(vec![("value", Arc::new(array1) as ArrayRef)]).unwrap();
5884 let batch2 =
5885 RecordBatch::try_from_iter(vec![("value", Arc::new(array2) as ArrayRef)]).unwrap();
5886 let batch3 =
5887 RecordBatch::try_from_iter(vec![("value", Arc::new(array3) as ArrayRef)]).unwrap();
5888
5889 let mut buffer = Vec::new();
5890 let options = WriterProperties::builder()
5891 .set_max_row_group_row_count(Some(2))
5892 .build();
5893 let mut writer = ArrowWriter::try_new(&mut buffer, batch1.schema(), Some(options)).unwrap();
5894 writer.write(&batch1).unwrap();
5895 writer.write(&batch2).unwrap();
5896 writer.write(&batch3).unwrap();
5897 writer.close().unwrap();
5898
5899 let file = Bytes::from(buffer);
5900 let row_group_index_field = Arc::new(
5901 Field::new("row_group_index", ArrowDataType::Int64, false)
5902 .with_extension_type(RowGroupIndex),
5903 );
5904
5905 let options = ArrowReaderOptions::new()
5906 .with_virtual_columns(vec![row_group_index_field.clone()])
5907 .unwrap();
5908 let mut arrow_reader =
5909 ParquetRecordBatchReaderBuilder::try_new_with_options(file.clone(), options)
5910 .expect("reader builder with virtual columns")
5911 .build()
5912 .expect("reader with virtual columns");
5913
5914 let batch = arrow_reader.next().unwrap().unwrap();
5915
5916 assert_eq!(batch.num_columns(), 2);
5917 assert_eq!(batch.num_rows(), 6);
5918
5919 assert_eq!(
5920 batch
5921 .column(0)
5922 .as_primitive::<types::Int64Type>()
5923 .iter()
5924 .collect::<Vec<_>>(),
5925 vec![Some(1), Some(2), Some(3), Some(4), Some(5), Some(6)]
5926 );
5927
5928 assert_eq!(
5929 batch
5930 .column(1)
5931 .as_primitive::<types::Int64Type>()
5932 .iter()
5933 .collect::<Vec<_>>(),
5934 vec![Some(0), Some(0), Some(1), Some(1), Some(2), Some(2)]
5935 );
5936 }
5937
5938 #[test]
5939 fn test_read_only_row_group_indices() {
5940 let array1 = Int64Array::from(vec![1, 2, 3]);
5941 let array2 = Int64Array::from(vec![4, 5]);
5942
5943 let batch1 =
5944 RecordBatch::try_from_iter(vec![("value", Arc::new(array1) as ArrayRef)]).unwrap();
5945 let batch2 =
5946 RecordBatch::try_from_iter(vec![("value", Arc::new(array2) as ArrayRef)]).unwrap();
5947
5948 let mut buffer = Vec::new();
5949 let options = WriterProperties::builder()
5950 .set_max_row_group_row_count(Some(3))
5951 .build();
5952 let mut writer = ArrowWriter::try_new(&mut buffer, batch1.schema(), Some(options)).unwrap();
5953 writer.write(&batch1).unwrap();
5954 writer.write(&batch2).unwrap();
5955 writer.close().unwrap();
5956
5957 let file = Bytes::from(buffer);
5958 let row_group_index_field = Arc::new(
5959 Field::new("row_group_index", ArrowDataType::Int64, false)
5960 .with_extension_type(RowGroupIndex),
5961 );
5962
5963 let options = ArrowReaderOptions::new()
5964 .with_virtual_columns(vec![row_group_index_field.clone()])
5965 .unwrap();
5966 let metadata = ArrowReaderMetadata::load(&file, options).unwrap();
5967 let num_columns = metadata
5968 .metadata
5969 .file_metadata()
5970 .schema_descr()
5971 .num_columns();
5972
5973 let mut arrow_reader = ParquetRecordBatchReaderBuilder::new_with_metadata(file, metadata)
5974 .with_projection(ProjectionMask::none(num_columns))
5975 .build()
5976 .expect("reader with virtual columns only");
5977
5978 let batch = arrow_reader.next().unwrap().unwrap();
5979 let schema = Arc::new(Schema::new(vec![(*row_group_index_field).clone()]));
5980
5981 assert_eq!(batch.schema(), schema);
5982 assert_eq!(batch.num_columns(), 1);
5983 assert_eq!(batch.num_rows(), 5);
5984
5985 assert_eq!(
5986 batch
5987 .column(0)
5988 .as_primitive::<types::Int64Type>()
5989 .iter()
5990 .collect::<Vec<_>>(),
5991 vec![Some(0), Some(0), Some(0), Some(1), Some(1)]
5992 );
5993 }
5994
5995 #[test]
5996 fn test_read_row_group_indices_with_selection() -> Result<()> {
5997 let mut buffer = Vec::new();
5998 let options = WriterProperties::builder()
5999 .set_max_row_group_row_count(Some(10))
6000 .build();
6001
6002 let schema = Arc::new(Schema::new(vec![Field::new(
6003 "value",
6004 ArrowDataType::Int64,
6005 false,
6006 )]));
6007
6008 let mut writer = ArrowWriter::try_new(&mut buffer, schema.clone(), Some(options))?;
6009
6010 for i in 0..3 {
6012 let start = i * 10;
6013 let array = Int64Array::from_iter_values(start..start + 10);
6014 let batch = RecordBatch::try_from_iter(vec![("value", Arc::new(array) as ArrayRef)])?;
6015 writer.write(&batch)?;
6016 }
6017 writer.close()?;
6018
6019 let file = Bytes::from(buffer);
6020 let row_group_index_field = Arc::new(
6021 Field::new("rg_idx", ArrowDataType::Int64, false).with_extension_type(RowGroupIndex),
6022 );
6023
6024 let options =
6025 ArrowReaderOptions::new().with_virtual_columns(vec![row_group_index_field])?;
6026
6027 let arrow_reader =
6029 ParquetRecordBatchReaderBuilder::try_new_with_options(file.clone(), options.clone())?
6030 .with_row_groups(vec![2, 1, 0])
6031 .build()?;
6032
6033 let batches: Vec<_> = arrow_reader.collect::<Result<Vec<_>, _>>()?;
6034 let combined = concat_batches(&batches[0].schema(), &batches)?;
6035
6036 let values = combined.column(0).as_primitive::<types::Int64Type>();
6037 let first_val = values.value(0);
6038 let last_val = values.value(combined.num_rows() - 1);
6039 assert_eq!(first_val, 20);
6041 assert_eq!(last_val, 9);
6043
6044 let rg_indices = combined.column(1).as_primitive::<types::Int64Type>();
6045 assert_eq!(rg_indices.value(0), 2);
6046 assert_eq!(rg_indices.value(10), 1);
6047 assert_eq!(rg_indices.value(20), 0);
6048
6049 Ok(())
6050 }
6051
6052 pub(crate) fn test_row_numbers_with_multiple_row_groups_helper<F>(
6053 use_filter: bool,
6054 test_case: F,
6055 ) where
6056 F: FnOnce(PathBuf, RowSelection, Option<RowFilter>, usize) -> Vec<RecordBatch>,
6057 {
6058 let seed: u64 = random();
6059 println!("test_row_numbers_with_multiple_row_groups seed: {seed}");
6060 let mut rng = StdRng::seed_from_u64(seed);
6061
6062 use tempfile::TempDir;
6063 let tempdir = TempDir::new().expect("Could not create temp dir");
6064
6065 let (bytes, metadata) = generate_file_with_row_numbers(&mut rng);
6066
6067 let path = tempdir.path().join("test.parquet");
6068 std::fs::write(&path, bytes).expect("Could not write file");
6069
6070 let mut case = vec![];
6071 let mut remaining = metadata.file_metadata().num_rows();
6072 while remaining > 0 {
6073 let row_count = rng.random_range(1..=remaining);
6074 remaining -= row_count;
6075 case.push(RowSelector {
6076 row_count: row_count as usize,
6077 skip: rng.random_bool(0.5),
6078 });
6079 }
6080
6081 let filter = use_filter.then(|| {
6082 let filter = (0..metadata.file_metadata().num_rows())
6083 .map(|_| rng.random_bool(0.99))
6084 .collect::<Vec<_>>();
6085 let mut filter_offset = 0;
6086 RowFilter::new(vec![Box::new(ArrowPredicateFn::new(
6087 ProjectionMask::all(),
6088 move |b| {
6089 let array = BooleanArray::from_iter(
6090 filter
6091 .iter()
6092 .skip(filter_offset)
6093 .take(b.num_rows())
6094 .map(|x| Some(*x)),
6095 );
6096 filter_offset += b.num_rows();
6097 Ok(array)
6098 },
6099 ))])
6100 });
6101
6102 let selection = RowSelection::from(case);
6103 let batches = test_case(path, selection.clone(), filter, rng.random_range(1..4096));
6104
6105 if selection.skipped_row_count() == metadata.file_metadata().num_rows() as usize {
6106 assert!(batches.into_iter().all(|batch| batch.num_rows() == 0));
6107 return;
6108 }
6109 let actual = concat_batches(batches.first().expect("No batches").schema_ref(), &batches)
6110 .expect("Failed to concatenate");
6111 let values = actual
6113 .column(0)
6114 .as_primitive::<types::Int64Type>()
6115 .iter()
6116 .collect::<Vec<_>>();
6117 let row_numbers = actual
6118 .column(1)
6119 .as_primitive::<types::Int64Type>()
6120 .iter()
6121 .collect::<Vec<_>>();
6122 assert_eq!(
6123 row_numbers
6124 .into_iter()
6125 .map(|number| number.map(|number| number + 1))
6126 .collect::<Vec<_>>(),
6127 values
6128 );
6129 }
6130
6131 fn generate_file_with_row_numbers(rng: &mut impl Rng) -> (Bytes, ParquetMetaData) {
6132 let schema = Arc::new(Schema::new(Fields::from(vec![Field::new(
6133 "value",
6134 ArrowDataType::Int64,
6135 false,
6136 )])));
6137
6138 let mut buf = Vec::with_capacity(1024);
6139 let mut writer =
6140 ArrowWriter::try_new(&mut buf, schema.clone(), None).expect("Could not create writer");
6141
6142 let mut values = 1..=rng.random_range(1..4096);
6143 while !values.is_empty() {
6144 let batch_values = values
6145 .by_ref()
6146 .take(rng.random_range(1..4096))
6147 .collect::<Vec<_>>();
6148 let array = Arc::new(Int64Array::from(batch_values)) as ArrayRef;
6149 let batch =
6150 RecordBatch::try_from_iter([("value", array)]).expect("Could not create batch");
6151 writer.write(&batch).expect("Could not write batch");
6152 writer.flush().expect("Could not flush");
6153 }
6154 let metadata = writer.close().expect("Could not close writer");
6155
6156 (Bytes::from(buf), metadata)
6157 }
6158}