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 6eef29ed fix(table): skip unrelated BLOB and vector files in row-range
reads (#629)
6eef29ed is described below
commit 6eef29ed1152c893eb9783e99af6e609f769199d
Author: XiaoHongbo <[email protected]>
AuthorDate: Thu Jul 30 09:12:14 2026 +0800
fix(table): skip unrelated BLOB and vector files in row-range reads (#629)
---
crates/paimon/src/table/data_evolution_reader.rs | 10 +-
crates/paimon/src/table/table_scan.rs | 375 ++++++++++++++++++++---
2 files changed, 340 insertions(+), 45 deletions(-)
diff --git a/crates/paimon/src/table/data_evolution_reader.rs
b/crates/paimon/src/table/data_evolution_reader.rs
index 648725df..3e62cd3e 100644
--- a/crates/paimon/src/table/data_evolution_reader.rs
+++ b/crates/paimon/src/table/data_evolution_reader.rs
@@ -5592,7 +5592,7 @@ mod tests {
}
#[tokio::test]
- async fn test_scan_and_read_retains_complete_rolled_dedicated_group() {
+ async fn test_scan_and_read_prunes_unselected_rolled_dedicated_sources() {
let tempdir = tempdir().unwrap();
let table_path = local_file_path(tempdir.path());
let bucket_dir = tempdir.path().join("bucket-0");
@@ -5688,15 +5688,11 @@ mod tests {
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);
+ assert_eq!(trace.manifest_entries_pruned_by_row_ranges, 4);
+ assert_eq!(trace.final_files, 3);
let batches = builder
.new_read()
diff --git a/crates/paimon/src/table/table_scan.rs
b/crates/paimon/src/table/table_scan.rs
index fa9b7100..632b6ee3 100644
--- a/crates/paimon/src/table/table_scan.rs
+++ b/crates/paimon/src/table/table_scan.rs
@@ -267,11 +267,12 @@ async fn read_all_manifest_entries(
counters.merge(manifest_counters);
all_entries.extend(entries);
}
+ let mut all_entries = merge_manifest_entries(all_entries);
+ let manifest_entries_after_merge = all_entries.len();
if let Some(index) = row_range_index {
let before = all_entries.len();
- all_entries = retain_manifest_entry_row_range_groups(all_entries,
index);
+ all_entries = retain_live_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;
@@ -282,6 +283,7 @@ async fn read_all_manifest_entries(
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;
+ trace.manifest_entries_after_merge = manifest_entries_after_merge;
}
Ok(all_entries)
}
@@ -367,10 +369,11 @@ fn data_file_overlaps_row_range_index(
.is_none_or(|(from, to)| row_range_index.intersects(from, to))
}
-fn retain_manifest_entry_row_range_groups(
+fn retain_live_manifest_entry_row_range_groups(
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
@@ -416,23 +419,28 @@ fn retain_manifest_entry_row_range_groups(
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();
- }
+ 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);
}
}
- if !component.is_empty() && row_range_index.intersects(component_from,
component_to) {
- for component_idx in component {
- keep[component_idx] = true;
- }
- }
+ retain_selected_row_range_component(
+ &entries,
+ &mut keep,
+ &component,
+ component_from,
+ component_to,
+ row_range_index,
+ );
}
drop(buckets);
@@ -443,6 +451,72 @@ fn retain_manifest_entry_row_range_groups(
.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]>,
@@ -1207,7 +1281,7 @@ impl<'a> PaimonTableScan<'a> {
&self,
snapshot: &Snapshot,
row_range_index: Option<&RowRangeIndex>,
- mut trace: Option<&mut ScanTrace>,
+ trace: Option<&mut ScanTrace>,
) -> crate::Result<Vec<ManifestEntry>> {
let file_io = self.table.file_io();
let table_path = self.table.location();
@@ -1281,14 +1355,10 @@ impl<'a> PaimonTableScan<'a> {
&bucket_key_fields,
bucket_function_type,
row_range_index,
- trace.as_deref_mut(),
+ trace,
)
.await?;
- let merged = merge_manifest_entries(entries);
- if let Some(trace) = trace {
- trace.manifest_entries_after_merge = merged.len();
- }
- Ok(merged)
+ Ok(entries)
}
fn can_push_down_limit_hint(&self, row_ranges: Option<&[RowRange]>) ->
bool {
@@ -1715,12 +1785,9 @@ impl<'a> PaimonTableScan<'a> {
)?;
entries.extend(manifest_entries);
}
- let entries = entries
- .into_iter()
- .filter(|entry| *entry.kind() == FileKind::Add)
- .collect::<Vec<_>>();
+ let entries = merge_manifest_entries(entries);
let entries = if let Some(index) = row_range_index {
- retain_manifest_entry_row_range_groups(entries, index)
+ retain_live_manifest_entry_row_range_groups(entries, index)
} else {
entries
};
@@ -2130,11 +2197,11 @@ 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,
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,
+ 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,
+ should_skip_level_zero_for_scan, split_row_ranges_for_files,
LimitPushdownAccumulator,
+ PaimonTableScan, RowRangeIndex, TableScan,
};
use crate::catalog::Identifier;
use crate::io::FileIOBuilder;
@@ -2281,7 +2348,7 @@ mod tests {
}
#[test]
- fn test_manifest_entry_row_range_pruning_retains_overlapping_group() {
+ fn
test_manifest_entry_row_range_pruning_drops_disjoint_files_from_wide_component()
{
let index = RowRangeIndex::create(vec![RowRange::new(2, 2)]);
let entry = |name: &str, first_row_id, row_count| {
ManifestEntry::new(
@@ -2300,14 +2367,157 @@ mod tests {
entry("other-group", 10, 2),
];
- let retained = retain_manifest_entry_row_range_groups(entries, &index);
+ 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!["anchor", "left-dedicated", "right-dedicated"]
+ vec!["anchor"]
+ );
+ }
+
+ #[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"]
);
}
@@ -2332,7 +2542,7 @@ mod tests {
entry("bucket-1-outside", 1, Some(10), 2),
];
- let retained = retain_manifest_entry_row_range_groups(entries, &index);
+ let retained = retain_live_manifest_entry_row_range_groups(entries,
&index);
let mut names = retained
.into_iter()
@@ -2895,8 +3105,8 @@ mod tests {
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.manifest_entries_after_manifest_filters, 4);
+ assert_eq!(trace.manifest_entries_after_merge, 4);
assert_eq!(trace.data_evolution_groups_before_stats, 1);
assert_eq!(trace.data_evolution_groups_pruned_by_row_ranges, 0);
@@ -2920,6 +3130,95 @@ mod tests {
assert_eq!(delta_files, vec!["a-new", "a-old"]);
}
+ #[tokio::test]
+ async fn test_row_range_trace_excludes_manifest_netting() {
+ let table_path = "memory:/de_row_range_trace_netting";
+ let table = data_evolution_test_table(table_path, two_column_schema(0,
"id", "name"));
+ setup_scan_trace_dirs(&table).await;
+
+ let deleted = make_evo_file_with_cols("deleted.parquet", 10, 1, 0,
&["id"]);
+ let retained = make_evo_file_with_cols("retained.parquet", 10, 2, 0,
&["id"]);
+ let pruned = make_evo_file_with_cols("pruned.parquet", 10, 3, 100,
&["id"]);
+ let partition = BinaryRowBuilder::new(0).build_serialized();
+ TableCommit::new(table.clone(), "row-range-netting-add".to_string())
+ .commit(vec![CommitMessage::new(
+ partition.clone(),
+ 0,
+ vec![deleted.clone(), retained, pruned],
+ )])
+ .await
+ .unwrap();
+
+ let mut delete = CommitMessage::new(partition, 0, Vec::new());
+ delete.deleted_files = vec![deleted];
+ TableCommit::new(table.clone(), "row-range-netting-delete".to_string())
+ .commit(vec![delete])
+ .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();
+
+ assert_eq!(
+ plan.splits()[0].data_files()[0].file_name,
+ "retained.parquet"
+ );
+ assert_eq!(trace.manifest_entries_read, 4);
+ assert_eq!(trace.manifest_entries_after_manifest_filters, 4);
+ assert_eq!(trace.manifest_entries_after_merge, 2);
+ assert_eq!(trace.manifest_entries_pruned_by_row_ranges, 1);
+ }
+
+ #[tokio::test]
+ async fn
test_snapshot_delta_row_range_pruning_does_not_resurrect_deleted_witness() {
+ let table_path = "memory:/de_delta_row_range_netting";
+ let table = data_evolution_test_table(table_path, two_column_schema(0,
"id", "payload"));
+ setup_scan_trace_dirs(&table).await;
+
+ let deleted = make_evo_file_with_cols("deleted.blob", 100, 1, 0,
&["payload"]);
+ let partition = BinaryRowBuilder::new(0).build_serialized();
+ TableCommit::new(table.clone(), "delta-netting-base".to_string())
+ .commit(vec![CommitMessage::new(
+ partition.clone(),
+ 0,
+ vec![deleted.clone()],
+ )])
+ .await
+ .unwrap();
+
+ let anchor = make_evo_file_with_cols("anchor.parquet", 1_000, 2, 0,
&["id"]);
+ let mut replacement = CommitMessage::new(partition, 0, vec![anchor,
deleted.clone()]);
+ replacement.deleted_files = vec![deleted];
+ TableCommit::new(table.clone(),
"delta-netting-replacement".to_string())
+ .commit(vec![replacement])
+ .await
+ .unwrap();
+
+ let snapshot = table
+ .snapshot_manager()
+ .get_latest_snapshot()
+ .await
+ .unwrap()
+ .unwrap();
+ let mut read_builder = table.new_read_builder();
+ read_builder.with_row_ranges(vec![RowRange::new(120, 129)]);
+ let plan = read_builder
+ .new_scan()
+ .plan_snapshot_delta(&snapshot)
+ .await
+ .unwrap();
+
+ assert_eq!(
+ plan.splits()
+ .iter()
+ .flat_map(|split| split.data_files())
+ .map(|file| file.file_name.as_str())
+ .collect::<Vec<_>>(),
+ vec!["anchor.parquet"]
+ );
+ }
+
#[tokio::test]
async fn test_data_evolution_prunes_files_without_projected_columns() {
let table =