1mod cache;
21use cache::*;
22mod context;
23use context::*;
24mod task;
25
26use std::collections::HashMap;
27use std::future::Future;
28use std::sync::Arc;
29
30use arrow_array::RecordBatch;
31use futures::channel::mpsc::{Sender, channel};
32use futures::stream::BoxStream;
33use futures::{SinkExt, StreamExt, TryStreamExt};
34pub use task::*;
35
36use crate::arrow::ArrowReaderBuilder;
37pub use crate::arrow::{ScanMetrics, ScanResult};
38use crate::delete_file_index::DeleteFileIndex;
39use crate::error::invalid_data;
40use crate::expr::visitors::inclusive_metrics_evaluator::InclusiveMetricsEvaluator;
41use crate::expr::{Bind, BoundPredicate, Predicate};
42use crate::io::FileIO;
43use crate::metadata_columns::{
44 RESERVED_FIELD_ID_PARTITION, get_metadata_field_id, is_metadata_column_name,
45};
46use crate::runtime::Runtime;
47use crate::spec::{DataContentType, Schema, SchemaRef, SnapshotRef, SortOrderRef, StructType};
48use crate::table::Table;
49use crate::util::available_parallelism;
50use crate::{Error, ErrorKind, Result};
51
52pub type ArrowRecordBatchStream = BoxStream<'static, Result<RecordBatch>>;
54
55fn resolve_field_id(schema: &Schema, column_name: &str, case_sensitive: bool) -> Option<i32> {
57 if case_sensitive {
58 schema.field_id_by_name(column_name)
59 } else {
60 schema
61 .field_by_name_case_insensitive(column_name)
62 .map(|field| field.id)
63 }
64}
65
66fn collect_scan_field_ids(
67 schema: &Schema,
68 column_names: Option<&[String]>,
69 case_sensitive: bool,
70) -> Result<Vec<i32>> {
71 let Some(column_names) = column_names else {
72 return Ok(schema.as_struct().fields().iter().map(|f| f.id).collect());
73 };
74
75 column_names
76 .iter()
77 .map(|column_name| {
78 if is_metadata_column_name(column_name) {
79 return get_metadata_field_id(column_name);
80 }
81
82 let field_id = resolve_field_id(schema, column_name, case_sensitive).ok_or_else(|| {
83 invalid_data!("Column {column_name} not found in table. Schema: {schema}")
84 })?;
85
86 schema
87 .as_struct()
88 .field_by_id(field_id)
89 .ok_or_else(|| {
90 Error::new(
91 ErrorKind::FeatureUnsupported,
92 format!(
93 "Column {column_name} is not a direct child of schema but a nested field, which is not supported now. Schema: {schema}"
94 ),
95 )
96 })?;
97
98 Ok(field_id)
99 })
100 .collect()
101}
102
103fn bind_scan_predicate(
104 schema: &SchemaRef,
105 predicate: Option<&Predicate>,
106 case_sensitive: bool,
107) -> Result<Option<Arc<BoundPredicate>>> {
108 predicate
109 .map(|predicate| predicate.bind(schema.clone(), case_sensitive))
110 .transpose()
111 .map(|predicate| predicate.map(Arc::new))
112}
113
114fn projected_partition_type(
115 table: &Table,
116 schema: &Schema,
117 field_ids: &[i32],
118) -> Result<Option<Arc<StructType>>> {
119 if !field_ids.contains(&RESERVED_FIELD_ID_PARTITION) {
120 return Ok(None);
121 }
122
123 table
124 .metadata()
125 .unified_partition_type(schema)
126 .map(Arc::new)
127 .map(Some)
128}
129
130pub struct TableScanBuilder<'a> {
132 table: &'a Table,
133 column_names: Option<Vec<String>>,
135 snapshot_id: Option<i64>,
136 batch_size: Option<usize>,
137 case_sensitive: bool,
138 filter: Option<Predicate>,
139 concurrency_limit_data_files: usize,
140 concurrency_limit_manifest_entries: usize,
141 concurrency_limit_manifest_files: usize,
142 row_group_filtering_enabled: bool,
143 row_selection_enabled: bool,
144 bloom_filter_enabled: bool,
145}
146
147impl<'a> TableScanBuilder<'a> {
148 pub(crate) fn new(table: &'a Table) -> Self {
149 let num_cpus = available_parallelism().get();
150
151 Self {
152 table,
153 column_names: None,
154 snapshot_id: None,
155 batch_size: None,
156 case_sensitive: true,
157 filter: None,
158 concurrency_limit_data_files: num_cpus,
159 concurrency_limit_manifest_entries: num_cpus,
160 concurrency_limit_manifest_files: num_cpus,
161 row_group_filtering_enabled: true,
162 row_selection_enabled: false,
163 bloom_filter_enabled: false,
164 }
165 }
166
167 pub fn with_batch_size(mut self, batch_size: Option<usize>) -> Self {
170 self.batch_size = batch_size;
171 self
172 }
173
174 pub fn with_case_sensitive(mut self, case_sensitive: bool) -> Self {
176 self.case_sensitive = case_sensitive;
177 self
178 }
179
180 pub fn with_filter(mut self, predicate: Predicate) -> Self {
182 self.filter = Some(predicate.rewrite_not());
185 self
186 }
187
188 pub fn select_all(mut self) -> Self {
190 self.column_names = None;
191 self
192 }
193
194 pub fn select_empty(mut self) -> Self {
196 self.column_names = Some(vec![]);
197 self
198 }
199
200 pub fn select(mut self, column_names: impl IntoIterator<Item = impl ToString>) -> Self {
202 self.column_names = Some(
203 column_names
204 .into_iter()
205 .map(|item| item.to_string())
206 .collect(),
207 );
208 self
209 }
210
211 pub fn snapshot_id(mut self, snapshot_id: i64) -> Self {
213 self.snapshot_id = Some(snapshot_id);
214 self
215 }
216
217 pub fn with_concurrency_limit(mut self, limit: usize) -> Self {
220 self.concurrency_limit_manifest_files = limit;
221 self.concurrency_limit_manifest_entries = limit;
222 self.concurrency_limit_data_files = limit;
223 self
224 }
225
226 pub fn with_data_file_concurrency_limit(mut self, limit: usize) -> Self {
228 self.concurrency_limit_data_files = limit;
229 self
230 }
231
232 pub fn with_manifest_entry_concurrency_limit(mut self, limit: usize) -> Self {
234 self.concurrency_limit_manifest_entries = limit;
235 self
236 }
237
238 pub fn with_row_group_filtering_enabled(mut self, row_group_filtering_enabled: bool) -> Self {
247 self.row_group_filtering_enabled = row_group_filtering_enabled;
248 self
249 }
250
251 pub fn with_row_selection_enabled(mut self, row_selection_enabled: bool) -> Self {
266 self.row_selection_enabled = row_selection_enabled;
267 self
268 }
269
270 pub fn with_bloom_filter_enabled(mut self, bloom_filter_enabled: bool) -> Self {
280 self.bloom_filter_enabled = bloom_filter_enabled;
281 self
282 }
283
284 pub fn build(self) -> Result<TableScan> {
286 let snapshot = match self.snapshot_id {
287 Some(snapshot_id) => self
288 .table
289 .metadata()
290 .snapshot_by_id(snapshot_id)
291 .ok_or_else(|| invalid_data!("Snapshot with id {snapshot_id} not found"))?
292 .clone(),
293 None => {
294 let Some(current_snapshot_id) = self.table.metadata().current_snapshot() else {
295 return Ok(TableScan {
296 batch_size: self.batch_size,
297 column_names: self.column_names,
298 file_io: self.table.file_io().clone(),
299 plan_context: None,
300 concurrency_limit_data_files: self.concurrency_limit_data_files,
301 concurrency_limit_manifest_entries: self.concurrency_limit_manifest_entries,
302 concurrency_limit_manifest_files: self.concurrency_limit_manifest_files,
303 row_group_filtering_enabled: self.row_group_filtering_enabled,
304 row_selection_enabled: self.row_selection_enabled,
305 bloom_filter_enabled: self.bloom_filter_enabled,
306 runtime: self.table.runtime().clone(),
307 });
308 };
309 current_snapshot_id.clone()
310 }
311 };
312
313 let schema = snapshot.schema(self.table.metadata())?;
314 let field_ids =
315 collect_scan_field_ids(&schema, self.column_names.as_deref(), self.case_sensitive)?;
316 let snapshot_bound_predicate =
317 bind_scan_predicate(&schema, self.filter.as_ref(), self.case_sensitive)?;
318 let name_mapping = self
319 .table
320 .metadata()
321 .table_properties()
322 .default_name_mapping()?
323 .map(Arc::new);
324 let unified_partition_type = projected_partition_type(self.table, &schema, &field_ids)?;
325
326 let sort_orders = Arc::new(
329 self.table
330 .metadata()
331 .sort_orders_iter()
332 .map(|order| (order.order_id, order.clone()))
333 .collect::<HashMap<i64, SortOrderRef>>(),
334 );
335
336 let plan_context = PlanContext {
337 snapshot,
338 table_metadata: self.table.metadata_ref(),
339 snapshot_schema: schema,
340 case_sensitive: self.case_sensitive,
341 predicate: self.filter.map(Arc::new),
342 snapshot_bound_predicate,
343 object_cache: self.table.object_cache(),
344 field_ids: Arc::new(field_ids),
345 name_mapping,
346 partition_filter_cache: Arc::new(PartitionFilterCache::new()),
347 manifest_evaluator_cache: Arc::new(ManifestEvaluatorCache::new()),
348 expression_evaluator_cache: Arc::new(ExpressionEvaluatorCache::new()),
349 unified_partition_type,
350 sort_orders,
351 };
352
353 Ok(TableScan {
354 batch_size: self.batch_size,
355 column_names: self.column_names,
356 file_io: self.table.file_io().clone(),
357 plan_context: Some(plan_context),
358 concurrency_limit_data_files: self.concurrency_limit_data_files,
359 concurrency_limit_manifest_entries: self.concurrency_limit_manifest_entries,
360 concurrency_limit_manifest_files: self.concurrency_limit_manifest_files,
361 row_group_filtering_enabled: self.row_group_filtering_enabled,
362 row_selection_enabled: self.row_selection_enabled,
363 bloom_filter_enabled: self.bloom_filter_enabled,
364 runtime: self.table.runtime().clone(),
365 })
366 }
367}
368
369#[derive(Debug)]
371pub struct TableScan {
372 plan_context: Option<PlanContext>,
376 batch_size: Option<usize>,
377 file_io: FileIO,
378 column_names: Option<Vec<String>>,
379 concurrency_limit_manifest_files: usize,
382
383 concurrency_limit_manifest_entries: usize,
386
387 concurrency_limit_data_files: usize,
390
391 row_group_filtering_enabled: bool,
392 row_selection_enabled: bool,
393 bloom_filter_enabled: bool,
394
395 runtime: Runtime,
396}
397
398impl TableScan {
399 pub async fn plan_files(&self) -> Result<FileScanTaskStream> {
401 self.plan_data_files(|ctx| ctx.into_file_scan_task()).await
402 }
403
404 pub(crate) async fn plan_cow_rewrite_files(
405 &self,
406 ) -> Result<BoxStream<'static, Result<crate::cow_rewrite::CowRewriteFile>>> {
407 self.plan_data_files(|ctx| ctx.into_cow_rewrite_file())
408 .await
409 }
410
411 async fn plan_data_files<T, F, Fut>(
412 &self,
413 build_result: F,
414 ) -> Result<BoxStream<'static, Result<T>>>
415 where
416 T: Send + 'static,
417 F: Fn(ManifestEntryContext) -> Fut + Copy + Send + Sync + 'static,
418 Fut: Future<Output = Result<T>> + Send + 'static,
419 {
420 let Some(plan_context) = self.plan_context.as_ref() else {
421 return Ok(Box::pin(futures::stream::empty()));
422 };
423
424 let concurrency_limit_manifest_files = self.concurrency_limit_manifest_files;
425 let concurrency_limit_manifest_entries = self.concurrency_limit_manifest_entries;
426
427 let (manifest_entry_data_ctx_tx, manifest_entry_data_ctx_rx) =
429 channel(concurrency_limit_manifest_files);
430 let (manifest_entry_delete_ctx_tx, manifest_entry_delete_ctx_rx) =
431 channel(concurrency_limit_manifest_files);
432
433 let (planned_file_tx, planned_file_rx) = channel(concurrency_limit_manifest_entries);
435
436 let (delete_file_idx, delete_file_tx) = DeleteFileIndex::new(self.runtime.clone());
437
438 let manifest_list = plan_context.get_manifest_list().await?;
439
440 let manifest_file_contexts = plan_context.build_manifest_file_contexts(
444 manifest_list,
445 manifest_entry_data_ctx_tx,
446 delete_file_idx.clone(),
447 manifest_entry_delete_ctx_tx,
448 )?;
449
450 let mut channel_for_manifest_error = planned_file_tx.clone();
451 let mut channel_for_data_manifest_entry_error = planned_file_tx.clone();
452 let mut channel_for_delete_manifest_entry_error = planned_file_tx.clone();
453
454 let rt = self.runtime.clone();
455
456 rt.io().spawn(async move {
458 let result = futures::stream::iter(manifest_file_contexts)
459 .try_for_each_concurrent(concurrency_limit_manifest_files, |ctx| async move {
460 ctx.fetch_manifest_and_stream_manifest_entries().await
461 })
462 .await;
463
464 if let Err(error) = result {
465 let _ = channel_for_manifest_error.send(Err(error)).await;
466 }
467 });
468
469 {
471 let rt = rt.clone();
472 let rt_inner = rt.clone();
473 rt.cpu().spawn(async move {
474 let result = manifest_entry_delete_ctx_rx
475 .map(|me_ctx| Ok((me_ctx, delete_file_tx.clone())))
476 .try_for_each_concurrent(
477 concurrency_limit_manifest_entries,
478 |(manifest_entry_context, tx)| {
479 let rt_inner = rt_inner.clone();
480 async move {
481 rt_inner
482 .cpu()
483 .spawn(async move {
484 Self::process_delete_manifest_entry(
485 manifest_entry_context,
486 tx,
487 )
488 .await
489 })
490 .await?
491 }
492 },
493 )
494 .await;
495
496 if let Err(error) = result {
497 let _ = channel_for_delete_manifest_entry_error
498 .send(Err(error))
499 .await;
500 }
501 });
502 }
503
504 {
506 let rt_inner = rt.clone();
507 rt.cpu().spawn(async move {
508 let result = manifest_entry_data_ctx_rx
509 .map(|me_ctx| Ok((me_ctx, planned_file_tx.clone())))
510 .try_for_each_concurrent(
511 concurrency_limit_manifest_entries,
512 |(manifest_entry_context, tx)| {
513 let rt_inner = rt_inner.clone();
514 async move {
515 rt_inner
516 .cpu()
517 .spawn(async move {
518 Self::process_data_manifest_entry(
519 manifest_entry_context,
520 tx,
521 build_result,
522 )
523 .await
524 })
525 .await?
526 }
527 },
528 )
529 .await;
530
531 if let Err(error) = result {
532 let _ = channel_for_data_manifest_entry_error.send(Err(error)).await;
533 }
534 });
535 }
536
537 Ok(planned_file_rx.boxed())
538 }
539
540 pub async fn to_arrow(&self) -> Result<ArrowRecordBatchStream> {
542 let mut arrow_reader_builder =
543 ArrowReaderBuilder::new(self.file_io.clone(), self.runtime.clone())
544 .with_data_file_concurrency_limit(self.concurrency_limit_data_files)
545 .with_row_group_filtering_enabled(self.row_group_filtering_enabled)
546 .with_row_selection_enabled(self.row_selection_enabled)
547 .with_bloom_filter_enabled(self.bloom_filter_enabled);
548
549 if let Some(batch_size) = self.batch_size {
550 arrow_reader_builder = arrow_reader_builder.with_batch_size(batch_size);
551 }
552
553 arrow_reader_builder
554 .build()
555 .read(self.plan_files().await?)
556 .map(|result| result.stream())
557 }
558
559 pub fn column_names(&self) -> Option<&[String]> {
561 self.column_names.as_deref()
562 }
563
564 pub fn snapshot(&self) -> Option<&SnapshotRef> {
566 self.plan_context.as_ref().map(|x| &x.snapshot)
567 }
568
569 async fn process_data_manifest_entry<T, F, Fut>(
570 manifest_entry_context: ManifestEntryContext,
571 mut planned_file_tx: Sender<Result<T>>,
572 build_result: F,
573 ) -> Result<()>
574 where
575 F: Fn(ManifestEntryContext) -> Fut + Send + Sync + 'static,
576 Fut: Future<Output = Result<T>> + Send + 'static,
577 {
578 if !manifest_entry_context.manifest_entry.is_alive() {
580 return Ok(());
581 }
582
583 if manifest_entry_context.manifest_entry.content_type() != DataContentType::Data {
585 return Err(Error::new(
586 ErrorKind::FeatureUnsupported,
587 "Encountered an entry for a delete file in a data file manifest",
588 ));
589 }
590
591 if let Some(ref bound_predicates) = manifest_entry_context.bound_predicates {
592 let BoundPredicates {
593 snapshot_bound_predicate,
594 partition_bound_predicate,
595 } = bound_predicates.as_ref();
596
597 let expression_evaluator_cache =
598 manifest_entry_context.expression_evaluator_cache.as_ref();
599
600 let expression_evaluator = expression_evaluator_cache.get(
601 manifest_entry_context.partition_spec_id,
602 partition_bound_predicate,
603 )?;
604
605 if !expression_evaluator.eval(manifest_entry_context.manifest_entry.data_file())? {
608 return Ok(());
609 }
610
611 if !InclusiveMetricsEvaluator::eval(
613 snapshot_bound_predicate,
614 manifest_entry_context.manifest_entry.data_file(),
615 false,
616 )? {
617 return Ok(());
618 }
619 }
620
621 planned_file_tx
625 .send(Ok(build_result(manifest_entry_context).await?))
626 .await?;
627
628 Ok(())
629 }
630
631 async fn process_delete_manifest_entry(
632 manifest_entry_context: ManifestEntryContext,
633 mut delete_file_ctx_tx: Sender<DeleteFileContext>,
634 ) -> Result<()> {
635 if !manifest_entry_context.manifest_entry.is_alive() {
637 return Ok(());
638 }
639
640 if manifest_entry_context.manifest_entry.content_type() == DataContentType::Data {
642 return Err(Error::new(
643 ErrorKind::FeatureUnsupported,
644 "Encountered an entry for a data file in a delete manifest",
645 ));
646 }
647
648 if let Some(ref bound_predicates) = manifest_entry_context.bound_predicates {
649 let expression_evaluator_cache =
650 manifest_entry_context.expression_evaluator_cache.as_ref();
651
652 let expression_evaluator = expression_evaluator_cache.get(
653 manifest_entry_context.partition_spec_id,
654 &bound_predicates.partition_bound_predicate,
655 )?;
656
657 if !expression_evaluator.eval(manifest_entry_context.manifest_entry.data_file())? {
660 return Ok(());
661 }
662 }
663
664 delete_file_ctx_tx
665 .send(DeleteFileContext {
666 manifest_entry: manifest_entry_context.manifest_entry.clone(),
667 partition_spec_id: manifest_entry_context.partition_spec_id,
668 })
669 .await?;
670
671 Ok(())
672 }
673}
674
675pub(crate) struct BoundPredicates {
676 partition_bound_predicate: BoundPredicate,
677 snapshot_bound_predicate: BoundPredicate,
678}
679
680#[cfg(test)]
681mod tests {
682 use std::collections::HashMap;
685 use std::sync::Arc;
686
687 use arrow_array::cast::AsArray;
688 use arrow_array::types::Int32Type;
689 use arrow_array::{
690 Array, ArrayRef, BooleanArray, Float64Array, Int32Array, Int64Array, RunArray, StringArray,
691 };
692 use futures::{TryStreamExt, stream};
693 use uuid::Uuid;
694
695 use crate::arrow::ArrowReaderBuilder;
696 use crate::expr::{BoundPredicate, Reference};
697 use crate::io::FileIO;
698 use crate::metadata_columns::{
699 RESERVED_COL_NAME_FILE, RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER,
700 RESERVED_COL_NAME_POS, RESERVED_COL_NAME_SPEC_ID, RESERVED_FIELD_ID_POS,
701 };
702 use crate::scan::{FileScanTask, FileScanTaskDeleteFile};
703 use crate::spec::{
704 DataContentType, DataFileFormat, Datum, Literal, MAIN_BRANCH, MappedField, NameMapping,
705 NestedField, NullOrder, Operation, PartitionSpec, PrimitiveType, Schema, Snapshot,
706 SortDirection, SortField, SortOrder, Struct, Summary, TableMetadataBuilder,
707 TableProperties, Transform, Type, UnboundPartitionSpec,
708 };
709 use crate::table::Table;
710 use crate::test_utils::scan::{TableTestFixture, assert_last_updated_seq_all};
711 use crate::test_utils::test_runtime;
712 use crate::{ErrorKind, TableIdent};
713
714 #[tokio::test]
715 async fn test_table_scan_columns() {
716 let table = TableTestFixture::new().table;
717
718 let table_scan = table.scan().select(["x", "y"]).build().unwrap();
719 assert_eq!(
720 Some(vec!["x".to_string(), "y".to_string()]),
721 table_scan.column_names
722 );
723
724 let table_scan = table
725 .scan()
726 .select(["x", "y"])
727 .select(["z"])
728 .build()
729 .unwrap();
730 assert_eq!(Some(vec!["z".to_string()]), table_scan.column_names);
731 }
732
733 #[tokio::test]
734 async fn test_select_all() {
735 let table = TableTestFixture::new().table;
736
737 let table_scan = table.scan().select_all().build().unwrap();
738 assert!(table_scan.column_names.is_none());
739 }
740
741 #[test]
742 fn test_select_no_exist_column() {
743 let table = TableTestFixture::new().table;
744
745 let table_scan = table.scan().select(["x", "y", "z", "a", "b"]).build();
746 assert!(table_scan.is_err());
747 }
748
749 #[test]
750 fn test_case_sensitive_scan_rejects_mismatched_column_case() {
751 let table = TableTestFixture::new().table;
752
753 assert!(table.scan().select(["X"]).build().is_err());
755 assert!(
756 table
757 .scan()
758 .with_filter(Reference::new("X").greater_than(Datum::long(1)))
759 .build()
760 .is_err()
761 );
762 }
763
764 #[tokio::test]
765 async fn test_case_insensitive_scan_resolves_mismatched_column_case() {
766 let mut fixture = TableTestFixture::new();
767 fixture.setup_manifest_files().await;
768
769 let table_scan = fixture
772 .table
773 .scan()
774 .with_case_sensitive(false)
775 .select(["X", "Z"])
776 .with_filter(Reference::new("Y").greater_than(Datum::long(1)))
777 .build()
778 .unwrap();
779
780 let batches: Vec<_> = table_scan
781 .to_arrow()
782 .await
783 .unwrap()
784 .try_collect()
785 .await
786 .unwrap();
787
788 assert_eq!(batches[0].num_columns(), 2);
789 assert_eq!(
790 batches[0]
791 .column_by_name("x")
792 .unwrap()
793 .as_any()
794 .downcast_ref::<Int64Array>()
795 .unwrap()
796 .value(0),
797 1
798 );
799 assert_eq!(
800 batches[0]
801 .column_by_name("z")
802 .unwrap()
803 .as_any()
804 .downcast_ref::<Int64Array>()
805 .unwrap()
806 .value(0),
807 3
808 );
809 }
810
811 #[tokio::test]
812 async fn test_table_scan_default_snapshot_id() {
813 let table = TableTestFixture::new().table;
814
815 let table_scan = table.scan().build().unwrap();
816 assert_eq!(
817 table.metadata().current_snapshot().unwrap().snapshot_id(),
818 table_scan.snapshot().unwrap().snapshot_id()
819 );
820 }
821
822 #[test]
823 fn test_table_scan_non_exist_snapshot_id() {
824 let table = TableTestFixture::new().table;
825
826 let table_scan = table.scan().snapshot_id(1024).build();
827 assert!(table_scan.is_err());
828 }
829
830 #[tokio::test]
831 async fn test_table_scan_with_snapshot_id() {
832 let table = TableTestFixture::new().table;
833
834 let table_scan = table
835 .scan()
836 .snapshot_id(3051729675574597004)
837 .with_row_selection_enabled(true)
838 .build()
839 .unwrap();
840 assert_eq!(
841 table_scan.snapshot().unwrap().snapshot_id(),
842 3051729675574597004
843 );
844 }
845
846 fn table_with_property(key: &str, value: &str) -> Table {
847 let fixture = TableTestFixture::new();
848 let mut metadata = fixture.table.metadata().clone();
849 metadata
850 .properties
851 .insert(key.to_string(), value.to_string());
852 Table::builder()
853 .metadata(metadata)
854 .identifier(fixture.table.identifier().clone())
855 .file_io(fixture.table.file_io().clone())
856 .metadata_location(fixture.table.metadata_location().unwrap().to_string())
857 .runtime(test_runtime())
858 .build()
859 .unwrap()
860 }
861
862 #[test]
863 fn test_table_scan_without_name_mapping_property() {
864 let table = TableTestFixture::new().table;
865
866 let table_scan = table.scan().build().unwrap();
867 assert!(
868 table_scan
869 .plan_context
870 .as_ref()
871 .unwrap()
872 .name_mapping
873 .is_none()
874 );
875 }
876
877 #[test]
878 fn test_table_scan_with_name_mapping_property() {
879 let mapping_json = r#"[{"field-id":1,"names":["id","record_id"]}]"#;
880 let table =
881 table_with_property(TableProperties::PROPERTY_DEFAULT_NAME_MAPPING, mapping_json);
882
883 let table_scan = table.scan().build().unwrap();
884 let mapping = table_scan
885 .plan_context
886 .as_ref()
887 .unwrap()
888 .name_mapping
889 .as_ref()
890 .expect("name_mapping should be parsed from the table property");
891 let fields = mapping.fields();
892 assert_eq!(fields.len(), 1);
893 assert_eq!(fields[0].field_id(), Some(1));
894 assert_eq!(fields[0].names(), &[
895 "id".to_string(),
896 "record_id".to_string()
897 ]);
898 }
899
900 #[test]
901 fn test_table_scan_with_malformed_name_mapping_property() {
902 let table = table_with_property(
903 TableProperties::PROPERTY_DEFAULT_NAME_MAPPING,
904 "{ not valid json",
905 );
906
907 let err = table
908 .scan()
909 .build()
910 .expect_err("malformed name mapping should fail to parse");
911 assert_eq!(err.kind(), ErrorKind::DataInvalid);
912 }
913
914 #[tokio::test]
915 async fn test_plan_files_carries_name_mapping_into_file_scan_task() {
916 let mut fixture = TableTestFixture::new();
917 fixture.setup_manifest_files().await;
918
919 let mapping_json = r#"[{"field-id":1,"names":["id","record_id"]}]"#;
920 let mut metadata = fixture.table.metadata().clone();
921 metadata.properties.insert(
922 TableProperties::PROPERTY_DEFAULT_NAME_MAPPING.to_string(),
923 mapping_json.to_string(),
924 );
925 let table = Table::builder()
926 .metadata(metadata)
927 .identifier(fixture.table.identifier().clone())
928 .file_io(fixture.table.file_io().clone())
929 .metadata_location(fixture.table.metadata_location().unwrap().to_string())
930 .runtime(test_runtime())
931 .build()
932 .unwrap();
933
934 let tasks: Vec<_> = table
935 .scan()
936 .build()
937 .unwrap()
938 .plan_files()
939 .await
940 .unwrap()
941 .try_collect()
942 .await
943 .unwrap();
944
945 assert!(!tasks.is_empty(), "expected at least one FileScanTask");
946 for task in &tasks {
947 let mapping = task
948 .name_mapping()
949 .expect("name_mapping should reach the FileScanTask");
950 assert_eq!(mapping.fields().len(), 1);
951 assert_eq!(mapping.fields()[0].field_id(), Some(1));
952 }
953 }
954
955 #[tokio::test]
956 async fn test_plan_files_carries_sort_order_into_file_scan_task() {
957 let mut fixture = TableTestFixture::new();
958
959 let mut metadata = fixture.table.metadata().clone();
963 metadata
964 .sort_orders
965 .insert(0, Arc::new(SortOrder::unsorted_order()));
966 fixture.table = fixture.table.with_metadata(Arc::new(metadata));
967
968 let expected_sort_order = fixture
969 .table
970 .metadata()
971 .sort_order_by_id(3)
972 .unwrap()
973 .clone();
974
975 fixture
977 .setup_manifest_files_with_sort_order_ids([Some(3), None, Some(99), Some(0)])
978 .await;
979
980 let tasks: Vec<_> = fixture
981 .table
982 .scan()
983 .build()
984 .unwrap()
985 .plan_files()
986 .await
987 .unwrap()
988 .try_collect()
989 .await
990 .unwrap();
991
992 assert_eq!(tasks.len(), 4, "expected all four FileScanTasks");
993
994 assert_eq!(
997 tasks.iter().filter(|t| t.sort_order().is_some()).count(),
998 1,
999 "exactly one file resolves to a sort order"
1000 );
1001 assert_eq!(
1002 tasks.iter().filter(|t| t.sort_order_id().is_some()).count(),
1003 3,
1004 "three files carry a raw sort_order_id"
1005 );
1006
1007 let resolved = tasks
1008 .iter()
1009 .find(|t| t.data_file_path().ends_with("1.parquet"))
1010 .unwrap();
1011 assert_eq!(resolved.sort_order_id(), Some(3));
1012 assert_eq!(
1013 resolved.sort_order(),
1014 Some(&expected_sort_order),
1015 "sort_order_id 3 should resolve to the table's sort order at id 3"
1016 );
1017
1018 let missing = tasks
1019 .iter()
1020 .find(|t| t.data_file_path().ends_with("2.parquet"))
1021 .unwrap();
1022 assert_eq!(missing.sort_order_id(), None);
1023 assert!(
1024 missing.sort_order().is_none(),
1025 "a file with no sort_order_id carries no sort_order"
1026 );
1027
1028 let unresolvable = tasks
1029 .iter()
1030 .find(|t| t.data_file_path().ends_with("3.parquet"))
1031 .unwrap();
1032 assert_eq!(
1033 unresolvable.sort_order_id(),
1034 Some(99),
1035 "the raw id is preserved even when it does not resolve"
1036 );
1037 assert!(
1038 unresolvable.sort_order().is_none(),
1039 "an unresolvable sort_order_id resolves to no sort_order"
1040 );
1041
1042 let unsorted = tasks
1043 .iter()
1044 .find(|t| t.data_file_path().ends_with("4.parquet"))
1045 .unwrap();
1046 assert_eq!(unsorted.sort_order_id(), Some(0));
1047 assert!(
1048 unsorted.sort_order().is_none(),
1049 "the reserved unsorted order (id 0) resolves to no sort_order"
1050 );
1051 }
1052
1053 #[tokio::test]
1054 async fn test_plan_files_on_table_without_any_snapshots() {
1055 let table = TableTestFixture::new_empty().table;
1056 let batch_stream = table.scan().build().unwrap().to_arrow().await.unwrap();
1057 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1058 assert!(batches.is_empty());
1059 }
1060
1061 #[tokio::test]
1062 async fn test_plan_files_no_deletions() {
1063 let mut fixture = TableTestFixture::new();
1064 fixture.setup_manifest_files().await;
1065
1066 let table_scan = fixture
1068 .table
1069 .scan()
1070 .with_row_selection_enabled(true)
1071 .build()
1072 .unwrap();
1073
1074 let mut tasks = table_scan
1075 .plan_files()
1076 .await
1077 .unwrap()
1078 .try_fold(vec![], |mut acc, task| async move {
1079 acc.push(task);
1080 Ok(acc)
1081 })
1082 .await
1083 .unwrap();
1084
1085 assert_eq!(tasks.len(), 2);
1086
1087 tasks.sort_by_key(|t| t.data_file_path().to_string());
1088
1089 assert_eq!(
1091 tasks[0].data_file_path(),
1092 format!("{}/1.parquet", &fixture.table_location)
1093 );
1094
1095 assert_eq!(
1097 tasks[1].data_file_path(),
1098 format!("{}/3.parquet", &fixture.table_location)
1099 );
1100 }
1101
1102 #[tokio::test]
1103 async fn test_plan_files_carries_row_lineage_into_file_scan_task() {
1104 let mut fixture = TableTestFixture::new();
1105 fixture.setup_manifest_files().await;
1106
1107 let mut tasks: Vec<_> = fixture
1108 .table
1109 .scan()
1110 .build()
1111 .unwrap()
1112 .plan_files()
1113 .await
1114 .unwrap()
1115 .try_collect()
1116 .await
1117 .unwrap();
1118
1119 assert_eq!(tasks.len(), 2);
1120 tasks.sort_by_key(|task| task.data_file_path().to_string());
1121
1122 assert_eq!(
1125 tasks[0].data_file_path(),
1126 format!("{}/1.parquet", &fixture.table_location)
1127 );
1128 assert_eq!(tasks[0].data_sequence_number(), Some(1));
1129 assert_eq!(
1130 tasks[1].data_file_path(),
1131 format!("{}/3.parquet", &fixture.table_location)
1132 );
1133 assert_eq!(tasks[1].data_sequence_number(), Some(0));
1134
1135 assert!(tasks.iter().all(|task| task.first_row_id().is_none()));
1137 }
1138
1139 #[tokio::test]
1140 async fn test_plan_files_carries_row_lineage_from_v3_manifest() {
1141 let mut fixture = TableTestFixture::new();
1142 fixture.setup_v3_manifest_files().await;
1143
1144 let task = fixture
1145 .table
1146 .scan()
1147 .build()
1148 .unwrap()
1149 .plan_files()
1150 .await
1151 .unwrap()
1152 .try_collect::<Vec<_>>()
1153 .await
1154 .unwrap()
1155 .into_iter()
1156 .next()
1157 .expect("expected one FileScanTask");
1158
1159 assert_eq!(task.first_row_id(), Some(42));
1162 assert_eq!(task.data_sequence_number(), Some(1));
1164 }
1165
1166 #[tokio::test]
1167 async fn test_filtered_scan_with_dropped_partition_source_column() {
1168 let mut fixture = TableTestFixture::new();
1169 fixture.setup_manifest_files().await;
1170
1171 let baseline = scan_y_gte_5(&fixture.table).await;
1173 assert!(!baseline.is_empty());
1174 assert!(baseline.iter().all(|y| *y >= 5));
1175
1176 let current_schema = fixture.table.metadata().current_schema();
1180 let evolved_schema = Schema::builder()
1181 .with_fields(
1182 current_schema
1183 .as_struct()
1184 .fields()
1185 .iter()
1186 .filter(|field| field.id != 1)
1187 .cloned(),
1188 )
1189 .with_identifier_field_ids(vec![2])
1190 .build()
1191 .unwrap();
1192 let evolved =
1193 TableMetadataBuilder::new_from_metadata(fixture.table.metadata().clone(), None)
1194 .add_default_partition_spec(UnboundPartitionSpec::builder().build())
1195 .unwrap()
1196 .add_current_schema(evolved_schema)
1197 .unwrap()
1198 .build()
1199 .unwrap()
1200 .metadata;
1201
1202 let parent = evolved.current_snapshot().unwrap().clone();
1205 let snapshot = Snapshot::builder()
1206 .with_snapshot_id(parent.snapshot_id() + 1)
1207 .with_parent_snapshot_id(Some(parent.snapshot_id()))
1208 .with_sequence_number(evolved.last_sequence_number() + 1)
1209 .with_timestamp_ms(evolved.last_updated_ms + 1)
1210 .with_schema_id(evolved.current_schema_id())
1211 .with_manifest_list(parent.manifest_list())
1212 .with_summary(Summary {
1213 operation: Operation::Append,
1214 additional_properties: HashMap::new(),
1215 })
1216 .build();
1217 let metadata = TableMetadataBuilder::new_from_metadata(evolved, None)
1218 .set_branch_snapshot(snapshot, MAIN_BRANCH)
1219 .unwrap()
1220 .build()
1221 .unwrap()
1222 .metadata;
1223 let table = fixture.table.clone().with_metadata(Arc::new(metadata));
1224
1225 let evolved = scan_y_gte_5(&table).await;
1228 assert_eq!(evolved, baseline);
1229 }
1230
1231 async fn scan_y_gte_5(table: &Table) -> Vec<i64> {
1232 let table_scan = table
1233 .scan()
1234 .select(["y"])
1235 .with_filter(Reference::new("y").greater_than_or_equal_to(Datum::long(5)))
1236 .build()
1237 .unwrap();
1238 let batches: Vec<_> = table_scan
1239 .to_arrow()
1240 .await
1241 .unwrap()
1242 .try_collect()
1243 .await
1244 .unwrap();
1245
1246 let mut values: Vec<i64> = batches
1247 .iter()
1248 .flat_map(|batch| {
1249 let col = batch.column_by_name("y").unwrap();
1250 let arr = col.as_any().downcast_ref::<Int64Array>().unwrap();
1251 (0..arr.len()).map(|i| arr.value(i)).collect::<Vec<_>>()
1252 })
1253 .collect();
1254 values.sort_unstable();
1255 values
1256 }
1257
1258 #[tokio::test]
1259 async fn test_open_parquet_no_deletions() {
1260 let mut fixture = TableTestFixture::new();
1261 fixture.setup_manifest_files().await;
1262
1263 let table_scan = fixture
1265 .table
1266 .scan()
1267 .with_row_selection_enabled(true)
1268 .build()
1269 .unwrap();
1270
1271 let batch_stream = table_scan.to_arrow().await.unwrap();
1272
1273 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1274
1275 let col = batches[0].column_by_name("x").unwrap();
1276
1277 let int64_arr = col.as_any().downcast_ref::<Int64Array>().unwrap();
1278 assert_eq!(int64_arr.value(0), 1);
1279 }
1280
1281 #[tokio::test]
1282 async fn test_open_parquet_no_deletions_by_separate_reader() {
1283 let mut fixture = TableTestFixture::new();
1284 fixture.setup_manifest_files().await;
1285
1286 let table_scan = fixture
1288 .table
1289 .scan()
1290 .with_row_selection_enabled(true)
1291 .build()
1292 .unwrap();
1293
1294 let mut plan_task: Vec<_> = table_scan
1295 .plan_files()
1296 .await
1297 .unwrap()
1298 .try_collect()
1299 .await
1300 .unwrap();
1301 assert_eq!(plan_task.len(), 2);
1302
1303 let reader = ArrowReaderBuilder::new(
1304 fixture.table.file_io().clone(),
1305 fixture.table.runtime().clone(),
1306 )
1307 .build();
1308 let batch_stream = reader
1309 .clone()
1310 .read(Box::pin(stream::iter(vec![Ok(plan_task.remove(0))])))
1311 .unwrap()
1312 .stream();
1313 let batch_1: Vec<_> = batch_stream.try_collect().await.unwrap();
1314
1315 let reader = ArrowReaderBuilder::new(
1316 fixture.table.file_io().clone(),
1317 fixture.table.runtime().clone(),
1318 )
1319 .build();
1320 let batch_stream = reader
1321 .read(Box::pin(stream::iter(vec![Ok(plan_task.remove(0))])))
1322 .unwrap()
1323 .stream();
1324 let batch_2: Vec<_> = batch_stream.try_collect().await.unwrap();
1325
1326 assert_eq!(batch_1, batch_2);
1327 }
1328
1329 #[tokio::test]
1330 async fn test_open_parquet_with_projection() {
1331 let mut fixture = TableTestFixture::new();
1332 fixture.setup_manifest_files().await;
1333
1334 let table_scan = fixture
1336 .table
1337 .scan()
1338 .select(["x", "z"])
1339 .with_row_selection_enabled(true)
1340 .build()
1341 .unwrap();
1342
1343 let batch_stream = table_scan.to_arrow().await.unwrap();
1344
1345 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1346
1347 assert_eq!(batches[0].num_columns(), 2);
1348
1349 let col1 = batches[0].column_by_name("x").unwrap();
1350 let int64_arr = col1.as_any().downcast_ref::<Int64Array>().unwrap();
1351 assert_eq!(int64_arr.value(0), 1);
1352
1353 let col2 = batches[0].column_by_name("z").unwrap();
1354 let int64_arr = col2.as_any().downcast_ref::<Int64Array>().unwrap();
1355 assert_eq!(int64_arr.value(0), 3);
1356
1357 let table_scan = fixture.table.scan().select_empty().build().unwrap();
1359 let batch_stream = table_scan.to_arrow().await.unwrap();
1360 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1361
1362 assert_eq!(batches[0].num_columns(), 0);
1363 assert_eq!(batches[0].num_rows(), 1024);
1364 }
1365
1366 #[tokio::test]
1367 async fn test_filter_on_arrow_lt() {
1368 let mut fixture = TableTestFixture::new();
1369 fixture.setup_manifest_files().await;
1370
1371 let mut builder = fixture.table.scan();
1373 let predicate = Reference::new("y").less_than(Datum::long(3));
1374 builder = builder
1375 .with_filter(predicate)
1376 .with_row_selection_enabled(true);
1377 let table_scan = builder.build().unwrap();
1378
1379 let batch_stream = table_scan.to_arrow().await.unwrap();
1380
1381 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1382
1383 assert_eq!(batches[0].num_rows(), 512);
1384
1385 let col = batches[0].column_by_name("x").unwrap();
1386 let int64_arr = col.as_any().downcast_ref::<Int64Array>().unwrap();
1387 assert_eq!(int64_arr.value(0), 1);
1388
1389 let col = batches[0].column_by_name("y").unwrap();
1390 let int64_arr = col.as_any().downcast_ref::<Int64Array>().unwrap();
1391 assert_eq!(int64_arr.value(0), 2);
1392 }
1393
1394 #[tokio::test]
1395 async fn test_filter_on_arrow_gt_eq() {
1396 let mut fixture = TableTestFixture::new();
1397 fixture.setup_manifest_files().await;
1398
1399 let mut builder = fixture.table.scan();
1401 let predicate = Reference::new("y").greater_than_or_equal_to(Datum::long(5));
1402 builder = builder
1403 .with_filter(predicate)
1404 .with_row_selection_enabled(true);
1405 let table_scan = builder.build().unwrap();
1406
1407 let batch_stream = table_scan.to_arrow().await.unwrap();
1408
1409 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1410
1411 assert_eq!(batches[0].num_rows(), 12);
1412
1413 let col = batches[0].column_by_name("x").unwrap();
1414 let int64_arr = col.as_any().downcast_ref::<Int64Array>().unwrap();
1415 assert_eq!(int64_arr.value(0), 1);
1416
1417 let col = batches[0].column_by_name("y").unwrap();
1418 let int64_arr = col.as_any().downcast_ref::<Int64Array>().unwrap();
1419 assert_eq!(int64_arr.value(0), 5);
1420 }
1421
1422 #[tokio::test]
1423 async fn test_filter_double_eq() {
1424 let mut fixture = TableTestFixture::new();
1425 fixture.setup_manifest_files().await;
1426
1427 let mut builder = fixture.table.scan();
1429 let predicate = Reference::new("dbl").equal_to(Datum::double(150.0f64));
1430 builder = builder
1431 .with_filter(predicate)
1432 .with_row_selection_enabled(true);
1433 let table_scan = builder.build().unwrap();
1434
1435 let batch_stream = table_scan.to_arrow().await.unwrap();
1436
1437 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1438
1439 assert_eq!(batches.len(), 2);
1440 assert_eq!(batches[0].num_rows(), 12);
1441
1442 let col = batches[0].column_by_name("dbl").unwrap();
1443 let f64_arr = col.as_any().downcast_ref::<Float64Array>().unwrap();
1444 assert_eq!(f64_arr.value(1), 150.0f64);
1445 }
1446
1447 #[tokio::test]
1448 async fn test_filter_int_eq() {
1449 let mut fixture = TableTestFixture::new();
1450 fixture.setup_manifest_files().await;
1451
1452 let mut builder = fixture.table.scan();
1454 let predicate = Reference::new("i32").equal_to(Datum::int(150i32));
1455 builder = builder
1456 .with_filter(predicate)
1457 .with_row_selection_enabled(true);
1458 let table_scan = builder.build().unwrap();
1459
1460 let batch_stream = table_scan.to_arrow().await.unwrap();
1461
1462 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1463
1464 assert_eq!(batches.len(), 2);
1465 assert_eq!(batches[0].num_rows(), 12);
1466
1467 let col = batches[0].column_by_name("i32").unwrap();
1468 let i32_arr = col.as_any().downcast_ref::<Int32Array>().unwrap();
1469 assert_eq!(i32_arr.value(1), 150i32);
1470 }
1471
1472 #[tokio::test]
1473 async fn test_filter_long_eq() {
1474 let mut fixture = TableTestFixture::new();
1475 fixture.setup_manifest_files().await;
1476
1477 let mut builder = fixture.table.scan();
1479 let predicate = Reference::new("i64").equal_to(Datum::long(150i64));
1480 builder = builder
1481 .with_filter(predicate)
1482 .with_row_selection_enabled(true);
1483 let table_scan = builder.build().unwrap();
1484
1485 let batch_stream = table_scan.to_arrow().await.unwrap();
1486
1487 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1488
1489 assert_eq!(batches.len(), 2);
1490 assert_eq!(batches[0].num_rows(), 12);
1491
1492 let col = batches[0].column_by_name("i64").unwrap();
1493 let i64_arr = col.as_any().downcast_ref::<Int64Array>().unwrap();
1494 assert_eq!(i64_arr.value(1), 150i64);
1495 }
1496
1497 #[tokio::test]
1498 async fn test_filter_bool_eq() {
1499 let mut fixture = TableTestFixture::new();
1500 fixture.setup_manifest_files().await;
1501
1502 let mut builder = fixture.table.scan();
1504 let predicate = Reference::new("bool").equal_to(Datum::bool(true));
1505 builder = builder
1506 .with_filter(predicate)
1507 .with_row_selection_enabled(true);
1508 let table_scan = builder.build().unwrap();
1509
1510 let batch_stream = table_scan.to_arrow().await.unwrap();
1511
1512 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1513
1514 assert_eq!(batches.len(), 2);
1515 assert_eq!(batches[0].num_rows(), 512);
1516
1517 let col = batches[0].column_by_name("bool").unwrap();
1518 let bool_arr = col.as_any().downcast_ref::<BooleanArray>().unwrap();
1519 assert!(bool_arr.value(1));
1520 }
1521
1522 #[tokio::test]
1523 async fn test_filter_on_arrow_is_null() {
1524 let mut fixture = TableTestFixture::new();
1525 fixture.setup_manifest_files().await;
1526
1527 let mut builder = fixture.table.scan();
1529 let predicate = Reference::new("y").is_null();
1530 builder = builder
1531 .with_filter(predicate)
1532 .with_row_selection_enabled(true);
1533 let table_scan = builder.build().unwrap();
1534
1535 let batch_stream = table_scan.to_arrow().await.unwrap();
1536
1537 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1538 assert_eq!(batches.len(), 0);
1539 }
1540
1541 #[tokio::test]
1542 async fn test_filter_on_arrow_is_not_null() {
1543 let mut fixture = TableTestFixture::new();
1544 fixture.setup_manifest_files().await;
1545
1546 let mut builder = fixture.table.scan();
1548 let predicate = Reference::new("y").is_not_null();
1549 builder = builder
1550 .with_filter(predicate)
1551 .with_row_selection_enabled(true);
1552 let table_scan = builder.build().unwrap();
1553
1554 let batch_stream = table_scan.to_arrow().await.unwrap();
1555
1556 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1557 assert_eq!(batches[0].num_rows(), 1024);
1558 }
1559
1560 #[tokio::test]
1561 async fn test_filter_on_arrow_lt_and_gt() {
1562 let mut fixture = TableTestFixture::new();
1563 fixture.setup_manifest_files().await;
1564
1565 let mut builder = fixture.table.scan();
1567 let predicate = Reference::new("y")
1568 .less_than(Datum::long(5))
1569 .and(Reference::new("z").greater_than_or_equal_to(Datum::long(4)));
1570 builder = builder
1571 .with_filter(predicate)
1572 .with_row_selection_enabled(true);
1573 let table_scan = builder.build().unwrap();
1574
1575 let batch_stream = table_scan.to_arrow().await.unwrap();
1576
1577 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1578 assert_eq!(batches[0].num_rows(), 500);
1579
1580 let col = batches[0].column_by_name("x").unwrap();
1581 let expected_x = Arc::new(Int64Array::from_iter_values(vec![1; 500])) as ArrayRef;
1582 assert_eq!(col, &expected_x);
1583
1584 let col = batches[0].column_by_name("y").unwrap();
1585 let mut values = vec![];
1586 values.append(vec![3; 200].as_mut());
1587 values.append(vec![4; 300].as_mut());
1588 let expected_y = Arc::new(Int64Array::from_iter_values(values)) as ArrayRef;
1589 assert_eq!(col, &expected_y);
1590
1591 let col = batches[0].column_by_name("z").unwrap();
1592 let expected_z = Arc::new(Int64Array::from_iter_values(vec![4; 500])) as ArrayRef;
1593 assert_eq!(col, &expected_z);
1594 }
1595
1596 #[tokio::test]
1597 async fn test_filter_on_arrow_lt_or_gt() {
1598 let mut fixture = TableTestFixture::new();
1599 fixture.setup_manifest_files().await;
1600
1601 let mut builder = fixture.table.scan();
1603 let predicate = Reference::new("y")
1604 .less_than(Datum::long(5))
1605 .or(Reference::new("z").greater_than_or_equal_to(Datum::long(4)));
1606 builder = builder
1607 .with_filter(predicate)
1608 .with_row_selection_enabled(true);
1609 let table_scan = builder.build().unwrap();
1610
1611 let batch_stream = table_scan.to_arrow().await.unwrap();
1612
1613 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1614 assert_eq!(batches[0].num_rows(), 1024);
1615
1616 let col = batches[0].column_by_name("x").unwrap();
1617 let expected_x = Arc::new(Int64Array::from_iter_values(vec![1; 1024])) as ArrayRef;
1618 assert_eq!(col, &expected_x);
1619
1620 let col = batches[0].column_by_name("y").unwrap();
1621 let mut values = vec![2; 512];
1622 values.append(vec![3; 200].as_mut());
1623 values.append(vec![4; 300].as_mut());
1624 values.append(vec![5; 12].as_mut());
1625 let expected_y = Arc::new(Int64Array::from_iter_values(values)) as ArrayRef;
1626 assert_eq!(col, &expected_y);
1627
1628 let col = batches[0].column_by_name("z").unwrap();
1629 let mut values = vec![3; 512];
1630 values.append(vec![4; 512].as_mut());
1631 let expected_z = Arc::new(Int64Array::from_iter_values(values)) as ArrayRef;
1632 assert_eq!(col, &expected_z);
1633 }
1634
1635 #[tokio::test]
1636 async fn test_filter_on_arrow_startswith() {
1637 let mut fixture = TableTestFixture::new();
1638 fixture.setup_manifest_files().await;
1639
1640 let mut builder = fixture.table.scan();
1642 let predicate = Reference::new("a").starts_with(Datum::string("Ice"));
1643 builder = builder
1644 .with_filter(predicate)
1645 .with_row_selection_enabled(true);
1646 let table_scan = builder.build().unwrap();
1647
1648 let batch_stream = table_scan.to_arrow().await.unwrap();
1649
1650 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1651
1652 assert_eq!(batches[0].num_rows(), 512);
1653
1654 let col = batches[0].column_by_name("a").unwrap();
1655 let string_arr = col.as_any().downcast_ref::<StringArray>().unwrap();
1656 assert_eq!(string_arr.value(0), "Iceberg");
1657 }
1658
1659 #[tokio::test]
1660 async fn test_filter_on_arrow_not_startswith() {
1661 let mut fixture = TableTestFixture::new();
1662 fixture.setup_manifest_files().await;
1663
1664 let mut builder = fixture.table.scan();
1666 let predicate = Reference::new("a").not_starts_with(Datum::string("Ice"));
1667 builder = builder
1668 .with_filter(predicate)
1669 .with_row_selection_enabled(true);
1670 let table_scan = builder.build().unwrap();
1671
1672 let batch_stream = table_scan.to_arrow().await.unwrap();
1673
1674 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1675
1676 assert_eq!(batches[0].num_rows(), 512);
1677
1678 let col = batches[0].column_by_name("a").unwrap();
1679 let string_arr = col.as_any().downcast_ref::<StringArray>().unwrap();
1680 assert_eq!(string_arr.value(0), "Apache");
1681 }
1682
1683 #[tokio::test]
1684 async fn test_filter_on_arrow_in() {
1685 let mut fixture = TableTestFixture::new();
1686 fixture.setup_manifest_files().await;
1687
1688 let mut builder = fixture.table.scan();
1690 let predicate =
1691 Reference::new("a").is_in([Datum::string("Sioux"), Datum::string("Iceberg")]);
1692 builder = builder
1693 .with_filter(predicate)
1694 .with_row_selection_enabled(true);
1695 let table_scan = builder.build().unwrap();
1696
1697 let batch_stream = table_scan.to_arrow().await.unwrap();
1698
1699 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1700
1701 assert_eq!(batches[0].num_rows(), 512);
1702
1703 let col = batches[0].column_by_name("a").unwrap();
1704 let string_arr = col.as_any().downcast_ref::<StringArray>().unwrap();
1705 assert_eq!(string_arr.value(0), "Iceberg");
1706 }
1707
1708 #[tokio::test]
1709 async fn test_filter_on_arrow_not_in() {
1710 let mut fixture = TableTestFixture::new();
1711 fixture.setup_manifest_files().await;
1712
1713 let mut builder = fixture.table.scan();
1715 let predicate =
1716 Reference::new("a").is_not_in([Datum::string("Sioux"), Datum::string("Iceberg")]);
1717 builder = builder
1718 .with_filter(predicate)
1719 .with_row_selection_enabled(true);
1720 let table_scan = builder.build().unwrap();
1721
1722 let batch_stream = table_scan.to_arrow().await.unwrap();
1723
1724 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1725
1726 assert_eq!(batches[0].num_rows(), 512);
1727
1728 let col = batches[0].column_by_name("a").unwrap();
1729 let string_arr = col.as_any().downcast_ref::<StringArray>().unwrap();
1730 assert_eq!(string_arr.value(0), "Apache");
1731 }
1732
1733 fn file_scan_task_test_schema(primitive_type: PrimitiveType) -> Arc<Schema> {
1734 Arc::new(
1735 Schema::builder()
1736 .with_fields(vec![Arc::new(NestedField::required(
1737 1,
1738 "x",
1739 Type::Primitive(primitive_type),
1740 ))])
1741 .build()
1742 .unwrap(),
1743 )
1744 }
1745
1746 fn assert_file_scan_task_serde_round_trip(task: FileScanTask) {
1747 let serialized = serde_json::to_string(&task).unwrap();
1749 let deserialized: FileScanTask = serde_json::from_str(&serialized).unwrap();
1750
1751 assert_eq!(task, deserialized);
1752 }
1753
1754 fn file_scan_task_with_partition(
1755 primitive_type: PrimitiveType,
1756 transform: Transform,
1757 partition_value: Literal,
1758 ) -> FileScanTask {
1759 let schema = file_scan_task_test_schema(primitive_type);
1760 let partition_spec = Arc::new(
1761 PartitionSpec::builder(schema.clone())
1762 .add_partition_field("x", "x_partition", transform)
1763 .unwrap()
1764 .build()
1765 .unwrap(),
1766 );
1767 FileScanTask::builder()
1768 .with_data_file_path("data_file_path".to_string())
1769 .with_file_size_in_bytes(123)
1770 .with_start(10)
1771 .with_length(100)
1772 .with_project_field_ids(vec![1])
1773 .with_schema(schema)
1774 .with_data_file_format(DataFileFormat::Parquet)
1775 .with_partition(Some(Struct::from_iter([Some(partition_value)])))
1776 .with_partition_spec(Some(partition_spec))
1777 .with_case_sensitive(true)
1778 .build()
1779 .unwrap()
1780 }
1781
1782 #[test]
1783 fn test_file_scan_task_serde_without_predicate() {
1784 let task = FileScanTask::builder()
1785 .with_data_file_path("data_file_path".to_string())
1786 .with_file_size_in_bytes(0)
1787 .with_start(0)
1788 .with_length(100)
1789 .with_project_field_ids(vec![1, 2, 3])
1790 .with_schema(file_scan_task_test_schema(PrimitiveType::Binary))
1791 .with_record_count(Some(100))
1792 .with_first_row_id(Some(1000))
1793 .with_data_sequence_number(Some(5))
1794 .with_data_file_format(DataFileFormat::Parquet)
1795 .with_case_sensitive(false)
1796 .build()
1797 .unwrap();
1798 assert_file_scan_task_serde_round_trip(task);
1799 }
1800
1801 #[test]
1802 fn test_file_scan_task_serde_with_predicate() {
1803 let task = FileScanTask::builder()
1804 .with_data_file_path("data_file_path".to_string())
1805 .with_file_size_in_bytes(0)
1806 .with_start(0)
1807 .with_length(100)
1808 .with_project_field_ids(vec![1, 2, 3])
1809 .with_predicate(Some(BoundPredicate::AlwaysTrue))
1810 .with_schema(file_scan_task_test_schema(PrimitiveType::Binary))
1811 .with_data_file_format(DataFileFormat::Avro)
1812 .with_case_sensitive(false)
1813 .build()
1814 .unwrap();
1815
1816 let serialized = serde_json::to_value(&task).unwrap();
1817 assert!(serialized.get("record_count").is_none());
1818 assert_file_scan_task_serde_round_trip(task);
1819 }
1820
1821 #[test]
1822 fn test_unpartitioned_file_scan_task_serde() {
1823 let task = FileScanTask::builder()
1824 .with_data_file_path("data_file_path".to_string())
1825 .with_file_size_in_bytes(0)
1826 .with_start(0)
1827 .with_length(100)
1828 .with_project_field_ids(vec![1, 2, 3])
1829 .with_schema(file_scan_task_test_schema(PrimitiveType::Binary))
1830 .with_data_file_format(DataFileFormat::Parquet)
1831 .with_partition(Some(Struct::empty()))
1832 .with_case_sensitive(false)
1833 .build()
1834 .unwrap();
1835 assert_file_scan_task_serde_round_trip(task);
1836 }
1837
1838 #[test]
1839 fn test_file_scan_task_serde_with_all_optional_fields() {
1840 let schema = file_scan_task_test_schema(PrimitiveType::Long);
1841 let partition_spec = Arc::new(
1842 PartitionSpec::builder(schema.clone())
1843 .add_partition_field("x", "x", Transform::Identity)
1844 .unwrap()
1845 .build()
1846 .unwrap(),
1847 );
1848 let unified_partition_type = Arc::new(partition_spec.partition_type(&schema).unwrap());
1849 let sort_order = Arc::new(
1850 SortOrder::builder()
1851 .with_order_id(1)
1852 .with_sort_field(
1853 SortField::builder()
1854 .source_id(1)
1855 .transform(Transform::Identity)
1856 .direction(SortDirection::Ascending)
1857 .null_order(NullOrder::First)
1858 .build(),
1859 )
1860 .build(&schema)
1861 .unwrap(),
1862 );
1863 let task = FileScanTask::builder()
1864 .with_data_file_path("data_file_path".to_string())
1865 .with_file_size_in_bytes(123)
1866 .with_start(10)
1867 .with_length(100)
1868 .with_project_field_ids(vec![1])
1869 .with_schema(schema)
1870 .with_data_file_format(DataFileFormat::Parquet)
1871 .with_deletes(vec![
1872 FileScanTaskDeleteFile::builder()
1873 .with_file_path("delete_file_path".to_string())
1874 .with_file_size_in_bytes(23)
1875 .with_file_type(DataContentType::EqualityDeletes)
1876 .with_file_format(DataFileFormat::Parquet)
1877 .with_partition_spec_id(0)
1878 .with_equality_ids(Some(vec![1]))
1879 .with_referenced_data_file(Some("data_file_path".to_string()))
1880 .with_content_offset(Some(12))
1881 .with_content_size_in_bytes(Some(34))
1882 .with_record_count(Some(5))
1883 .with_key_metadata(Some(vec![4, 5, 6].into_boxed_slice()))
1884 .build()
1885 .unwrap(),
1886 ])
1887 .with_partition(Some(Struct::from_iter([Some(Literal::long(42))])))
1888 .with_partition_spec(Some(partition_spec))
1889 .with_name_mapping(Some(Arc::new(NameMapping::new(vec![MappedField::new(
1890 Some(1),
1891 vec!["x".to_string()],
1892 vec![],
1893 )]))))
1894 .with_unified_partition_type(Some(unified_partition_type))
1895 .with_sort_order_id(Some(1))
1896 .with_sort_order(Some(sort_order))
1897 .with_case_sensitive(true)
1898 .with_key_metadata(Some(vec![1, 2, 3].into_boxed_slice()))
1899 .build()
1900 .unwrap();
1901 assert_file_scan_task_serde_round_trip(task);
1902 }
1903
1904 #[test]
1905 fn test_file_scan_task_serde_with_date_partition() {
1906 let task = file_scan_task_with_partition(
1907 PrimitiveType::Date,
1908 Transform::Identity,
1909 Literal::date(19_000),
1910 );
1911 assert_file_scan_task_serde_round_trip(task);
1912 }
1913
1914 #[test]
1915 fn test_file_scan_task_serde_with_timestamp_ns_partition() {
1916 let task = file_scan_task_with_partition(
1917 PrimitiveType::TimestampNs,
1918 Transform::Identity,
1919 Literal::timestamp_nano(1_510_871_468_123_456_789),
1920 );
1921 assert_file_scan_task_serde_round_trip(task);
1922 }
1923
1924 #[test]
1925 fn test_file_scan_task_serde_with_timestamptz_ns_partition() {
1926 let task = file_scan_task_with_partition(
1927 PrimitiveType::TimestamptzNs,
1928 Transform::Identity,
1929 Literal::timestamptz_nano(1_510_871_468_123_456_789),
1930 );
1931 assert_file_scan_task_serde_round_trip(task);
1932 }
1933
1934 #[test]
1935 fn test_file_scan_task_serde_with_decimal_partition() {
1936 let task = file_scan_task_with_partition(
1937 PrimitiveType::Decimal {
1938 precision: 9,
1939 scale: 2,
1940 },
1941 Transform::Identity,
1942 Literal::decimal(12_345),
1943 );
1944 assert_file_scan_task_serde_round_trip(task);
1945 }
1946
1947 #[test]
1948 fn test_file_scan_task_serde_with_uuid_partition() {
1949 let task = file_scan_task_with_partition(
1950 PrimitiveType::Uuid,
1951 Transform::Identity,
1952 Literal::uuid(Uuid::from_u128(0x12345678_90ab_cdef_1234_567890abcdef)),
1953 );
1954 assert_file_scan_task_serde_round_trip(task);
1955 }
1956
1957 #[test]
1958 fn test_file_scan_task_serde_with_fixed_partition() {
1959 let task = file_scan_task_with_partition(
1960 PrimitiveType::Fixed(4),
1961 Transform::Identity,
1962 Literal::fixed([1, 2, 3, 4]),
1963 );
1964 assert_file_scan_task_serde_round_trip(task);
1965 }
1966
1967 #[test]
1968 fn test_file_scan_task_serde_with_bucket_partition() {
1969 let task = file_scan_task_with_partition(
1970 PrimitiveType::String,
1971 Transform::Bucket(4),
1972 Literal::int(2),
1973 );
1974 assert_file_scan_task_serde_round_trip(task);
1975 }
1976
1977 #[tokio::test]
1978 async fn test_select_with_file_column() {
1979 let mut fixture = TableTestFixture::new();
1980 fixture.setup_manifest_files().await;
1981
1982 let table_scan = fixture
1984 .table
1985 .scan()
1986 .select(["x", RESERVED_COL_NAME_FILE])
1987 .with_row_selection_enabled(true)
1988 .build()
1989 .unwrap();
1990
1991 let batch_stream = table_scan.to_arrow().await.unwrap();
1992 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1993
1994 assert_eq!(batches[0].num_columns(), 2);
1996
1997 let x_col = batches[0].column_by_name("x").unwrap();
1999 let x_arr = x_col.as_primitive::<arrow_array::types::Int64Type>();
2000 assert_eq!(x_arr.value(0), 1);
2001
2002 let file_col = batches[0].column_by_name(RESERVED_COL_NAME_FILE);
2004 assert!(
2005 file_col.is_some(),
2006 "_file column should be present in the batch"
2007 );
2008
2009 let file_col = file_col.unwrap();
2011 assert!(
2012 matches!(
2013 file_col.data_type(),
2014 arrow_schema::DataType::RunEndEncoded(_, _)
2015 ),
2016 "_file column should use RunEndEncoded type"
2017 );
2018
2019 let run_array = file_col
2021 .as_any()
2022 .downcast_ref::<RunArray<Int32Type>>()
2023 .expect("_file column should be a RunArray");
2024
2025 let values = run_array.values();
2026 let string_values = values.as_string::<i32>();
2027 assert_eq!(string_values.len(), 1, "Should have a single file path");
2028
2029 let file_path = string_values.value(0);
2030 assert!(
2031 file_path.ends_with(".parquet"),
2032 "File path should end with .parquet, got: {file_path}"
2033 );
2034 }
2035
2036 #[tokio::test]
2037 async fn test_select_file_column_position() {
2038 let mut fixture = TableTestFixture::new();
2039 fixture.setup_manifest_files().await;
2040
2041 let table_scan = fixture
2043 .table
2044 .scan()
2045 .select(["x", RESERVED_COL_NAME_FILE, "z"])
2046 .with_row_selection_enabled(true)
2047 .build()
2048 .unwrap();
2049
2050 let batch_stream = table_scan.to_arrow().await.unwrap();
2051 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
2052
2053 assert_eq!(batches[0].num_columns(), 3);
2054
2055 let schema = batches[0].schema();
2057 assert_eq!(schema.field(0).name(), "x");
2058 assert_eq!(schema.field(1).name(), RESERVED_COL_NAME_FILE);
2059 assert_eq!(schema.field(2).name(), "z");
2060
2061 assert!(batches[0].column_by_name("x").is_some());
2063 assert!(batches[0].column_by_name(RESERVED_COL_NAME_FILE).is_some());
2064 assert!(batches[0].column_by_name("z").is_some());
2065 }
2066
2067 #[tokio::test]
2068 async fn test_select_file_column_only() {
2069 let mut fixture = TableTestFixture::new();
2070 fixture.setup_manifest_files().await;
2071
2072 let table_scan = fixture
2074 .table
2075 .scan()
2076 .select([RESERVED_COL_NAME_FILE])
2077 .with_row_selection_enabled(true)
2078 .build()
2079 .unwrap();
2080
2081 let batch_stream = table_scan.to_arrow().await.unwrap();
2082 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
2083
2084 assert_eq!(batches[0].num_columns(), 1);
2086
2087 let schema = batches[0].schema();
2089 assert_eq!(schema.field(0).name(), RESERVED_COL_NAME_FILE);
2090
2091 let total_rows: usize = batches.iter().map(|b| b.num_rows()).sum();
2095 assert_eq!(total_rows, 2048);
2096 }
2097
2098 #[tokio::test]
2099 async fn test_file_column_with_multiple_files() {
2100 use std::collections::HashSet;
2101
2102 let mut fixture = TableTestFixture::new();
2103 fixture.setup_manifest_files().await;
2104
2105 let table_scan = fixture
2107 .table
2108 .scan()
2109 .select(["x", RESERVED_COL_NAME_FILE])
2110 .with_row_selection_enabled(true)
2111 .build()
2112 .unwrap();
2113
2114 let batch_stream = table_scan.to_arrow().await.unwrap();
2115 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
2116
2117 let mut file_paths = HashSet::new();
2119 for batch in &batches {
2120 let file_col = batch.column_by_name(RESERVED_COL_NAME_FILE).unwrap();
2121 let run_array = file_col
2122 .as_any()
2123 .downcast_ref::<RunArray<Int32Type>>()
2124 .expect("_file column should be a RunArray");
2125
2126 let values = run_array.values();
2127 let string_values = values.as_string::<i32>();
2128 for i in 0..string_values.len() {
2129 file_paths.insert(string_values.value(i).to_string());
2130 }
2131 }
2132
2133 assert!(!file_paths.is_empty(), "Should have at least one file path");
2135
2136 for path in &file_paths {
2138 assert!(
2139 path.ends_with(".parquet"),
2140 "All file paths should end with .parquet, got: {path}"
2141 );
2142 }
2143 }
2144
2145 #[tokio::test]
2146 async fn test_file_column_at_start() {
2147 let mut fixture = TableTestFixture::new();
2148 fixture.setup_manifest_files().await;
2149
2150 let table_scan = fixture
2152 .table
2153 .scan()
2154 .select([RESERVED_COL_NAME_FILE, "x", "y"])
2155 .with_row_selection_enabled(true)
2156 .build()
2157 .unwrap();
2158
2159 let batch_stream = table_scan.to_arrow().await.unwrap();
2160 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
2161
2162 assert_eq!(batches[0].num_columns(), 3);
2163
2164 let schema = batches[0].schema();
2166 assert_eq!(schema.field(0).name(), RESERVED_COL_NAME_FILE);
2167 assert_eq!(schema.field(1).name(), "x");
2168 assert_eq!(schema.field(2).name(), "y");
2169 }
2170
2171 #[tokio::test]
2172 async fn test_file_column_at_end() {
2173 let mut fixture = TableTestFixture::new();
2174 fixture.setup_manifest_files().await;
2175
2176 let table_scan = fixture
2178 .table
2179 .scan()
2180 .select(["x", "y", RESERVED_COL_NAME_FILE])
2181 .with_row_selection_enabled(true)
2182 .build()
2183 .unwrap();
2184
2185 let batch_stream = table_scan.to_arrow().await.unwrap();
2186 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
2187
2188 assert_eq!(batches[0].num_columns(), 3);
2189
2190 let schema = batches[0].schema();
2192 assert_eq!(schema.field(0).name(), "x");
2193 assert_eq!(schema.field(1).name(), "y");
2194 assert_eq!(schema.field(2).name(), RESERVED_COL_NAME_FILE);
2195 }
2196
2197 #[tokio::test]
2198 async fn test_select_with_repeated_column_names() {
2199 let mut fixture = TableTestFixture::new();
2200 fixture.setup_manifest_files().await;
2201
2202 let table_scan = fixture
2205 .table
2206 .scan()
2207 .select([
2208 "x",
2209 RESERVED_COL_NAME_FILE,
2210 "x", "y",
2212 RESERVED_COL_NAME_FILE, "y", ])
2215 .with_row_selection_enabled(true)
2216 .build()
2217 .unwrap();
2218
2219 let batch_stream = table_scan.to_arrow().await.unwrap();
2220 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
2221
2222 assert_eq!(
2224 batches[0].num_columns(),
2225 6,
2226 "Should have exactly 6 columns with duplicates"
2227 );
2228
2229 let schema = batches[0].schema();
2230
2231 assert_eq!(schema.field(0).name(), "x", "Column 0 should be x");
2233 assert_eq!(
2234 schema.field(1).name(),
2235 RESERVED_COL_NAME_FILE,
2236 "Column 1 should be _file"
2237 );
2238 assert_eq!(
2239 schema.field(2).name(),
2240 "x",
2241 "Column 2 should be x (duplicate)"
2242 );
2243 assert_eq!(schema.field(3).name(), "y", "Column 3 should be y");
2244 assert_eq!(
2245 schema.field(4).name(),
2246 RESERVED_COL_NAME_FILE,
2247 "Column 4 should be _file (duplicate)"
2248 );
2249 assert_eq!(
2250 schema.field(5).name(),
2251 "y",
2252 "Column 5 should be y (duplicate)"
2253 );
2254
2255 assert!(
2257 matches!(schema.field(0).data_type(), arrow_schema::DataType::Int64),
2258 "Column x should be Int64"
2259 );
2260 assert!(
2261 matches!(schema.field(2).data_type(), arrow_schema::DataType::Int64),
2262 "Column x (duplicate) should be Int64"
2263 );
2264 assert!(
2265 matches!(schema.field(3).data_type(), arrow_schema::DataType::Int64),
2266 "Column y should be Int64"
2267 );
2268 assert!(
2269 matches!(schema.field(5).data_type(), arrow_schema::DataType::Int64),
2270 "Column y (duplicate) should be Int64"
2271 );
2272 assert!(
2273 matches!(
2274 schema.field(1).data_type(),
2275 arrow_schema::DataType::RunEndEncoded(_, _)
2276 ),
2277 "_file column should use RunEndEncoded type"
2278 );
2279 assert!(
2280 matches!(
2281 schema.field(4).data_type(),
2282 arrow_schema::DataType::RunEndEncoded(_, _)
2283 ),
2284 "_file column (duplicate) should use RunEndEncoded type"
2285 );
2286 }
2287
2288 fn table_with_data_column(column_name: &str) -> Table {
2294 use crate::spec::{
2295 FormatVersion, MAIN_BRANCH, Operation, Snapshot, SnapshotReference, SnapshotRetention,
2296 SortOrder, Summary, TableMetadataBuilder, UnboundPartitionSpec,
2297 };
2298
2299 let schema = Schema::builder()
2300 .with_schema_id(0)
2301 .with_fields(vec![
2302 NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
2303 NestedField::required(2, column_name, Type::Primitive(PrimitiveType::Int)).into(),
2304 ])
2305 .build()
2306 .unwrap();
2307
2308 let snapshot = Snapshot::builder()
2309 .with_snapshot_id(1)
2310 .with_timestamp_ms(1)
2311 .with_sequence_number(0)
2312 .with_schema_id(0)
2313 .with_manifest_list("/snap-1.avro")
2314 .with_summary(Summary {
2315 operation: Operation::Append,
2316 additional_properties: HashMap::new(),
2317 })
2318 .build();
2319
2320 let metadata = TableMetadataBuilder::new(
2321 schema,
2322 UnboundPartitionSpec::builder().with_spec_id(0).build(),
2323 SortOrder::unsorted_order(),
2324 "s3://bucket/table".to_string(),
2325 FormatVersion::V2,
2326 HashMap::new(),
2327 )
2328 .unwrap()
2329 .add_snapshot(snapshot)
2330 .unwrap()
2331 .set_ref(MAIN_BRANCH, SnapshotReference {
2332 snapshot_id: 1,
2333 retention: SnapshotRetention::Branch {
2334 min_snapshots_to_keep: None,
2335 max_snapshot_age_ms: None,
2336 max_ref_age_ms: None,
2337 },
2338 })
2339 .unwrap()
2340 .build()
2341 .unwrap()
2342 .metadata;
2343
2344 Table::builder()
2345 .metadata(metadata)
2346 .identifier(TableIdent::from_strs(["db", "table1"]).unwrap())
2347 .file_io(FileIO::new_with_fs())
2348 .runtime(test_runtime())
2349 .build()
2350 .unwrap()
2351 }
2352
2353 #[test]
2357 fn test_scan_projects_data_column_named_like_delete_file_column() {
2358 for column_name in ["pos", "file_path"] {
2359 let table = table_with_data_column(column_name);
2360
2361 let table_scan = table
2364 .scan()
2365 .select([column_name])
2366 .build()
2367 .unwrap_or_else(|e| panic!("scan of data column `{column_name}` failed: {e}"));
2368
2369 assert_eq!(
2370 table_scan.plan_context.as_ref().unwrap().field_ids.as_ref(),
2371 &[2]
2372 );
2373
2374 let default_scan = table.scan().build().unwrap();
2377 assert_eq!(
2378 default_scan
2379 .plan_context
2380 .as_ref()
2381 .unwrap()
2382 .field_ids
2383 .as_ref(),
2384 &[1, 2]
2385 );
2386 }
2387 }
2388
2389 #[test]
2392 fn test_scan_rejects_unknown_column_named_like_delete_file_column() {
2393 let table = table_with_data_column("file_path");
2395
2396 let err = table
2397 .scan()
2398 .select(["pos"])
2399 .build()
2400 .expect_err("projecting an absent column should fail");
2401 assert_eq!(err.kind(), ErrorKind::DataInvalid);
2402 assert!(err.to_string().contains("not found"));
2403 }
2404
2405 #[tokio::test]
2406 async fn test_scan_deadlock() {
2407 let mut fixture = TableTestFixture::new();
2408 fixture.setup_deadlock_manifests().await;
2409
2410 let table_scan = fixture
2417 .table
2418 .scan()
2419 .with_concurrency_limit(1)
2420 .build()
2421 .unwrap();
2422
2423 let result = tokio::time::timeout(std::time::Duration::from_secs(5), async {
2426 table_scan
2427 .plan_files()
2428 .await
2429 .unwrap()
2430 .try_collect::<Vec<_>>()
2431 .await
2432 })
2433 .await;
2434
2435 assert!(result.is_ok(), "Scan timed out - deadlock detected");
2437 }
2438
2439 #[tokio::test]
2440 async fn test_select_with_spec_id_column() {
2441 let mut fixture = TableTestFixture::new();
2442 fixture.setup_manifest_files().await;
2443
2444 let table_scan = fixture
2446 .table
2447 .scan()
2448 .select(["x", RESERVED_COL_NAME_SPEC_ID, "z"])
2449 .with_row_selection_enabled(true)
2450 .build()
2451 .unwrap();
2452
2453 let batch_stream = table_scan.to_arrow().await.unwrap();
2454 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
2455
2456 assert_eq!(batches[0].num_columns(), 3);
2458
2459 let col1 = batches[0].column_by_name("x").unwrap();
2461 let int64_arr = col1.as_any().downcast_ref::<Int64Array>().unwrap();
2462 assert_eq!(int64_arr.value(0), 1);
2463
2464 let spec_id_col = batches[0].column_by_name(RESERVED_COL_NAME_SPEC_ID);
2466 assert!(
2467 spec_id_col.is_some(),
2468 "_spec_id column should be present in the batch"
2469 );
2470
2471 let spec_id_col = spec_id_col.unwrap();
2473 assert!(
2474 matches!(
2475 spec_id_col.data_type(),
2476 arrow_schema::DataType::RunEndEncoded(_, _)
2477 ),
2478 "_spec_id column should use RunEndEncoded type"
2479 );
2480
2481 let run_array = spec_id_col
2483 .as_any()
2484 .downcast_ref::<RunArray<Int32Type>>()
2485 .expect("_spec_id column should be a RunArray");
2486
2487 let values = run_array.values();
2488 let int_values = values.as_primitive::<Int32Type>();
2489 assert_eq!(int_values.len(), 1, "Should have a single _spec_id");
2490
2491 let spec_id = int_values.value(0);
2492 assert_eq!(spec_id, 0, "_spec_id should be 0, got: {spec_id}");
2493
2494 assert!(batches[0].column_by_name("z").is_some());
2496 }
2497
2498 #[tokio::test]
2499 async fn test_select_with_last_updated_sequence_number_column() {
2500 let mut fixture = TableTestFixture::new();
2505 fixture.setup_manifest_files().await;
2506
2507 let table_scan = fixture
2508 .table
2509 .scan()
2510 .select(["x", RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER])
2511 .with_row_selection_enabled(true)
2512 .build()
2513 .unwrap();
2514
2515 let batches: Vec<_> = table_scan
2516 .to_arrow()
2517 .await
2518 .unwrap()
2519 .try_collect()
2520 .await
2521 .unwrap();
2522
2523 assert_last_updated_seq_all(&batches, None);
2525 }
2526
2527 #[tokio::test]
2528 async fn test_select_with_last_updated_sequence_number_column_v3() {
2529 let mut fixture = TableTestFixture::new();
2535 fixture.setup_v3_manifest_files().await;
2536
2537 let expected_seq = fixture
2540 .table
2541 .metadata()
2542 .current_snapshot()
2543 .unwrap()
2544 .sequence_number();
2545
2546 let table_scan = fixture
2547 .table
2548 .scan()
2549 .select(["x", RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER])
2550 .with_row_selection_enabled(true)
2551 .build()
2552 .unwrap();
2553
2554 let batches: Vec<_> = table_scan
2555 .to_arrow()
2556 .await
2557 .unwrap()
2558 .try_collect()
2559 .await
2560 .unwrap();
2561
2562 assert_last_updated_seq_all(&batches, Some(expected_seq));
2563 }
2564
2565 #[tokio::test]
2566 async fn test_select_with_spec_id_column_from_unpartitioned_table() {
2567 let mut fixture = TableTestFixture::new_unpartitioned();
2568 fixture.setup_unpartitioned_manifest_files().await;
2569
2570 let table_scan = fixture
2572 .table
2573 .scan()
2574 .select(["x", RESERVED_COL_NAME_SPEC_ID])
2575 .with_row_selection_enabled(true)
2576 .build()
2577 .unwrap();
2578
2579 let batch_stream = table_scan.to_arrow().await.unwrap();
2580 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
2581
2582 assert_eq!(batches[0].num_columns(), 2);
2584
2585 let spec_id_col = batches[0].column_by_name(RESERVED_COL_NAME_SPEC_ID);
2587 assert!(
2588 spec_id_col.is_some(),
2589 "_spec_id column should be present in the batch"
2590 );
2591
2592 let spec_id_col = spec_id_col.unwrap();
2594 assert!(
2595 matches!(
2596 spec_id_col.data_type(),
2597 arrow_schema::DataType::RunEndEncoded(_, _)
2598 ),
2599 "_spec_id column should use RunEndEncoded type"
2600 );
2601
2602 let run_array = spec_id_col
2604 .as_any()
2605 .downcast_ref::<RunArray<Int32Type>>()
2606 .expect("_spec_id column should be a RunArray");
2607
2608 let values = run_array.values();
2609 let int_values = values.as_primitive::<Int32Type>();
2610 assert_eq!(int_values.len(), 1, "Should have a single _spec_id");
2611
2612 let spec_id = int_values.value(0);
2613 assert_eq!(spec_id, 0, "_spec_id should be 0, got: {spec_id}");
2614 }
2615
2616 #[tokio::test]
2617 async fn test_select_with_spec_id_column_with_partition_evolution() {
2618 let mut fixture = TableTestFixture::new_with_partition_evolution();
2619 fixture
2620 .setup_manifest_files_with_partition_evolution()
2621 .await;
2622
2623 let table_scan = fixture
2625 .table
2626 .scan()
2627 .select(["x", RESERVED_COL_NAME_SPEC_ID, "z"])
2628 .with_row_selection_enabled(true)
2629 .build()
2630 .unwrap();
2631
2632 let batch_stream = table_scan.to_arrow().await.unwrap();
2633 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
2634
2635 let col1 = batches[0].column_by_name("x").unwrap();
2637 let int64_arr = col1.as_any().downcast_ref::<Int64Array>().unwrap();
2638 assert_eq!(int64_arr.value(0), 1);
2639
2640 let spec_id_col = batches[0].column_by_name(RESERVED_COL_NAME_SPEC_ID);
2642 assert!(
2643 spec_id_col.is_some(),
2644 "_spec_id column should be present in the batch"
2645 );
2646
2647 let spec_id_col = spec_id_col.unwrap();
2649 assert!(
2650 matches!(
2651 spec_id_col.data_type(),
2652 arrow_schema::DataType::RunEndEncoded(_, _)
2653 ),
2654 "_spec_id column should use RunEndEncoded type"
2655 );
2656
2657 let run_array = spec_id_col
2659 .as_any()
2660 .downcast_ref::<RunArray<Int32Type>>()
2661 .expect("_spec_id column should be a RunArray");
2662
2663 let values = run_array.values();
2664 let int_values = values.as_primitive::<Int32Type>();
2665 assert_eq!(int_values.len(), 1, "Should have a single _spec_id");
2666
2667 let spec_id = int_values.value(0);
2668 assert_eq!(spec_id, 2, "_spec_id should be 2, got: {spec_id}");
2669 }
2670
2671 #[tokio::test]
2672 async fn test_select_with_pos_and_file_columns() {
2673 use arrow_array::cast::AsArray;
2674
2675 let mut fixture = TableTestFixture::new();
2676 fixture.setup_manifest_files().await;
2677
2678 let table_scan = fixture
2680 .table
2681 .scan()
2682 .select(["x", RESERVED_COL_NAME_POS, RESERVED_COL_NAME_FILE])
2683 .with_row_selection_enabled(true)
2684 .build()
2685 .unwrap();
2686
2687 let batch_stream = table_scan.to_arrow().await.unwrap();
2688 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
2689 assert_eq!(batches.len(), 2);
2690
2691 for batch in batches.iter() {
2693 assert_eq!(batch.num_columns(), 3);
2695
2696 let x_col = batch.column_by_name("x").unwrap();
2698 let x_arr = x_col.as_primitive::<arrow_array::types::Int64Type>();
2699 assert_eq!(x_arr.value(0), 1);
2700
2701 let pos_col = batch.column(1);
2703 let pos_array: &Int64Array = pos_col
2704 .as_any()
2705 .downcast_ref::<Int64Array>()
2706 .expect("_pos column should be a Int64Array");
2707 assert_eq!(*pos_array, Int64Array::from_iter_values(0i64..1024));
2708
2709 let file_col = batch.column_by_name(RESERVED_COL_NAME_FILE);
2711 assert!(
2712 file_col.is_some(),
2713 "_file column should be present in the batch"
2714 );
2715 }
2716 }
2717
2718 #[tokio::test]
2719 async fn test_pos_column_at_start_with_filters() {
2720 let mut fixture = TableTestFixture::new();
2721 fixture.setup_manifest_files().await;
2722
2723 let predicate = Reference::new("y")
2725 .greater_than(Datum::long(4i64))
2726 .and(Reference::new("y").less_than_or_equal_to(Datum::long(5i64)));
2727 let table_scan = fixture
2729 .table
2730 .scan()
2731 .select([RESERVED_COL_NAME_POS, "x", "y"])
2732 .with_filter(predicate)
2733 .with_row_selection_enabled(true)
2734 .build()
2735 .unwrap();
2736
2737 let batch_stream = table_scan.to_arrow().await.unwrap();
2738 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
2739 assert_eq!(batches.len(), 2);
2740
2741 for batch in batches.iter() {
2743 assert_eq!(batch.num_columns(), 3);
2744 assert_eq!(batch.num_rows(), 12);
2745
2746 let schema = batch.schema();
2748 assert_eq!(schema.field(0).name(), RESERVED_COL_NAME_POS);
2749 assert_eq!(schema.field(1).name(), "x");
2750 assert_eq!(schema.field(2).name(), "y");
2751
2752 let pos_col = batch.column(0);
2753 let pos_array: &Int64Array = pos_col
2754 .as_any()
2755 .downcast_ref::<Int64Array>()
2756 .expect("_pos column should be a Int64Array");
2757 assert_eq!(*pos_array, Int64Array::from_iter_values(1012i64..1024));
2758 }
2759 }
2760
2761 #[tokio::test]
2762 async fn test_repeated_pos_column_with_filter() {
2763 let mut fixture = TableTestFixture::new();
2764 fixture.setup_manifest_files().await;
2765
2766 let predicate = Reference::new("a").not_starts_with(Datum::string("Apa"));
2768 let table_scan = fixture
2770 .table
2771 .scan()
2772 .select([RESERVED_COL_NAME_POS, "a", RESERVED_COL_NAME_POS, "x"])
2773 .with_row_selection_enabled(true)
2774 .with_filter(predicate)
2775 .build()
2776 .unwrap();
2777
2778 let batch_stream = table_scan.to_arrow().await.unwrap();
2779
2780 let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
2781 assert_eq!(batches.len(), 2);
2782
2783 for batch in batches.iter() {
2785 assert_eq!(batch.num_rows(), 512);
2786
2787 let pos_col = batch
2789 .column_by_name("_pos")
2790 .expect("_pos column should be present in the batch");
2791 let pos_array: &Int64Array = pos_col
2792 .as_any()
2793 .downcast_ref::<Int64Array>()
2794 .expect("_pos column should be a Int64Array");
2795 assert_eq!(*pos_array, Int64Array::from_iter_values(512i64..1024));
2796
2797 let pos_col = batch.column(2);
2799 let pos_array: &Int64Array = pos_col
2800 .as_any()
2801 .downcast_ref::<Int64Array>()
2802 .expect("_pos column should be a Int64Array");
2803 assert_eq!(*pos_array, Int64Array::from_iter_values(512i64..1024));
2804 }
2805 }
2806
2807 #[tokio::test]
2811 async fn test_pos_across_row_groups_via_table_scan() {
2812 let mut fixture = TableTestFixture::new();
2813 fixture.setup_multi_row_group_manifest(&[]).await;
2814
2815 let tasks: Vec<_> = fixture
2818 .table
2819 .scan()
2820 .select(["x", RESERVED_COL_NAME_POS])
2821 .build()
2822 .unwrap()
2823 .plan_files()
2824 .await
2825 .unwrap()
2826 .try_collect()
2827 .await
2828 .unwrap();
2829 assert_eq!(tasks.len(), 1, "expected a single FileScanTask");
2830 let task = &tasks[0];
2831 assert!(
2832 task.project_field_ids().contains(&RESERVED_FIELD_ID_POS),
2833 "_pos field id must be projected into the FileScanTask"
2834 );
2835 assert_eq!(task.start(), 0, "TableScan should plan whole-file tasks");
2836 assert_eq!(task.length(), task.file_size_in_bytes());
2837 assert!(task.deletes().is_empty());
2838
2839 let batches: Vec<_> = fixture
2841 .table
2842 .scan()
2843 .select(["x", RESERVED_COL_NAME_POS])
2844 .build()
2845 .unwrap()
2846 .to_arrow()
2847 .await
2848 .unwrap()
2849 .try_collect()
2850 .await
2851 .unwrap();
2852
2853 let pos: Vec<i64> = batches
2854 .iter()
2855 .flat_map(|b| {
2856 b.column_by_name(RESERVED_COL_NAME_POS)
2857 .expect("_pos column should be present")
2858 .as_any()
2859 .downcast_ref::<Int64Array>()
2860 .expect("_pos column should be a Int64Array")
2861 .values()
2862 .to_vec()
2863 })
2864 .collect();
2865 assert_eq!(pos, (0..300).collect::<Vec<i64>>());
2866
2867 let x: Vec<i64> = batches
2869 .iter()
2870 .flat_map(|b| {
2871 b.column_by_name("x")
2872 .unwrap()
2873 .as_primitive::<arrow_array::types::Int64Type>()
2874 .values()
2875 .to_vec()
2876 })
2877 .collect();
2878 assert_eq!(x, (1000..1300).collect::<Vec<i64>>());
2879 }
2880
2881 #[tokio::test]
2884 async fn test_pos_with_positional_deletes_via_table_scan() {
2885 let mut fixture = TableTestFixture::new();
2886 fixture.setup_multi_row_group_manifest(&[150, 299]).await;
2888
2889 let tasks: Vec<_> = fixture
2891 .table
2892 .scan()
2893 .select(["x", RESERVED_COL_NAME_POS])
2894 .build()
2895 .unwrap()
2896 .plan_files()
2897 .await
2898 .unwrap()
2899 .try_collect()
2900 .await
2901 .unwrap();
2902 assert_eq!(tasks.len(), 1);
2903 assert_eq!(
2904 tasks[0].deletes().len(),
2905 1,
2906 "positional delete file should be planned into the task"
2907 );
2908 assert_eq!(
2909 tasks[0].deletes()[0].file_type(),
2910 DataContentType::PositionDeletes
2911 );
2912
2913 let batches: Vec<_> = fixture
2915 .table
2916 .scan()
2917 .select(["x", RESERVED_COL_NAME_POS])
2918 .build()
2919 .unwrap()
2920 .to_arrow()
2921 .await
2922 .unwrap()
2923 .try_collect()
2924 .await
2925 .unwrap();
2926
2927 let pos: Vec<i64> = batches
2928 .iter()
2929 .flat_map(|b| {
2930 b.column_by_name(RESERVED_COL_NAME_POS)
2931 .expect("_pos column should be present")
2932 .as_any()
2933 .downcast_ref::<Int64Array>()
2934 .expect("_pos column should be a Int64Array")
2935 .values()
2936 .to_vec()
2937 })
2938 .collect();
2939
2940 let total: usize = batches.iter().map(|b| b.num_rows()).sum();
2941 assert_eq!(
2942 total, 298,
2943 "two rows should be removed by positional deletes"
2944 );
2945 assert!(!pos.contains(&150) && !pos.contains(&299), "got {pos:?}");
2946 let expected: Vec<i64> = (0..150).chain(151..299).collect();
2947 assert_eq!(pos, expected);
2948 }
2949
2950 #[tokio::test]
2958 async fn test_pos_reads_only_middle_row_group_via_filter() {
2959 let mut fixture = TableTestFixture::new();
2960 fixture.setup_multi_row_group_manifest(&[]).await;
2961
2962 let predicate = Reference::new("y")
2964 .greater_than_or_equal_to(Datum::long(1100))
2965 .and(Reference::new("y").less_than(Datum::long(1200)));
2966
2967 let batches: Vec<_> = fixture
2968 .table
2969 .scan()
2970 .select(["y", RESERVED_COL_NAME_POS])
2971 .with_filter(predicate)
2972 .with_row_group_filtering_enabled(true)
2973 .build()
2974 .unwrap()
2975 .to_arrow()
2976 .await
2977 .unwrap()
2978 .try_collect()
2979 .await
2980 .unwrap();
2981
2982 let total: usize = batches.iter().map(|b| b.num_rows()).sum();
2983 assert_eq!(total, 100, "only the middle row group should be read");
2984
2985 let pos: Vec<i64> = batches
2986 .iter()
2987 .flat_map(|b| {
2988 b.column_by_name(RESERVED_COL_NAME_POS)
2989 .expect("_pos column should be present")
2990 .as_any()
2991 .downcast_ref::<Int64Array>()
2992 .expect("_pos column should be a Int64Array")
2993 .values()
2994 .to_vec()
2995 })
2996 .collect();
2997 assert_eq!(
2998 pos,
2999 (100..200).collect::<Vec<i64>>(),
3000 "_pos must be file-absolute for the middle row group"
3001 );
3002
3003 let y: Vec<i64> = batches
3005 .iter()
3006 .flat_map(|b| {
3007 b.column_by_name("y")
3008 .unwrap()
3009 .as_primitive::<arrow_array::types::Int64Type>()
3010 .values()
3011 .to_vec()
3012 })
3013 .collect();
3014 assert_eq!(y, (1100..1200).collect::<Vec<i64>>());
3015 }
3016}