1use std::collections::HashMap;
24use std::sync::Arc;
25use std::sync::atomic::AtomicU64;
26
27use arrow_schema::{DataType, Field};
28use futures::{StreamExt, TryStreamExt};
29use parquet::arrow::arrow_reader::{ArrowReaderMetadata, ArrowReaderOptions};
30use parquet::arrow::{
31 PARQUET_FIELD_ID_META_KEY, ParquetRecordBatchStreamBuilder, ProjectionMask, RowNumber,
32};
33use parquet::encryption::decrypt::FileDecryptionProperties;
34
35use super::row_lineage::synthesize_row_id_column;
36use super::{
37 ArrowFileReader, ArrowReader, ParquetReadOptions, add_fallback_field_ids_to_arrow_schema,
38 apply_name_mapping_to_arrow_schema, find_leaf_by_field_id,
39};
40use crate::arrow::build_partition_constant;
41use crate::arrow::caching_delete_file_loader::CachingDeleteFileLoader;
42use crate::arrow::int96::coerce_int96_timestamps;
43use crate::arrow::record_batch_transformer::RecordBatchTransformerBuilder;
44use crate::arrow::scan_metrics::{CountingFileRead, ScanMetrics, ScanResult};
45use crate::encryption::StandardKeyMetadata;
46use crate::error::{Result, invalid_data};
47use crate::expr::BoundPredicate;
48use crate::expr::visitors::bloom_filter_evaluator::{
49 BloomFilterEvaluator, ColumnBloomFilter, collect_bloom_filter_field_ids,
50};
51use crate::io::{FileIO, FileMetadata, FileRead};
52use crate::metadata_columns::{
53 RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER, RESERVED_COL_NAME_POS,
54 RESERVED_COL_NAME_ROW_ID, RESERVED_FIELD_ID_FILE,
55 RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER, RESERVED_FIELD_ID_PARTITION,
56 RESERVED_FIELD_ID_POS, RESERVED_FIELD_ID_ROW_ID, RESERVED_FIELD_ID_SPEC_ID, is_metadata_field,
57};
58use crate::scan::{ArrowRecordBatchStream, FileScanTask, FileScanTaskStream};
59use crate::spec::{Datum, PartitionSpec, Struct};
60use crate::{Error, ErrorKind};
61
62impl ArrowReader {
63 pub fn read(self, tasks: FileScanTaskStream) -> Result<ScanResult> {
66 let concurrency_limit_data_files = self.concurrency_limit_data_files;
67 let scan_metrics = ScanMetrics::new();
68
69 let task_reader = FileScanTaskReader {
70 batch_size: self.batch_size,
71 file_io: self.file_io,
72 delete_file_loader: self
73 .delete_file_loader
74 .with_scan_metrics(scan_metrics.clone()),
75 row_group_filtering_enabled: self.row_group_filtering_enabled,
76 row_selection_enabled: self.row_selection_enabled,
77 bloom_filter_enabled: self.bloom_filter_enabled,
78 parquet_read_options: self.parquet_read_options,
79 scan_metrics: scan_metrics.clone(),
80 };
81
82 let stream: ArrowRecordBatchStream = if concurrency_limit_data_files == 1 {
84 Box::pin(
85 tasks
86 .and_then(move |task| task_reader.clone().process(task))
87 .map_err(|err| {
88 Error::new(ErrorKind::Unexpected, "file scan task generate failed")
89 .with_source(err)
90 })
91 .try_flatten(),
92 )
93 } else {
94 Box::pin(
95 tasks
96 .map_ok(move |task| task_reader.clone().process(task))
97 .map_err(|err| {
98 Error::new(ErrorKind::Unexpected, "file scan task generate failed")
99 .with_source(err)
100 })
101 .try_buffer_unordered(concurrency_limit_data_files)
102 .try_flatten_unordered(concurrency_limit_data_files),
103 )
104 };
105
106 Ok(ScanResult::new(stream, scan_metrics))
107 }
108}
109
110const PRUNABLE_METADATA_FIELDS: &[i32] = &[
114 RESERVED_FIELD_ID_FILE,
115 RESERVED_FIELD_ID_SPEC_ID,
116 RESERVED_FIELD_ID_PARTITION,
117 RESERVED_FIELD_ID_POS,
118 RESERVED_FIELD_ID_ROW_ID,
119 RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
120];
121
122#[derive(Clone)]
125struct FileScanTaskReader {
126 batch_size: Option<usize>,
127 file_io: FileIO,
128 delete_file_loader: CachingDeleteFileLoader,
129 row_group_filtering_enabled: bool,
130 row_selection_enabled: bool,
131 bloom_filter_enabled: bool,
132 parquet_read_options: ParquetReadOptions,
133 scan_metrics: ScanMetrics,
134}
135
136impl FileScanTaskReader {
137 async fn process(self, task: FileScanTask) -> Result<ArrowRecordBatchStream> {
138 let should_load_page_index = (self.row_selection_enabled && task.predicate().is_some())
139 || !task.deletes().is_empty();
140 let mut parquet_read_options = self.parquet_read_options;
141 parquet_read_options.preload_page_index = should_load_page_index;
142
143 let delete_filter_rx = self
144 .delete_file_loader
145 .load_deletes(task.deletes(), task.schema_ref());
146
147 let (parquet_file_reader, arrow_metadata) = ArrowReader::open_parquet_file(
149 task.data_file_path(),
150 &self.file_io,
151 task.file_size_in_bytes(),
152 parquet_read_options,
153 self.scan_metrics.bytes_read_counter(),
154 task.key_metadata(),
155 )
156 .await?;
157
158 let missing_field_ids = arrow_metadata
162 .schema()
163 .fields()
164 .iter()
165 .next()
166 .is_some_and(|f| f.metadata().get(PARQUET_FIELD_ID_META_KEY).is_none());
167
168 let use_position_fallback = missing_field_ids && task.name_mapping().is_none();
173
174 let project_pos = task.project_field_ids().contains(&RESERVED_FIELD_ID_POS);
175 let project_row_id = task.project_field_ids().contains(&RESERVED_FIELD_ID_ROW_ID);
176
177 let need_row_number = project_pos || (project_row_id && task.first_row_id().is_some());
182
183 let field_ids = task.project_field_ids();
184 let metadata_only_projection = !field_ids.is_empty()
185 && field_ids
186 .iter()
187 .all(|id| PRUNABLE_METADATA_FIELDS.contains(id));
188
189 let install_row_number = need_row_number || metadata_only_projection;
190
191 let arrow_metadata = Self::configure_arrow_reader_metadata(
192 arrow_metadata,
193 &task,
194 missing_field_ids,
195 install_row_number,
196 )?;
197
198 let mut record_batch_stream_builder =
200 ParquetRecordBatchStreamBuilder::new_with_metadata(parquet_file_reader, arrow_metadata);
201
202 let project_last_updated_seq = task
206 .project_field_ids()
207 .contains(&RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER);
208
209 let phys_last_updated_seq_leaf = if project_last_updated_seq {
215 find_leaf_by_field_id(
216 record_batch_stream_builder.parquet_schema(),
217 RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
218 )
219 } else {
220 None
221 };
222
223 let last_updated_seq_present_by_name_only = project_last_updated_seq
226 && phys_last_updated_seq_leaf.is_none()
227 && record_batch_stream_builder
228 .schema()
229 .fields()
230 .iter()
231 .any(|f| f.name() == RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER);
232
233 let coalesce_last_updated_seq_leaf = phys_last_updated_seq_leaf
238 .filter(|_| task.first_row_id().is_some() && task.data_sequence_number().is_some());
239
240 let phys_row_id_leaf = if project_row_id {
241 find_leaf_by_field_id(
242 record_batch_stream_builder.parquet_schema(),
243 RESERVED_FIELD_ID_ROW_ID,
244 )
245 } else {
246 None
247 };
248
249 let row_id_present_by_name_only = project_row_id
255 && !use_position_fallback
256 && phys_row_id_leaf.is_none()
257 && record_batch_stream_builder
258 .schema()
259 .fields()
260 .iter()
261 .any(|f| f.name() == RESERVED_COL_NAME_ROW_ID);
262
263 let coalesce_row_id_leaf = phys_row_id_leaf.filter(|_| task.first_row_id().is_some());
266
267 let project_field_ids_without_metadata: Vec<i32> = task
269 .project_field_ids()
270 .iter()
271 .filter(|&&id| !is_metadata_field(id))
272 .copied()
273 .collect();
274
275 let mut projection_mask = ArrowReader::get_arrow_projection_mask(
281 &project_field_ids_without_metadata,
282 task.schema(),
283 record_batch_stream_builder.parquet_schema(),
284 record_batch_stream_builder.schema(),
285 use_position_fallback, )?;
287
288 if project_field_ids_without_metadata.is_empty() && install_row_number {
299 projection_mask =
300 ProjectionMask::none(record_batch_stream_builder.parquet_schema().num_columns());
301 }
302
303 for leaf in [coalesce_last_updated_seq_leaf, coalesce_row_id_leaf]
308 .into_iter()
309 .flatten()
310 {
311 let phys_mask =
312 ProjectionMask::leaves(record_batch_stream_builder.parquet_schema(), vec![leaf]);
313 projection_mask.union(&phys_mask);
314 }
315
316 record_batch_stream_builder =
317 record_batch_stream_builder.with_projection(projection_mask.clone());
318
319 let mut record_batch_transformer_builder =
323 RecordBatchTransformerBuilder::new(task.schema_ref(), task.project_field_ids());
324
325 if task.project_field_ids().contains(&RESERVED_FIELD_ID_FILE) {
327 let file_datum = Datum::string(task.data_file_path().to_string());
328 record_batch_transformer_builder =
329 record_batch_transformer_builder.with_constant(RESERVED_FIELD_ID_FILE, file_datum);
330 }
331
332 if task
333 .project_field_ids()
334 .contains(&RESERVED_FIELD_ID_SPEC_ID)
335 {
336 let partition_spec = task
337 .partition_spec()
338 .ok_or_else(|| Error::new(ErrorKind::Unexpected, "Partition spec is missing"))?;
339
340 let spec_id_datum = Datum::int(partition_spec.spec_id());
341 record_batch_transformer_builder = record_batch_transformer_builder
342 .with_constant(RESERVED_FIELD_ID_SPEC_ID, spec_id_datum);
343 }
344
345 if project_last_updated_seq {
346 record_batch_transformer_builder =
351 match (task.first_row_id(), task.data_sequence_number()) {
352 (Some(_), Some(seq)) => {
353 let datum = Datum::long(seq);
354 if coalesce_last_updated_seq_leaf.is_some() {
355 record_batch_transformer_builder.with_coalesced_last_updated_seq_column(
358 RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
359 datum,
360 )
361 } else if last_updated_seq_present_by_name_only {
362 return Err(Error::new(
370 ErrorKind::FeatureUnsupported,
371 "Reading a physically-stored _last_updated_sequence_number column \
372 without an embedded field id is not supported",
373 ));
374 } else {
375 record_batch_transformer_builder.with_constant(
377 RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
378 datum,
379 )
380 }
381 }
382 (None, _) => record_batch_transformer_builder.with_null_metadata_column(
384 RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
385 )?,
386 (Some(_), None) => {
390 return Err(invalid_data!(
391 "Data file {} has a first_row_id but no data sequence number",
392 task.data_file_path()
393 ));
394 }
395 };
396 }
397
398 if project_row_id {
399 if task.first_row_id().is_some() && row_id_present_by_name_only {
405 return Err(Error::new(
406 ErrorKind::FeatureUnsupported,
407 "Reading a physically-stored _row_id column without an embedded field id \
408 is not supported",
409 ));
410 }
411
412 record_batch_transformer_builder =
416 record_batch_transformer_builder.with_virtual_field(RESERVED_FIELD_ID_ROW_ID);
417 }
418
419 if let (Some(partition_spec), Some(partition_data)) =
420 (task.partition_spec(), task.partition())
421 {
422 record_batch_transformer_builder = record_batch_transformer_builder
423 .with_partition(Arc::clone(partition_spec), partition_data.clone())?;
424 }
425
426 if project_pos {
427 record_batch_transformer_builder =
428 record_batch_transformer_builder.with_virtual_field(RESERVED_FIELD_ID_POS);
429 }
430
431 if task
434 .project_field_ids()
435 .contains(&RESERVED_FIELD_ID_PARTITION)
436 && let Some(unified_type) = task.unified_partition_type()
437 {
438 let (spec, partition_data) = match (task.partition_spec(), task.partition()) {
439 (Some(spec), Some(data)) => (Arc::clone(spec), data.clone()),
440 _ if unified_type.fields().is_empty() => {
445 (Arc::new(PartitionSpec::unpartition_spec()), Struct::empty())
446 }
447 _ => {
448 return Err(Error::new(
449 ErrorKind::Unexpected,
450 "cannot build _partition column: unified partition type has fields \
451 but the scan task is missing its partition spec or data",
452 ));
453 }
454 };
455 let constant = build_partition_constant(unified_type, &spec, &partition_data)?;
456 record_batch_transformer_builder =
457 record_batch_transformer_builder.with_partition_constant(constant);
458 }
459
460 let mut record_batch_transformer = record_batch_transformer_builder.build();
461
462 if let Some(batch_size) = self.batch_size {
463 record_batch_stream_builder = record_batch_stream_builder.with_batch_size(batch_size);
464 }
465
466 let delete_filter = delete_filter_rx.await.unwrap()?;
467 let delete_predicate = delete_filter.build_equality_delete_predicate(&task).await?;
468
469 let final_predicate = match (task.predicate(), delete_predicate) {
474 (None, None) => None,
475 (Some(predicate), None) => Some(predicate.clone()),
476 (None, Some(ref predicate)) => Some(predicate.clone()),
477 (Some(filter_predicate), Some(delete_predicate)) => {
478 Some(filter_predicate.clone().and(delete_predicate))
479 }
480 };
481
482 let mut selected_row_group_indices = None;
498 let mut row_selection = None;
499
500 if task.start() != 0 || task.length() != 0 {
503 let byte_range_filtered_row_groups = ArrowReader::filter_row_groups_by_byte_range(
504 record_batch_stream_builder.metadata(),
505 task.start(),
506 task.length(),
507 )?;
508 selected_row_group_indices = Some(byte_range_filtered_row_groups);
509 }
510
511 if let Some(predicate) = final_predicate {
512 let (iceberg_field_ids, field_id_map) = ArrowReader::build_field_id_set_and_map(
513 record_batch_stream_builder.parquet_schema(),
514 record_batch_stream_builder.schema(),
515 &predicate,
516 use_position_fallback,
517 )?;
518
519 let row_filter = ArrowReader::get_row_filter(
520 &predicate,
521 record_batch_stream_builder.parquet_schema(),
522 &iceberg_field_ids,
523 &field_id_map,
524 )?;
525 record_batch_stream_builder = record_batch_stream_builder.with_row_filter(row_filter);
526
527 if self.row_group_filtering_enabled {
528 let predicate_filtered_row_groups = ArrowReader::get_selected_row_group_indices(
529 &predicate,
530 record_batch_stream_builder.metadata(),
531 &field_id_map,
532 task.schema(),
533 )?;
534
535 selected_row_group_indices = match selected_row_group_indices {
538 Some(byte_range_filtered) => {
539 let intersection: Vec<usize> = byte_range_filtered
541 .into_iter()
542 .filter(|idx| predicate_filtered_row_groups.contains(idx))
543 .collect();
544 Some(intersection)
545 }
546 None => Some(predicate_filtered_row_groups),
547 };
548 }
549
550 if self.bloom_filter_enabled {
551 let all_rgs;
552 let candidate_rgs = match &selected_row_group_indices {
553 Some(indices) => indices.as_slice(),
554 None => {
555 all_rgs = (0..record_batch_stream_builder.metadata().num_row_groups())
556 .collect::<Vec<_>>();
557 &all_rgs
558 }
559 };
560
561 let bloom_filtered = Self::filter_row_groups_by_bloom_filter(
562 &predicate,
563 &mut record_batch_stream_builder,
564 candidate_rgs,
565 &field_id_map,
566 )
567 .await?;
568
569 if bloom_filtered.len() < candidate_rgs.len() {
570 selected_row_group_indices = Some(bloom_filtered);
571 }
572 }
573
574 if self.row_selection_enabled {
575 row_selection = ArrowReader::get_row_selection_for_filter_predicate(
576 &predicate,
577 record_batch_stream_builder.metadata(),
578 &selected_row_group_indices,
579 &field_id_map,
580 task.schema(),
581 )?;
582 }
583 }
584
585 let positional_delete_indexes = delete_filter.get_delete_vector(&task);
586
587 if let Some(positional_delete_indexes) = positional_delete_indexes {
588 let delete_row_selection = {
589 let positional_delete_indexes = positional_delete_indexes.lock().unwrap();
590
591 ArrowReader::build_deletes_row_selection(
592 record_batch_stream_builder.metadata().row_groups(),
593 &selected_row_group_indices,
594 &positional_delete_indexes,
595 )
596 }?;
597
598 row_selection = match row_selection {
601 None => Some(delete_row_selection),
602 Some(filter_row_selection) => {
603 Some(filter_row_selection.intersection(&delete_row_selection))
604 }
605 };
606 }
607
608 if let Some(row_selection) = row_selection {
609 record_batch_stream_builder =
610 record_batch_stream_builder.with_row_selection(row_selection);
611 }
612
613 if let Some(selected_row_group_indices) = selected_row_group_indices {
614 record_batch_stream_builder =
615 record_batch_stream_builder.with_row_groups(selected_row_group_indices);
616 }
617
618 let first_row_id = task.first_row_id();
623 let record_batch_stream = record_batch_stream_builder.build()?.map(move |batch| {
624 let mut batch = batch.map_err(|err| -> Error { err.into() })?;
625 if project_row_id {
626 batch = synthesize_row_id_column(batch, first_row_id)?;
627 }
628 record_batch_transformer.process_record_batch(batch)
630 });
631
632 Ok(Box::pin(record_batch_stream) as ArrowRecordBatchStream)
633 }
634
635 fn configure_arrow_reader_metadata(
638 arrow_metadata: ArrowReaderMetadata,
639 task: &FileScanTask,
640 missing_field_ids: bool,
641 install_row_number: bool,
642 ) -> Result<ArrowReaderMetadata> {
643 let arrow_schema = if missing_field_ids {
652 let schema = if let Some(name_mapping) = task.name_mapping() {
653 apply_name_mapping_to_arrow_schema(
654 Arc::clone(arrow_metadata.schema()),
655 name_mapping,
656 )?
657 } else {
658 add_fallback_field_ids_to_arrow_schema(arrow_metadata.schema())
659 };
660 coerce_int96_timestamps(&schema, task.schema()).unwrap_or(schema)
664 } else if let Some(coerced) =
665 coerce_int96_timestamps(arrow_metadata.schema(), task.schema())
666 {
667 coerced
668 } else if install_row_number {
669 Arc::clone(arrow_metadata.schema())
670 } else {
671 return Ok(arrow_metadata);
672 };
673
674 let mut options = ArrowReaderOptions::new().with_schema(Arc::clone(&arrow_schema));
675 if install_row_number {
676 let row_number_field = Arc::new(
677 Field::new(RESERVED_COL_NAME_POS, DataType::Int64, false)
678 .with_metadata(HashMap::from([(
679 PARQUET_FIELD_ID_META_KEY.to_string(),
680 RESERVED_FIELD_ID_POS.to_string(),
681 )]))
682 .with_extension_type(RowNumber),
683 );
684 options = options.with_virtual_columns(vec![row_number_field])?;
685 }
686
687 ArrowReaderMetadata::try_new(Arc::clone(arrow_metadata.metadata()), options).map_err(|e| {
688 Error::new(
689 ErrorKind::Unexpected,
690 format!(
691 "Failed to create ArrowReaderMetadata with the configured reader options \
692 (missing_field_ids: {missing_field_ids}, \
693 install_row_number: {install_row_number}, schema: {arrow_schema})"
694 ),
695 )
696 .with_source(e)
697 })
698 }
699
700 async fn filter_row_groups_by_bloom_filter(
703 predicate: &BoundPredicate,
704 builder: &mut ParquetRecordBatchStreamBuilder<ArrowFileReader>,
705 candidate_row_groups: &[usize],
706 field_id_map: &HashMap<i32, usize>,
707 ) -> Result<Vec<usize>> {
708 let bloom_filter_field_ids: Vec<i32> = collect_bloom_filter_field_ids(predicate)?
711 .into_iter()
712 .filter(|id| field_id_map.contains_key(id))
713 .collect();
714
715 if bloom_filter_field_ids.is_empty() {
716 return Ok(candidate_row_groups.to_vec());
717 }
718
719 let mut result = Vec::with_capacity(candidate_row_groups.len());
720
721 for &rg_idx in candidate_row_groups {
722 let mut bloom_filters: HashMap<i32, ColumnBloomFilter> = HashMap::new();
723
724 for &field_id in &bloom_filter_field_ids {
725 let col_idx = field_id_map[&field_id];
726 let col_meta = builder.metadata().row_group(rg_idx).column(col_idx);
727
728 if col_meta.bloom_filter_offset().is_none() {
730 continue;
731 }
732
733 let physical_type = col_meta.column_type();
734 let type_length = col_meta.column_descr().type_length();
735
736 match builder
737 .get_row_group_column_bloom_filter(rg_idx, col_idx)
738 .await
739 {
740 Ok(Some(sbbf)) => {
741 bloom_filters.insert(
742 field_id,
743 ColumnBloomFilter::new(sbbf, physical_type, type_length),
744 );
745 }
746 Ok(None) => {}
747 Err(e) => {
748 tracing::debug!(
751 "Bloom filter for field {field_id} in row group {rg_idx} could not be read: {e}"
752 );
753 }
754 }
755 }
756
757 match BloomFilterEvaluator::eval(predicate, &bloom_filters) {
758 Ok(true) => result.push(rg_idx),
759 Ok(false) => { }
760 Err(e) => {
761 tracing::debug!(
762 "Bloom filter evaluation failed for row group {rg_idx}, including it: {e}"
763 );
764 result.push(rg_idx);
765 }
766 }
767 }
768
769 Ok(result)
770 }
771}
772
773impl ArrowReader {
774 pub(crate) async fn open_parquet_file(
777 data_file_path: &str,
778 file_io: &FileIO,
779 file_size_in_bytes: u64,
780 parquet_read_options: ParquetReadOptions,
781 bytes_read: &Arc<AtomicU64>,
782 key_metadata: Option<&[u8]>,
783 ) -> Result<(ArrowFileReader, ArrowReaderMetadata)> {
784 let parquet_file = file_io.new_input(data_file_path)?;
785 let counting_reader =
786 CountingFileRead::new(parquet_file.reader().await?, Arc::clone(bytes_read));
787 Self::build_parquet_reader(
788 Box::new(counting_reader),
789 file_size_in_bytes,
790 parquet_read_options,
791 key_metadata,
792 )
793 .await
794 }
795
796 async fn build_parquet_reader(
797 parquet_reader: Box<dyn FileRead>,
798 file_size_in_bytes: u64,
799 parquet_read_options: ParquetReadOptions,
800 key_metadata: Option<&[u8]>,
801 ) -> Result<(ArrowFileReader, ArrowReaderMetadata)> {
802 let mut reader = ArrowFileReader::new(
803 FileMetadata {
804 size: file_size_in_bytes,
805 },
806 parquet_reader,
807 )
808 .with_parquet_read_options(parquet_read_options);
809
810 let arrow_reader_options = Self::build_arrow_reader_options(key_metadata)?;
811
812 let arrow_metadata = ArrowReaderMetadata::load_async(&mut reader, arrow_reader_options)
813 .await
814 .map_err(|e| {
815 Error::new(ErrorKind::Unexpected, "Failed to load Parquet metadata").with_source(e)
816 })?;
817
818 Ok((reader, arrow_metadata))
819 }
820
821 fn build_arrow_reader_options(key_metadata: Option<&[u8]>) -> Result<ArrowReaderOptions> {
824 match key_metadata {
825 Some(km) => {
826 let standard_key_metadata = StandardKeyMetadata::decode(km)?;
827 let mut builder = FileDecryptionProperties::builder(
828 standard_key_metadata.encryption_key().as_bytes().to_vec(),
829 );
830 if let Some(aad) = standard_key_metadata.aad_prefix() {
831 builder = builder.with_aad_prefix(aad.to_vec());
832 }
833 let decryption_properties = builder.build().map_err(|e| {
834 Error::new(
835 ErrorKind::Unexpected,
836 "Failed to build Parquet file decryption properties",
837 )
838 .with_source(e)
839 })?;
840 Ok(
841 ArrowReaderOptions::new()
842 .with_file_decryption_properties(decryption_properties),
843 )
844 }
845 None => Ok(ArrowReaderOptions::default()),
846 }
847 }
848}
849
850#[cfg(test)]
851mod tests {
852 use std::collections::HashMap;
853 use std::fs::File;
854 use std::sync::Arc;
855
856 use arrow_array::cast::AsArray;
857 use arrow_array::{Array, ArrayRef, Int32Array, Int64Array, RecordBatch, StringArray};
858 use arrow_cast::cast;
859 use arrow_schema::{DataType, Field, Schema as ArrowSchema};
860 use futures::TryStreamExt;
861 use parquet::arrow::{ArrowWriter, PARQUET_FIELD_ID_META_KEY};
862 use parquet::basic::Compression;
863 use parquet::file::properties::WriterProperties;
864 use tempfile::TempDir;
865
866 use crate::Runtime;
867 use crate::arrow::ArrowReaderBuilder;
868 use crate::arrow::test_utils::write_encrypted_parquet;
869 use crate::io::FileIO;
870 use crate::metadata_columns::{
871 RESERVED_COL_NAME_POS, RESERVED_COL_NAME_ROW_ID, RESERVED_FIELD_ID_FILE,
872 RESERVED_FIELD_ID_POS, RESERVED_FIELD_ID_ROW_ID,
873 };
874 use crate::scan::{FileScanTask, FileScanTaskDeleteFile, FileScanTaskStream};
875 use crate::spec::{DataFileFormat, NestedField, PrimitiveType, Schema, SchemaRef, Type};
876
877 const UNIX_EPOCH_JULIAN: i64 = 2_440_588;
880 const MICROS_PER_DAY: i64 = 86_400_000_000;
881 const INT96_TEST_NANOS_WITHIN_DAY: u64 = 43_200_000_000_000;
883 const INT96_TEST_JULIAN_DAY: u32 = 2_953_529;
884
885 fn make_int96_test_value() -> (parquet::data_type::Int96, i64) {
886 let mut val = parquet::data_type::Int96::new();
887 val.set_data(
888 (INT96_TEST_NANOS_WITHIN_DAY & 0xFFFFFFFF) as u32,
889 (INT96_TEST_NANOS_WITHIN_DAY >> 32) as u32,
890 INT96_TEST_JULIAN_DAY,
891 );
892 let expected_micros = (INT96_TEST_JULIAN_DAY as i64 - UNIX_EPOCH_JULIAN) * MICROS_PER_DAY
893 + (INT96_TEST_NANOS_WITHIN_DAY / 1_000) as i64;
894 (val, expected_micros)
895 }
896
897 async fn read_int96_batches(
898 file_path: &str,
899 schema: SchemaRef,
900 project_field_ids: Vec<i32>,
901 ) -> Vec<RecordBatch> {
902 let file_io = FileIO::new_with_fs();
903 let reader = ArrowReaderBuilder::new(file_io, Runtime::current()).build();
904
905 let file_size = std::fs::metadata(file_path).unwrap().len();
906 let task = FileScanTask::builder()
907 .with_file_size_in_bytes(file_size)
908 .with_start(0)
909 .with_length(file_size)
910 .with_data_file_path(file_path.to_string())
911 .with_data_file_format(DataFileFormat::Parquet)
912 .with_schema(schema)
913 .with_project_field_ids(project_field_ids)
914 .with_case_sensitive(false)
915 .build()
916 .unwrap();
917
918 let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
919 reader
920 .read(tasks)
921 .unwrap()
922 .stream()
923 .try_collect()
924 .await
925 .unwrap()
926 }
927
928 fn write_int96_parquet_file(
930 table_location: &str,
931 filename: &str,
932 with_field_ids: bool,
933 ) -> (String, Vec<i64>) {
934 use parquet::basic::{Repetition, Type as PhysicalType};
935 use parquet::data_type::{Int32Type, Int96, Int96Type};
936 use parquet::file::writer::SerializedFileWriter;
937 use parquet::schema::types::Type as SchemaType;
938
939 let file_path = format!("{table_location}/{filename}");
940
941 let mut ts_builder = SchemaType::primitive_type_builder("ts", PhysicalType::INT96)
942 .with_repetition(Repetition::OPTIONAL);
943 let mut id_builder = SchemaType::primitive_type_builder("id", PhysicalType::INT32)
944 .with_repetition(Repetition::REQUIRED);
945
946 if with_field_ids {
947 ts_builder = ts_builder.with_id(Some(1));
948 id_builder = id_builder.with_id(Some(2));
949 }
950
951 let schema = SchemaType::group_type_builder("schema")
952 .with_fields(vec![
953 Arc::new(ts_builder.build().unwrap()),
954 Arc::new(id_builder.build().unwrap()),
955 ])
956 .build()
957 .unwrap();
958
959 const NOON_NANOS: u64 = INT96_TEST_NANOS_WITHIN_DAY;
961 const JULIAN_3333: u32 = INT96_TEST_JULIAN_DAY;
962 const JULIAN_2100: u32 = 2_488_070;
963
964 let test_data: Vec<(u32, u32, u32, i64)> = vec![
965 (
967 0,
968 0,
969 JULIAN_3333,
970 (JULIAN_3333 as i64 - UNIX_EPOCH_JULIAN) * MICROS_PER_DAY,
971 ),
972 (
974 (NOON_NANOS & 0xFFFFFFFF) as u32,
975 (NOON_NANOS >> 32) as u32,
976 JULIAN_3333,
977 (JULIAN_3333 as i64 - UNIX_EPOCH_JULIAN) * MICROS_PER_DAY
978 + (NOON_NANOS / 1_000) as i64,
979 ),
980 (
982 0,
983 0,
984 JULIAN_2100,
985 (JULIAN_2100 as i64 - UNIX_EPOCH_JULIAN) * MICROS_PER_DAY,
986 ),
987 ];
988
989 let int96_values: Vec<Int96> = test_data
990 .iter()
991 .map(|(lo, hi, day, _)| {
992 let mut v = Int96::new();
993 v.set_data(*lo, *hi, *day);
994 v
995 })
996 .collect();
997
998 let id_values: Vec<i32> = (0..test_data.len() as i32).collect();
999 let expected_micros: Vec<i64> = test_data.iter().map(|(_, _, _, m)| *m).collect();
1000
1001 let file = File::create(&file_path).unwrap();
1002 let mut writer =
1003 SerializedFileWriter::new(file, Arc::new(schema), Default::default()).unwrap();
1004
1005 let mut row_group = writer.next_row_group().unwrap();
1006 {
1007 let mut col = row_group.next_column().unwrap().unwrap();
1009 col.typed::<Int96Type>()
1010 .write_batch(&int96_values, Some(&vec![1; test_data.len()]), None)
1011 .unwrap();
1012 col.close().unwrap();
1013 }
1014 {
1015 let mut col = row_group.next_column().unwrap().unwrap();
1016 col.typed::<Int32Type>()
1017 .write_batch(&id_values, None, None)
1018 .unwrap();
1019 col.close().unwrap();
1020 }
1021 row_group.close().unwrap();
1022 writer.close().unwrap();
1023
1024 (file_path, expected_micros)
1025 }
1026
1027 async fn assert_int96_read_matches(
1028 file_path: &str,
1029 schema: SchemaRef,
1030 project_field_ids: Vec<i32>,
1031 expected_micros: &[i64],
1032 ) {
1033 use arrow_array::TimestampMicrosecondArray;
1034
1035 let batches = read_int96_batches(file_path, schema, project_field_ids).await;
1036
1037 assert_eq!(batches.len(), 1);
1038 let ts_array = batches[0]
1039 .column(0)
1040 .as_any()
1041 .downcast_ref::<TimestampMicrosecondArray>()
1042 .expect("Expected TimestampMicrosecondArray");
1043
1044 for (i, expected) in expected_micros.iter().enumerate() {
1045 assert_eq!(
1046 ts_array.value(i),
1047 *expected,
1048 "Row {i}: got {}, expected {expected}",
1049 ts_array.value(i)
1050 );
1051 }
1052 }
1053
1054 async fn assert_encrypted_parquet_roundtrip(encryption_key: &[u8]) {
1058 let aad_prefix = b"aad_prefix";
1059
1060 let schema = Arc::new(
1061 Schema::builder()
1062 .with_schema_id(1)
1063 .with_fields(vec![
1064 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1065 ])
1066 .build()
1067 .unwrap(),
1068 );
1069
1070 let arrow_schema = Arc::new(ArrowSchema::new(vec![
1071 Field::new("id", DataType::Int32, false).with_metadata(HashMap::from([(
1072 PARQUET_FIELD_ID_META_KEY.to_string(),
1073 "1".to_string(),
1074 )])),
1075 ]));
1076
1077 let tmp_dir = TempDir::new().unwrap();
1078 let table_location = tmp_dir.path().to_str().unwrap().to_string();
1079 let file_io = FileIO::new_with_fs();
1080
1081 let id_data = Arc::new(Int32Array::from(vec![10, 20, 30])) as ArrayRef;
1082 let batch = RecordBatch::try_new(arrow_schema.clone(), vec![id_data]).unwrap();
1083
1084 let file_path = format!("{table_location}/encrypted.parquet");
1085 write_encrypted_parquet(&file_path, &batch, encryption_key, Some(aad_prefix));
1086
1087 let key_metadata = crate::encryption::StandardKeyMetadata::try_new(encryption_key)
1088 .unwrap()
1089 .with_aad_prefix(aad_prefix)
1090 .encode()
1091 .unwrap();
1092
1093 let reader = ArrowReaderBuilder::new(file_io, Runtime::current()).build();
1094
1095 let task = FileScanTask::builder()
1096 .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
1097 .with_start(0)
1098 .with_length(0)
1099 .with_data_file_path(file_path)
1100 .with_data_file_format(DataFileFormat::Parquet)
1101 .with_schema(schema)
1102 .with_project_field_ids(vec![1])
1103 .with_case_sensitive(false)
1104 .with_key_metadata(Some(key_metadata))
1105 .build()
1106 .unwrap();
1107
1108 let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
1109 let batches: Vec<RecordBatch> = reader
1110 .read(tasks)
1111 .unwrap()
1112 .stream()
1113 .try_collect()
1114 .await
1115 .unwrap();
1116
1117 assert_eq!(batches.len(), 1);
1118 let ids = batches[0]
1119 .column(0)
1120 .as_any()
1121 .downcast_ref::<Int32Array>()
1122 .unwrap();
1123 assert_eq!(ids.values(), &[10, 20, 30]);
1124 }
1125
1126 #[tokio::test]
1127 async fn test_read_encrypted_parquet_aes_128() {
1128 assert_encrypted_parquet_roundtrip(b"0123456789abcdef").await;
1129 }
1130
1131 #[tokio::test]
1132 async fn test_read_encrypted_parquet_aes_256() {
1133 assert_encrypted_parquet_roundtrip(b"0123456789abcdef0123456789abcdef").await;
1134 }
1135
1136 #[tokio::test]
1137 async fn test_read_encrypted_parquet_without_key_metadata_fails() {
1138 let encryption_key = b"0123456789abcdef";
1139
1140 let schema = Arc::new(
1141 Schema::builder()
1142 .with_schema_id(1)
1143 .with_fields(vec![
1144 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1145 ])
1146 .build()
1147 .unwrap(),
1148 );
1149
1150 let arrow_schema = Arc::new(ArrowSchema::new(vec![
1151 Field::new("id", DataType::Int32, false).with_metadata(HashMap::from([(
1152 PARQUET_FIELD_ID_META_KEY.to_string(),
1153 "1".to_string(),
1154 )])),
1155 ]));
1156
1157 let tmp_dir = TempDir::new().unwrap();
1158 let table_location = tmp_dir.path().to_str().unwrap().to_string();
1159 let file_io = FileIO::new_with_fs();
1160
1161 let id_data = Arc::new(Int32Array::from(vec![1, 2, 3])) as ArrayRef;
1162 let batch = RecordBatch::try_new(arrow_schema.clone(), vec![id_data]).unwrap();
1163
1164 let file_path = format!("{table_location}/encrypted_no_key.parquet");
1165 write_encrypted_parquet(&file_path, &batch, encryption_key, None);
1166
1167 let reader = ArrowReaderBuilder::new(file_io, Runtime::current()).build();
1168
1169 let task = FileScanTask::builder()
1170 .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
1171 .with_start(0)
1172 .with_length(0)
1173 .with_data_file_path(file_path)
1174 .with_data_file_format(DataFileFormat::Parquet)
1175 .with_schema(schema)
1176 .with_project_field_ids(vec![1])
1177 .with_case_sensitive(false)
1178 .build()
1179 .unwrap();
1180
1181 let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
1182 let result: Result<Vec<RecordBatch>, _> =
1183 reader.read(tasks).unwrap().stream().try_collect().await;
1184
1185 let err = result.unwrap_err();
1186 assert_eq!(err.kind(), crate::ErrorKind::Unexpected);
1187 let err_str = format!("{err}");
1188 assert!(
1189 err_str.contains("encrypted footer"),
1190 "Expected error about encrypted footer, got: {err_str}"
1191 );
1192 assert!(
1193 err_str.contains("decryption properties were not provided"),
1194 "Expected error about missing decryption properties, got: {err_str}"
1195 );
1196 }
1197
1198 fn write_plain_parquet(
1201 dir: &str,
1202 name: &str,
1203 extra_fields: Vec<Field>,
1204 extra_columns: Vec<ArrayRef>,
1205 ) -> String {
1206 let mut fields =
1207 vec![
1208 Field::new("id", DataType::Int32, false).with_metadata(HashMap::from([(
1209 PARQUET_FIELD_ID_META_KEY.to_string(),
1210 "1".to_string(),
1211 )])),
1212 ];
1213 fields.extend(extra_fields);
1214 let arrow_schema = Arc::new(ArrowSchema::new(fields));
1215
1216 let mut columns: Vec<ArrayRef> = vec![Arc::new(Int32Array::from(vec![1, 2, 3]))];
1217 columns.extend(extra_columns);
1218 let batch = RecordBatch::try_new(arrow_schema.clone(), columns).unwrap();
1219
1220 let file_path = format!("{dir}/{name}");
1221 let file = File::create(&file_path).unwrap();
1222 let props = WriterProperties::builder()
1223 .set_compression(Compression::SNAPPY)
1224 .build();
1225 let mut writer = ArrowWriter::try_new(file, arrow_schema, Some(props)).unwrap();
1226 writer.write(&batch).unwrap();
1227 writer.close().unwrap();
1228 file_path
1229 }
1230
1231 fn last_updated_seq_task(
1232 file_path: String,
1233 first_row_id: Option<i64>,
1234 data_sequence_number: Option<i64>,
1235 ) -> FileScanTask {
1236 use crate::metadata_columns::RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER;
1237
1238 let schema = Arc::new(
1239 Schema::builder()
1240 .with_schema_id(1)
1241 .with_fields(vec![
1242 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1243 ])
1244 .build()
1245 .unwrap(),
1246 );
1247
1248 FileScanTask::builder()
1249 .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
1250 .with_start(0)
1251 .with_length(0)
1252 .with_data_file_path(file_path)
1253 .with_data_file_format(DataFileFormat::Parquet)
1254 .with_schema(schema)
1255 .with_project_field_ids(vec![1, RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER])
1256 .with_first_row_id(first_row_id)
1257 .with_data_sequence_number(data_sequence_number)
1258 .with_case_sensitive(false)
1259 .build()
1260 .unwrap()
1261 }
1262
1263 fn assert_last_updated_seq_column(batches: &[RecordBatch], expected: &[Option<i64>]) {
1266 use arrow_array::cast::AsArray;
1267 use arrow_cast::cast;
1268 use arrow_schema::DataType;
1269
1270 use crate::metadata_columns::RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER;
1271
1272 let mut actual = Vec::new();
1273 for batch in batches {
1274 let col = batch
1275 .column_by_name(RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER)
1276 .expect("_last_updated_sequence_number column should be present");
1277 let logical = cast(col, &DataType::Int64).unwrap();
1278 let values = logical.as_primitive::<arrow_array::types::Int64Type>();
1279 for i in 0..values.len() {
1280 actual.push((!values.is_null(i)).then(|| values.value(i)));
1281 }
1282 }
1283 assert_eq!(actual, expected);
1284 }
1285
1286 #[tokio::test]
1287 async fn test_last_updated_sequence_number_null_when_no_first_row_id() {
1288 let tmp_dir = TempDir::new().unwrap();
1289 let dir = tmp_dir.path().to_str().unwrap();
1290 let file_path = write_plain_parquet(dir, "no_first_row_id.parquet", vec![], vec![]);
1291
1292 let task = last_updated_seq_task(file_path, None, Some(9));
1296
1297 let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
1298 let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
1299 let batches: Vec<RecordBatch> = reader
1300 .read(tasks)
1301 .unwrap()
1302 .stream()
1303 .try_collect()
1304 .await
1305 .unwrap();
1306
1307 assert_last_updated_seq_column(&batches, &[None, None, None]);
1308 }
1309
1310 #[tokio::test]
1311 async fn test_last_updated_sequence_number_error_when_no_data_seq() {
1312 let tmp_dir = TempDir::new().unwrap();
1313 let dir = tmp_dir.path().to_str().unwrap();
1314 let file_path = write_plain_parquet(dir, "no_data_seq.parquet", vec![], vec![]);
1315
1316 let task = last_updated_seq_task(file_path, Some(42), None);
1320
1321 let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
1322 let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
1323 let result: Result<Vec<RecordBatch>, _> =
1324 reader.read(tasks).unwrap().stream().try_collect().await;
1325
1326 let err = result.unwrap_err();
1327 assert_eq!(err.kind(), crate::ErrorKind::DataInvalid);
1328 assert!(
1329 format!("{err}").contains("no data sequence number"),
1330 "unexpected error: {err}"
1331 );
1332 }
1333
1334 #[tokio::test]
1335 async fn test_last_updated_sequence_number_derived_from_data_seq() {
1336 let tmp_dir = TempDir::new().unwrap();
1337 let dir = tmp_dir.path().to_str().unwrap();
1338 let file_path = write_plain_parquet(dir, "with_first_row_id.parquet", vec![], vec![]);
1339
1340 let task = last_updated_seq_task(file_path, Some(42), Some(7));
1343
1344 let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
1345 let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
1346 let batches: Vec<RecordBatch> = reader
1347 .read(tasks)
1348 .unwrap()
1349 .stream()
1350 .try_collect()
1351 .await
1352 .unwrap();
1353
1354 assert_last_updated_seq_column(&batches, &[Some(7), Some(7), Some(7)]);
1355 }
1356
1357 #[tokio::test]
1358 async fn test_last_updated_sequence_number_mixed_files_share_schema() {
1359 use arrow_select::concat::concat_batches;
1360
1361 let tmp_dir = TempDir::new().unwrap();
1362 let dir = tmp_dir.path().to_str().unwrap();
1363
1364 let constant = last_updated_seq_task(
1370 write_plain_parquet(dir, "constant.parquet", vec![], vec![]),
1371 Some(42),
1372 Some(7),
1373 );
1374 let nulled = last_updated_seq_task(
1375 write_plain_parquet(dir, "nulled.parquet", vec![], vec![]),
1376 None,
1377 Some(7),
1378 );
1379 let coalesced = last_updated_seq_task(
1380 write_plain_parquet(
1381 dir,
1382 "coalesced.parquet",
1383 vec![physical_last_updated_seq_field()],
1384 vec![Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef],
1385 ),
1386 Some(50),
1387 Some(7),
1388 );
1389
1390 let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
1391 let tasks = Box::pin(futures::stream::iter(vec![
1392 Ok(constant),
1393 Ok(nulled),
1394 Ok(coalesced),
1395 ])) as FileScanTaskStream;
1396 let batches: Vec<RecordBatch> = reader
1397 .read(tasks)
1398 .unwrap()
1399 .stream()
1400 .try_collect()
1401 .await
1402 .unwrap();
1403
1404 assert_eq!(batches.len(), 3);
1405 let schema = batches[0].schema();
1407 concat_batches(&schema, &batches)
1408 .expect("constant, null and coalesce files must share one column type");
1409 }
1410
1411 fn physical_last_updated_seq_field() -> Field {
1413 use crate::metadata_columns::{
1414 RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER,
1415 RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
1416 };
1417 Field::new(
1418 RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER,
1419 DataType::Int64,
1420 true,
1421 )
1422 .with_metadata(HashMap::from([(
1423 PARQUET_FIELD_ID_META_KEY.to_string(),
1424 RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER.to_string(),
1425 )]))
1426 }
1427
1428 #[tokio::test]
1429 async fn test_last_updated_sequence_number_physical_column_coalesced() {
1430 let tmp_dir = TempDir::new().unwrap();
1431 let dir = tmp_dir.path().to_str().unwrap();
1432 let seq_col = Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef;
1436 let file_path = write_plain_parquet(
1437 dir,
1438 "with_seq.parquet",
1439 vec![physical_last_updated_seq_field()],
1440 vec![seq_col],
1441 );
1442
1443 let task = last_updated_seq_task(file_path, Some(100), Some(9));
1444
1445 let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
1446 let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
1447 let batches: Vec<RecordBatch> = reader
1448 .read(tasks)
1449 .unwrap()
1450 .stream()
1451 .try_collect()
1452 .await
1453 .unwrap();
1454
1455 assert_last_updated_seq_column(&batches, &[Some(5), Some(9), Some(8)]);
1457 }
1458
1459 #[tokio::test]
1460 async fn test_last_updated_sequence_number_coalesced_with_pos_column() {
1461 use crate::metadata_columns::{
1462 RESERVED_COL_NAME_POS, RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
1463 RESERVED_FIELD_ID_POS,
1464 };
1465
1466 let tmp_dir = TempDir::new().unwrap();
1467 let dir = tmp_dir.path().to_str().unwrap();
1468 let seq_col = Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef;
1469 let file_path = write_plain_parquet(
1470 dir,
1471 "with_seq_and_pos.parquet",
1472 vec![physical_last_updated_seq_field()],
1473 vec![seq_col],
1474 );
1475
1476 let schema = Arc::new(
1481 Schema::builder()
1482 .with_schema_id(1)
1483 .with_fields(vec![
1484 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1485 ])
1486 .build()
1487 .unwrap(),
1488 );
1489 let task = FileScanTask::builder()
1490 .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
1491 .with_start(0)
1492 .with_length(0)
1493 .with_data_file_path(file_path)
1494 .with_data_file_format(DataFileFormat::Parquet)
1495 .with_schema(schema)
1496 .with_project_field_ids(vec![
1497 1,
1498 RESERVED_FIELD_ID_POS,
1499 RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
1500 ])
1501 .with_first_row_id(Some(100))
1502 .with_data_sequence_number(Some(9))
1503 .with_case_sensitive(false)
1504 .build()
1505 .unwrap();
1506
1507 let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
1508 let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
1509 let batches: Vec<RecordBatch> = reader
1510 .read(tasks)
1511 .unwrap()
1512 .stream()
1513 .try_collect()
1514 .await
1515 .unwrap();
1516
1517 assert_last_updated_seq_column(&batches, &[Some(5), Some(9), Some(8)]);
1519 let pos_col = batches[0]
1521 .column_by_name(RESERVED_COL_NAME_POS)
1522 .expect("_pos column should be present")
1523 .as_primitive::<arrow_array::types::Int64Type>();
1524 assert_eq!(pos_col.values(), &[0, 1, 2]);
1525 }
1526
1527 #[tokio::test]
1528 async fn test_last_updated_sequence_number_physical_column_nulled_without_first_row_id() {
1529 let tmp_dir = TempDir::new().unwrap();
1530 let dir = tmp_dir.path().to_str().unwrap();
1531 let seq_col = Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef;
1532 let file_path = write_plain_parquet(
1533 dir,
1534 "with_seq_no_first_row_id.parquet",
1535 vec![physical_last_updated_seq_field()],
1536 vec![seq_col],
1537 );
1538
1539 let task = last_updated_seq_task(file_path, None, Some(9));
1542
1543 let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
1544 let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
1545 let batches: Vec<RecordBatch> = reader
1546 .read(tasks)
1547 .unwrap()
1548 .stream()
1549 .try_collect()
1550 .await
1551 .unwrap();
1552
1553 assert_last_updated_seq_column(&batches, &[None, None, None]);
1554 }
1555
1556 #[tokio::test]
1557 async fn test_last_updated_sequence_number_present_by_name_without_id_unsupported() {
1558 let tmp_dir = TempDir::new().unwrap();
1559 let dir = tmp_dir.path().to_str().unwrap();
1560 let seq_field = Field::new(
1564 crate::metadata_columns::RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER,
1565 DataType::Int64,
1566 true,
1567 );
1568 let seq_col = Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef;
1569 let file_path =
1570 write_plain_parquet(dir, "with_seq_by_name.parquet", vec![seq_field], vec![
1571 seq_col,
1572 ]);
1573
1574 let task = last_updated_seq_task(file_path, Some(100), Some(9));
1575
1576 let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
1577 let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
1578 let result: Result<Vec<RecordBatch>, _> =
1579 reader.read(tasks).unwrap().stream().try_collect().await;
1580
1581 let err = result.unwrap_err();
1582 assert_eq!(err.kind(), crate::ErrorKind::FeatureUnsupported);
1583 assert!(
1584 format!("{err}").contains("without an embedded field id"),
1585 "unexpected error: {err}"
1586 );
1587 }
1588
1589 #[tokio::test]
1590 async fn test_last_updated_sequence_number_physical_column_first_row_id_without_data_seq() {
1591 let tmp_dir = TempDir::new().unwrap();
1592 let dir = tmp_dir.path().to_str().unwrap();
1593 let seq_col = Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef;
1594 let file_path = write_plain_parquet(
1595 dir,
1596 "with_seq_no_data_seq.parquet",
1597 vec![physical_last_updated_seq_field()],
1598 vec![seq_col],
1599 );
1600
1601 let task = last_updated_seq_task(file_path, Some(100), None);
1605
1606 let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
1607 let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
1608 let result: Result<Vec<RecordBatch>, _> =
1609 reader.read(tasks).unwrap().stream().try_collect().await;
1610
1611 let err = result.unwrap_err();
1612 assert_eq!(err.kind(), crate::ErrorKind::DataInvalid);
1613 assert!(
1614 format!("{err}").contains("no data sequence number"),
1615 "unexpected error: {err}"
1616 );
1617 }
1618
1619 fn row_id_task(file_path: String, first_row_id: Option<i64>) -> FileScanTask {
1621 row_id_task_with_options(file_path, first_row_id, 0, 0, vec![])
1622 }
1623
1624 fn row_id_task_with_options(
1625 file_path: String,
1626 first_row_id: Option<i64>,
1627 start: u64,
1628 length: u64,
1629 deletes: Vec<FileScanTaskDeleteFile>,
1630 ) -> FileScanTask {
1631 let schema = Arc::new(
1632 Schema::builder()
1633 .with_schema_id(1)
1634 .with_fields(vec![
1635 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1636 ])
1637 .build()
1638 .unwrap(),
1639 );
1640
1641 FileScanTask::builder()
1642 .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
1643 .with_start(start)
1644 .with_length(length)
1645 .with_data_file_path(file_path)
1646 .with_data_file_format(DataFileFormat::Parquet)
1647 .with_schema(schema)
1648 .with_project_field_ids(vec![1, RESERVED_FIELD_ID_ROW_ID])
1649 .with_first_row_id(first_row_id)
1650 .with_deletes(deletes)
1651 .with_case_sensitive(false)
1652 .build()
1653 .unwrap()
1654 }
1655
1656 fn assert_row_id_column(batches: &[RecordBatch], expected: &[Option<i64>]) {
1659 use arrow_array::cast::AsArray;
1660 use arrow_cast::cast;
1661 use arrow_schema::DataType;
1662
1663 let mut actual = Vec::new();
1664 for batch in batches {
1665 let col = batch
1666 .column_by_name(RESERVED_COL_NAME_ROW_ID)
1667 .expect("_row_id column should be present");
1668 let logical = cast(col, &DataType::Int64).unwrap();
1669 let values = logical.as_primitive::<arrow_array::types::Int64Type>();
1670 for i in 0..values.len() {
1671 actual.push((!values.is_null(i)).then(|| values.value(i)));
1672 }
1673 }
1674 assert_eq!(actual, expected);
1675 }
1676
1677 fn physical_row_id_field() -> Field {
1679 Field::new(RESERVED_COL_NAME_ROW_ID, DataType::Int64, true).with_metadata(HashMap::from([
1680 (
1681 PARQUET_FIELD_ID_META_KEY.to_string(),
1682 RESERVED_FIELD_ID_ROW_ID.to_string(),
1683 ),
1684 ]))
1685 }
1686
1687 #[tokio::test]
1688 async fn test_row_id_synthesized_from_first_row_id_and_pos() {
1689 let tmp_dir = TempDir::new().unwrap();
1690 let dir = tmp_dir.path().to_str().unwrap();
1691 let file_path = write_plain_parquet(dir, "row_id_synth.parquet", vec![], vec![]);
1692
1693 let task = row_id_task(file_path, Some(100));
1695
1696 let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
1697 let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
1698 let batches: Vec<RecordBatch> = reader
1699 .read(tasks)
1700 .unwrap()
1701 .stream()
1702 .try_collect()
1703 .await
1704 .unwrap();
1705
1706 assert_row_id_column(&batches, &[Some(100), Some(101), Some(102)]);
1707 }
1708
1709 #[tokio::test]
1710 async fn test_row_id_physical_column_coalesced() {
1711 let tmp_dir = TempDir::new().unwrap();
1712 let dir = tmp_dir.path().to_str().unwrap();
1713 let id_col = Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef;
1716 let file_path = write_plain_parquet(
1717 dir,
1718 "row_id_phys.parquet",
1719 vec![physical_row_id_field()],
1720 vec![id_col],
1721 );
1722
1723 let task = row_id_task(file_path, Some(100));
1724
1725 let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
1726 let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
1727 let batches: Vec<RecordBatch> = reader
1728 .read(tasks)
1729 .unwrap()
1730 .stream()
1731 .try_collect()
1732 .await
1733 .unwrap();
1734
1735 assert_row_id_column(&batches, &[Some(5), Some(101), Some(8)]);
1737 }
1738
1739 #[tokio::test]
1740 async fn test_row_id_only_synthesis_reads_no_data_columns() {
1741 let tmp_dir = TempDir::new().unwrap();
1746 let dir = tmp_dir.path().to_str().unwrap();
1747
1748 let meta_only = metadata_projection_task_with_first_row_id(
1749 write_parquet_with_wide_column(dir, "row_id_only.parquet", vec![], vec![]),
1750 id_and_wide_schema(),
1751 vec![RESERVED_FIELD_ID_ROW_ID],
1752 Some(100),
1753 );
1754 let (batches, meta_only_bytes) = scan_task(meta_only).await;
1755
1756 assert_eq!(batches[0].num_columns(), 1);
1757 assert_row_id_column(&batches, &[Some(100), Some(101), Some(102)]);
1758
1759 let with_data = metadata_projection_task_with_first_row_id(
1761 write_parquet_with_wide_column(dir, "row_id_only_ref.parquet", vec![], vec![]),
1762 id_and_wide_schema(),
1763 vec![2, RESERVED_FIELD_ID_ROW_ID],
1764 Some(100),
1765 );
1766 let (_, with_data_bytes) = scan_task(with_data).await;
1767
1768 assert!(
1769 meta_only_bytes < with_data_bytes,
1770 "_row_id-only synthesis should read fewer bytes than a scan of the wide column: \
1771 {meta_only_bytes} vs {with_data_bytes}"
1772 );
1773 }
1774
1775 #[tokio::test]
1776 async fn test_row_id_only_null_first_row_id_reads_no_data_columns() {
1777 let tmp_dir = TempDir::new().unwrap();
1782 let dir = tmp_dir.path().to_str().unwrap();
1783
1784 let meta_only = metadata_projection_task_with_first_row_id(
1785 write_parquet_with_wide_column(dir, "row_id_null.parquet", vec![], vec![]),
1786 id_and_wide_schema(),
1787 vec![RESERVED_FIELD_ID_ROW_ID],
1788 None,
1789 );
1790 let (batches, meta_only_bytes) = scan_task(meta_only).await;
1791
1792 assert_eq!(batches[0].num_columns(), 1);
1793 assert_row_id_column(&batches, &[None, None, None]);
1794
1795 let with_data = metadata_projection_task_with_first_row_id(
1796 write_parquet_with_wide_column(dir, "row_id_null_ref.parquet", vec![], vec![]),
1797 id_and_wide_schema(),
1798 vec![2, RESERVED_FIELD_ID_ROW_ID],
1799 None,
1800 );
1801 let (_, with_data_bytes) = scan_task(with_data).await;
1802
1803 assert!(
1804 meta_only_bytes < with_data_bytes,
1805 "_row_id-only scan with a null first_row_id should read fewer bytes than a scan \
1806 of the wide column: {meta_only_bytes} vs {with_data_bytes}"
1807 );
1808 }
1809
1810 #[tokio::test]
1811 async fn test_last_updated_seq_only_reads_no_data_columns() {
1812 use crate::metadata_columns::RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER;
1813
1814 let tmp_dir = TempDir::new().unwrap();
1819 let dir = tmp_dir.path().to_str().unwrap();
1820
1821 let lusn_only_task = |file_path: String, project_field_ids: Vec<i32>| {
1822 FileScanTask::builder()
1823 .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
1824 .with_start(0)
1825 .with_length(0)
1826 .with_data_file_path(file_path)
1827 .with_data_file_format(DataFileFormat::Parquet)
1828 .with_schema(id_and_wide_schema())
1829 .with_project_field_ids(project_field_ids)
1830 .with_first_row_id(Some(42))
1831 .with_data_sequence_number(Some(7))
1832 .with_case_sensitive(false)
1833 .build()
1834 .unwrap()
1835 };
1836
1837 let meta_only = lusn_only_task(
1838 write_parquet_with_wide_column(dir, "lusn_only.parquet", vec![], vec![]),
1839 vec![RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER],
1840 );
1841 let (batches, meta_only_bytes) = scan_task(meta_only).await;
1842
1843 assert_eq!(batches[0].num_columns(), 1);
1844 assert_last_updated_seq_column(&batches, &[Some(7), Some(7), Some(7)]);
1845
1846 let with_data = lusn_only_task(
1847 write_parquet_with_wide_column(dir, "lusn_only_ref.parquet", vec![], vec![]),
1848 vec![2, RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER],
1849 );
1850 let (_, with_data_bytes) = scan_task(with_data).await;
1851
1852 assert!(
1853 meta_only_bytes < with_data_bytes,
1854 "_last_updated_sequence_number-only scan should read fewer bytes than a scan of \
1855 the wide column: {meta_only_bytes} vs {with_data_bytes}"
1856 );
1857 }
1858
1859 #[tokio::test]
1860 async fn test_partition_only_reads_no_data_columns() {
1861 use arrow_array::StructArray;
1862
1863 use crate::metadata_columns::{RESERVED_COL_NAME_PARTITION, RESERVED_FIELD_ID_PARTITION};
1864 use crate::spec::{Literal, PartitionSpec, Struct, Transform};
1865
1866 let tmp_dir = TempDir::new().unwrap();
1870 let dir = tmp_dir.path().to_str().unwrap();
1871 let schema = id_and_wide_schema();
1872 let spec = Arc::new(
1873 PartitionSpec::builder(schema.clone())
1874 .with_spec_id(7)
1875 .add_partition_field("id", "id", Transform::Identity)
1876 .unwrap()
1877 .build()
1878 .unwrap(),
1879 );
1880 let unified_type = Arc::new(spec.partition_type(&schema).unwrap());
1881 let partition_data = Struct::from_iter(vec![Some(Literal::int(42))]);
1882
1883 let partition_task = |file_path: String, project_field_ids: Vec<i32>| {
1884 FileScanTask::builder()
1885 .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
1886 .with_start(0)
1887 .with_length(0)
1888 .with_data_file_path(file_path)
1889 .with_data_file_format(DataFileFormat::Parquet)
1890 .with_schema(schema.clone())
1891 .with_project_field_ids(project_field_ids)
1892 .with_partition_spec(Some(spec.clone()))
1893 .with_partition(Some(partition_data.clone()))
1894 .with_unified_partition_type(Some(unified_type.clone()))
1895 .with_case_sensitive(false)
1896 .build()
1897 .unwrap()
1898 };
1899
1900 let meta_only = partition_task(
1901 write_parquet_with_wide_column(dir, "partition_only.parquet", vec![], vec![]),
1902 vec![RESERVED_FIELD_ID_PARTITION],
1903 );
1904 let (batches, meta_only_bytes) = scan_task(meta_only).await;
1905
1906 assert_eq!(batches[0].num_columns(), 1);
1907 let partition_col = batches[0]
1908 .column_by_name(RESERVED_COL_NAME_PARTITION)
1909 .expect("_partition column should be present")
1910 .as_any()
1911 .downcast_ref::<StructArray>()
1912 .unwrap();
1913 assert_eq!(partition_col.len(), 3);
1914 let inner = partition_col
1915 .column(0)
1916 .as_any()
1917 .downcast_ref::<Int32Array>()
1918 .unwrap();
1919 assert_eq!(inner.values(), &[42, 42, 42]);
1920
1921 let with_data = partition_task(
1922 write_parquet_with_wide_column(dir, "partition_only_ref.parquet", vec![], vec![]),
1923 vec![2, RESERVED_FIELD_ID_PARTITION],
1924 );
1925 let (_, with_data_bytes) = scan_task(with_data).await;
1926
1927 assert!(
1928 meta_only_bytes < with_data_bytes,
1929 "_partition-only scan should read fewer bytes than a scan of the wide column: \
1930 {meta_only_bytes} vs {with_data_bytes}"
1931 );
1932 }
1933
1934 #[tokio::test]
1935 async fn test_row_id_resolves_alongside_id_less_leaf() {
1936 let tmp_dir = TempDir::new().unwrap();
1941 let dir = tmp_dir.path().to_str().unwrap();
1942 let idless_field = Field::new("variant_internal", DataType::Utf8, true);
1943 let idless_col = Arc::new(StringArray::from(vec!["a", "b", "c"])) as ArrayRef;
1944 let row_id_col = Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef;
1945 let file_path = write_plain_parquet(
1946 dir,
1947 "row_id_with_idless_leaf.parquet",
1948 vec![idless_field, physical_row_id_field()],
1949 vec![idless_col, row_id_col],
1950 );
1951
1952 let task = row_id_task(file_path, Some(100));
1953
1954 let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
1955 let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
1956 let batches: Vec<RecordBatch> = reader
1957 .read(tasks)
1958 .unwrap()
1959 .stream()
1960 .try_collect()
1961 .await
1962 .unwrap();
1963
1964 assert_row_id_column(&batches, &[Some(5), Some(101), Some(8)]);
1966 }
1967
1968 #[tokio::test]
1969 async fn test_row_id_and_last_updated_seq_co_projected() {
1970 use crate::metadata_columns::RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER;
1971
1972 let tmp_dir = TempDir::new().unwrap();
1976 let dir = tmp_dir.path().to_str().unwrap();
1977 let row_id_col = Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef;
1978 let seq_col = Arc::new(Int64Array::from(vec![Some(50), None, Some(70)])) as ArrayRef;
1979 let file_path = write_plain_parquet(
1980 dir,
1981 "row_id_and_seq.parquet",
1982 vec![physical_row_id_field(), physical_last_updated_seq_field()],
1983 vec![row_id_col, seq_col],
1984 );
1985
1986 let schema = Arc::new(
1987 Schema::builder()
1988 .with_schema_id(1)
1989 .with_fields(vec![
1990 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1991 ])
1992 .build()
1993 .unwrap(),
1994 );
1995 let task = FileScanTask::builder()
1996 .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
1997 .with_start(0)
1998 .with_length(0)
1999 .with_data_file_path(file_path)
2000 .with_data_file_format(DataFileFormat::Parquet)
2001 .with_schema(schema)
2002 .with_project_field_ids(vec![
2003 1,
2004 RESERVED_FIELD_ID_ROW_ID,
2005 RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
2006 ])
2007 .with_first_row_id(Some(100))
2008 .with_data_sequence_number(Some(9))
2009 .with_case_sensitive(false)
2010 .build()
2011 .unwrap();
2012
2013 let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
2014 let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
2015 let batches: Vec<RecordBatch> = reader
2016 .read(tasks)
2017 .unwrap()
2018 .stream()
2019 .try_collect()
2020 .await
2021 .unwrap();
2022
2023 assert_row_id_column(&batches, &[Some(5), Some(101), Some(8)]);
2025 assert_last_updated_seq_column(&batches, &[Some(50), Some(9), Some(70)]);
2027 }
2028
2029 #[tokio::test]
2030 async fn test_row_id_null_when_no_first_row_id() {
2031 let tmp_dir = TempDir::new().unwrap();
2032 let dir = tmp_dir.path().to_str().unwrap();
2033 let id_col = Arc::new(Int64Array::from(vec![Some(5), Some(6), Some(7)])) as ArrayRef;
2035 let file_path = write_plain_parquet(
2036 dir,
2037 "row_id_no_first.parquet",
2038 vec![physical_row_id_field()],
2039 vec![id_col],
2040 );
2041
2042 let task = row_id_task(file_path, None);
2044
2045 let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
2046 let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
2047 let batches: Vec<RecordBatch> = reader
2048 .read(tasks)
2049 .unwrap()
2050 .stream()
2051 .try_collect()
2052 .await
2053 .unwrap();
2054
2055 assert_row_id_column(&batches, &[None, None, None]);
2056 }
2057
2058 #[tokio::test]
2059 async fn test_row_id_with_pos_column() {
2060 use crate::metadata_columns::RESERVED_COL_NAME_POS;
2061
2062 let tmp_dir = TempDir::new().unwrap();
2063 let dir = tmp_dir.path().to_str().unwrap();
2064 let id_col = Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef;
2065 let file_path = write_plain_parquet(
2066 dir,
2067 "row_id_and_pos.parquet",
2068 vec![physical_row_id_field()],
2069 vec![id_col],
2070 );
2071
2072 let schema = Arc::new(
2076 Schema::builder()
2077 .with_schema_id(1)
2078 .with_fields(vec![
2079 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
2080 ])
2081 .build()
2082 .unwrap(),
2083 );
2084 let task = FileScanTask::builder()
2085 .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
2086 .with_start(0)
2087 .with_length(0)
2088 .with_data_file_path(file_path)
2089 .with_data_file_format(DataFileFormat::Parquet)
2090 .with_schema(schema)
2091 .with_project_field_ids(vec![1, RESERVED_FIELD_ID_POS, RESERVED_FIELD_ID_ROW_ID])
2092 .with_first_row_id(Some(100))
2093 .with_case_sensitive(false)
2094 .build()
2095 .unwrap();
2096
2097 let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
2098 let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
2099 let batches: Vec<RecordBatch> = reader
2100 .read(tasks)
2101 .unwrap()
2102 .stream()
2103 .try_collect()
2104 .await
2105 .unwrap();
2106
2107 assert_row_id_column(&batches, &[Some(5), Some(101), Some(8)]);
2109 let pos_col = batches[0]
2111 .column_by_name(RESERVED_COL_NAME_POS)
2112 .expect("_pos column should be present")
2113 .as_primitive::<arrow_array::types::Int64Type>();
2114 assert_eq!(pos_col.values(), &[0, 1, 2]);
2115 }
2116
2117 #[tokio::test]
2118 async fn test_row_id_mixed_files_share_schema() {
2119 use arrow_select::concat::concat_batches;
2120
2121 let tmp_dir = TempDir::new().unwrap();
2122 let dir = tmp_dir.path().to_str().unwrap();
2123
2124 let synth = row_id_task(
2130 write_plain_parquet(dir, "row_id_synth2.parquet", vec![], vec![]),
2131 Some(42),
2132 );
2133 let nulled = row_id_task(
2134 write_plain_parquet(dir, "row_id_null2.parquet", vec![], vec![]),
2135 None,
2136 );
2137 let coalesced = row_id_task(
2138 write_plain_parquet(
2139 dir,
2140 "row_id_coalesced2.parquet",
2141 vec![physical_row_id_field()],
2142 vec![Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef],
2143 ),
2144 Some(50),
2145 );
2146
2147 let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
2148 let tasks = Box::pin(futures::stream::iter(vec![
2149 Ok(synth),
2150 Ok(nulled),
2151 Ok(coalesced),
2152 ])) as FileScanTaskStream;
2153 let batches: Vec<RecordBatch> = reader
2154 .read(tasks)
2155 .unwrap()
2156 .stream()
2157 .try_collect()
2158 .await
2159 .unwrap();
2160
2161 assert_eq!(batches.len(), 3);
2162 let schema = batches[0].schema();
2163 concat_batches(&schema, &batches)
2164 .expect("synthesis, null and coalesce files must share one column type");
2165 }
2166
2167 #[tokio::test]
2168 async fn test_row_id_present_by_name_without_id_unsupported() {
2169 let tmp_dir = TempDir::new().unwrap();
2170 let dir = tmp_dir.path().to_str().unwrap();
2171 let id_field = Field::new(RESERVED_COL_NAME_ROW_ID, DataType::Int64, true);
2174 let id_col = Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef;
2175 let file_path =
2176 write_plain_parquet(dir, "row_id_by_name.parquet", vec![id_field], vec![id_col]);
2177
2178 let task = row_id_task(file_path, Some(100));
2179
2180 let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
2181 let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
2182 let result: Result<Vec<RecordBatch>, _> =
2183 reader.read(tasks).unwrap().stream().try_collect().await;
2184
2185 let err = result.unwrap_err();
2186 assert_eq!(err.kind(), crate::ErrorKind::FeatureUnsupported);
2187 assert!(
2188 format!("{err}").contains("without an embedded field id"),
2189 "unexpected error: {err}"
2190 );
2191 }
2192
2193 #[tokio::test]
2194 async fn test_row_id_present_by_name_without_id_nulls_when_no_lineage() {
2195 let tmp_dir = TempDir::new().unwrap();
2196 let dir = tmp_dir.path().to_str().unwrap();
2197 let id_field = Field::new(RESERVED_COL_NAME_ROW_ID, DataType::Int64, true);
2202 let id_col = Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef;
2203 let file_path = write_plain_parquet(
2204 dir,
2205 "row_id_by_name_no_lineage.parquet",
2206 vec![id_field],
2207 vec![id_col],
2208 );
2209
2210 let task = row_id_task(file_path, None);
2211
2212 let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
2213 let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
2214 let batches: Vec<RecordBatch> = reader
2215 .read(tasks)
2216 .unwrap()
2217 .stream()
2218 .try_collect()
2219 .await
2220 .unwrap();
2221
2222 assert_row_id_column(&batches, &[None, None, None]);
2223 }
2224
2225 fn row_id_task_with_predicate(
2228 file_path: String,
2229 first_row_id: Option<i64>,
2230 extra_project_field_ids: Vec<i32>,
2231 predicate: crate::expr::Predicate,
2232 ) -> FileScanTask {
2233 use crate::expr::Bind;
2234
2235 let schema = Arc::new(
2236 Schema::builder()
2237 .with_schema_id(1)
2238 .with_fields(vec![
2239 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
2240 ])
2241 .build()
2242 .unwrap(),
2243 );
2244 let bound = predicate.bind(Arc::clone(&schema), false).unwrap();
2245
2246 let mut project_field_ids = vec![1];
2247 project_field_ids.extend(extra_project_field_ids);
2248 project_field_ids.push(RESERVED_FIELD_ID_ROW_ID);
2249
2250 FileScanTask::builder()
2251 .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
2252 .with_start(0)
2253 .with_length(0)
2254 .with_data_file_path(file_path)
2255 .with_data_file_format(DataFileFormat::Parquet)
2256 .with_schema(schema)
2257 .with_project_field_ids(project_field_ids)
2258 .with_predicate(Some(bound))
2259 .with_first_row_id(first_row_id)
2260 .with_case_sensitive(false)
2261 .build()
2262 .unwrap()
2263 }
2264
2265 #[tokio::test]
2266 async fn test_row_id_stable_under_row_selection() {
2267 use crate::expr::Reference;
2268 use crate::spec::Datum;
2269
2270 let tmp_dir = TempDir::new().unwrap();
2271 let dir = tmp_dir.path().to_str().unwrap();
2272 let file_path = write_plain_parquet(dir, "row_id_selection.parquet", vec![], vec![]);
2274
2275 let task = row_id_task_with_predicate(
2276 file_path,
2277 Some(100),
2278 vec![],
2279 Reference::new("id").not_equal_to(Datum::int(2)),
2280 );
2281
2282 let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current())
2284 .with_row_selection_enabled(true)
2285 .build();
2286 let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
2287 let batches: Vec<RecordBatch> = reader
2288 .read(tasks)
2289 .unwrap()
2290 .stream()
2291 .try_collect()
2292 .await
2293 .unwrap();
2294
2295 assert_row_id_column(&batches, &[Some(100), Some(102)]);
2298 }
2299
2300 #[tokio::test]
2301 async fn test_row_id_coalesce_stable_under_row_selection() {
2302 use crate::expr::Reference;
2303 use crate::spec::Datum;
2304
2305 let tmp_dir = TempDir::new().unwrap();
2306 let dir = tmp_dir.path().to_str().unwrap();
2307 let id_col = Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef;
2310 let file_path = write_plain_parquet(
2311 dir,
2312 "row_id_coalesce_selection.parquet",
2313 vec![physical_row_id_field()],
2314 vec![id_col],
2315 );
2316
2317 let task = row_id_task_with_predicate(
2318 file_path,
2319 Some(100),
2320 vec![],
2321 Reference::new("id").not_equal_to(Datum::int(2)),
2322 );
2323
2324 let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current())
2325 .with_row_selection_enabled(true)
2326 .build();
2327 let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
2328 let batches: Vec<RecordBatch> = reader
2329 .read(tasks)
2330 .unwrap()
2331 .stream()
2332 .try_collect()
2333 .await
2334 .unwrap();
2335
2336 assert_row_id_column(&batches, &[Some(5), Some(8)]);
2340 }
2341
2342 #[tokio::test]
2343 async fn test_row_id_global_across_row_groups() {
2344 let tmp_dir = TempDir::new().unwrap();
2345 let dir = tmp_dir.path().to_str().unwrap();
2346
2347 let file_path = format!("{dir}/row_id_multi_rg.parquet");
2351 let field = Field::new("id", DataType::Int32, false).with_metadata(HashMap::from([(
2352 PARQUET_FIELD_ID_META_KEY.to_string(),
2353 "1".to_string(),
2354 )]));
2355 let arrow_schema = Arc::new(ArrowSchema::new(vec![field]));
2356 let batch = RecordBatch::try_new(arrow_schema.clone(), vec![Arc::new(Int32Array::from(
2357 vec![1, 2, 3, 4, 5],
2358 ))])
2359 .unwrap();
2360 let props = WriterProperties::builder()
2361 .set_compression(Compression::SNAPPY)
2362 .set_max_row_group_row_count(Some(2))
2363 .build();
2364 let file = File::create(&file_path).unwrap();
2365 let mut writer = ArrowWriter::try_new(file, arrow_schema, Some(props)).unwrap();
2366 writer.write(&batch).unwrap();
2367 writer.close().unwrap();
2368
2369 let task = row_id_task(file_path, Some(0));
2370 let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
2371 let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
2372 let batches: Vec<RecordBatch> = reader
2373 .read(tasks)
2374 .unwrap()
2375 .stream()
2376 .try_collect()
2377 .await
2378 .unwrap();
2379
2380 assert_row_id_column(&batches, &[Some(0), Some(1), Some(2), Some(3), Some(4)]);
2381 }
2382
2383 #[tokio::test]
2384 async fn test_row_id_global_when_first_row_group_pruned() {
2385 use parquet::file::reader::{FileReader, SerializedFileReader};
2386
2387 let tmp_dir = TempDir::new().unwrap();
2388 let dir = tmp_dir.path().to_str().unwrap();
2389
2390 let file_path = format!("{dir}/row_id_prune_rg0.parquet");
2395 let field = Field::new("id", DataType::Int32, false).with_metadata(HashMap::from([(
2396 PARQUET_FIELD_ID_META_KEY.to_string(),
2397 "1".to_string(),
2398 )]));
2399 let arrow_schema = Arc::new(ArrowSchema::new(vec![field]));
2400 let batch = RecordBatch::try_new(arrow_schema.clone(), vec![Arc::new(Int32Array::from(
2401 vec![1, 2, 3, 4, 5, 6],
2402 ))])
2403 .unwrap();
2404 let props = WriterProperties::builder()
2405 .set_compression(Compression::SNAPPY)
2406 .set_max_row_group_row_count(Some(2))
2407 .build();
2408 let file = File::create(&file_path).unwrap();
2409 let mut writer = ArrowWriter::try_new(file, arrow_schema, Some(props)).unwrap();
2410 writer.write(&batch).unwrap();
2411 writer.close().unwrap();
2412
2413 let metadata = SerializedFileReader::new(File::open(&file_path).unwrap())
2416 .unwrap()
2417 .metadata()
2418 .clone();
2419 assert_eq!(metadata.num_row_groups(), 3);
2420 let start = 4 + metadata.row_group(0).compressed_size() as u64;
2421 let file_size = std::fs::metadata(&file_path).unwrap().len();
2422
2423 let task = row_id_task_with_options(file_path, Some(100), start, file_size - start, vec![]);
2424
2425 let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
2426 let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
2427 let batches: Vec<RecordBatch> = reader
2428 .read(tasks)
2429 .unwrap()
2430 .stream()
2431 .try_collect()
2432 .await
2433 .unwrap();
2434
2435 assert_row_id_column(&batches, &[Some(102), Some(103), Some(104), Some(105)]);
2438 }
2439
2440 fn write_positional_delete(
2443 dir: &str,
2444 name: &str,
2445 data_file_path: &str,
2446 positions: &[i64],
2447 ) -> String {
2448 use arrow_array::StringArray;
2449
2450 let file_path_field =
2451 Field::new("file_path", DataType::Utf8, false).with_metadata(HashMap::from([(
2452 PARQUET_FIELD_ID_META_KEY.to_string(),
2453 "2147483546".to_string(),
2454 )]));
2455 let pos_field = Field::new("pos", DataType::Int64, false).with_metadata(HashMap::from([(
2456 PARQUET_FIELD_ID_META_KEY.to_string(),
2457 "2147483545".to_string(),
2458 )]));
2459 let arrow_schema = Arc::new(ArrowSchema::new(vec![file_path_field, pos_field]));
2460 let batch = RecordBatch::try_new(arrow_schema.clone(), vec![
2461 Arc::new(StringArray::from(vec![data_file_path; positions.len()])),
2462 Arc::new(Int64Array::from(positions.to_vec())),
2463 ])
2464 .unwrap();
2465
2466 let path = format!("{dir}/{name}");
2467 let file = File::create(&path).unwrap();
2468 let props = WriterProperties::builder()
2469 .set_compression(Compression::SNAPPY)
2470 .build();
2471 let mut writer = ArrowWriter::try_new(file, arrow_schema, Some(props)).unwrap();
2472 writer.write(&batch).unwrap();
2473 writer.close().unwrap();
2474 path
2475 }
2476
2477 #[tokio::test]
2478 async fn test_row_id_survives_positional_delete() {
2479 use crate::spec::DataContentType;
2480
2481 let tmp_dir = TempDir::new().unwrap();
2482 let dir = tmp_dir.path().to_str().unwrap();
2483 let data_path = write_plain_parquet(dir, "row_id_posdel_data.parquet", vec![], vec![]);
2485 let del_path = write_positional_delete(dir, "row_id_posdel.parquet", &data_path, &[1]);
2486
2487 let delete = FileScanTaskDeleteFile::builder()
2488 .with_file_path(del_path.clone())
2489 .with_file_size_in_bytes(std::fs::metadata(&del_path).unwrap().len())
2490 .with_file_type(DataContentType::PositionDeletes)
2491 .with_file_format(DataFileFormat::Parquet)
2492 .with_partition_spec_id(0)
2493 .build()
2494 .unwrap();
2495 let task = row_id_task_with_options(data_path, Some(100), 0, 0, vec![delete]);
2496
2497 let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
2498 let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
2499 let batches: Vec<RecordBatch> = reader
2500 .read(tasks)
2501 .unwrap()
2502 .stream()
2503 .try_collect()
2504 .await
2505 .unwrap();
2506
2507 assert_row_id_column(&batches, &[Some(100), Some(102)]);
2511 }
2512
2513 #[tokio::test]
2514 async fn test_row_id_name_collision_under_positional_fallback() {
2515 let tmp_dir = TempDir::new().unwrap();
2516 let dir = tmp_dir.path().to_str().unwrap();
2517
2518 let file_path = format!("{dir}/fallback_row_id_name.parquet");
2523 let arrow_schema = Arc::new(ArrowSchema::new(vec![
2524 Field::new("id", DataType::Int32, false),
2525 Field::new(RESERVED_COL_NAME_ROW_ID, DataType::Int64, true),
2526 ]));
2527 let batch = RecordBatch::try_new(arrow_schema.clone(), vec![
2528 Arc::new(Int32Array::from(vec![1, 2, 3])),
2529 Arc::new(Int64Array::from(vec![7i64, 8, 9])),
2530 ])
2531 .unwrap();
2532 let props = WriterProperties::builder()
2533 .set_compression(Compression::SNAPPY)
2534 .build();
2535 let file = File::create(&file_path).unwrap();
2536 let mut writer = ArrowWriter::try_new(file, arrow_schema, Some(props)).unwrap();
2537 writer.write(&batch).unwrap();
2538 writer.close().unwrap();
2539
2540 let task = row_id_task(file_path, Some(100));
2541 let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
2542 let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
2543 let batches: Vec<RecordBatch> = reader
2544 .read(tasks)
2545 .unwrap()
2546 .stream()
2547 .try_collect()
2548 .await
2549 .unwrap();
2550
2551 assert_row_id_column(&batches, &[Some(100), Some(101), Some(102)]);
2553 }
2554
2555 #[tokio::test]
2556 async fn test_read_encrypted_parquet_with_wrong_key_fails() {
2557 let encryption_key = b"0123456789abcdef";
2558 let wrong_key = b"fedcba9876543210";
2559
2560 let schema = Arc::new(
2561 Schema::builder()
2562 .with_schema_id(1)
2563 .with_fields(vec![
2564 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
2565 ])
2566 .build()
2567 .unwrap(),
2568 );
2569
2570 let arrow_schema = Arc::new(ArrowSchema::new(vec![
2571 Field::new("id", DataType::Int32, false).with_metadata(HashMap::from([(
2572 PARQUET_FIELD_ID_META_KEY.to_string(),
2573 "1".to_string(),
2574 )])),
2575 ]));
2576
2577 let tmp_dir = TempDir::new().unwrap();
2578 let table_location = tmp_dir.path().to_str().unwrap().to_string();
2579 let file_io = FileIO::new_with_fs();
2580
2581 let id_data = Arc::new(Int32Array::from(vec![1, 2, 3])) as ArrayRef;
2582 let batch = RecordBatch::try_new(arrow_schema.clone(), vec![id_data]).unwrap();
2583
2584 let file_path = format!("{table_location}/encrypted_wrong_key.parquet");
2585 write_encrypted_parquet(&file_path, &batch, encryption_key, None);
2586
2587 let wrong_key_metadata = crate::encryption::StandardKeyMetadata::try_new(wrong_key)
2588 .unwrap()
2589 .encode()
2590 .unwrap();
2591
2592 let reader = ArrowReaderBuilder::new(file_io, Runtime::current()).build();
2593
2594 let task = FileScanTask::builder()
2595 .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
2596 .with_start(0)
2597 .with_length(0)
2598 .with_data_file_path(file_path)
2599 .with_data_file_format(DataFileFormat::Parquet)
2600 .with_schema(schema)
2601 .with_project_field_ids(vec![1])
2602 .with_case_sensitive(false)
2603 .with_key_metadata(Some(wrong_key_metadata))
2604 .build()
2605 .unwrap();
2606
2607 let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
2608 let result: Result<Vec<RecordBatch>, _> =
2609 reader.read(tasks).unwrap().stream().try_collect().await;
2610
2611 let err = result.unwrap_err();
2612 assert_eq!(err.kind(), crate::ErrorKind::Unexpected);
2613 let err_str = format!("{err}");
2614 assert!(
2615 err_str.contains("unable to decrypt parquet footer"),
2616 "Expected error about decryption failure, got: {err_str}"
2617 );
2618 }
2619
2620 #[tokio::test]
2623 async fn test_read_with_concurrency_one() {
2624 use arrow_array::Int32Array;
2625
2626 let schema = Arc::new(
2627 Schema::builder()
2628 .with_schema_id(1)
2629 .with_fields(vec![
2630 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
2631 NestedField::required(2, "file_num", Type::Primitive(PrimitiveType::Int))
2632 .into(),
2633 ])
2634 .build()
2635 .unwrap(),
2636 );
2637
2638 let arrow_schema = Arc::new(ArrowSchema::new(vec![
2639 Field::new("id", DataType::Int32, false).with_metadata(HashMap::from([(
2640 PARQUET_FIELD_ID_META_KEY.to_string(),
2641 "1".to_string(),
2642 )])),
2643 Field::new("file_num", DataType::Int32, false).with_metadata(HashMap::from([(
2644 PARQUET_FIELD_ID_META_KEY.to_string(),
2645 "2".to_string(),
2646 )])),
2647 ]));
2648
2649 let tmp_dir = TempDir::new().unwrap();
2650 let table_location = tmp_dir.path().to_str().unwrap().to_string();
2651 let file_io = FileIO::new_with_fs();
2652
2653 let props = WriterProperties::builder()
2655 .set_compression(Compression::SNAPPY)
2656 .build();
2657
2658 for file_num in 0..3 {
2659 let id_data = Arc::new(Int32Array::from_iter_values(
2660 file_num * 10..(file_num + 1) * 10,
2661 )) as ArrayRef;
2662 let file_num_data = Arc::new(Int32Array::from(vec![file_num; 10])) as ArrayRef;
2663
2664 let to_write =
2665 RecordBatch::try_new(arrow_schema.clone(), vec![id_data, file_num_data]).unwrap();
2666
2667 let file = File::create(format!("{table_location}/file_{file_num}.parquet")).unwrap();
2668 let mut writer =
2669 ArrowWriter::try_new(file, to_write.schema(), Some(props.clone())).unwrap();
2670 writer.write(&to_write).expect("Writing batch");
2671 writer.close().unwrap();
2672 }
2673
2674 let reader = ArrowReaderBuilder::new(file_io, Runtime::current())
2676 .with_data_file_concurrency_limit(1)
2677 .build();
2678
2679 let tasks = vec![
2681 Ok(FileScanTask::builder()
2682 .with_file_size_in_bytes(
2683 std::fs::metadata(format!("{table_location}/file_0.parquet"))
2684 .unwrap()
2685 .len(),
2686 )
2687 .with_start(0)
2688 .with_length(0)
2689 .with_data_file_path(format!("{table_location}/file_0.parquet"))
2690 .with_data_file_format(DataFileFormat::Parquet)
2691 .with_schema(schema.clone())
2692 .with_project_field_ids(vec![1, 2])
2693 .with_case_sensitive(false)
2694 .build()
2695 .unwrap()),
2696 Ok(FileScanTask::builder()
2697 .with_file_size_in_bytes(
2698 std::fs::metadata(format!("{table_location}/file_1.parquet"))
2699 .unwrap()
2700 .len(),
2701 )
2702 .with_start(0)
2703 .with_length(0)
2704 .with_data_file_path(format!("{table_location}/file_1.parquet"))
2705 .with_data_file_format(DataFileFormat::Parquet)
2706 .with_schema(schema.clone())
2707 .with_project_field_ids(vec![1, 2])
2708 .with_case_sensitive(false)
2709 .build()
2710 .unwrap()),
2711 Ok(FileScanTask::builder()
2712 .with_file_size_in_bytes(
2713 std::fs::metadata(format!("{table_location}/file_2.parquet"))
2714 .unwrap()
2715 .len(),
2716 )
2717 .with_start(0)
2718 .with_length(0)
2719 .with_data_file_path(format!("{table_location}/file_2.parquet"))
2720 .with_data_file_format(DataFileFormat::Parquet)
2721 .with_schema(schema.clone())
2722 .with_project_field_ids(vec![1, 2])
2723 .with_case_sensitive(false)
2724 .build()
2725 .unwrap()),
2726 ];
2727
2728 let tasks_stream = Box::pin(futures::stream::iter(tasks)) as FileScanTaskStream;
2729
2730 let result = reader
2731 .read(tasks_stream)
2732 .unwrap()
2733 .stream()
2734 .try_collect::<Vec<RecordBatch>>()
2735 .await
2736 .unwrap();
2737
2738 let total_rows: usize = result.iter().map(|b| b.num_rows()).sum();
2740 assert_eq!(total_rows, 30, "Should have 30 total rows");
2741
2742 let mut all_ids = Vec::new();
2744 let mut all_file_nums = Vec::new();
2745
2746 for batch in &result {
2747 let id_col = batch
2748 .column(0)
2749 .as_primitive::<arrow_array::types::Int32Type>();
2750 let file_num_col = batch
2751 .column(1)
2752 .as_primitive::<arrow_array::types::Int32Type>();
2753
2754 for i in 0..batch.num_rows() {
2755 all_ids.push(id_col.value(i));
2756 all_file_nums.push(file_num_col.value(i));
2757 }
2758 }
2759
2760 assert_eq!(all_ids.len(), 30);
2761 assert_eq!(all_file_nums.len(), 30);
2762
2763 for i in 0..10 {
2768 assert_eq!(all_file_nums[i], 0, "First 10 rows should be from file_0");
2769 assert_eq!(all_ids[i], i as i32, "IDs should be 0-9");
2770 }
2771 for i in 10..20 {
2772 assert_eq!(all_file_nums[i], 1, "Next 10 rows should be from file_1");
2773 assert_eq!(all_ids[i], i as i32, "IDs should be 10-19");
2774 }
2775 for i in 20..30 {
2776 assert_eq!(all_file_nums[i], 2, "Last 10 rows should be from file_2");
2777 assert_eq!(all_ids[i], i as i32, "IDs should be 20-29");
2778 }
2779 }
2780
2781 #[tokio::test]
2782 async fn test_read_int96_timestamps_with_field_ids() {
2783 let schema = Arc::new(
2784 Schema::builder()
2785 .with_schema_id(1)
2786 .with_fields(vec![
2787 NestedField::optional(1, "ts", Type::Primitive(PrimitiveType::Timestamp))
2788 .into(),
2789 NestedField::required(2, "id", Type::Primitive(PrimitiveType::Int)).into(),
2790 ])
2791 .build()
2792 .unwrap(),
2793 );
2794
2795 let tmp_dir = TempDir::new().unwrap();
2796 let table_location = tmp_dir.path().to_str().unwrap().to_string();
2797 let (file_path, expected_micros) =
2798 write_int96_parquet_file(&table_location, "with_ids.parquet", true);
2799
2800 assert_int96_read_matches(&file_path, schema, vec![1, 2], &expected_micros).await;
2801 }
2802
2803 #[tokio::test]
2804 async fn test_read_int96_timestamps_without_field_ids() {
2805 let schema = Arc::new(
2806 Schema::builder()
2807 .with_schema_id(1)
2808 .with_fields(vec![
2809 NestedField::optional(1, "ts", Type::Primitive(PrimitiveType::Timestamp))
2810 .into(),
2811 NestedField::required(2, "id", Type::Primitive(PrimitiveType::Int)).into(),
2812 ])
2813 .build()
2814 .unwrap(),
2815 );
2816
2817 let tmp_dir = TempDir::new().unwrap();
2818 let table_location = tmp_dir.path().to_str().unwrap().to_string();
2819 let (file_path, expected_micros) =
2820 write_int96_parquet_file(&table_location, "no_ids.parquet", false);
2821
2822 assert_int96_read_matches(&file_path, schema, vec![1, 2], &expected_micros).await;
2823 }
2824
2825 #[tokio::test]
2826 async fn test_read_int96_timestamps_with_fallback_ids_and_pos() {
2827 use arrow_array::TimestampMicrosecondArray;
2828
2829 let schema = Arc::new(
2834 Schema::builder()
2835 .with_schema_id(1)
2836 .with_fields(vec![
2837 NestedField::optional(1, "ts", Type::Primitive(PrimitiveType::Timestamp))
2838 .into(),
2839 NestedField::required(2, "id", Type::Primitive(PrimitiveType::Int)).into(),
2840 ])
2841 .build()
2842 .unwrap(),
2843 );
2844
2845 let tmp_dir = TempDir::new().unwrap();
2846 let table_location = tmp_dir.path().to_str().unwrap().to_string();
2847 let (file_path, expected_micros) =
2848 write_int96_parquet_file(&table_location, "no_ids_with_pos.parquet", false);
2849
2850 let batches =
2851 read_int96_batches(&file_path, schema, vec![1, 2, RESERVED_FIELD_ID_POS]).await;
2852
2853 assert_eq!(batches.len(), 1);
2854 let ts_col = batches[0]
2856 .column_by_name("ts")
2857 .expect("ts column should be present")
2858 .as_any()
2859 .downcast_ref::<TimestampMicrosecondArray>()
2860 .expect("Expected TimestampMicrosecondArray");
2861 for (i, expected) in expected_micros.iter().enumerate() {
2862 assert_eq!(ts_col.value(i), *expected, "Row {i}");
2863 }
2864 let pos_col = batches[0]
2866 .column_by_name(RESERVED_COL_NAME_POS)
2867 .expect("_pos column should be present")
2868 .as_primitive::<arrow_array::types::Int64Type>();
2869 assert_eq!(pos_col.values(), &[0, 1, 2]);
2870 }
2871
2872 #[tokio::test]
2873 async fn test_read_int96_timestamps_with_field_ids_and_pos() {
2874 use arrow_array::TimestampMicrosecondArray;
2875
2876 let schema = Arc::new(
2877 Schema::builder()
2878 .with_schema_id(1)
2879 .with_fields(vec![
2880 NestedField::optional(1, "ts", Type::Primitive(PrimitiveType::Timestamp))
2881 .into(),
2882 NestedField::required(2, "id", Type::Primitive(PrimitiveType::Int)).into(),
2883 ])
2884 .build()
2885 .unwrap(),
2886 );
2887
2888 let tmp_dir = TempDir::new().unwrap();
2889 let table_location = tmp_dir.path().to_str().unwrap().to_string();
2890 let (file_path, expected_micros) =
2891 write_int96_parquet_file(&table_location, "with_ids_with_pos.parquet", true);
2892
2893 let batches =
2894 read_int96_batches(&file_path, schema, vec![1, 2, RESERVED_FIELD_ID_POS]).await;
2895
2896 assert_eq!(batches.len(), 1);
2897 let ts_col = batches[0]
2898 .column_by_name("ts")
2899 .expect("ts column should be present")
2900 .as_any()
2901 .downcast_ref::<TimestampMicrosecondArray>()
2902 .expect("Expected TimestampMicrosecondArray");
2903
2904 for (i, expected) in expected_micros.iter().enumerate() {
2905 assert_eq!(ts_col.value(i), *expected, "Row {i}");
2906 }
2907
2908 let pos_col = batches[0]
2909 .column_by_name(RESERVED_COL_NAME_POS)
2910 .expect("_pos column should be present")
2911 .as_primitive::<arrow_array::types::Int64Type>();
2912
2913 assert_eq!(pos_col.values(), &[0, 1, 2]);
2914 }
2915
2916 #[tokio::test]
2917 async fn test_read_int96_timestamps_in_struct() {
2918 use arrow_array::{StructArray, TimestampMicrosecondArray};
2919 use parquet::basic::{Repetition, Type as PhysicalType};
2920 use parquet::data_type::Int96Type;
2921 use parquet::file::writer::SerializedFileWriter;
2922 use parquet::schema::types::Type as SchemaType;
2923
2924 let tmp_dir = TempDir::new().unwrap();
2925 let table_location = tmp_dir.path().to_str().unwrap().to_string();
2926 let file_path = format!("{table_location}/struct_int96.parquet");
2927
2928 let ts_type = SchemaType::primitive_type_builder("ts", PhysicalType::INT96)
2929 .with_repetition(Repetition::OPTIONAL)
2930 .with_id(Some(2))
2931 .build()
2932 .unwrap();
2933
2934 let struct_type = SchemaType::group_type_builder("data")
2935 .with_repetition(Repetition::REQUIRED)
2936 .with_id(Some(1))
2937 .with_fields(vec![Arc::new(ts_type)])
2938 .build()
2939 .unwrap();
2940
2941 let parquet_schema = SchemaType::group_type_builder("schema")
2942 .with_fields(vec![Arc::new(struct_type)])
2943 .build()
2944 .unwrap();
2945
2946 let (int96_val, expected_micros) = make_int96_test_value();
2947
2948 let file = File::create(&file_path).unwrap();
2949 let mut writer =
2950 SerializedFileWriter::new(file, Arc::new(parquet_schema), Default::default()).unwrap();
2951
2952 let mut row_group = writer.next_row_group().unwrap();
2955 {
2956 let mut col = row_group.next_column().unwrap().unwrap();
2957 col.typed::<Int96Type>()
2958 .write_batch(&[int96_val], Some(&[1]), None)
2959 .unwrap();
2960 col.close().unwrap();
2961 }
2962 row_group.close().unwrap();
2963 writer.close().unwrap();
2964
2965 let iceberg_schema = Arc::new(
2966 Schema::builder()
2967 .with_schema_id(1)
2968 .with_fields(vec![
2969 NestedField::required(
2970 1,
2971 "data",
2972 Type::Struct(crate::spec::StructType::new(vec![
2973 NestedField::optional(
2974 2,
2975 "ts",
2976 Type::Primitive(PrimitiveType::Timestamp),
2977 )
2978 .into(),
2979 ])),
2980 )
2981 .into(),
2982 ])
2983 .build()
2984 .unwrap(),
2985 );
2986
2987 let batches = read_int96_batches(&file_path, iceberg_schema, vec![1]).await;
2988
2989 assert_eq!(batches.len(), 1);
2990 let struct_array = batches[0]
2991 .column(0)
2992 .as_any()
2993 .downcast_ref::<StructArray>()
2994 .expect("Expected StructArray");
2995 let ts_array = struct_array
2996 .column(0)
2997 .as_any()
2998 .downcast_ref::<TimestampMicrosecondArray>()
2999 .expect("Expected TimestampMicrosecondArray inside struct");
3000
3001 assert_eq!(
3002 ts_array.value(0),
3003 expected_micros,
3004 "INT96 in struct: got {}, expected {expected_micros}",
3005 ts_array.value(0)
3006 );
3007 }
3008
3009 #[tokio::test]
3010 async fn test_read_int96_timestamps_in_list() {
3011 use arrow_array::{ListArray, TimestampMicrosecondArray};
3012 use parquet::basic::{Repetition, Type as PhysicalType};
3013 use parquet::data_type::Int96Type;
3014 use parquet::file::writer::SerializedFileWriter;
3015 use parquet::schema::types::Type as SchemaType;
3016
3017 let tmp_dir = TempDir::new().unwrap();
3018 let table_location = tmp_dir.path().to_str().unwrap().to_string();
3019 let file_path = format!("{table_location}/list_int96.parquet");
3020
3021 let element_type = SchemaType::primitive_type_builder("element", PhysicalType::INT96)
3028 .with_repetition(Repetition::OPTIONAL)
3029 .with_id(Some(2))
3030 .build()
3031 .unwrap();
3032
3033 let list_group = SchemaType::group_type_builder("list")
3034 .with_repetition(Repetition::REPEATED)
3035 .with_fields(vec![Arc::new(element_type)])
3036 .build()
3037 .unwrap();
3038
3039 let list_type = SchemaType::group_type_builder("timestamps")
3040 .with_repetition(Repetition::OPTIONAL)
3041 .with_id(Some(1))
3042 .with_logical_type(Some(parquet::basic::LogicalType::List))
3043 .with_fields(vec![Arc::new(list_group)])
3044 .build()
3045 .unwrap();
3046
3047 let parquet_schema = SchemaType::group_type_builder("schema")
3048 .with_fields(vec![Arc::new(list_type)])
3049 .build()
3050 .unwrap();
3051
3052 let (int96_val, expected_micros) = make_int96_test_value();
3053
3054 let file = File::create(&file_path).unwrap();
3055 let mut writer =
3056 SerializedFileWriter::new(file, Arc::new(parquet_schema), Default::default()).unwrap();
3057
3058 let mut row_group = writer.next_row_group().unwrap();
3062 {
3063 let mut col = row_group.next_column().unwrap().unwrap();
3064 col.typed::<Int96Type>()
3065 .write_batch(&[int96_val], Some(&[3]), Some(&[0]))
3066 .unwrap();
3067 col.close().unwrap();
3068 }
3069 row_group.close().unwrap();
3070 writer.close().unwrap();
3071
3072 let iceberg_schema = Arc::new(
3073 Schema::builder()
3074 .with_schema_id(1)
3075 .with_fields(vec![
3076 NestedField::optional(
3077 1,
3078 "timestamps",
3079 Type::List(crate::spec::ListType {
3080 element_field: NestedField::optional(
3081 2,
3082 "element",
3083 Type::Primitive(PrimitiveType::Timestamp),
3084 )
3085 .into(),
3086 }),
3087 )
3088 .into(),
3089 ])
3090 .build()
3091 .unwrap(),
3092 );
3093
3094 let batches = read_int96_batches(&file_path, iceberg_schema, vec![1]).await;
3095
3096 assert_eq!(batches.len(), 1);
3097 let list_array = batches[0]
3098 .column(0)
3099 .as_any()
3100 .downcast_ref::<ListArray>()
3101 .expect("Expected ListArray");
3102 let ts_array = list_array
3103 .values()
3104 .as_any()
3105 .downcast_ref::<TimestampMicrosecondArray>()
3106 .expect("Expected TimestampMicrosecondArray inside list");
3107
3108 assert_eq!(
3109 ts_array.value(0),
3110 expected_micros,
3111 "INT96 in list: got {}, expected {expected_micros}",
3112 ts_array.value(0)
3113 );
3114 }
3115
3116 #[tokio::test]
3117 async fn test_read_int96_timestamps_in_map() {
3118 use arrow_array::{MapArray, TimestampMicrosecondArray};
3119 use parquet::basic::{Repetition, Type as PhysicalType};
3120 use parquet::data_type::{ByteArrayType, Int96Type};
3121 use parquet::file::writer::SerializedFileWriter;
3122 use parquet::schema::types::Type as SchemaType;
3123
3124 let tmp_dir = TempDir::new().unwrap();
3125 let table_location = tmp_dir.path().to_str().unwrap().to_string();
3126 let file_path = format!("{table_location}/map_int96.parquet");
3127
3128 let key_type = SchemaType::primitive_type_builder("key", PhysicalType::BYTE_ARRAY)
3136 .with_repetition(Repetition::REQUIRED)
3137 .with_logical_type(Some(parquet::basic::LogicalType::String))
3138 .with_id(Some(2))
3139 .build()
3140 .unwrap();
3141
3142 let value_type = SchemaType::primitive_type_builder("value", PhysicalType::INT96)
3143 .with_repetition(Repetition::OPTIONAL)
3144 .with_id(Some(3))
3145 .build()
3146 .unwrap();
3147
3148 let key_value_group = SchemaType::group_type_builder("key_value")
3149 .with_repetition(Repetition::REPEATED)
3150 .with_fields(vec![Arc::new(key_type), Arc::new(value_type)])
3151 .build()
3152 .unwrap();
3153
3154 let map_type = SchemaType::group_type_builder("ts_map")
3155 .with_repetition(Repetition::OPTIONAL)
3156 .with_id(Some(1))
3157 .with_logical_type(Some(parquet::basic::LogicalType::Map))
3158 .with_fields(vec![Arc::new(key_value_group)])
3159 .build()
3160 .unwrap();
3161
3162 let parquet_schema = SchemaType::group_type_builder("schema")
3163 .with_fields(vec![Arc::new(map_type)])
3164 .build()
3165 .unwrap();
3166
3167 let (int96_val, expected_micros) = make_int96_test_value();
3168
3169 let file = File::create(&file_path).unwrap();
3170 let mut writer =
3171 SerializedFileWriter::new(file, Arc::new(parquet_schema), Default::default()).unwrap();
3172
3173 let mut row_group = writer.next_row_group().unwrap();
3178 {
3179 let mut col = row_group.next_column().unwrap().unwrap();
3180 col.typed::<ByteArrayType>()
3181 .write_batch(
3182 &[parquet::data_type::ByteArray::from("event_time")],
3183 Some(&[2]),
3184 Some(&[0]),
3185 )
3186 .unwrap();
3187 col.close().unwrap();
3188 }
3189 {
3190 let mut col = row_group.next_column().unwrap().unwrap();
3191 col.typed::<Int96Type>()
3192 .write_batch(&[int96_val], Some(&[3]), Some(&[0]))
3193 .unwrap();
3194 col.close().unwrap();
3195 }
3196 row_group.close().unwrap();
3197 writer.close().unwrap();
3198
3199 let iceberg_schema = Arc::new(
3200 Schema::builder()
3201 .with_schema_id(1)
3202 .with_fields(vec![
3203 NestedField::optional(
3204 1,
3205 "ts_map",
3206 Type::Map(crate::spec::MapType {
3207 key_field: NestedField::required(
3208 2,
3209 "key",
3210 Type::Primitive(PrimitiveType::String),
3211 )
3212 .into(),
3213 value_field: NestedField::optional(
3214 3,
3215 "value",
3216 Type::Primitive(PrimitiveType::Timestamp),
3217 )
3218 .into(),
3219 }),
3220 )
3221 .into(),
3222 ])
3223 .build()
3224 .unwrap(),
3225 );
3226
3227 let batches = read_int96_batches(&file_path, iceberg_schema, vec![1]).await;
3228
3229 assert_eq!(batches.len(), 1);
3230 let map_array = batches[0]
3231 .column(0)
3232 .as_any()
3233 .downcast_ref::<MapArray>()
3234 .expect("Expected MapArray");
3235 let ts_array = map_array
3236 .values()
3237 .as_any()
3238 .downcast_ref::<TimestampMicrosecondArray>()
3239 .expect("Expected TimestampMicrosecondArray as map values");
3240
3241 assert_eq!(
3242 ts_array.value(0),
3243 expected_micros,
3244 "INT96 in map: got {}, expected {expected_micros}",
3245 ts_array.value(0)
3246 );
3247 }
3248
3249 fn write_parquet_with_wide_column(
3255 dir: &str,
3256 name: &str,
3257 extra_fields: Vec<Field>,
3258 extra_columns: Vec<ArrayRef>,
3259 ) -> String {
3260 let wide_field =
3261 Field::new("wide", DataType::Utf8, false).with_metadata(HashMap::from([(
3262 PARQUET_FIELD_ID_META_KEY.to_string(),
3263 "2".to_string(),
3264 )]));
3265 let wide_values: Vec<String> = (0..3)
3268 .map(|i| {
3269 (0..2048)
3270 .map(|j| ((i * 2048 + j) % 251) as u8 as char)
3271 .collect()
3272 })
3273 .collect();
3274
3275 let mut fields = vec![wide_field];
3276 fields.extend(extra_fields);
3277 let mut columns: Vec<ArrayRef> = vec![Arc::new(StringArray::from(wide_values))];
3278 columns.extend(extra_columns);
3279 write_plain_parquet(dir, name, fields, columns)
3280 }
3281
3282 fn id_and_wide_schema() -> SchemaRef {
3285 Arc::new(
3286 Schema::builder()
3287 .with_schema_id(1)
3288 .with_fields(vec![
3289 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
3290 NestedField::required(2, "wide", Type::Primitive(PrimitiveType::String)).into(),
3291 ])
3292 .build()
3293 .unwrap(),
3294 )
3295 }
3296
3297 fn metadata_projection_task(
3299 file_path: String,
3300 schema: SchemaRef,
3301 project_field_ids: Vec<i32>,
3302 ) -> FileScanTask {
3303 metadata_projection_task_with_first_row_id(file_path, schema, project_field_ids, None)
3304 }
3305
3306 fn metadata_projection_task_with_first_row_id(
3307 file_path: String,
3308 schema: SchemaRef,
3309 project_field_ids: Vec<i32>,
3310 first_row_id: Option<i64>,
3311 ) -> FileScanTask {
3312 FileScanTask::builder()
3313 .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
3314 .with_start(0)
3315 .with_length(0)
3316 .with_data_file_path(file_path)
3317 .with_data_file_format(DataFileFormat::Parquet)
3318 .with_schema(schema)
3319 .with_project_field_ids(project_field_ids)
3320 .with_first_row_id(first_row_id)
3321 .with_case_sensitive(false)
3322 .build()
3323 .unwrap()
3324 }
3325
3326 async fn scan_task(task: FileScanTask) -> (Vec<RecordBatch>, u64) {
3328 let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
3329 let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
3330 let scan = reader.read(tasks).unwrap();
3331 let metrics = scan.metrics().clone();
3332 let batches = scan.stream().try_collect().await.unwrap();
3333 (batches, metrics.bytes_read())
3334 }
3335
3336 #[tokio::test]
3337 async fn test_pos_only_projection_reads_no_data_columns() {
3338 let tmp_dir = TempDir::new().unwrap();
3339 let dir = tmp_dir.path().to_str().unwrap();
3340
3341 let pos_only = metadata_projection_task(
3342 write_parquet_with_wide_column(dir, "pos_only.parquet", vec![], vec![]),
3343 id_and_wide_schema(),
3344 vec![RESERVED_FIELD_ID_POS],
3345 );
3346 let (batches, pos_only_bytes) = scan_task(pos_only).await;
3347
3348 assert_eq!(batches[0].num_columns(), 1);
3350 let pos_col = batches[0]
3351 .column_by_name(RESERVED_COL_NAME_POS)
3352 .expect("_pos column should be present")
3353 .as_primitive::<arrow_array::types::Int64Type>();
3354 assert_eq!(pos_col.values(), &[0, 1, 2]);
3355
3356 let with_data = metadata_projection_task(
3359 write_parquet_with_wide_column(dir, "pos_only_ref.parquet", vec![], vec![]),
3360 id_and_wide_schema(),
3361 vec![2, RESERVED_FIELD_ID_POS],
3362 );
3363 let (_, with_data_bytes) = scan_task(with_data).await;
3364
3365 assert!(
3366 pos_only_bytes < with_data_bytes,
3367 "_pos-only scan should read fewer bytes than a scan of the wide column: \
3368 {pos_only_bytes} vs {with_data_bytes}"
3369 );
3370 }
3371
3372 #[tokio::test]
3373 async fn test_pos_only_projection_keeps_absolute_pos_under_predicate() {
3374 use crate::expr::{Bind, Reference};
3375 use crate::spec::Datum;
3376
3377 let tmp_dir = TempDir::new().unwrap();
3378 let dir = tmp_dir.path().to_str().unwrap();
3379 let file_path = write_plain_parquet(dir, "pos_only_predicate.parquet", vec![], vec![]);
3381
3382 let schema = Arc::new(
3383 Schema::builder()
3384 .with_schema_id(1)
3385 .with_fields(vec![
3386 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
3387 ])
3388 .build()
3389 .unwrap(),
3390 );
3391 let bound = Reference::new("id")
3392 .not_equal_to(Datum::int(2))
3393 .bind(Arc::clone(&schema), false)
3394 .unwrap();
3395 let task = FileScanTask::builder()
3396 .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
3397 .with_start(0)
3398 .with_length(0)
3399 .with_data_file_path(file_path)
3400 .with_data_file_format(DataFileFormat::Parquet)
3401 .with_schema(schema)
3402 .with_project_field_ids(vec![RESERVED_FIELD_ID_POS])
3403 .with_predicate(Some(bound))
3404 .with_case_sensitive(false)
3405 .build()
3406 .unwrap();
3407
3408 let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current())
3412 .with_row_selection_enabled(true)
3413 .build();
3414 let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
3415 let batches: Vec<RecordBatch> = reader
3416 .read(tasks)
3417 .unwrap()
3418 .stream()
3419 .try_collect()
3420 .await
3421 .unwrap();
3422
3423 let pos: Vec<i64> = batches
3424 .iter()
3425 .flat_map(|b| {
3426 b.column_by_name(RESERVED_COL_NAME_POS)
3427 .expect("_pos column should be present")
3428 .as_primitive::<arrow_array::types::Int64Type>()
3429 .values()
3430 .to_vec()
3431 })
3432 .collect();
3433 assert_eq!(pos, vec![0, 2]);
3434 }
3435
3436 #[tokio::test]
3437 async fn test_pos_and_file_projection() {
3438 use crate::metadata_columns::RESERVED_COL_NAME_FILE;
3439
3440 let tmp_dir = TempDir::new().unwrap();
3441 let dir = tmp_dir.path().to_str().unwrap();
3442 let file_path = write_parquet_with_wide_column(dir, "pos_and_file.parquet", vec![], vec![]);
3445 let task = metadata_projection_task(file_path.clone(), id_and_wide_schema(), vec![
3446 RESERVED_FIELD_ID_POS,
3447 RESERVED_FIELD_ID_FILE,
3448 ]);
3449 let (batches, _) = scan_task(task).await;
3450
3451 assert_eq!(batches[0].num_columns(), 2);
3453 let pos_col = batches[0]
3454 .column_by_name(RESERVED_COL_NAME_POS)
3455 .expect("_pos column should be present")
3456 .as_primitive::<arrow_array::types::Int64Type>();
3457 assert_eq!(pos_col.values(), &[0, 1, 2]);
3458 let file_col = batches[0]
3459 .column_by_name(RESERVED_COL_NAME_FILE)
3460 .expect("_file column should be present");
3461 let file_col = cast(file_col, &DataType::Utf8).unwrap();
3462 let file_col = file_col.as_any().downcast_ref::<StringArray>().unwrap();
3463 assert_eq!(file_col.value(0), file_path);
3464 }
3465
3466 #[tokio::test]
3467 async fn test_pos_and_physical_seq_projection_reads_only_the_leaf() {
3468 use crate::metadata_columns::{
3469 RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER,
3470 RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
3471 };
3472
3473 let tmp_dir = TempDir::new().unwrap();
3477 let dir = tmp_dir.path().to_str().unwrap();
3478
3479 let write = |name: &str| {
3481 write_parquet_with_wide_column(
3482 dir,
3483 name,
3484 vec![physical_last_updated_seq_field()],
3485 vec![Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef],
3486 )
3487 };
3488
3489 let seq_task = |path: String, ids: Vec<i32>| {
3490 FileScanTask::builder()
3491 .with_file_size_in_bytes(std::fs::metadata(&path).unwrap().len())
3492 .with_start(0)
3493 .with_length(0)
3494 .with_data_file_path(path)
3495 .with_data_file_format(DataFileFormat::Parquet)
3496 .with_schema(id_and_wide_schema())
3497 .with_project_field_ids(ids)
3498 .with_first_row_id(Some(100))
3499 .with_data_sequence_number(Some(9))
3500 .with_case_sensitive(false)
3501 .build()
3502 .unwrap()
3503 };
3504
3505 let meta_only = seq_task(write("pos_seq.parquet"), vec![
3506 RESERVED_FIELD_ID_POS,
3507 RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
3508 ]);
3509 let (batches, meta_only_bytes) = scan_task(meta_only).await;
3510
3511 let pos_col = batches[0]
3513 .column_by_name(RESERVED_COL_NAME_POS)
3514 .expect("_pos column should be present")
3515 .as_primitive::<arrow_array::types::Int64Type>();
3516 assert_eq!(pos_col.values(), &[0, 1, 2]);
3517 let seq_col = batches[0]
3518 .column_by_name(RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER)
3519 .expect("_last_updated_sequence_number column should be present");
3520 let seq_col = cast(seq_col, &DataType::Int64).unwrap();
3521 let seq_col = seq_col.as_any().downcast_ref::<Int64Array>().unwrap();
3522 assert_eq!(seq_col.value(0), 5);
3524 assert_eq!(seq_col.value(1), 9);
3525 assert_eq!(seq_col.value(2), 8);
3526 assert!(batches[0].column_by_name("wide").is_none());
3527
3528 let with_data = seq_task(write("pos_seq_ref.parquet"), vec![
3531 2,
3532 RESERVED_FIELD_ID_POS,
3533 RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
3534 ]);
3535 let (_, with_data_bytes) = scan_task(with_data).await;
3536
3537 assert!(
3538 meta_only_bytes < with_data_bytes,
3539 "_pos + physical sequence scan should read fewer bytes than one that also \
3540 reads the wide column: {meta_only_bytes} vs {with_data_bytes}"
3541 );
3542 }
3543
3544 #[tokio::test]
3545 async fn test_seq_only_projection_reads_only_the_leaf() {
3546 use crate::metadata_columns::{
3547 RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER,
3548 RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
3549 };
3550
3551 let tmp_dir = TempDir::new().unwrap();
3555 let dir = tmp_dir.path().to_str().unwrap();
3556 let write = |name: &str| {
3557 write_parquet_with_wide_column(
3558 dir,
3559 name,
3560 vec![physical_last_updated_seq_field()],
3561 vec![Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef],
3562 )
3563 };
3564 let seq_task = |path: String, ids: Vec<i32>| {
3565 FileScanTask::builder()
3566 .with_file_size_in_bytes(std::fs::metadata(&path).unwrap().len())
3567 .with_start(0)
3568 .with_length(0)
3569 .with_data_file_path(path)
3570 .with_data_file_format(DataFileFormat::Parquet)
3571 .with_schema(id_and_wide_schema())
3572 .with_project_field_ids(ids)
3573 .with_first_row_id(Some(100))
3574 .with_data_sequence_number(Some(9))
3575 .with_case_sensitive(false)
3576 .build()
3577 .unwrap()
3578 };
3579
3580 let meta_only = seq_task(write("seq_only.parquet"), vec![
3581 RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
3582 ]);
3583 let (batches, meta_only_bytes) = scan_task(meta_only).await;
3584
3585 let seq_col = batches[0]
3586 .column_by_name(RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER)
3587 .expect("_last_updated_sequence_number column should be present");
3588 let seq_col = cast(seq_col, &DataType::Int64).unwrap();
3589 let seq_col = seq_col.as_any().downcast_ref::<Int64Array>().unwrap();
3590 assert_eq!(seq_col.value(0), 5);
3591 assert_eq!(seq_col.value(1), 9);
3592 assert_eq!(seq_col.value(2), 8);
3593 let total_rows: usize = batches.iter().map(|b| b.num_rows()).sum();
3594 assert_eq!(total_rows, 3);
3595 assert!(batches[0].column_by_name("wide").is_none());
3596
3597 let with_data = seq_task(write("seq_only_ref.parquet"), vec![
3598 2,
3599 RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
3600 ]);
3601 let (_, with_data_bytes) = scan_task(with_data).await;
3602
3603 assert!(
3604 meta_only_bytes < with_data_bytes,
3605 "seq-only scan should read fewer bytes than one that also reads the wide \
3606 column: {meta_only_bytes} vs {with_data_bytes}"
3607 );
3608 }
3609
3610 #[tokio::test]
3611 async fn test_seq_only_null_first_row_id_reads_no_data_columns() {
3612 use crate::metadata_columns::{
3613 RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER,
3614 RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
3615 };
3616
3617 let tmp_dir = TempDir::new().unwrap();
3622 let dir = tmp_dir.path().to_str().unwrap();
3623 let seq_task = |name: &str, ids: Vec<i32>| {
3624 let file_path = write_parquet_with_wide_column(
3625 dir,
3626 name,
3627 vec![physical_last_updated_seq_field()],
3628 vec![Arc::new(Int64Array::from(vec![Some(5), Some(6), Some(7)])) as ArrayRef],
3629 );
3630 FileScanTask::builder()
3631 .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
3632 .with_start(0)
3633 .with_length(0)
3634 .with_data_file_path(file_path)
3635 .with_data_file_format(DataFileFormat::Parquet)
3636 .with_schema(id_and_wide_schema())
3637 .with_project_field_ids(ids)
3638 .with_first_row_id(None)
3639 .with_data_sequence_number(Some(9))
3640 .with_case_sensitive(false)
3641 .build()
3642 .unwrap()
3643 };
3644
3645 let meta_only = seq_task("seq_only_null_first.parquet", vec![
3646 RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
3647 ]);
3648 let (batches, meta_only_bytes) = scan_task(meta_only).await;
3649
3650 assert_eq!(batches[0].num_columns(), 1);
3651 let total_rows: usize = batches.iter().map(|b| b.num_rows()).sum();
3652 assert_eq!(total_rows, 3);
3653 let seq_col = batches[0]
3654 .column_by_name(RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER)
3655 .expect("_last_updated_sequence_number column should be present");
3656 let seq_col = cast(seq_col, &DataType::Int64).unwrap();
3657 let seq_col = seq_col.as_any().downcast_ref::<Int64Array>().unwrap();
3658 assert!((0..3).all(|i| seq_col.is_null(i)));
3659
3660 let with_data = seq_task("seq_only_null_first_ref.parquet", vec![
3661 2,
3662 RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
3663 ]);
3664 let (_, with_data_bytes) = scan_task(with_data).await;
3665
3666 assert!(
3667 meta_only_bytes < with_data_bytes,
3668 "seq-only scan with a null first_row_id should read fewer bytes than a scan of \
3669 the wide column: {meta_only_bytes} vs {with_data_bytes}"
3670 );
3671 }
3672
3673 #[tokio::test]
3674 async fn test_file_only_reads_no_data_columns() {
3675 use crate::metadata_columns::RESERVED_COL_NAME_FILE;
3676
3677 let tmp_dir = TempDir::new().unwrap();
3681 let dir = tmp_dir.path().to_str().unwrap();
3682
3683 let file_path = write_parquet_with_wide_column(dir, "file_only.parquet", vec![], vec![]);
3684 let meta_only = metadata_projection_task(file_path.clone(), id_and_wide_schema(), vec![
3685 RESERVED_FIELD_ID_FILE,
3686 ]);
3687 let (batches, meta_only_bytes) = scan_task(meta_only).await;
3688
3689 assert_eq!(batches[0].num_columns(), 1);
3690 let total_rows: usize = batches.iter().map(|b| b.num_rows()).sum();
3691 assert_eq!(total_rows, 3);
3692 let file_col = batches[0]
3693 .column_by_name(RESERVED_COL_NAME_FILE)
3694 .expect("_file column should be present");
3695 let file_col = cast(file_col, &DataType::Utf8).unwrap();
3696 let file_col = file_col.as_any().downcast_ref::<StringArray>().unwrap();
3697 assert_eq!(file_col.value(0), file_path);
3698
3699 let with_data = metadata_projection_task(
3700 write_parquet_with_wide_column(dir, "file_only_ref.parquet", vec![], vec![]),
3701 id_and_wide_schema(),
3702 vec![2, RESERVED_FIELD_ID_FILE],
3703 );
3704 let (_, with_data_bytes) = scan_task(with_data).await;
3705
3706 assert!(
3707 meta_only_bytes < with_data_bytes,
3708 "_file-only scan should read fewer bytes than a scan of the wide column: \
3709 {meta_only_bytes} vs {with_data_bytes}"
3710 );
3711 }
3712
3713 #[tokio::test]
3714 async fn test_spec_id_only_reads_no_data_columns() {
3715 use crate::metadata_columns::{RESERVED_COL_NAME_SPEC_ID, RESERVED_FIELD_ID_SPEC_ID};
3716 use crate::spec::{Literal, PartitionSpec, Struct, Transform};
3717
3718 let tmp_dir = TempDir::new().unwrap();
3722 let dir = tmp_dir.path().to_str().unwrap();
3723 let schema = id_and_wide_schema();
3724 let spec = Arc::new(
3725 PartitionSpec::builder(schema.clone())
3726 .with_spec_id(7)
3727 .add_partition_field("id", "id", Transform::Identity)
3728 .unwrap()
3729 .build()
3730 .unwrap(),
3731 );
3732
3733 let spec_id_task = |file_path: String, project_field_ids: Vec<i32>| {
3734 FileScanTask::builder()
3735 .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
3736 .with_start(0)
3737 .with_length(0)
3738 .with_data_file_path(file_path)
3739 .with_data_file_format(DataFileFormat::Parquet)
3740 .with_schema(schema.clone())
3741 .with_project_field_ids(project_field_ids)
3742 .with_partition(Some(Struct::from_iter([Some(Literal::int(42))])))
3743 .with_partition_spec(Some(spec.clone()))
3744 .with_case_sensitive(false)
3745 .build()
3746 .unwrap()
3747 };
3748
3749 let meta_only = spec_id_task(
3750 write_parquet_with_wide_column(dir, "spec_id_only.parquet", vec![], vec![]),
3751 vec![RESERVED_FIELD_ID_SPEC_ID],
3752 );
3753 let (batches, meta_only_bytes) = scan_task(meta_only).await;
3754
3755 assert_eq!(batches[0].num_columns(), 1);
3756 let total_rows: usize = batches.iter().map(|b| b.num_rows()).sum();
3757 assert_eq!(total_rows, 3);
3758 let spec_id_col = batches[0]
3759 .column_by_name(RESERVED_COL_NAME_SPEC_ID)
3760 .expect("_spec_id column should be present");
3761 let spec_id_col = cast(spec_id_col, &DataType::Int32).unwrap();
3762 let spec_id_col = spec_id_col.as_any().downcast_ref::<Int32Array>().unwrap();
3763 assert_eq!(spec_id_col.values(), &[7, 7, 7]);
3764
3765 let with_data = spec_id_task(
3766 write_parquet_with_wide_column(dir, "spec_id_only_ref.parquet", vec![], vec![]),
3767 vec![2, RESERVED_FIELD_ID_SPEC_ID],
3768 );
3769 let (_, with_data_bytes) = scan_task(with_data).await;
3770
3771 assert!(
3772 meta_only_bytes < with_data_bytes,
3773 "_spec_id-only scan should read fewer bytes than a scan of the wide column: \
3774 {meta_only_bytes} vs {with_data_bytes}"
3775 );
3776 }
3777
3778 #[tokio::test]
3779 async fn test_empty_projection_preserves_row_count() {
3780 let tmp_dir = TempDir::new().unwrap();
3781 let dir = tmp_dir.path().to_str().unwrap();
3782 let file_path = write_plain_parquet(dir, "empty_projection.parquet", vec![], vec![]);
3783 let schema = Arc::new(
3784 Schema::builder()
3785 .with_schema_id(1)
3786 .with_fields(vec![
3787 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
3788 ])
3789 .build()
3790 .unwrap(),
3791 );
3792 let task = metadata_projection_task(file_path, schema, vec![]);
3793 let (batches, _) = scan_task(task).await;
3794
3795 let total_rows: usize = batches.iter().map(|b| b.num_rows()).sum();
3797 assert_eq!(total_rows, 3);
3798 }
3799}