This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-rust.git
The following commit(s) were added to refs/heads/main by this push:
new b7d7ee2c perf(table): push row-range pruning into manifest scan (#590)
b7d7ee2c is described below
commit b7d7ee2ce2c64499ffdb3873f9dd50193e7837a6
Author: XiaoHongbo <[email protected]>
AuthorDate: Thu Jul 23 20:58:26 2026 +0800
perf(table): push row-range pruning into manifest scan (#590)
---
crates/paimon/src/spec/data_file.rs | 20 +-
.../src/table/btree_global_index_build_builder.rs | 126 +++
crates/paimon/src/table/data_evolution_reader.rs | 259 +++++-
crates/paimon/src/table/read_builder.rs | 43 +-
crates/paimon/src/table/scan_trace.rs | 6 +-
crates/paimon/src/table/table_scan.rs | 975 ++++++++++++++++++---
6 files changed, 1295 insertions(+), 134 deletions(-)
diff --git a/crates/paimon/src/spec/data_file.rs
b/crates/paimon/src/spec/data_file.rs
index 43a23b43..d2b901bd 100644
--- a/crates/paimon/src/spec/data_file.rs
+++ b/crates/paimon/src/spec/data_file.rs
@@ -257,7 +257,12 @@ impl DataFileMeta {
/// Returns the row ID range `[first_row_id, first_row_id + row_count -
1]` if `first_row_id` is set.
pub fn row_id_range(&self) -> Option<(i64, i64)> {
- self.first_row_id.map(|fid| (fid, fid + self.row_count - 1))
+ let from = self.first_row_id?;
+ if self.row_count <= 0 {
+ return None;
+ }
+ let to = from.checked_add(self.row_count - 1)?;
+ Some((from, to))
}
/// Serialize as a `DataFileMeta.SCHEMA` (version 8) BinaryRow, raw data
without the
@@ -584,6 +589,19 @@ mod tests {
}
}
+ #[test]
+ fn row_id_range_rejects_non_positive_count_and_overflow() {
+ let mut file = data_file("data.parquet");
+ assert_eq!(file.row_id_range(), Some((100, 106)));
+
+ file.row_count = 0;
+ assert_eq!(file.row_id_range(), None);
+
+ file.row_count = 2;
+ file.first_row_id = Some(i64::MAX);
+ assert_eq!(file.row_id_range(), None);
+ }
+
#[test]
fn value_stats_for_field_resolves_dense_columns_and_schema() {
let fields = vec![int_field(0, "id"), int_field(1, "v")];
diff --git a/crates/paimon/src/table/btree_global_index_build_builder.rs
b/crates/paimon/src/table/btree_global_index_build_builder.rs
index da0b6081..f9ce468e 100644
--- a/crates/paimon/src/table/btree_global_index_build_builder.rs
+++ b/crates/paimon/src/table/btree_global_index_build_builder.rs
@@ -1369,6 +1369,132 @@ mod tests {
);
}
+ #[tokio::test]
+ async fn test_global_index_prunes_during_manifest_read() {
+ for (search_mode, expected_manifest_pruned, expected_entries_read) in
+ [("fast", 1, 1), ("full", 1, 1)]
+ {
+ let table_path =
format!("memory:/test_global_index_manifest_pruning_{search_mode}");
+ let mut options = table_options("2");
+ options.insert(
+ "global-index.search-mode".to_string(),
+ search_mode.to_string(),
+ );
+ let table = test_table_with_path(&table_path, options);
+ setup_dirs(&table).await;
+
+ for (user, ids, names) in [
+ ("writer-1", vec![1, 2], vec!["alice", "bob"]),
+ ("writer-2", vec![3, 4], vec!["carol", "dave"]),
+ ] {
+ let mut table_write = TableWrite::new(&table,
user.to_string()).unwrap();
+ table_write
+ .write_arrow_batch(&data_batch(ids, names))
+ .await
+ .unwrap();
+ TableCommit::new(table.clone(), user.to_string())
+ .commit(table_write.prepare_commit().await.unwrap())
+ .await
+ .unwrap();
+ }
+
+ table
+ .new_btree_global_index_build_builder()
+ .with_index_column("name")
+ .execute()
+ .await
+ .unwrap();
+
+ let predicate = PredicateBuilder::new(table.schema().fields())
+ .equal("name", crate::spec::Datum::String("alice".to_string()))
+ .unwrap();
+ let mut read_builder = table.new_read_builder();
+ read_builder.with_filter(predicate);
+ let (plan, trace) =
read_builder.new_scan().plan_with_trace().await.unwrap();
+
+ assert_eq!(
+ plan.splits()
+ .iter()
+ .flat_map(|split| split.data_files())
+ .count(),
+ 1
+ );
+ assert_eq!(
+ trace.manifest_files_pruned_by_row_ranges,
+ expected_manifest_pruned
+ );
+ assert_eq!(trace.manifest_entries_read, expected_entries_read);
+ assert_eq!(trace.manifest_entries_pruned_by_row_ranges, 0);
+ assert_eq!(
+ trace.manifest_entries_after_manifest_filters,
+ expected_entries_read
+ );
+ }
+ }
+
+ #[tokio::test]
+ async fn test_detail_mode_defers_manifest_pruning_for_unindexed_ranges() {
+ let table_path = "memory:/test_detail_manifest_pruning";
+ let mut options = table_options("2");
+ options.insert("global-index.search-mode".to_string(),
"detail".to_string());
+ let table = test_table_with_path(table_path, options);
+ setup_dirs(&table).await;
+
+ let mut table_write = TableWrite::new(&table,
"writer-1".to_string()).unwrap();
+ table_write
+ .write_arrow_batch(&data_batch(vec![1, 2], vec!["alice", "bob"]))
+ .await
+ .unwrap();
+ TableCommit::new(table.clone(), "writer-1".to_string())
+ .commit(table_write.prepare_commit().await.unwrap())
+ .await
+ .unwrap();
+ table
+ .new_btree_global_index_build_builder()
+ .with_index_column("name")
+ .execute()
+ .await
+ .unwrap();
+
+ let mut table_write = TableWrite::new(&table,
"writer-2".to_string()).unwrap();
+ table_write
+ .write_arrow_batch(&data_batch(vec![3, 4], vec!["alice", "dave"]))
+ .await
+ .unwrap();
+ TableCommit::new(table.clone(), "writer-2".to_string())
+ .commit(table_write.prepare_commit().await.unwrap())
+ .await
+ .unwrap();
+
+ let predicate = PredicateBuilder::new(table.schema().fields())
+ .equal("name", crate::spec::Datum::String("alice".to_string()))
+ .unwrap();
+ let mut read_builder = table.new_read_builder();
+ read_builder.with_filter(predicate);
+ let (plan, trace) =
read_builder.new_scan().plan_with_trace().await.unwrap();
+ let planned_ranges = merge_row_ranges(
+ plan.splits()
+ .iter()
+ .flat_map(|split| split.row_ranges().unwrap_or_default())
+ .cloned()
+ .collect(),
+ );
+
+ assert_eq!(
+ plan.splits()
+ .iter()
+ .flat_map(|split| split.data_files())
+ .count(),
+ 2
+ );
+ assert_eq!(
+ planned_ranges,
+ vec![RowRange::new(0, 0), RowRange::new(2, 3)]
+ );
+ assert_eq!(trace.manifest_files_pruned_by_row_ranges, 0);
+ assert_eq!(trace.manifest_entries_read, 2);
+ }
+
#[tokio::test]
async fn test_execute_writes_bitmap_index_manifest_and_java_file() {
let table_path = "memory:/test_bitmap_global_index_builder_e2e";
diff --git a/crates/paimon/src/table/data_evolution_reader.rs
b/crates/paimon/src/table/data_evolution_reader.rs
index 6a792d17..648725df 100644
--- a/crates/paimon/src/table/data_evolution_reader.rs
+++ b/crates/paimon/src/table/data_evolution_reader.rs
@@ -2379,10 +2379,12 @@ mod tests {
use crate::io::FileIOBuilder;
use crate::spec::stats::BinaryTableStats;
use crate::spec::{
- ArrayType, BinaryRow, BlobType, Datum, FloatType, IntType,
PredicateBuilder, Schema,
- TableSchema, VectorType,
+ ArrayType, BinaryRow, BinaryRowBuilder, BlobType, Datum, FloatType,
IntType,
+ PredicateBuilder, Schema, TableSchema, VectorType,
+ };
+ use crate::table::{
+ CommitMessage, DataSplitBuilder, DeletionFile, Table, TableCommit,
TableRead,
};
- use crate::table::{DataSplitBuilder, DeletionFile, Table, TableRead};
use arrow_array::{
Array, BinaryArray, FixedSizeListArray, Float32Array, Int32Array,
Int64Array, ListArray,
RecordBatch,
@@ -5589,6 +5591,257 @@ mod tests {
);
}
+ #[tokio::test]
+ async fn test_scan_and_read_retains_complete_rolled_dedicated_group() {
+ let tempdir = tempdir().unwrap();
+ let table_path = local_file_path(tempdir.path());
+ let bucket_dir = tempdir.path().join("bucket-0");
+ fs::create_dir_all(&bucket_dir).unwrap();
+ fs::create_dir_all(tempdir.path().join("snapshot")).unwrap();
+ fs::create_dir_all(tempdir.path().join("manifest")).unwrap();
+
+ let normal_path = bucket_dir.join("data.parquet");
+ write_int_parquet_file(&normal_path, vec![("id", vec![1, 2, 3, 4, 5,
6])], None);
+
+ let vector_rows = [
+ [Some(vec![1.0, 1.0]), Some(vec![2.0, 2.0])],
+ [Some(vec![3.0, 3.0]), Some(vec![4.0, 4.0])],
+ [Some(vec![5.0, 5.0]), Some(vec![6.0, 6.0])],
+ ];
+ let blob_rows: [[Option<&[u8]>; 2]; 3] = [
+ [Some(b"b1"), Some(b"b2")],
+ [Some(b"b3"), Some(b"b4")],
+ [Some(b"b5"), Some(b"b6")],
+ ];
+
+ let mut files = vec![data_file_meta_with_path(
+ "data.parquet",
+ 0,
+ 6,
+ 1,
+ normal_path.metadata().unwrap().len() as i64,
+ Some(vec!["id"]),
+ )];
+ for segment in 0..3 {
+ let vector_name = format!("emb-{}.vector.parquet", segment + 1);
+ let vector_path = bucket_dir.join(&vector_name);
+ write_fixed_size_list_parquet(&vector_path, "embedding", 2,
&vector_rows[segment]);
+ files.push(data_file_meta_with_path(
+ &vector_name,
+ (segment * 2) as i64,
+ 2,
+ 1,
+ vector_path.metadata().unwrap().len() as i64,
+ Some(vec!["embedding"]),
+ ));
+
+ let blob_name = format!("payload-{}.blob", segment + 1);
+ let blob_path = bucket_dir.join(&blob_name);
+ write_blob_file(&blob_path, &blob_rows[segment]);
+ files.push(data_file_meta_with_path(
+ &blob_name,
+ (segment * 2) as i64,
+ 2,
+ 1,
+ blob_path.metadata().unwrap().len() as i64,
+ Some(vec!["payload"]),
+ ));
+ }
+
+ let file_io = FileIOBuilder::new("file").build().unwrap();
+ let table_schema = TableSchema::new(
+ 0,
+ &Schema::builder()
+ .column("id", DataType::Int(IntType::new()))
+ .column("embedding", vector_float_type(2))
+ .column("payload", DataType::Blob(BlobType::new()))
+ .option("data-evolution.enabled", "true")
+ .build()
+ .unwrap(),
+ );
+ let table = Table::new(
+ file_io,
+ Identifier::new("default", "rolled_dedicated_scan_t"),
+ table_path,
+ table_schema,
+ None,
+ );
+ TableCommit::new(table.clone(), "rolled-dedicated-test".to_string())
+ .commit(vec![CommitMessage::new(
+ BinaryRowBuilder::new(0).build_serialized(),
+ 0,
+ files,
+ )])
+ .await
+ .unwrap();
+
+ let mut builder = table.new_read_builder();
+ builder.with_row_ranges(vec![RowRange::new(0, 0)]);
+ let (plan, trace) =
builder.new_scan().plan_with_trace().await.unwrap();
+ let mut planned_files = plan
+ .splits()
+ .iter()
+ .flat_map(|split| split.data_files())
+ .map(|file| file.file_name.clone())
+ .collect::<Vec<_>>();
+ planned_files.sort();
+ let expected_planned_files = vec![
+ "data.parquet".to_string(),
+ "emb-1.vector.parquet".to_string(),
+ "emb-2.vector.parquet".to_string(),
+ "emb-3.vector.parquet".to_string(),
+ "payload-1.blob".to_string(),
+ "payload-2.blob".to_string(),
+ "payload-3.blob".to_string(),
+ ];
+ assert_eq!(planned_files, expected_planned_files);
+ assert_eq!(trace.manifest_entries_pruned_by_row_ranges, 0);
+ assert_eq!(trace.final_files, 7);
+
+ let batches = builder
+ .new_read()
+ .unwrap()
+ .to_arrow(plan.splits())
+ .unwrap()
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap();
+ assert_eq!(collect_int_values(&batches, "id"), vec![1]);
+ assert_fixed_size_list(&batches, "embedding", 2, &[Some(vec![1.0,
1.0])]);
+ assert_eq!(
+ collect_binary_values(&batches, "payload"),
+ vec![Some(b"b1".to_vec())]
+ );
+
+ for field in ["embedding", "payload"] {
+ let predicate = PredicateBuilder::new(table.schema().fields())
+ .is_null(field)
+ .unwrap();
+ let mut predicate_builder = table.new_read_builder();
+ predicate_builder.with_projection(&["id"]).unwrap();
+ predicate_builder.with_filter(predicate);
+ predicate_builder.with_row_ranges(vec![RowRange::new(0, 0)]);
+ let predicate_plan =
predicate_builder.new_scan().plan().await.unwrap();
+ let predicate_batches = predicate_builder
+ .new_read()
+ .unwrap()
+ .to_arrow(predicate_plan.splits())
+ .unwrap()
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap();
+ assert!(
+ collect_int_values(&predicate_batches, "id").is_empty(),
+ "predicate-only dedicated field '{field}' must not be treated
as missing/null"
+ );
+ }
+
+ let snapshot = table
+ .snapshot_manager()
+ .get_latest_snapshot()
+ .await
+ .unwrap()
+ .unwrap();
+ let delta_plan = builder
+ .new_scan()
+ .plan_snapshot_delta(&snapshot)
+ .await
+ .unwrap();
+ let mut delta_files = delta_plan
+ .splits()
+ .iter()
+ .flat_map(|split| split.data_files())
+ .map(|file| file.file_name.clone())
+ .collect::<Vec<_>>();
+ delta_files.sort();
+ assert_eq!(delta_files, expected_planned_files);
+ }
+
+ #[tokio::test]
+ async fn test_scan_and_read_rejects_selected_gap_after_dedicated_pruning()
{
+ let tempdir = tempdir().unwrap();
+ let table_path = local_file_path(tempdir.path());
+ let bucket_dir = tempdir.path().join("bucket-0");
+ fs::create_dir_all(&bucket_dir).unwrap();
+ fs::create_dir_all(tempdir.path().join("snapshot")).unwrap();
+ fs::create_dir_all(tempdir.path().join("manifest")).unwrap();
+
+ let normal_path = bucket_dir.join("data.parquet");
+ write_int_parquet_file(&normal_path, vec![("id", vec![1, 2, 3, 4, 5,
6])], None);
+
+ let files = vec![
+ data_file_meta_with_path(
+ "data.parquet",
+ 0,
+ 6,
+ 1,
+ normal_path.metadata().unwrap().len() as i64,
+ Some(vec!["id"]),
+ ),
+ data_file_meta_with_path(
+ "emb-left.vector.parquet",
+ 0,
+ 2,
+ 1,
+ 1,
+ Some(vec!["embedding"]),
+ ),
+ data_file_meta_with_path(
+ "emb-right.vector.parquet",
+ 4,
+ 2,
+ 1,
+ 1,
+ Some(vec!["embedding"]),
+ ),
+ data_file_meta_with_path("payload-left.blob", 0, 2, 1, 1,
Some(vec!["payload"])),
+ data_file_meta_with_path("payload-right.blob", 4, 2, 1, 1,
Some(vec!["payload"])),
+ ];
+
+ let file_io = FileIOBuilder::new("file").build().unwrap();
+ let table_schema = TableSchema::new(
+ 0,
+ &Schema::builder()
+ .column("id", DataType::Int(IntType::new()))
+ .column("embedding", vector_float_type(2))
+ .column("payload", DataType::Blob(BlobType::new()))
+ .option("data-evolution.enabled", "true")
+ .build()
+ .unwrap(),
+ );
+ let table = Table::new(
+ file_io,
+ Identifier::new("default", "missing_selected_dedicated_t"),
+ table_path,
+ table_schema,
+ None,
+ );
+ TableCommit::new(table.clone(),
"missing-selected-dedicated-test".to_string())
+ .commit(vec![CommitMessage::new(
+ BinaryRowBuilder::new(0).build_serialized(),
+ 0,
+ files,
+ )])
+ .await
+ .unwrap();
+
+ for (field, provider_kind) in [("embedding", "Vector"), ("payload",
"Blob")] {
+ let mut builder = table.new_read_builder();
+ builder.with_projection(&["id", field]).unwrap();
+ builder.with_row_ranges(vec![RowRange::new(2, 2)]);
+ let plan = builder.new_scan().plan().await.unwrap();
+
+ let mut stream =
builder.new_read().unwrap().to_arrow(plan.splits()).unwrap();
+ let error = stream.try_next().await.unwrap_err();
+ assert!(
+ matches!(error, Error::DataInvalid { ref message, .. }
+ if message.contains(provider_kind)
+ && message.contains("cover effective selected row
ranges")),
+ "expected missing {provider_kind} coverage error, got
{error:?}"
+ );
+ }
+ }
+
/// (6) Row-range mismatch: normal file row_count=3 but `.vector.parquet`
row_count=2
/// must surface as DataInvalid.
#[tokio::test]
diff --git a/crates/paimon/src/table/read_builder.rs
b/crates/paimon/src/table/read_builder.rs
index 4870b8f8..de432a5c 100644
--- a/crates/paimon/src/table/read_builder.rs
+++ b/crates/paimon/src/table/read_builder.rs
@@ -443,7 +443,11 @@ impl<'a> PaimonReadBuilder<'a> {
self.limit,
self.row_ranges.clone(),
)
- .with_projected_read_field_ids(projected_read_field_ids(&read_type))
+
.with_projected_read_field_ids(projected_read_field_ids_with_predicates(
+ &read_type,
+ &self.filter.data_predicates,
+ self.table.schema().fields(),
+ ))
}
/// Create a table read for consuming splits (e.g. from a scan plan).
@@ -621,6 +625,28 @@ fn projected_read_field_ids(read_type:
&Option<Vec<DataField>>) -> Option<HashSe
.map(|fields| projected_read_field_ids_from_fields(fields))
}
+fn projected_read_field_ids_with_predicates(
+ read_type: &Option<Vec<DataField>>,
+ predicates: &[Predicate],
+ table_fields: &[DataField],
+) -> Option<HashSet<i32>> {
+ let mut field_ids = projected_read_field_ids(read_type)?;
+ let mut predicate_indices = Vec::new();
+ for predicate in predicates {
+ crate::arrow::residual::collect_predicate_field_indices(predicate,
&mut predicate_indices);
+ }
+ for index in predicate_indices {
+ let Some(field) = table_fields.get(index) else {
+ // A malformed predicate must not make scan planning discard files.
+ return None;
+ };
+ if !is_system_projection_field(field.id()) {
+ field_ids.insert(field.id());
+ }
+ }
+ Some(field_ids)
+}
+
pub(super) fn is_system_projection_field(field_id: i32) -> bool {
matches!(
field_id,
@@ -840,6 +866,21 @@ mod tests {
);
}
+ #[test]
+ fn test_projected_read_field_ids_include_predicate_only_fields() {
+ let fields = vec![
+ DataField::new(1, "id".to_string(), DataType::Int(IntType::new())),
+ DataField::new(2, "payload".to_string(),
DataType::Int(IntType::new())),
+ ];
+ let read_type = Some(vec![fields[0].clone()]);
+ let predicate =
PredicateBuilder::new(&fields).is_null("payload").unwrap();
+
+ assert_eq!(
+ super::projected_read_field_ids_with_predicates(&read_type,
&[predicate], &fields,),
+ Some(HashSet::from([1, 2]))
+ );
+ }
+
#[test]
fn test_with_projection_validates_unknown_projection() {
// A column that cannot match under any case sensitivity is an obvious
diff --git a/crates/paimon/src/table/scan_trace.rs
b/crates/paimon/src/table/scan_trace.rs
index 9dd616ac..5aec3ac9 100644
--- a/crates/paimon/src/table/scan_trace.rs
+++ b/crates/paimon/src/table/scan_trace.rs
@@ -32,11 +32,13 @@ pub struct ScanTrace {
pub delta_manifest_files: usize,
pub manifest_files_before_partition_pruning: usize,
pub manifest_files_after_partition_pruning: usize,
+ pub manifest_files_pruned_by_row_ranges: usize,
pub manifest_entries_read: usize,
pub manifest_entries_pruned_by_bucket: usize,
pub manifest_entries_pruned_by_partition: usize,
pub manifest_entries_after_entry_pruning: usize,
pub manifest_entries_pruned_by_level: usize,
+ pub manifest_entries_pruned_by_row_ranges: usize,
pub manifest_entries_pruned_by_data_stats: usize,
pub manifest_entries_after_manifest_filters: usize,
pub manifest_entries_after_merge: usize,
@@ -97,13 +99,15 @@ impl fmt::Display for ScanTrace {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(
f,
- "snapshot={:?}, manifests={}/{}, entries_read={},
bucket_pruned={}, partition_pruned={}, data_stats_pruned={},
cross_schema_pruned={}, split_candidates_built={}, limit_early_stopped={},
splits_before_limit={}, splits_after_limit={}, files={}",
+ "snapshot={:?}, manifests={}/{}, manifest_row_range_pruned={},
entries_read={}, bucket_pruned={}, partition_pruned={},
entry_row_range_pruned={}, data_stats_pruned={}, cross_schema_pruned={},
split_candidates_built={}, limit_early_stopped={}, splits_before_limit={},
splits_after_limit={}, files={}",
self.snapshot_id,
self.manifest_files_after_partition_pruning,
self.manifest_files_before_partition_pruning,
+ self.manifest_files_pruned_by_row_ranges,
self.manifest_entries_read,
self.manifest_entries_pruned_by_bucket,
self.manifest_entries_pruned_by_partition,
+ self.manifest_entries_pruned_by_row_ranges,
self.manifest_entries_pruned_by_data_stats,
self.manifest_entries_pruned_by_cross_schema_stats,
self.split_candidates_built,
diff --git a/crates/paimon/src/table/table_scan.rs
b/crates/paimon/src/table/table_scan.rs
index af36b177..a9da04e9 100644
--- a/crates/paimon/src/table/table_scan.rs
+++ b/crates/paimon/src/table/table_scan.rs
@@ -22,6 +22,8 @@
use super::bucket_filter::compute_target_buckets;
use super::format_table_scan::FormatTableScan;
+use super::global_index_scanner::RowRangeIndex;
+use super::global_index_types::normalize_sorted_global_index_type;
use super::kv_file_reader::retain_primary_key_conjuncts;
use super::partition_filter::PartitionFilter;
use super::stats_filter::{
@@ -33,8 +35,8 @@ use super::Table;
use crate::io::FileIO;
use crate::spec::{
avro::SharedSchemaCache, bucket_dir_name, BinaryRow, BucketFunctionType,
CoreOptions,
- DataField, DataFileMeta, FileKind, GlobalIndexSearchMode, IndexManifest,
ManifestEntry,
- PartitionComputer, Predicate, Snapshot, ROW_ID_FIELD_ID, ROW_ID_FIELD_NAME,
+ DataField, DataFileMeta, FileKind, GlobalIndexSearchMode, IndexManifest,
IndexManifestEntry,
+ ManifestEntry, PartitionComputer, Predicate, Snapshot, ROW_ID_FIELD_ID,
ROW_ID_FIELD_NAME,
SEQUENCE_NUMBER_FIELD_ID, SEQUENCE_NUMBER_FIELD_NAME, VALUE_KIND_FIELD_ID,
VALUE_KIND_FIELD_NAME,
};
@@ -56,6 +58,7 @@ use std::sync::Arc;
const MANIFEST_DIR: &str = "manifest";
/// Path segment for index directory under table.
const INDEX_DIR: &str = "index";
+const DELETION_VECTORS_INDEX_TYPE: &str = "DELETION_VECTORS";
#[derive(Debug, Default)]
struct ManifestReadCounters {
@@ -64,6 +67,7 @@ struct ManifestReadCounters {
pruned_by_partition: usize,
after_entry_pruning: usize,
pruned_by_level: usize,
+ pruned_by_row_ranges: usize,
pruned_by_data_stats: usize,
after_manifest_filters: usize,
}
@@ -75,6 +79,7 @@ impl ManifestReadCounters {
self.pruned_by_partition += other.pruned_by_partition;
self.after_entry_pruning += other.after_entry_pruning;
self.pruned_by_level += other.pruned_by_level;
+ self.pruned_by_row_ranges += other.pruned_by_row_ranges;
self.pruned_by_data_stats += other.pruned_by_data_stats;
self.after_manifest_filters += other.after_manifest_filters;
}
@@ -122,6 +127,7 @@ async fn read_all_manifest_entries(
bucket_predicate: Option<&Predicate>,
bucket_key_fields: &[DataField],
bucket_function_type: BucketFunctionType,
+ row_range_index: Option<&RowRangeIndex>,
trace: Option<&mut ScanTrace>,
) -> crate::Result<Vec<ManifestEntry>> {
let (mut manifest_files, delta) = futures::try_join!(
@@ -159,6 +165,14 @@ async fn read_all_manifest_entries(
trace.manifest_files_after_partition_pruning = manifest_files.len();
}
+ if let Some(index) = row_range_index {
+ let before = manifest_files.len();
+ retain_manifest_row_range_components(&mut manifest_files, index);
+ if let Some(trace) = trace.as_deref_mut() {
+ trace.manifest_files_pruned_by_row_ranges = before -
manifest_files.len();
+ }
+ }
+
let manifest_path_prefix = format!("{}/{}",
table_path.trim_end_matches('/'), MANIFEST_DIR);
let shared_cache = SharedSchemaCache::new();
let manifest_results: Vec<(Vec<ManifestEntry>, ManifestReadCounters)> =
@@ -253,18 +267,257 @@ async fn read_all_manifest_entries(
counters.merge(manifest_counters);
all_entries.extend(entries);
}
+ if let Some(index) = row_range_index {
+ let before = all_entries.len();
+ all_entries = retain_manifest_entry_row_range_groups(all_entries,
index);
+ counters.pruned_by_row_ranges = before - all_entries.len();
+ counters.after_manifest_filters = all_entries.len();
+ }
if let Some(trace) = trace {
trace.manifest_entries_read = counters.entries_read;
trace.manifest_entries_pruned_by_bucket = counters.pruned_by_bucket;
trace.manifest_entries_pruned_by_partition =
counters.pruned_by_partition;
trace.manifest_entries_after_entry_pruning =
counters.after_entry_pruning;
trace.manifest_entries_pruned_by_level = counters.pruned_by_level;
+ trace.manifest_entries_pruned_by_row_ranges =
counters.pruned_by_row_ranges;
trace.manifest_entries_pruned_by_data_stats =
counters.pruned_by_data_stats;
trace.manifest_entries_after_manifest_filters =
counters.after_manifest_filters;
}
Ok(all_entries)
}
+#[cfg(test)]
+fn manifest_file_overlaps_row_range_index(
+ manifest: &crate::spec::ManifestFileMeta,
+ row_range_index: &RowRangeIndex,
+) -> bool {
+ manifest_row_id_range(manifest).is_none_or(|(min, max)|
row_range_index.intersects(min, max))
+}
+
+fn manifest_row_id_range(manifest: &crate::spec::ManifestFileMeta) ->
Option<(i64, i64)> {
+ match (manifest.min_row_id(), manifest.max_row_id()) {
+ (Some(min), Some(max)) if min <= max => Some((min, max)),
+ _ => None,
+ }
+}
+
+fn retain_manifest_row_range_components(
+ manifests: &mut Vec<crate::spec::ManifestFileMeta>,
+ row_range_index: &RowRangeIndex,
+) {
+ let ranges = manifests
+ .iter()
+ .map(manifest_row_id_range)
+ .collect::<Option<Vec<_>>>();
+ let Some(ranges) = ranges else {
+ // An unknown manifest range may contain the anchor that connects
otherwise
+ // disjoint dedicated-file ranges. Fail open for the whole list.
+ return;
+ };
+
+ let mut order = (0..manifests.len()).collect::<Vec<_>>();
+ order.sort_unstable_by_key(|&idx| ranges[idx]);
+ let mut keep = vec![false; manifests.len()];
+ let mut component = Vec::new();
+ let mut component_from = 0i64;
+ let mut component_to = 0i64;
+
+ for idx in order {
+ let (from, to) = ranges[idx];
+ if component.is_empty() {
+ component_from = from;
+ component_to = to;
+ component.push(idx);
+ } else if from <= component_to {
+ component_to = component_to.max(to);
+ component.push(idx);
+ } else {
+ if row_range_index.intersects(component_from, component_to) {
+ for component_idx in component.drain(..) {
+ keep[component_idx] = true;
+ }
+ } else {
+ component.clear();
+ }
+ component_from = from;
+ component_to = to;
+ component.push(idx);
+ }
+ }
+ if !component.is_empty() && row_range_index.intersects(component_from,
component_to) {
+ for component_idx in component {
+ keep[component_idx] = true;
+ }
+ }
+
+ let mut idx = 0usize;
+ manifests.retain(|_| {
+ let retain = keep[idx];
+ idx += 1;
+ retain
+ });
+}
+
+#[cfg(test)]
+fn data_file_overlaps_row_range_index(
+ file: &DataFileMeta,
+ row_range_index: &RowRangeIndex,
+) -> bool {
+ file.row_id_range()
+ .is_none_or(|(from, to)| row_range_index.intersects(from, to))
+}
+
+fn retain_manifest_entry_row_range_groups(
+ entries: Vec<ManifestEntry>,
+ row_range_index: &RowRangeIndex,
+) -> Vec<ManifestEntry> {
+ let mut buckets: HashMap<(&[u8], i32), Vec<usize>> = HashMap::new();
+ for (idx, entry) in entries.iter().enumerate() {
+ buckets
+ .entry((entry.partition(), entry.bucket()))
+ .or_default()
+ .push(idx);
+ }
+
+ let mut keep = vec![false; entries.len()];
+ for indices in buckets.values_mut() {
+ if indices
+ .iter()
+ .any(|&idx| entries[idx].file().row_id_range().is_none())
+ {
+ // Unknown file ranges may bridge otherwise disjoint row groups.
+ // Keep the whole bucket rather than risking a partial group.
+ for &idx in indices.iter() {
+ keep[idx] = true;
+ }
+ continue;
+ }
+
+ indices.sort_unstable_by_key(|&idx| {
+ entries[idx]
+ .file()
+ .row_id_range()
+ .expect("validated row-id range")
+ });
+ let mut component = Vec::new();
+ let mut component_from = 0i64;
+ let mut component_to = 0i64;
+
+ for &idx in indices.iter() {
+ let (from, to) = entries[idx]
+ .file()
+ .row_id_range()
+ .expect("validated row-id range");
+ if component.is_empty() {
+ component_from = from;
+ component_to = to;
+ component.push(idx);
+ } else if from <= component_to {
+ component_to = component_to.max(to);
+ component.push(idx);
+ } else {
+ if row_range_index.intersects(component_from, component_to) {
+ for component_idx in component.drain(..) {
+ keep[component_idx] = true;
+ }
+ } else {
+ component.clear();
+ }
+ component_from = from;
+ component_to = to;
+ component.push(idx);
+ }
+ }
+ if !component.is_empty() && row_range_index.intersects(component_from,
component_to) {
+ for component_idx in component {
+ keep[component_idx] = true;
+ }
+ }
+ }
+ drop(buckets);
+
+ entries
+ .into_iter()
+ .enumerate()
+ .filter_map(|(idx, entry)| keep[idx].then_some(entry))
+ .collect()
+}
+
+fn data_evolution_row_range_groups(
+ data_files: Vec<DataFileMeta>,
+ row_ranges: Option<&[RowRange]>,
+) -> (Vec<Vec<DataFileMeta>>, usize) {
+ if data_files.is_empty() {
+ return (Vec::new(), 0);
+ }
+ let Some(row_ranges) = row_ranges else {
+ return (group_by_overlapping_row_id(data_files), 0);
+ };
+ let all_ranges_known = data_files.iter().all(|file|
file.row_id_range().is_some());
+ if !all_ranges_known {
+ // Avoid unchecked row-range arithmetic in downstream grouping and keep
+ // the whole bucket as one non-raw group. The reader will then fail on
+ // invalid metadata instead of silently losing dedicated providers.
+ return (vec![data_files], 0);
+ }
+ let row_id_groups = group_by_overlapping_row_id(data_files);
+ let groups_before_pruning = row_id_groups.len();
+
+ let retained = row_id_groups
+ .into_iter()
+ .filter(|group| {
+ group
+ .iter()
+ .any(|file| any_range_overlaps_file(row_ranges, file))
+ })
+ .collect::<Vec<_>>();
+ let pruned = groups_before_pruning - retained.len();
+ (retained, pruned)
+}
+
+fn split_row_ranges_for_files(
+ effective_row_ranges: Option<&[RowRange]>,
+ files: &[DataFileMeta],
+) -> crate::Result<Option<Vec<RowRange>>> {
+ let Some(ranges) = effective_row_ranges else {
+ return Ok(None);
+ };
+ if let Some(file) = files.iter().find(|file|
file.row_id_range().is_none()) {
+ return Err(crate::Error::DataInvalid {
+ message: format!(
+ "Cannot apply selected row ranges to file '{}' with missing or
invalid row-id range",
+ file.file_name
+ ),
+ source: None,
+ });
+ }
+
+ let split_ranges = merge_row_ranges(
+ files
+ .iter()
+ .flat_map(|file| intersect_ranges_with_file(ranges, file))
+ .collect(),
+ );
+ if split_ranges.is_empty() {
+ return Err(crate::Error::DataInvalid {
+ message: "Planned data-evolution split does not overlap selected
row ranges"
+ .to_string(),
+ source: None,
+ });
+ }
+ Ok(Some(split_ranges))
+}
+
+fn retain_index_manifest_entry(
+ entry: &IndexManifestEntry,
+ global_index_needed: bool,
+ deletion_vectors_needed: bool,
+) -> bool {
+ (deletion_vectors_needed && entry.index_file.index_type ==
DELETION_VECTORS_INDEX_TYPE)
+ || (global_index_needed
+ &&
normalize_sorted_global_index_type(&entry.index_file.index_type).is_some())
+}
+
/// Builds a map from (partition, bucket) to (data_file_name -> DeletionFile)
from index manifest entries.
/// Only considers ADD entries with index_type "DELETION_VECTORS" and their
deletion_vectors_ranges.
fn build_deletion_files_map(
@@ -280,7 +533,7 @@ fn build_deletion_files_map(
if entry.kind != FileKind::Add {
continue;
}
- if entry.index_file.index_type != "DELETION_VECTORS" {
+ if entry.index_file.index_type != DELETION_VECTORS_INDEX_TYPE {
continue;
}
let ranges = match &entry.index_file.deletion_vectors_ranges {
@@ -408,16 +661,24 @@ impl LimitPushdownAccumulator {
type BucketDataFileGroups = HashMap<(Vec<u8>, i32), (i32, Vec<DataFileMeta>)>;
-fn global_index_detail_data_ranges(groups: &BucketDataFileGroups) ->
Vec<RowRange> {
- let mut ranges = Vec::new();
- for (_, data_files) in groups.values() {
- for file in data_files {
- if let Some((from, to)) = file.row_id_range() {
- ranges.push(RowRange::new(from, to));
- }
- }
- }
- merge_row_ranges(ranges)
+#[derive(Clone, Copy)]
+struct GlobalIndexScanSettings {
+ search_mode: GlobalIndexSearchMode,
+ thread_num: usize,
+}
+
+fn global_index_detail_data_ranges(entries: &[ManifestEntry]) -> Vec<RowRange>
{
+ merge_row_ranges(
+ entries
+ .iter()
+ .filter_map(|entry| {
+ entry
+ .file()
+ .row_id_range()
+ .map(|(from, to)| RowRange::new(from, to))
+ })
+ .collect(),
+ )
}
fn should_skip_level_zero_for_scan(
@@ -937,12 +1198,14 @@ impl<'a> PaimonTableScan<'a> {
&self,
snapshot: &Snapshot,
) -> crate::Result<Vec<ManifestEntry>> {
- self.plan_manifest_entries_with_trace(snapshot, None).await
+ self.plan_manifest_entries_with_trace(snapshot, None, None)
+ .await
}
async fn plan_manifest_entries_with_trace(
&self,
snapshot: &Snapshot,
+ row_range_index: Option<&RowRangeIndex>,
mut trace: Option<&mut ScanTrace>,
) -> crate::Result<Vec<ManifestEntry>> {
let file_io = self.table.file_io();
@@ -1016,6 +1279,7 @@ impl<'a> PaimonTableScan<'a> {
self.bucket_predicate.as_ref(),
&bucket_key_fields,
bucket_function_type,
+ row_range_index,
trace.as_deref_mut(),
)
.await?;
@@ -1030,6 +1294,121 @@ impl<'a> PaimonTableScan<'a> {
can_push_down_limit_hint_for_scan(&self.data_predicates, row_ranges)
}
+ fn global_index_scan_settings(
+ &self,
+ core_options: &CoreOptions,
+ data_evolution_enabled: bool,
+ ) -> crate::Result<Option<GlobalIndexScanSettings>> {
+ if data_evolution_enabled
+ && core_options.global_index_enabled()
+ && !self.data_predicates.is_empty()
+ {
+ Ok(Some(GlobalIndexScanSettings {
+ search_mode: core_options.global_index_search_mode()?,
+ thread_num: core_options.global_index_thread_num()?,
+ }))
+ } else {
+ Ok(None)
+ }
+ }
+
+ async fn read_index_manifest_entries(
+ &self,
+ snapshot: &Snapshot,
+ global_index_needed: bool,
+ deletion_vectors_needed: bool,
+ ) -> crate::Result<Option<Vec<IndexManifestEntry>>> {
+ if !global_index_needed && !deletion_vectors_needed {
+ return Ok(None);
+ }
+ let Some(index_manifest_name) = snapshot.index_manifest() else {
+ return Ok(None);
+ };
+ let table_path = self.table.location().trim_end_matches('/');
+ let path =
format!("{table_path}/{MANIFEST_DIR}/{index_manifest_name}");
+ let entries = IndexManifest::read(self.table.file_io(), &path)
+ .await?
+ .into_iter()
+ .filter(|entry| {
+ retain_index_manifest_entry(entry, global_index_needed,
deletion_vectors_needed)
+ })
+ .collect();
+ Ok(Some(entries))
+ }
+
+ async fn evaluate_global_index_row_ranges(
+ &self,
+ snapshot: &Snapshot,
+ index_entries: &[IndexManifestEntry],
+ settings: GlobalIndexScanSettings,
+ data_ranges: &[RowRange],
+ ) -> crate::Result<Option<Vec<RowRange>>> {
+ let core_options = CoreOptions::new(self.table.schema().options());
+ super::global_index_scanner::evaluate_global_index(
+ super::global_index_scanner::GlobalIndexEvaluation {
+ file_io: self.table.file_io(),
+ table_path: self.table.location().trim_end_matches('/'),
+ index_entries,
+ predicates: &self.data_predicates,
+ schema_fields: self.table.schema().fields(),
+ search_mode: settings.search_mode,
+ global_index_thread_num: settings.thread_num,
+ btree_fallback_scan_max_size:
core_options.btree_index_fallback_scan_max_size()?,
+ bitmap_fallback_scan_max_size: core_options
+ .bitmap_index_fallback_scan_max_size()?,
+ next_row_id: snapshot.next_row_id(),
+ data_ranges,
+ },
+ )
+ .await
+ }
+
+ async fn manifest_row_ranges(
+ &self,
+ snapshot: &Snapshot,
+ index_entries: Option<&[IndexManifestEntry]>,
+ settings: Option<GlobalIndexScanSettings>,
+ ) -> crate::Result<Option<Vec<RowRange>>> {
+ if self.row_ranges.is_some() {
+ return Ok(self.row_ranges.clone());
+ }
+ let (Some(index_entries), Some(settings)) = (index_entries, settings)
else {
+ return Ok(None);
+ };
+ if settings.search_mode == GlobalIndexSearchMode::Detail {
+ return Ok(None);
+ }
+ self.evaluate_global_index_row_ranges(snapshot, index_entries,
settings, &[])
+ .await
+ }
+
+ async fn effective_row_ranges(
+ &self,
+ snapshot: &Snapshot,
+ entries: &[ManifestEntry],
+ index_entries: Option<&[IndexManifestEntry]>,
+ settings: Option<GlobalIndexScanSettings>,
+ manifest_row_ranges: Option<Vec<RowRange>>,
+ ) -> crate::Result<Option<Vec<RowRange>>> {
+ if self.row_ranges.is_some()
+ || !matches!(
+ settings,
+ Some(GlobalIndexScanSettings {
+ search_mode: GlobalIndexSearchMode::Detail,
+ ..
+ })
+ )
+ {
+ return Ok(manifest_row_ranges);
+ }
+ let (Some(index_entries), Some(settings)) = (index_entries, settings)
else {
+ return Ok(None);
+ };
+ let data_ranges = global_index_detail_data_ranges(entries);
+ self.evaluate_global_index_row_ranges(snapshot, index_entries,
settings, &data_ranges)
+ .await
+ }
+
/// The predicate set that may prune WHOLE FILES by their stats.
///
/// For primary-key tables read by merging, only key conjuncts are safe: a
@@ -1074,15 +1453,11 @@ impl<'a> PaimonTableScan<'a> {
/// reads the delta manifest list and keeps ADD entries.
pub(crate) async fn plan_snapshot_delta(&self, snapshot: &Snapshot) ->
crate::Result<Plan> {
self.ensure_query_auth_allowed()?;
- let entries = self
- .plan_manifest_list_entries(snapshot.delta_manifest_list())
- .await?;
let data_evolution_read_field_ids = self.projected_read_field_ids()?;
- self.plan_snapshot_from_entries(
- snapshot.clone(),
- entries,
+ self.plan_snapshot_manifest_list(
+ snapshot,
+ snapshot.delta_manifest_list(),
data_evolution_read_field_ids.as_ref(),
- None,
)
.await
}
@@ -1097,12 +1472,61 @@ impl<'a> PaimonTableScan<'a> {
let Some(list_name) = snapshot.changelog_manifest_list() else {
return Ok(Plan::new(Vec::new()));
};
- let entries = self.plan_manifest_list_entries(list_name).await?;
let data_evolution_read_field_ids = self.projected_read_field_ids()?;
+ self.plan_snapshot_manifest_list(
+ snapshot,
+ list_name,
+ data_evolution_read_field_ids.as_ref(),
+ )
+ .await
+ }
+
+ async fn plan_snapshot_manifest_list(
+ &self,
+ snapshot: &Snapshot,
+ manifest_list_name: &str,
+ data_evolution_read_field_ids: Option<&HashSet<i32>>,
+ ) -> crate::Result<Plan> {
+ if matches!(self.limit, Some(0)) {
+ return Ok(Plan::new(Vec::new()));
+ }
+ let core_options = CoreOptions::new(self.table.schema().options());
+ let data_evolution_enabled = core_options.data_evolution_enabled();
+ let global_index_settings =
+ self.global_index_scan_settings(&core_options,
data_evolution_enabled)?;
+ let index_entries = self
+ .read_index_manifest_entries(
+ snapshot,
+ global_index_settings.is_some(),
+ core_options.deletion_vectors_enabled(),
+ )
+ .await?;
+ let manifest_row_ranges = self
+ .manifest_row_ranges(snapshot, index_entries.as_deref(),
global_index_settings)
+ .await?;
+ let row_range_index = if data_evolution_enabled {
+ manifest_row_ranges.clone().map(RowRangeIndex::create)
+ } else {
+ None
+ };
+ let entries = self
+ .plan_manifest_list_entries(manifest_list_name,
row_range_index.as_ref())
+ .await?;
+ let effective_row_ranges = self
+ .effective_row_ranges(
+ snapshot,
+ &entries,
+ index_entries.as_deref(),
+ global_index_settings,
+ manifest_row_ranges,
+ )
+ .await?;
self.plan_snapshot_from_entries(
snapshot.clone(),
entries,
- data_evolution_read_field_ids.as_ref(),
+ data_evolution_read_field_ids,
+ index_entries,
+ effective_row_ranges,
None,
)
.await
@@ -1113,6 +1537,7 @@ impl<'a> PaimonTableScan<'a> {
async fn plan_manifest_list_entries(
&self,
manifest_list_name: &str,
+ row_range_index: Option<&RowRangeIndex>,
) -> crate::Result<Vec<ManifestEntry>> {
let file_io = self.table.file_io();
let table_path = self.table.location();
@@ -1140,6 +1565,9 @@ impl<'a> PaimonTableScan<'a> {
});
}
}
+ if let Some(index) = row_range_index {
+ retain_manifest_row_range_components(&mut manifest_metas, index);
+ }
let bucket_key_fields: Vec<DataField> = if
self.bucket_predicate.is_none() {
Vec::new()
@@ -1214,7 +1642,12 @@ impl<'a> PaimonTableScan<'a> {
let entries = entries
.into_iter()
.filter(|entry| *entry.kind() == FileKind::Add)
- .collect();
+ .collect::<Vec<_>>();
+ let entries = if let Some(index) = row_range_index {
+ retain_manifest_entry_row_range_groups(entries, index)
+ } else {
+ entries
+ };
Ok(entries)
}
@@ -1224,11 +1657,56 @@ impl<'a> PaimonTableScan<'a> {
data_evolution_read_field_ids: Option<&HashSet<i32>>,
mut trace: Option<&mut ScanTrace>,
) -> crate::Result<Plan> {
+ if matches!(self.limit, Some(0)) {
+ if let Some(trace) = trace {
+ trace.record_final_plan_with_limit(0, 0, 0, 0, true);
+ }
+ return Ok(Plan::new(Vec::new()));
+ }
+ let core_options = CoreOptions::new(self.table.schema().options());
+ let data_evolution_enabled = core_options.data_evolution_enabled();
+ let global_index_settings =
+ self.global_index_scan_settings(&core_options,
data_evolution_enabled)?;
+ let index_entries = self
+ .read_index_manifest_entries(
+ &snapshot,
+ global_index_settings.is_some(),
+ core_options.deletion_vectors_enabled(),
+ )
+ .await?;
+ let manifest_row_ranges = self
+ .manifest_row_ranges(&snapshot, index_entries.as_deref(),
global_index_settings)
+ .await?;
+ let row_range_index = if data_evolution_enabled {
+ manifest_row_ranges.clone().map(RowRangeIndex::create)
+ } else {
+ None
+ };
let entries = self
- .plan_manifest_entries_with_trace(&snapshot, trace.as_deref_mut())
+ .plan_manifest_entries_with_trace(
+ &snapshot,
+ row_range_index.as_ref(),
+ trace.as_deref_mut(),
+ )
.await?;
- self.plan_snapshot_from_entries(snapshot, entries,
data_evolution_read_field_ids, trace)
- .await
+ let effective_row_ranges = self
+ .effective_row_ranges(
+ &snapshot,
+ &entries,
+ index_entries.as_deref(),
+ global_index_settings,
+ manifest_row_ranges,
+ )
+ .await?;
+ self.plan_snapshot_from_entries(
+ snapshot,
+ entries,
+ data_evolution_read_field_ids,
+ index_entries,
+ effective_row_ranges,
+ trace,
+ )
+ .await
}
async fn plan_snapshot_from_entries(
@@ -1236,9 +1714,10 @@ impl<'a> PaimonTableScan<'a> {
snapshot: Snapshot,
entries: Vec<ManifestEntry>,
data_evolution_read_field_ids: Option<&HashSet<i32>>,
+ index_entries: Option<Vec<IndexManifestEntry>>,
+ effective_row_ranges: Option<Vec<RowRange>>,
mut trace: Option<&mut ScanTrace>,
) -> crate::Result<Plan> {
- let file_io = self.table.file_io();
let table_path = self.table.location();
let table_schema_id = self.table.schema().id();
let table_fields = self.table.schema().fields();
@@ -1325,30 +1804,6 @@ impl<'a> PaimonTableScan<'a> {
entry.1.push(file);
}
- let global_index_settings = if data_evolution_enabled
- && core_options.global_index_enabled()
- && !self.data_predicates.is_empty()
- {
- Some((
- core_options.global_index_search_mode()?,
- core_options.global_index_thread_num()?,
- ))
- } else {
- None
- };
- let global_index_detail_data_ranges = if matches!(
- global_index_settings,
- Some((GlobalIndexSearchMode::Detail, _))
- ) {
- global_index_detail_data_ranges(&groups)
- } else {
- Vec::new()
- };
- let btree_index_fallback_scan_max_size =
- core_options.btree_index_fallback_scan_max_size()?;
- let bitmap_index_fallback_scan_max_size =
- core_options.bitmap_index_fallback_scan_max_size()?;
-
let snapshot_id = snapshot.id();
let base_path = table_path.trim_end_matches('/');
let mut splits = Vec::with_capacity(groups.len());
@@ -1382,42 +1837,11 @@ impl<'a> PaimonTableScan<'a> {
None
};
- // Read deletion vector index manifest once (like Java generateSplits
/ scanDvIndex).
- let (deletion_files_map, effective_row_ranges) =
- if let Some(index_manifest_name) = snapshot.index_manifest() {
- let index_manifest_path =
format!("{base_path}/{MANIFEST_DIR}");
- let path =
format!("{index_manifest_path}/{index_manifest_name}");
- let index_entries = IndexManifest::read(file_io, &path).await?;
- let dv_map = build_deletion_files_map(&index_entries,
base_path);
-
- // Use pushed-down row_ranges first; otherwise try global
index.
- let row_ranges = if self.row_ranges.is_some() {
- self.row_ranges.clone()
- } else if let Some((search_mode, global_index_thread_num)) =
global_index_settings {
- super::global_index_scanner::evaluate_global_index(
- super::global_index_scanner::GlobalIndexEvaluation {
- file_io,
- table_path: base_path,
- index_entries: &index_entries,
- predicates: &self.data_predicates,
- schema_fields: self.table.schema().fields(),
- search_mode,
- global_index_thread_num,
- btree_fallback_scan_max_size:
btree_index_fallback_scan_max_size,
- bitmap_fallback_scan_max_size:
bitmap_index_fallback_scan_max_size,
- next_row_id: snapshot.next_row_id(),
- data_ranges: &global_index_detail_data_ranges,
- },
- )
- .await?
- } else {
- None
- };
-
- (Some(dv_map), row_ranges)
- } else {
- (None, self.row_ranges.clone())
- };
+ // The index manifest was read before data manifests so global-index
row
+ // ranges can prune manifest I/O. Reuse it here for deletion vectors.
+ let deletion_files_map = index_entries
+ .as_deref()
+ .map(|entries| build_deletion_files_map(entries, base_path));
let mut data_file_field_ids_cache = DataFileFieldIdsCache::new();
let can_push_down_limit =
self.can_push_down_limit_hint(effective_row_ranges.as_deref());
@@ -1442,13 +1866,13 @@ impl<'a> PaimonTableScan<'a> {
.as_ref()
.and_then(|map| map.get(&PartitionBucket::new(partition,
bucket)));
- // Data-evolution tables merge overlapping row-id groups
column-wise during read.
- // Keep that split boundary intact and only bin-pack single-file
groups.
- // Apply group-level predicate filtering after grouping by row_id
range.
+ // Data-evolution reads merge overlapping row-id groups
column-wise.
let file_groups: Vec<SplitGroup> = if data_evolution_enabled {
- let row_id_groups = group_by_overlapping_row_id(data_files);
+ let (row_id_groups, groups_pruned_by_row_ranges) =
+ data_evolution_row_range_groups(data_files,
effective_row_ranges.as_deref());
if let Some(trace) = trace.as_deref_mut() {
trace.data_evolution_groups_before_stats +=
row_id_groups.len();
+ trace.data_evolution_groups_pruned_by_row_ranges +=
groups_pruned_by_row_ranges;
}
// Filter groups by merged stats before splitting.
@@ -1472,21 +1896,6 @@ impl<'a> PaimonTableScan<'a> {
groups
};
- // Filter groups by row ID ranges.
- let row_id_groups = if let Some(ref ranges) =
effective_row_ranges {
- let before = row_id_groups.len();
- let groups = row_id_groups
- .into_iter()
- .filter(|group| group.iter().any(|f|
any_range_overlaps_file(ranges, f)))
- .collect::<Vec<_>>();
- if let Some(trace) = trace.as_deref_mut() {
- trace.data_evolution_groups_pruned_by_row_ranges +=
before - groups.len();
- }
- groups
- } else {
- row_id_groups
- };
-
let row_id_groups = if let Some(read_field_ids) =
data_evolution_read_field_ids {
if read_field_ids.is_empty() {
row_id_groups
@@ -1579,21 +1988,8 @@ impl<'a> PaimonTableScan<'a> {
.collect::<Vec<Option<DeletionFile>>>()
});
- // Compute row_ranges before moving file_group to avoid clone
- let split_row_ranges = if let Some(ref ranges) =
effective_row_ranges {
- let mut split_ranges = Vec::new();
- for file in &file_group {
- split_ranges.extend(intersect_ranges_with_file(ranges,
file));
- }
- let split_ranges = merge_row_ranges(split_ranges);
- if split_ranges.is_empty() {
- None
- } else {
- Some(split_ranges)
- }
- } else {
- None
- };
+ let split_row_ranges =
+
split_row_ranges_for_files(effective_row_ranges.as_deref(), &file_group)?;
let mut builder = DataSplitBuilder::new()
.with_snapshot(snapshot_id)
@@ -1651,20 +2047,25 @@ impl<'a> PaimonTableScan<'a> {
#[cfg(test)]
mod tests {
use super::{
- prune_data_evolution_group_by_read_fields,
should_skip_level_zero_for_scan,
- LimitPushdownAccumulator, TableScan,
+ data_evolution_row_range_groups, data_file_overlaps_row_range_index,
+ manifest_file_overlaps_row_range_index,
prune_data_evolution_group_by_read_fields,
+ retain_index_manifest_entry, retain_manifest_entry_row_range_groups,
+ retain_manifest_row_range_components, should_skip_level_zero_for_scan,
+ split_row_ranges_for_files, LimitPushdownAccumulator, PaimonTableScan,
RowRangeIndex,
+ TableScan,
};
use crate::catalog::Identifier;
use crate::io::FileIOBuilder;
use crate::spec::{
stats::BinaryTableStats, ArrayType, BinaryRow, BinaryRowBuilder,
BucketFunctionType,
- DataField, DataFileMeta, DataType, Datum, DeletionVectorMeta,
FileKind, IndexFileMeta,
- IndexManifestEntry, IntType, Predicate, PredicateBuilder,
PredicateOperator,
- Schema as PaimonSchema, TableSchema, VarCharType,
+ CommitKind, DataField, DataFileMeta, DataType, Datum,
DeletionVectorMeta, FileKind,
+ IndexFileMeta, IndexManifestEntry, IntType, ManifestEntry,
ManifestFileMeta, Predicate,
+ PredicateBuilder, PredicateOperator, Schema as PaimonSchema, Snapshot,
TableSchema,
+ VarCharType,
};
use crate::table::bucket_filter::{compute_target_buckets,
extract_predicate_for_keys};
use crate::table::partition_filter::PartitionFilter;
- use crate::table::source::{DataSplit, DataSplitBuilder, DeletionFile};
+ use crate::table::source::{DataSplit, DataSplitBuilder, DeletionFile,
RowRange};
use crate::table::stats_filter::{
data_evolution_group_matches_predicates, data_file_matches_predicates,
group_by_overlapping_row_id,
@@ -1719,6 +2120,189 @@ mod tests {
file
}
+ #[test]
+ fn test_row_range_manifest_filters_fail_open_without_valid_metadata() {
+ let index = RowRangeIndex::create(vec![RowRange::new(10, 20)]);
+ let stats = BinaryTableStats::new(Vec::new(), Vec::new(), Vec::new());
+
+ let outside = ManifestFileMeta::new("outside".to_string(), 1, 1, 0,
stats.clone(), 0)
+ .with_row_id_stats(Some(30), Some(40));
+ let missing = ManifestFileMeta::new("missing".to_string(), 1, 1, 0,
stats.clone(), 0);
+ let invalid = ManifestFileMeta::new("invalid".to_string(), 1, 1, 0,
stats, 0)
+ .with_row_id_stats(Some(40), Some(30));
+
+ assert!(!manifest_file_overlaps_row_range_index(&outside, &index));
+ assert!(manifest_file_overlaps_row_range_index(&missing, &index));
+ assert!(manifest_file_overlaps_row_range_index(&invalid, &index));
+
+ assert!(!data_file_overlaps_row_range_index(
+ &make_evo_file("outside", 1, 5, 0, Some(30)),
+ &index
+ ));
+ assert!(data_file_overlaps_row_range_index(
+ &make_evo_file("missing", 1, 5, 0, None),
+ &index
+ ));
+ assert!(data_file_overlaps_row_range_index(
+ &make_evo_file("invalid", 1, 0, 0, Some(30)),
+ &index
+ ));
+ assert!(data_file_overlaps_row_range_index(
+ &make_evo_file("overflow", 1, 2, 0, Some(i64::MAX)),
+ &index
+ ));
+ }
+
+ #[test]
+ fn test_manifest_row_range_pruning_retains_overlapping_component() {
+ let index = RowRangeIndex::create(vec![RowRange::new(2, 2)]);
+ let stats = BinaryTableStats::new(Vec::new(), Vec::new(), Vec::new());
+ let manifest = |name: &str, min, max| {
+ ManifestFileMeta::new(name.to_string(), 1, 1, 0, stats.clone(), 0)
+ .with_row_id_stats(Some(min), Some(max))
+ };
+ let mut manifests = vec![
+ manifest("anchor", 0, 5),
+ manifest("left-dedicated", 0, 1),
+ manifest("right-dedicated", 4, 5),
+ manifest("other-group", 10, 15),
+ ];
+
+ retain_manifest_row_range_components(&mut manifests, &index);
+
+ assert_eq!(
+ manifests
+ .iter()
+ .map(ManifestFileMeta::file_name)
+ .collect::<Vec<_>>(),
+ vec!["anchor", "left-dedicated", "right-dedicated"]
+ );
+ }
+
+ #[test]
+ fn test_manifest_row_range_component_pruning_fails_open_on_unknown_range()
{
+ let index = RowRangeIndex::create(vec![RowRange::new(2, 2)]);
+ let stats = BinaryTableStats::new(Vec::new(), Vec::new(), Vec::new());
+ let mut manifests = vec![
+ ManifestFileMeta::new("unknown".to_string(), 1, 1, 0,
stats.clone(), 0),
+ ManifestFileMeta::new("inverted".to_string(), 1, 1, 0,
stats.clone(), 0)
+ .with_row_id_stats(Some(20), Some(10)),
+ ManifestFileMeta::new("one-sided".to_string(), 1, 1, 0,
stats.clone(), 0)
+ .with_row_id_stats(Some(0), None),
+ ManifestFileMeta::new("outside".to_string(), 1, 1, 0, stats, 0)
+ .with_row_id_stats(Some(10), Some(15)),
+ ];
+
+ retain_manifest_row_range_components(&mut manifests, &index);
+
+ assert_eq!(manifests.len(), 4);
+ }
+
+ #[test]
+ fn test_manifest_entry_row_range_pruning_retains_overlapping_group() {
+ let index = RowRangeIndex::create(vec![RowRange::new(2, 2)]);
+ let entry = |name: &str, first_row_id, row_count| {
+ ManifestEntry::new(
+ FileKind::Add,
+ Vec::new(),
+ 0,
+ 1,
+ make_evo_file(name, 1, row_count, 0, Some(first_row_id)),
+ 3,
+ )
+ };
+ let entries = vec![
+ entry("anchor", 0, 6),
+ entry("left-dedicated", 0, 2),
+ entry("right-dedicated", 4, 2),
+ entry("other-group", 10, 2),
+ ];
+
+ let retained = retain_manifest_entry_row_range_groups(entries, &index);
+
+ assert_eq!(
+ retained
+ .iter()
+ .map(|entry| entry.file().file_name.as_str())
+ .collect::<Vec<_>>(),
+ vec!["anchor", "left-dedicated", "right-dedicated"]
+ );
+ }
+
+ #[test]
+ fn test_manifest_entry_row_range_pruning_fails_open_per_bucket() {
+ let index = RowRangeIndex::create(vec![RowRange::new(2, 2)]);
+ let entry = |name: &str, bucket, first_row_id, row_count| {
+ ManifestEntry::new(
+ FileKind::Add,
+ Vec::new(),
+ bucket,
+ 2,
+ make_evo_file(name, 1, row_count, 0, first_row_id),
+ 3,
+ )
+ };
+ let entries = vec![
+ entry("unknown-anchor", 0, None, 6),
+ entry("zero-count", 0, Some(10), 0),
+ entry("overflow", 0, Some(i64::MAX), 2),
+ entry("bucket-0-outside", 0, Some(10), 2),
+ entry("bucket-1-outside", 1, Some(10), 2),
+ ];
+
+ let retained = retain_manifest_entry_row_range_groups(entries, &index);
+
+ let mut names = retained
+ .into_iter()
+ .map(|entry| entry.file().file_name.clone())
+ .collect::<Vec<_>>();
+ names.sort();
+ assert_eq!(
+ names,
+ vec![
+ "bucket-0-outside",
+ "overflow",
+ "unknown-anchor",
+ "zero-count",
+ ]
+ );
+ }
+
+ #[test]
+ fn
test_data_evolution_row_range_group_pruning_fails_open_on_unknown_range() {
+ let files = vec![
+ make_evo_file("unknown-anchor", 1, 6, 0, None),
+ make_evo_file("left.blob", 1, 2, 0, Some(0)),
+ make_evo_file("right.blob", 1, 2, 0, Some(4)),
+ ];
+ let ranges = [RowRange::new(2, 2)];
+
+ let (groups, pruned) = data_evolution_row_range_groups(files,
Some(&ranges));
+
+ assert_eq!(pruned, 0);
+ let mut names = groups
+ .into_iter()
+ .flatten()
+ .map(|file| file.file_name)
+ .collect::<Vec<_>>();
+ names.sort();
+ assert_eq!(names, vec!["left.blob", "right.blob", "unknown-anchor"]);
+ }
+
+ #[test]
+ fn test_split_row_ranges_rejects_unknown_or_non_overlapping_files() {
+ let ranges = [RowRange::new(2, 2)];
+ let unknown = make_evo_file("unknown", 1, 6, 0, None);
+ let error = split_row_ranges_for_files(Some(&ranges),
&[unknown]).unwrap_err();
+ assert!(matches!(error, Error::DataInvalid { message, .. }
+ if message.contains("missing or invalid row-id range")));
+
+ let outside = make_evo_file("outside", 1, 2, 0, Some(10));
+ let error = split_row_ranges_for_files(Some(&ranges),
&[outside]).unwrap_err();
+ assert!(matches!(error, Error::DataInvalid { message, .. }
+ if message.contains("does not overlap selected row ranges")));
+ }
+
fn data_evolution_test_table(table_path: &str, schema: TableSchema) ->
Table {
let file_io = FileIOBuilder::new("memory").build().unwrap();
let schema = schema.copy_with_options(HashMap::from([(
@@ -2196,6 +2780,64 @@ mod tests {
assert_eq!(file_names(&groups), vec![vec!["a", "b"]]);
}
+ #[tokio::test]
+ async fn test_data_evolution_row_ranges_prune_normal_groups() {
+ let table_path = "memory:/de_row_range_prune_normal_groups";
+ let table = data_evolution_test_table(table_path, two_column_schema(0,
"id", "name"));
+ setup_scan_trace_dirs(&table).await;
+
+ TableCommit::new(table.clone(), "row-range-prune-test".to_string())
+ .commit(vec![CommitMessage::new(
+ BinaryRowBuilder::new(0).build_serialized(),
+ 0,
+ vec![
+ make_evo_file("a-new", 10, 101, 2, Some(0)),
+ make_evo_file("a-old", 10, 101, 1, Some(0)),
+ make_evo_file("b-new", 10, 101, 4, Some(400)),
+ make_evo_file("b-old", 10, 101, 3, Some(400)),
+ ],
+ )])
+ .await
+ .unwrap();
+
+ let mut read_builder = table.new_read_builder();
+ read_builder.with_row_ranges(vec![RowRange::new(0, 0)]);
+ let (plan, trace) =
read_builder.new_scan().plan_with_trace().await.unwrap();
+ let planned_files = plan
+ .splits()
+ .iter()
+ .flat_map(|split| split.data_files())
+ .map(|file| file.file_name.as_str())
+ .collect::<Vec<_>>();
+
+ assert_eq!(planned_files, vec!["a-new", "a-old"]);
+ assert_eq!(trace.manifest_entries_read, 4);
+ assert_eq!(trace.manifest_entries_pruned_by_row_ranges, 2);
+ assert_eq!(trace.manifest_entries_after_manifest_filters, 2);
+ assert_eq!(trace.manifest_entries_after_merge, 2);
+ assert_eq!(trace.data_evolution_groups_before_stats, 1);
+ assert_eq!(trace.data_evolution_groups_pruned_by_row_ranges, 0);
+
+ let snapshot = table
+ .snapshot_manager()
+ .get_latest_snapshot()
+ .await
+ .unwrap()
+ .unwrap();
+ let delta_plan = read_builder
+ .new_scan()
+ .plan_snapshot_delta(&snapshot)
+ .await
+ .unwrap();
+ let delta_files = delta_plan
+ .splits()
+ .iter()
+ .flat_map(|split| split.data_files())
+ .map(|file| file.file_name.as_str())
+ .collect::<Vec<_>>();
+ assert_eq!(delta_files, vec!["a-new", "a-old"]);
+ }
+
#[tokio::test]
async fn test_data_evolution_prunes_files_without_projected_columns() {
let table =
@@ -3136,6 +3778,83 @@ mod tests {
);
}
+ #[test]
+ fn test_retain_index_manifest_entries_for_active_consumers() {
+ let entry = |index_type: &str| IndexManifestEntry {
+ version: 1,
+ kind: FileKind::Add,
+ partition: Vec::new(),
+ bucket: 0,
+ index_file: IndexFileMeta {
+ index_type: index_type.to_string(),
+ file_name: format!("{index_type}.idx"),
+ file_size: 1,
+ row_count: 1,
+ deletion_vectors_ranges: None,
+ global_index_meta: None,
+ },
+ };
+ let entries = [
+ entry("DELETION_VECTORS"),
+ entry("btree"),
+ entry("bitmap"),
+ entry("HASH"),
+ ];
+ let retained = |global_index_needed, deletion_vectors_needed| {
+ entries
+ .iter()
+ .filter(|entry| {
+ retain_index_manifest_entry(entry, global_index_needed,
deletion_vectors_needed)
+ })
+ .map(|entry| entry.index_file.index_type.as_str())
+ .collect::<Vec<_>>()
+ };
+
+ assert!(retained(false, false).is_empty());
+ assert_eq!(retained(false, true), vec!["DELETION_VECTORS"]);
+ assert_eq!(retained(true, false), vec!["btree", "bitmap"]);
+ assert_eq!(
+ retained(true, true),
+ vec!["DELETION_VECTORS", "btree", "bitmap"]
+ );
+ }
+
+ #[tokio::test]
+ async fn test_skip_index_manifest_without_active_consumer() {
+ let table = Table::new(
+ FileIOBuilder::new("memory").build().unwrap(),
+ Identifier::new("test_db", "index_manifest_gate"),
+ "memory:/index_manifest_gate".to_string(),
+ TableSchema::new(
+ 0,
+ &PaimonSchema::builder()
+ .column("id", DataType::Int(IntType::new()))
+ .build()
+ .unwrap(),
+ ),
+ None,
+ );
+ let snapshot = Snapshot::builder()
+ .version(3)
+ .id(1)
+ .schema_id(0)
+ .base_manifest_list(String::new())
+ .delta_manifest_list(String::new())
+ .index_manifest(Some("missing-index-manifest".to_string()))
+ .commit_user("test-user".to_string())
+ .commit_identifier(1)
+ .commit_kind(CommitKind::APPEND)
+ .time_millis(1)
+ .build();
+ let scan = PaimonTableScan::new(&table, None, Vec::new(), None, None,
None);
+
+ let entries = scan
+ .read_index_manifest_entries(&snapshot, false, false)
+ .await
+ .unwrap();
+ assert!(entries.is_none());
+ }
+
// ======================== Bucket predicate filtering
========================
fn bucket_key_fields() -> Vec<DataField> {