This is an automated email from the ASF dual-hosted git repository.
jerry-024 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 c6665663 fix(table): align Data Evolution row-range pruning with Java
(#628)
c6665663 is described below
commit c666566339326ca3fdcdb48f286d1d0819991f72
Author: jerry <[email protected]>
AuthorDate: Mon Aug 3 10:55:43 2026 +0800
fix(table): align Data Evolution row-range pruning with Java (#628)
---
crates/paimon/src/table/data_evolution_reader.rs | 425 +++++++++++++++++----
.../table/data_evolution_reader/blob_fallback.rs | 27 +-
crates/paimon/src/table/source.rs | 8 +
crates/paimon/src/table/table_scan.rs | 390 ++-----------------
4 files changed, 401 insertions(+), 449 deletions(-)
diff --git a/crates/paimon/src/table/data_evolution_reader.rs
b/crates/paimon/src/table/data_evolution_reader.rs
index 3e62cd3e..03c76732 100644
--- a/crates/paimon/src/table/data_evolution_reader.rs
+++ b/crates/paimon/src/table/data_evolution_reader.rs
@@ -1147,13 +1147,27 @@ fn open_source_stream(
let uncovered_ranges =
crate::table::source::exclude_row_ranges(selected_ranges,
bunch.logical_ranges());
if !uncovered_ranges.is_empty() {
- return Err(Error::DataInvalid {
- message: format!(
- "Blob bunch logical row ranges {:?} do not cover
effective selected row ranges {selected_ranges:?}; uncovered ranges are
{uncovered_ranges:?}",
- bunch.logical_ranges()
- ),
- source: None,
- });
+ if read_fields
+ .iter()
+ .any(|field| !field.data_type().is_nullable())
+ {
+ return Err(Error::DataInvalid {
+ message: format!(
+ "Cannot NULL-fill uncovered ranges
{uncovered_ranges:?} for non-nullable BLOB field"
+ ),
+ source: None,
+ });
+ }
+ return blob_fallback::read(
+ split,
+ bunch.clone(),
+ read_fields.clone(),
+ row_ranges,
+ batch_size,
+ file_io,
+ blob_as_descriptor,
+ anchor_deletion_vector.cloned(),
+ );
}
}
@@ -1164,6 +1178,7 @@ fn open_source_stream(
bunch.clone(),
read_fields.clone(),
row_ranges,
+ batch_size,
file_io,
blob_as_descriptor,
anchor_deletion_vector.cloned(),
@@ -1233,45 +1248,33 @@ fn open_source_stream(
let covered_ranges = bunch
.files
.iter()
- .map(|file| {
- let first_row_id =
- file.first_row_id.ok_or_else(||
Error::DataInvalid {
- message: format!(
- "Vector file '{}' is missing
first_row_id",
- file.file_name
- ),
- source: None,
- })?;
- if file.row_count <= 0 {
- return Err(Error::DataInvalid {
- message: format!(
- "Vector file '{}' row count must be
positive, got {}",
- file.file_name, file.row_count
- ),
- source: None,
- });
- }
- let last_row_id = first_row_id
- .checked_add(file.row_count - 1)
- .ok_or_else(|| Error::DataInvalid {
- message: format!(
- "Vector file '{}' row range overflows
i64",
- file.file_name
- ),
- source: None,
- })?;
- Ok(RowRange::new(first_row_id, last_row_id))
- })
+ .map(vector_file_row_range)
.collect::<crate::Result<Vec<_>>>()?;
let uncovered_ranges =
crate::table::source::exclude_row_ranges(ranges,
&covered_ranges);
if !uncovered_ranges.is_empty() {
- return Err(Error::DataInvalid {
- message: format!(
- "Vector bunch does not cover effective
selected row ranges {uncovered_ranges:?}"
- ),
- source: None,
- });
+ if source
+ .read_fields()
+ .iter()
+ .any(|field| !field.data_type().is_nullable())
+ {
+ return Err(Error::DataInvalid {
+ message: format!(
+ "Cannot NULL-fill uncovered ranges
{uncovered_ranges:?} for non-nullable vector field"
+ ),
+ source: None,
+ });
+ }
+ return read_vector_bunch_with_null_gaps_stream(
+ file_reader,
+ split,
+ bunch.files_overlapping(ranges),
+ data_fields.clone(),
+ ranges.to_vec(),
+ anchor_deletion_vector.cloned(),
+ source.read_fields(),
+ batch_size,
+ );
}
bunch.files_overlapping(ranges)
}
@@ -1289,6 +1292,102 @@ fn open_source_stream(
}
}
+#[allow(clippy::too_many_arguments)]
+fn read_vector_bunch_with_null_gaps_stream(
+ file_reader: DataFileReader,
+ split: &DataSplit,
+ files: Vec<DataFileMeta>,
+ data_fields: Option<Vec<DataField>>,
+ selected_ranges: Vec<RowRange>,
+ anchor_deletion_vector: Option<DeletionVectorContext>,
+ read_fields: &[DataField],
+ batch_size: Option<usize>,
+) -> crate::Result<ArrowRecordBatchStream> {
+ let files = files
+ .into_iter()
+ .map(|file| vector_file_row_range(&file).map(|range| (file, range)))
+ .collect::<crate::Result<Vec<_>>>()?;
+ let mut reads: Vec<(Option<usize>, Vec<RowRange>)> = Vec::new();
+ let mut file_idx = 0usize;
+ for selected in selected_ranges {
+ let mut cursor = selected.from();
+ loop {
+ while files
+ .get(file_idx)
+ .is_some_and(|(_, range)| range.to() < cursor)
+ {
+ file_idx += 1;
+ }
+ let (provider, range_to) = match files
+ .get(file_idx)
+ .filter(|(_, range)| range.from() <= cursor)
+ {
+ Some((_, range)) => (Some(file_idx),
selected.to().min(range.to())),
+ None => (
+ None,
+ files
+ .get(file_idx)
+ .map(|(_, range)| selected.to().min(range.from() - 1))
+ .unwrap_or_else(|| selected.to()),
+ ),
+ };
+ let range = RowRange::new(cursor, range_to);
+ match reads.last_mut() {
+ Some((last_provider, ranges)) if *last_provider == provider =>
ranges.push(range),
+ _ => reads.push((provider, vec![range])),
+ }
+ if range_to == selected.to() {
+ break;
+ }
+ cursor = range_to + 1;
+ }
+ }
+
+ let target_schema = build_target_arrow_schema(read_fields)?;
+ let null_batch_size = batch_size.unwrap_or(1024).max(1);
+ let split = split.clone();
+
+ Ok(try_stream! {
+ for (provider, ranges) in reads {
+ if let Some(file_idx) = provider {
+ let file = &files[file_idx].0;
+ let deletion_vector =
+ shifted_deletion_vector_for_file(file,
anchor_deletion_vector.as_ref())?;
+ let mut stream = file_reader.read_single_file_stream(
+ &split,
+ file.clone(),
+ data_fields.clone(),
+ deletion_vector,
+ Some(ranges),
+ )?;
+ while let Some(batch) = stream.next().await {
+ yield batch?;
+ }
+ } else {
+ let mut gap_rows =
ranges.iter().map(RowRange::count).sum::<i64>();
+ while gap_rows > 0 {
+ let rows = usize::try_from(gap_rows)
+ .unwrap_or(usize::MAX)
+ .min(null_batch_size);
+ let columns = target_schema
+ .fields()
+ .iter()
+ .map(|field|
arrow_array::new_null_array(field.data_type(), rows))
+ .collect::<Vec<_>>();
+ yield RecordBatch::try_new(target_schema.clone(),
columns).map_err(|e| {
+ Error::UnexpectedError {
+ message: format!("Failed to build NULL-filled
vector batch: {e}"),
+ source: Some(Box::new(e)),
+ }
+ })?;
+ gap_rows -= i64::try_from(rows).unwrap_or(i64::MAX);
+ }
+ }
+ }
+ }
+ .boxed())
+}
+
fn read_bunch_files_stream(
file_reader: DataFileReader,
split: &DataSplit,
@@ -1419,16 +1518,18 @@ fn selected_absolute_row_ranges_for_file(
};
if let Some(ranges) = row_ranges {
- let selected = ranges
- .iter()
- .filter_map(|range| {
- range
- .intersect_inclusive(first_row_id, first_row_id +
row_count - 1)
- .map(|range| {
- RowRange::new(range.from() - first_row_id, range.to()
- first_row_id)
- })
- })
- .collect::<Vec<_>>();
+ let selected = crate::table::merge_row_ranges(
+ ranges
+ .iter()
+ .filter_map(|range| {
+ range
+ .intersect_inclusive(first_row_id, first_row_id +
row_count - 1)
+ .map(|range| {
+ RowRange::new(range.from() - first_row_id,
range.to() - first_row_id)
+ })
+ })
+ .collect(),
+ );
local_ranges = intersect_local_ranges(&local_ranges, &selected);
}
@@ -1808,6 +1909,14 @@ fn build_source_plan_with_row_id_pushdown(
if let Some(source_idx) = source_idx {
let field_offset =
sources[source_idx].add_read_field(field.clone());
column_plan.push(Some((source_idx, field_offset)));
+ } else if !field.data_type().is_nullable() {
+ return Err(Error::DataInvalid {
+ message: format!(
+ "Cannot read non-nullable field '{}' without a provider",
+ field.name()
+ ),
+ source: None,
+ });
} else {
column_plan.push(None);
}
@@ -2153,6 +2262,29 @@ fn blob_file_row_range(file: &DataFileMeta) ->
crate::Result<RowRange> {
Ok(RowRange::new(first_row_id, last_row_id))
}
+fn vector_file_row_range(file: &DataFileMeta) -> crate::Result<RowRange> {
+ let first_row_id = file.first_row_id.ok_or_else(|| Error::DataInvalid {
+ message: format!("Vector file '{}' is missing first_row_id",
file.file_name),
+ source: None,
+ })?;
+ if file.row_count <= 0 {
+ return Err(Error::DataInvalid {
+ message: format!(
+ "Vector file '{}' row count must be positive, got {}",
+ file.file_name, file.row_count
+ ),
+ source: None,
+ });
+ }
+ let last_row_id = first_row_id
+ .checked_add(file.row_count - 1)
+ .ok_or_else(|| Error::DataInvalid {
+ message: format!("Vector file '{}' row range overflows i64",
file.file_name),
+ source: None,
+ })?;
+ Ok(RowRange::new(first_row_id, last_row_id))
+}
+
/// Aggregates rolled `.vector.<format>` segments belonging to one logical
vector
/// source, mirroring upstream `VectorFileBunch` non-pushdown semantics. Unlike
/// `BlobBunch`, the expected row count is taken directly from the prepared
group's
@@ -2407,6 +2539,25 @@ mod tests {
use blob_test_utils::{write_blob_file, write_blob_file_with_values,
BlobFixtureValue};
use test_utils::{local_file_path, write_int_parquet_file};
+ #[test]
+ fn test_selected_absolute_row_ranges_normalizes_before_intersection() {
+ let selected = selected_absolute_row_ranges_for_file(
+ 0,
+ 6,
+ Some(&[
+ RowRange::new(4, 4),
+ RowRange::new(2, 3),
+ RowRange::new(0, 2),
+ RowRange::new(4, 4),
+ ]),
+ None,
+ )
+ .unwrap()
+ .unwrap();
+
+ assert_eq!(selected, vec![RowRange::new(0, 4)]);
+ }
+
#[tokio::test]
async fn test_descriptor_columns_resolve_concurrently_and_preserve_order()
{
let schema = Arc::new(arrow_schema::Schema::new(vec![
@@ -3203,6 +3354,33 @@ mod tests {
assert_eq!(source_plan.column_plan, vec![Some((0, 0)), Some((2, 0))]);
}
+ #[test]
+ fn test_row_id_pushdown_rejects_missing_non_nullable_provider() {
+ let files = vec![data_file("data.parquet", 0, 4, 1, None)];
+ let prepared_group = PreparedMergeGroup::new(&files).unwrap();
+ let file_infos = vec![resolved_info(vec![1])];
+ let read_type = vec![
+ DataField::new(1, "id".to_string(), DataType::Int(IntType::new())),
+ DataField::new(
+ 2,
+ "payload".to_string(),
+ DataType::Blob(BlobType::with_nullable(false)),
+ ),
+ ];
+
+ let err = build_source_plan_with_row_id_pushdown(
+ &prepared_group,
+ &file_infos,
+ &read_type,
+ &HashSet::new(),
+ true,
+ )
+ .unwrap_err();
+
+ assert!(matches!(err, Error::DataInvalid { message, .. }
+ if message.contains("non-nullable field 'payload'")));
+ }
+
#[test]
fn
test_build_source_plan_validates_complete_range_only_for_read_blob_source() {
let files = vec![
@@ -4065,7 +4243,7 @@ mod tests {
}
#[tokio::test]
- async fn test_selected_blob_fallback_rejects_uncovered_non_deleted_range()
{
+ async fn test_selected_blob_fallback_fills_uncovered_nullable_range() {
use BlobFixtureValue::{Placeholder, Value};
let tempdir = tempdir().unwrap();
@@ -4138,8 +4316,8 @@ mod tests {
.build()
.unwrap();
- // Row 0 is outside the selection and row 1 is deleted, so only the
- // uncovered, selected, non-deleted row 3 must make the read fail.
+ // Row 0 is outside the selection and row 1 is deleted. Row 2 is
covered
+ // by the BLOB provider, while the uncovered row 3 must be NULL-filled.
for blob_as_descriptor in [false, true] {
let mode_table = table.copy_with_options(HashMap::from([(
"blob-as-descriptor".to_string(),
@@ -4150,18 +4328,20 @@ mod tests {
mode_table.schema().fields().to_vec(),
Vec::new(),
);
- let mut stream =
read.to_arrow(std::slice::from_ref(&split)).unwrap();
- let first = stream.try_next().await;
- assert!(
- matches!(
- &first,
- Err(Error::DataInvalid { message, .. })
- if message.contains(
- "uncovered ranges are [RowRange { from: 3, to: 3
}]"
- )
- ),
- "blob_as_descriptor={blob_as_descriptor}: expected uncovered
selected BLOB range error, got {first:?}"
- );
+ let batches = read
+ .to_arrow(std::slice::from_ref(&split))
+ .unwrap()
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap();
+ assert_eq!(collect_int_values(&batches, "id"), vec![3, 4]);
+ let payloads = collect_binary_values(&batches, "payload");
+ assert_eq!(payloads.len(), 2);
+ assert!(payloads[0].is_some());
+ assert_eq!(payloads[1], None);
+ if !blob_as_descriptor {
+ assert_eq!(payloads[0], Some(b"covered".to_vec()));
+ }
}
}
@@ -5429,12 +5609,14 @@ mod tests {
.with_row_ranges(vec![RowRange::new(2, 2), RowRange::new(4, 4)])
.build()
.unwrap();
- let mut stream = read.to_arrow(&[uncovered_split]).unwrap();
- let error = stream.try_next().await.unwrap_err();
- assert!(matches!(error, Error::DataInvalid { message, .. }
- if message.contains(
- "does not cover effective selected row ranges [RowRange { from: 2,
to: 2 }]"
- )));
+ let batches = read
+ .to_arrow(&[uncovered_split])
+ .unwrap()
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap();
+ assert_eq!(collect_int_values(&batches, "id"), vec![3, 5]);
+ assert_fixed_size_list(&batches, "embedding", 2, &[None,
Some(vec![5.0, 5.0])]);
let full_split = DataSplitBuilder::new()
.with_snapshot(1)
@@ -5592,7 +5774,7 @@ mod tests {
}
#[tokio::test]
- async fn test_scan_and_read_prunes_unselected_rolled_dedicated_sources() {
+ async fn test_scan_and_read_prunes_rolled_dedicated_ranges() {
let tempdir = tempdir().unwrap();
let table_path = local_file_path(tempdir.path());
let bucket_dir = tempdir.path().join("bucket-0");
@@ -5754,7 +5936,7 @@ mod tests {
}
#[tokio::test]
- async fn test_scan_and_read_rejects_selected_gap_after_dedicated_pruning()
{
+ async fn test_scan_and_read_null_fills_selected_dedicated_gap() {
let tempdir = tempdir().unwrap();
let table_path = local_file_path(tempdir.path());
let bucket_dir = tempdir.path().join("bucket-0");
@@ -5828,16 +6010,99 @@ mod tests {
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();
+ let batch = stream.try_next().await.unwrap().unwrap();
+ assert_eq!(
+ collect_int_values(std::slice::from_ref(&batch), "id"),
+ vec![3]
+ );
+ let column = batch.column_by_name(field).unwrap();
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:?}"
+ column.is_null(0),
+ "missing {provider_kind} range should be null-filled"
);
}
}
+ #[tokio::test]
+ async fn test_btree_hit_null_fills_missing_nullable_blob() {
+ 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])],
None);
+ let mut normal_file = data_file_meta_with_path(
+ "data.parquet",
+ 0,
+ 3,
+ 1,
+ normal_path.metadata().unwrap().len() as i64,
+ Some(vec!["id"]),
+ );
+ normal_file.first_row_id = None;
+ normal_file.file_source = Some(0);
+
+ let table = Table::new(
+ FileIOBuilder::new("file").build().unwrap(),
+ Identifier::new("default", "btree_missing_blob_t"),
+ table_path,
+ TableSchema::new(
+ 0,
+ &Schema::builder()
+ .column("id", DataType::Int(IntType::new()))
+ .column("payload", DataType::Blob(BlobType::new()))
+ .option("data-evolution.enabled", "true")
+ .option("row-tracking.enabled", "true")
+ .option("global-index.enabled", "true")
+ .option("sorted-index.records-per-range", "10")
+ .build()
+ .unwrap(),
+ ),
+ None,
+ );
+ TableCommit::new(table.clone(), "btree-missing-blob-test".to_string())
+ .commit(vec![CommitMessage::new(
+ BinaryRowBuilder::new(0).build_serialized(),
+ 0,
+ vec![normal_file],
+ )])
+ .await
+ .unwrap();
+ table
+ .new_btree_global_index_build_builder()
+ .with_index_column("id")
+ .execute()
+ .await
+ .unwrap();
+
+ let mut builder = table.new_read_builder();
+ builder.with_projection(&["id", "payload"]).unwrap();
+ builder.with_filter(
+ PredicateBuilder::new(table.schema().fields())
+ .equal("id", Datum::Int(2))
+ .unwrap(),
+ );
+ let plan = builder.new_scan().plan().await.unwrap();
+ assert_eq!(
+ plan.splits()[0].row_ranges(),
+ Some(&[RowRange::new(1, 1)][..])
+ );
+
+ let batches = builder
+ .new_read()
+ .unwrap()
+ .to_arrow(plan.splits())
+ .unwrap()
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap();
+ assert_eq!(collect_int_values(&batches, "id"), vec![2]);
+ assert_eq!(collect_binary_values(&batches, "payload"), vec![None]);
+ }
+
/// (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/data_evolution_reader/blob_fallback.rs
b/crates/paimon/src/table/data_evolution_reader/blob_fallback.rs
index 4331b4d6..bc538df1 100644
--- a/crates/paimon/src/table/data_evolution_reader/blob_fallback.rs
+++ b/crates/paimon/src/table/data_evolution_reader/blob_fallback.rs
@@ -106,11 +106,13 @@ impl LazyBlobFile {
}
}
+#[allow(clippy::too_many_arguments)]
pub(super) fn read(
split: &DataSplit,
bunch: BlobBunch,
read_fields: Vec<DataField>,
row_ranges: Option<Vec<RowRange>>,
+ batch_size: Option<usize>,
file_io: FileIO,
blob_as_descriptor: bool,
anchor_deletion_vector: Option<DeletionVectorContext>,
@@ -124,6 +126,7 @@ pub(super) fn read(
let target_schema = build_target_arrow_schema(&read_fields)?;
let array_field = matches!(read_fields[0].data_type(), DataType::Array(_));
+ let batch_size = batch_size.unwrap_or(BATCH_SIZE).max(1);
let split = split.clone();
Ok(try_stream! {
@@ -161,8 +164,8 @@ pub(super) fn read(
}
}
- let mut row_cursor = RowIdBatchCursor::new(selected_ranges);
- while let Some(row_ids) = row_cursor.next_batch(BATCH_SIZE) {
+ let mut row_cursor = RowIdBatchCursor::new(selected_ranges,
batch_size);
+ while let Some(row_ids) = row_cursor.next_batch() {
yield resolve_batch(
&mut sequence_groups,
&row_ids,
@@ -270,21 +273,23 @@ struct RowIdBatchCursor {
ranges: Vec<RowRange>,
range_index: usize,
next_row_id: Option<i64>,
+ batch_size: usize,
}
impl RowIdBatchCursor {
- fn new(ranges: Vec<RowRange>) -> Self {
+ fn new(ranges: Vec<RowRange>, batch_size: usize) -> Self {
let next_row_id = ranges.first().map(RowRange::from);
Self {
ranges,
range_index: 0,
next_row_id,
+ batch_size,
}
}
- fn next_batch(&mut self, batch_size: usize) -> Option<Vec<i64>> {
- let mut row_ids = Vec::with_capacity(batch_size);
- while row_ids.len() < batch_size {
+ fn next_batch(&mut self) -> Option<Vec<i64>> {
+ let mut row_ids = Vec::with_capacity(self.batch_size);
+ while row_ids.len() < self.batch_size {
let Some(row_id) = self.next_row_id else {
break;
};
@@ -312,6 +317,16 @@ mod tests {
use std::ops::Range;
use std::sync::atomic::{AtomicUsize, Ordering};
+ #[test]
+ fn test_row_id_batch_cursor_honors_batch_size() {
+ let mut cursor = RowIdBatchCursor::new(vec![RowRange::new(0, 4)], 2);
+
+ assert_eq!(cursor.next_batch(), Some(vec![0, 1]));
+ assert_eq!(cursor.next_batch(), Some(vec![2, 3]));
+ assert_eq!(cursor.next_batch(), Some(vec![4]));
+ assert_eq!(cursor.next_batch(), None);
+ }
+
#[allow(dead_code)]
mod blob_test_utils {
include!(concat!(env!("CARGO_MANIFEST_DIR"), "/blob_test_utils.rs"));
diff --git a/crates/paimon/src/table/source.rs
b/crates/paimon/src/table/source.rs
index 933cb441..a47dc582 100644
--- a/crates/paimon/src/table/source.rs
+++ b/crates/paimon/src/table/source.rs
@@ -114,6 +114,14 @@ pub fn merge_row_ranges(mut ranges: Vec<RowRange>) ->
Vec<RowRange> {
if ranges.len() <= 1 {
return ranges;
}
+ if ranges.windows(2).all(|pair| {
+ pair[0]
+ .to
+ .checked_add(1)
+ .is_some_and(|next| next < pair[1].from)
+ }) {
+ return ranges;
+ }
ranges.sort_by_key(|r| r.from);
let mut merged: Vec<RowRange> = Vec::with_capacity(ranges.len());
let mut iter = ranges.into_iter();
diff --git a/crates/paimon/src/table/table_scan.rs
b/crates/paimon/src/table/table_scan.rs
index 632b6ee3..ded0cab0 100644
--- a/crates/paimon/src/table/table_scan.rs
+++ b/crates/paimon/src/table/table_scan.rs
@@ -167,7 +167,7 @@ async fn read_all_manifest_entries(
if let Some(index) = row_range_index {
let before = manifest_files.len();
- retain_manifest_row_range_components(&mut manifest_files, index);
+ retain_manifest_row_ranges(&mut manifest_files, index);
if let Some(trace) = trace.as_deref_mut() {
trace.manifest_files_pruned_by_row_ranges = before -
manifest_files.len();
}
@@ -271,7 +271,7 @@ async fn read_all_manifest_entries(
let manifest_entries_after_merge = all_entries.len();
if let Some(index) = row_range_index {
let before = all_entries.len();
- all_entries = retain_live_manifest_entry_row_range_groups(all_entries,
index);
+ all_entries = retain_manifest_entry_row_ranges(all_entries, index);
counters.pruned_by_row_ranges = before - all_entries.len();
}
if let Some(trace) = trace {
@@ -288,7 +288,6 @@ async fn read_all_manifest_entries(
Ok(all_entries)
}
-#[cfg(test)]
fn manifest_file_overlaps_row_range_index(
manifest: &crate::spec::ManifestFileMeta,
row_range_index: &RowRangeIndex,
@@ -303,64 +302,13 @@ fn manifest_row_id_range(manifest:
&crate::spec::ManifestFileMeta) -> Option<(i6
}
}
-fn retain_manifest_row_range_components(
+fn retain_manifest_row_ranges(
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
- });
+ manifests.retain(|manifest|
manifest_file_overlaps_row_range_index(manifest, row_range_index));
}
-#[cfg(test)]
fn data_file_overlaps_row_range_index(
file: &DataFileMeta,
row_range_index: &RowRangeIndex,
@@ -369,154 +317,16 @@ fn data_file_overlaps_row_range_index(
.is_none_or(|(from, to)| row_range_index.intersects(from, to))
}
-fn retain_live_manifest_entry_row_range_groups(
+fn retain_manifest_entry_row_ranges(
entries: Vec<ManifestEntry>,
row_range_index: &RowRangeIndex,
) -> Vec<ManifestEntry> {
- debug_assert!(entries.iter().all(|entry| *entry.kind() == FileKind::Add));
- 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 {
- retain_selected_row_range_component(
- &entries,
- &mut keep,
- &component,
- component_from,
- component_to,
- row_range_index,
- );
- component.clear();
- component_from = from;
- component_to = to;
- component.push(idx);
- }
- }
- retain_selected_row_range_component(
- &entries,
- &mut keep,
- &component,
- component_from,
- component_to,
- row_range_index,
- );
- }
- drop(buckets);
-
entries
.into_iter()
- .enumerate()
- .filter_map(|(idx, entry)| keep[idx].then_some(entry))
+ .filter(|entry| data_file_overlaps_row_range_index(entry.file(),
row_range_index))
.collect()
}
-#[derive(Debug, Clone, PartialEq, Eq, Hash)]
-enum DataEvolutionProviderKey<'a> {
- Normal,
- // Blob bunches span schema versions and resolve to one write column.
- Blob(Option<&'a [String]>),
- Vector(i64, String, Option<Vec<&'a str>>),
-}
-
-fn data_evolution_provider_key(file: &DataFileMeta) ->
DataEvolutionProviderKey<'_> {
- if
crate::table::dedicated_format_file_writer::is_blob_file_name(&file.file_name) {
- DataEvolutionProviderKey::Blob(file.write_cols.as_deref())
- } else if is_vector_store_file_name(&file.file_name) {
- let write_cols = file.write_cols.as_ref().map(|cols| {
- let mut cols = cols.iter().map(String::as_str).collect::<Vec<_>>();
- cols.sort_unstable();
- cols
- });
- DataEvolutionProviderKey::Vector(
- file.schema_id,
- file.file_name
- .rsplit('.')
- .next()
- .unwrap_or("")
- .to_ascii_lowercase(),
- write_cols,
- )
- } else {
- DataEvolutionProviderKey::Normal
- }
-}
-
-fn retain_selected_row_range_component(
- entries: &[ManifestEntry],
- keep: &mut [bool],
- component: &[usize],
- component_from: i64,
- component_to: i64,
- row_range_index: &RowRangeIndex,
-) {
- if component.is_empty() || !row_range_index.intersects(component_from,
component_to) {
- return;
- }
-
- // Keep files that can contribute to the selected rows.
- let mut selected_providers = HashSet::new();
- for &idx in component {
- let file = entries[idx].file();
- let (from, to) = file.row_id_range().expect("validated row-id range");
- if row_range_index.intersects(from, to) {
- keep[idx] = true;
- selected_providers.insert(data_evolution_provider_key(file));
- }
- }
-
- // Keep one missing-provider witness for gap validation.
- for &idx in component {
- if keep[idx] {
- continue;
- }
- let provider = data_evolution_provider_key(entries[idx].file());
- if selected_providers.insert(provider) {
- keep[idx] = true;
- }
- }
-}
-
fn data_evolution_row_range_groups(
data_files: Vec<DataFileMeta>,
row_ranges: Option<&[RowRange]>,
@@ -1712,7 +1522,7 @@ impl<'a> PaimonTableScan<'a> {
}
}
if let Some(index) = row_range_index {
- retain_manifest_row_range_components(&mut manifest_metas, index);
+ retain_manifest_row_ranges(&mut manifest_metas, index);
}
let bucket_key_fields: Vec<DataField> = if
self.bucket_predicate.is_none() {
@@ -1787,7 +1597,7 @@ impl<'a> PaimonTableScan<'a> {
}
let entries = merge_manifest_entries(entries);
let entries = if let Some(index) = row_range_index {
- retain_live_manifest_entry_row_range_groups(entries, index)
+ retain_manifest_entry_row_ranges(entries, index)
} else {
entries
};
@@ -2197,9 +2007,8 @@ impl<'a> PaimonTableScan<'a> {
mod tests {
use super::{
data_evolution_row_range_groups, data_file_overlaps_row_range_index,
- manifest_file_overlaps_row_range_index, merge_manifest_entries,
- prune_data_evolution_group_by_read_fields, retain_index_manifest_entry,
- retain_live_manifest_entry_row_range_groups,
retain_manifest_row_range_components,
+ manifest_file_overlaps_row_range_index,
prune_data_evolution_group_by_read_fields,
+ retain_index_manifest_entry, retain_manifest_entry_row_ranges,
retain_manifest_row_ranges,
should_skip_level_zero_for_scan, split_row_ranges_for_files,
LimitPushdownAccumulator,
PaimonTableScan, RowRangeIndex, TableScan,
};
@@ -2303,7 +2112,7 @@ mod tests {
}
#[test]
- fn test_manifest_row_range_pruning_retains_overlapping_component() {
+ fn test_manifest_row_range_pruning_uses_each_envelope() {
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| {
@@ -2317,19 +2126,19 @@ mod tests {
manifest("other-group", 10, 15),
];
- retain_manifest_row_range_components(&mut manifests, &index);
+ retain_manifest_row_ranges(&mut manifests, &index);
assert_eq!(
manifests
.iter()
.map(ManifestFileMeta::file_name)
.collect::<Vec<_>>(),
- vec!["anchor", "left-dedicated", "right-dedicated"]
+ vec!["anchor"]
);
}
#[test]
- fn test_manifest_row_range_component_pruning_fails_open_on_unknown_range()
{
+ fn test_manifest_row_range_pruning_fails_open_per_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![
@@ -2342,13 +2151,19 @@ mod tests {
.with_row_id_stats(Some(10), Some(15)),
];
- retain_manifest_row_range_components(&mut manifests, &index);
+ retain_manifest_row_ranges(&mut manifests, &index);
- assert_eq!(manifests.len(), 4);
+ assert_eq!(
+ manifests
+ .iter()
+ .map(ManifestFileMeta::file_name)
+ .collect::<Vec<_>>(),
+ vec!["unknown", "inverted", "one-sided"]
+ );
}
#[test]
- fn
test_manifest_entry_row_range_pruning_drops_disjoint_files_from_wide_component()
{
+ fn test_manifest_entry_row_range_pruning_uses_each_file_range() {
let index = RowRangeIndex::create(vec![RowRange::new(2, 2)]);
let entry = |name: &str, first_row_id, row_count| {
ManifestEntry::new(
@@ -2367,7 +2182,7 @@ mod tests {
entry("other-group", 10, 2),
];
- let retained = retain_live_manifest_entry_row_range_groups(entries,
&index);
+ let retained = retain_manifest_entry_row_ranges(entries, &index);
assert_eq!(
retained
@@ -2379,150 +2194,7 @@ mod tests {
}
#[test]
- fn
test_manifest_entry_row_range_pruning_prunes_rolled_blob_sidecars_behind_wide_anchor()
{
- let index = RowRangeIndex::create(vec![RowRange::new(120, 129)]);
- let entry = |name: &str, first_row_id, row_count, schema_id,
write_cols: &[&str]| {
- let mut file = make_evo_file_with_cols(name, row_count, 0,
first_row_id, write_cols);
- file.schema_id = schema_id;
- ManifestEntry::new(FileKind::Add, Vec::new(), 0, 1, file, 3)
- };
- let entries = vec![
- entry("base.parquet", 0, 1_000, 3, &["record_index"]),
- entry("image-0.blob", 0, 100, 1, &["image"]),
- entry("image-1.blob", 100, 100, 3, &["image"]),
- entry("image-2.blob", 200, 100, 2, &["image"]),
- entry("image-3.blob", 300, 100, 3, &["image"]),
- ];
-
- let retained = retain_live_manifest_entry_row_range_groups(entries,
&index);
-
- assert_eq!(
- retained
- .iter()
- .map(|entry| entry.file().file_name.as_str())
- .collect::<Vec<_>>(),
- vec!["base.parquet", "image-1.blob"]
- );
- }
-
- #[test]
- fn test_vector_provider_key_distinguishes_schema_ids() {
- let index = RowRangeIndex::create(vec![RowRange::new(120, 129)]);
- let entry = |name: &str, first_row_id, row_count, schema_id,
write_cols: &[&str]| {
- let mut file = make_evo_file_with_cols(name, row_count, 0,
first_row_id, write_cols);
- file.schema_id = schema_id;
- ManifestEntry::new(FileKind::Add, Vec::new(), 0, 1, file, 3)
- };
- let entries = vec![
- entry("base.parquet", 0, 1_000, 3, &["id"]),
- entry("old.vector.parquet", 0, 100, 1, &["embedding"]),
- entry("new.vector.parquet", 100, 100, 3, &["embedding"]),
- ];
-
- let retained = retain_live_manifest_entry_row_range_groups(entries,
&index);
-
- assert_eq!(
- retained
- .iter()
- .map(|entry| entry.file().file_name.as_str())
- .collect::<Vec<_>>(),
- vec!["base.parquet", "old.vector.parquet", "new.vector.parquet"]
- );
- }
-
- #[test]
- fn test_vector_provider_key_normalizes_write_col_order() {
- let index = RowRangeIndex::create(vec![RowRange::new(120, 129)]);
- let entry = |name: &str, first_row_id, row_count, write_cols: &[&str]|
{
- ManifestEntry::new(
- FileKind::Add,
- Vec::new(),
- 0,
- 1,
- make_evo_file_with_cols(name, row_count, 0, first_row_id,
write_cols),
- 3,
- )
- };
- let entries = vec![
- entry("base.parquet", 0, 1_000, &["id"]),
- entry("old.vector.parquet", 0, 100, &["velocity", "embedding"]),
- entry("new.vector.parquet", 100, 100, &["embedding", "velocity"]),
- ];
-
- let retained = retain_live_manifest_entry_row_range_groups(entries,
&index);
-
- assert_eq!(
- retained
- .iter()
- .map(|entry| entry.file().file_name.as_str())
- .collect::<Vec<_>>(),
- vec!["base.parquet", "new.vector.parquet"]
- );
- }
-
- #[test]
- fn
test_manifest_entry_row_range_pruning_keeps_provider_witness_for_selected_gap()
{
- let index = RowRangeIndex::create(vec![RowRange::new(2, 2)]);
- let entry = |name: &str, first_row_id, row_count, write_cols: &[&str]|
{
- ManifestEntry::new(
- FileKind::Add,
- Vec::new(),
- 0,
- 1,
- make_evo_file_with_cols(name, row_count, 0, first_row_id,
write_cols),
- 3,
- )
- };
- let entries = vec![
- entry("base.parquet", 0, 6, &["id"]),
- entry("payload-left.blob", 0, 2, &["payload"]),
- entry("payload-right.blob", 4, 2, &["payload"]),
- ];
-
- let retained = retain_live_manifest_entry_row_range_groups(entries,
&index);
-
- assert_eq!(
- retained
- .iter()
- .map(|entry| entry.file().file_name.as_str())
- .collect::<Vec<_>>(),
- vec!["base.parquet", "payload-left.blob"]
- );
- }
-
- #[test]
- fn
test_manifest_entry_row_range_pruning_does_not_resurrect_deleted_provider_witness()
{
- let index = RowRangeIndex::create(vec![RowRange::new(2, 2)]);
- let entry = |kind, name: &str, first_row_id, row_count, write_cols:
&[&str]| {
- ManifestEntry::new(
- kind,
- Vec::new(),
- 0,
- 1,
- make_evo_file_with_cols(name, row_count, 0, first_row_id,
write_cols),
- 3,
- )
- };
- let entries = vec![
- entry(FileKind::Add, "base.parquet", 0, 6, &["id"]),
- entry(FileKind::Add, "deleted-payload.blob", 0, 2, &["payload"]),
- entry(FileKind::Delete, "deleted-payload.blob", 0, 2,
&["payload"]),
- ];
-
- let live_entries = merge_manifest_entries(entries);
- let retained =
retain_live_manifest_entry_row_range_groups(live_entries, &index);
-
- assert_eq!(
- retained
- .iter()
- .map(|entry| entry.file().file_name.as_str())
- .collect::<Vec<_>>(),
- vec!["base.parquet"]
- );
- }
-
- #[test]
- fn test_manifest_entry_row_range_pruning_fails_open_per_bucket() {
+ fn test_manifest_entry_row_range_pruning_fails_open_per_file() {
let index = RowRangeIndex::create(vec![RowRange::new(2, 2)]);
let entry = |name: &str, bucket, first_row_id, row_count| {
ManifestEntry::new(
@@ -2542,22 +2214,14 @@ mod tests {
entry("bucket-1-outside", 1, Some(10), 2),
];
- let retained = retain_live_manifest_entry_row_range_groups(entries,
&index);
+ let retained = retain_manifest_entry_row_ranges(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",
- ]
- );
+ assert_eq!(names, vec!["overflow", "unknown-anchor", "zero-count"]);
}
#[test]