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 4dcf2fa fix(scan): prune data-evolution files by projection (#442)
4dcf2fa is described below
commit 4dcf2fa3345a6bdaf4d2a47bc3816c9e2f0311cb
Author: QuakeWang <[email protected]>
AuthorDate: Fri Jul 3 11:23:44 2026 +0800
fix(scan): prune data-evolution files by projection (#442)
---
crates/integrations/datafusion/src/table/mod.rs | 223 +++++++++++
crates/paimon/src/table/read_builder.rs | 174 +++++++--
crates/paimon/src/table/table_scan.rs | 476 +++++++++++++++++++++++-
3 files changed, 832 insertions(+), 41 deletions(-)
diff --git a/crates/integrations/datafusion/src/table/mod.rs
b/crates/integrations/datafusion/src/table/mod.rs
index 3508362..ec8893b 100644
--- a/crates/integrations/datafusion/src/table/mod.rs
+++ b/crates/integrations/datafusion/src/table/mod.rs
@@ -177,6 +177,16 @@ impl TableProvider for PaimonTableProvider {
// Plan splits eagerly so we know partition count upfront.
let filter_analysis = analyze_filters(filters,
self.table.schema().fields());
let mut read_builder = self.table.new_read_builder();
+ let projected_columns = projection.map(|indices| {
+ indices
+ .iter()
+ .map(|&i| self.schema.field(i).name().clone())
+ .collect::<Vec<_>>()
+ });
+ if let Some(ref columns) = projected_columns {
+ let col_refs: Vec<&str> = columns.iter().map(|s|
s.as_str()).collect();
+ read_builder.with_projection(&col_refs);
+ }
if let Some(filter) = filter_analysis.pushed_predicate.clone() {
read_builder.with_filter(filter);
}
@@ -365,6 +375,183 @@ mod tests {
.collect()
}
+ fn empty_binary_stats_json() -> serde_json::Value {
+ let row = paimon::spec::EMPTY_SERIALIZED_ROW.as_slice().to_vec();
+ serde_json::json!({
+ "_MIN_VALUES": row,
+ "_MAX_VALUES": row,
+ "_NULL_COUNTS": [],
+ })
+ }
+
+ fn data_evolution_file(
+ file_name: &str,
+ file_size: i64,
+ row_count: i64,
+ first_row_id: i64,
+ write_cols: &[&str],
+ ) -> paimon::spec::DataFileMeta {
+ serde_json::from_value(serde_json::json!({
+ "_FILE_NAME": file_name,
+ "_FILE_SIZE": file_size,
+ "_ROW_COUNT": row_count,
+ "_MIN_KEY": [],
+ "_MAX_KEY": [],
+ "_KEY_STATS": empty_binary_stats_json(),
+ "_VALUE_STATS": empty_binary_stats_json(),
+ "_MIN_SEQUENCE_NUMBER": 0,
+ "_MAX_SEQUENCE_NUMBER": 0,
+ "_SCHEMA_ID": 0,
+ "_LEVEL": 1,
+ "_EXTRA_FILES": [],
+ "_CREATION_TIME": null,
+ "_DELETE_ROW_COUNT": null,
+ "_EMBEDDED_FILE_INDEX": null,
+ "_FILE_SOURCE": null,
+ "_VALUE_STATS_COLS": null,
+ "_FIRST_ROW_ID": first_row_id,
+ "_WRITE_COLS": write_cols,
+ "_EXTERNAL_PATH": null,
+ }))
+ .expect("test data file should deserialize")
+ }
+
+ fn manifest_file_meta(
+ file_name: &str,
+ file_size: i64,
+ num_added_files: i64,
+ ) -> paimon::spec::ManifestFileMeta {
+ serde_json::from_value(serde_json::json!({
+ "_VERSION": 2,
+ "_FILE_NAME": file_name,
+ "_FILE_SIZE": file_size,
+ "_NUM_ADDED_FILES": num_added_files,
+ "_NUM_DELETED_FILES": 0,
+ "_PARTITION_STATS": empty_binary_stats_json(),
+ "_SCHEMA_ID": 0,
+ }))
+ .expect("test manifest file meta should deserialize")
+ }
+
+ async fn data_evolution_projection_pruning_provider() ->
PaimonTableProvider {
+ use paimon::io::FileIOBuilder;
+ use paimon::spec::{
+ CommitKind, DataType, FileKind, IntType, Manifest, ManifestEntry,
ManifestList,
+ Schema as PaimonSchema, Snapshot, TableSchema,
+ };
+ use paimon::table::{SnapshotManager, Table};
+
+ let file_io = FileIOBuilder::new("memory").build().unwrap();
+ let table_path = format!("memory:/df_de_projection_pruning_{}",
uuid::Uuid::new_v4());
+ file_io
+ .mkdirs(&format!("{table_path}/snapshot/"))
+ .await
+ .unwrap();
+ file_io
+ .mkdirs(&format!("{table_path}/manifest/"))
+ .await
+ .unwrap();
+
+ let schema = PaimonSchema::builder()
+ .column("id", DataType::Int(IntType::new()))
+ .column("name", DataType::Int(IntType::new()))
+ .option("data-evolution.enabled", "true")
+ .build()
+ .unwrap();
+ let table_schema = TableSchema::new(0, &schema);
+ let table = Table::new(
+ file_io.clone(),
+ Identifier::new("default", "df_de_projection_pruning"),
+ table_path.clone(),
+ table_schema,
+ None,
+ );
+
+ let partition = paimon::spec::EMPTY_SERIALIZED_ROW.as_slice().to_vec();
+ let entries = vec![
+ ManifestEntry::new(
+ FileKind::Add,
+ partition.clone(),
+ 0,
+ 1,
+ data_evolution_file("id.parquet", 11, 10, 0, &["id"]),
+ 2,
+ ),
+ ManifestEntry::new(
+ FileKind::Add,
+ partition,
+ 0,
+ 1,
+ data_evolution_file("name.parquet", 13, 10, 0, &["name"]),
+ 2,
+ ),
+ ];
+
+ let manifest_name = "manifest-de-projection-0";
+ let manifest_path = format!("{table_path}/manifest/{manifest_name}");
+ Manifest::write(&file_io, &manifest_path, &entries)
+ .await
+ .unwrap();
+ let manifest_size = file_io
+ .new_input(&manifest_path)
+ .unwrap()
+ .metadata()
+ .await
+ .unwrap()
+ .size;
+
+ let base_list_name = "base-list-de-projection";
+ let delta_list_name = "delta-list-de-projection";
+ ManifestList::write(
+ &file_io,
+ &format!("{table_path}/manifest/{base_list_name}"),
+ &[manifest_file_meta(
+ manifest_name,
+ manifest_size as i64,
+ entries.len() as i64,
+ )],
+ )
+ .await
+ .unwrap();
+ ManifestList::write(
+ &file_io,
+ &format!("{table_path}/manifest/{delta_list_name}"),
+ &[],
+ )
+ .await
+ .unwrap();
+
+ let snapshot = Snapshot::builder()
+ .version(3)
+ .id(1)
+ .schema_id(0)
+ .base_manifest_list(base_list_name.to_string())
+ .delta_manifest_list(delta_list_name.to_string())
+ .commit_user("test-user".to_string())
+ .commit_identifier(1)
+ .commit_kind(CommitKind::APPEND)
+ .time_millis(1)
+ .total_record_count(Some(10))
+ .delta_record_count(Some(10))
+ .build();
+ let snapshot_manager = SnapshotManager::new(file_io, table_path);
+ assert!(snapshot_manager.commit_snapshot(&snapshot).await.unwrap());
+
+ PaimonTableProvider::try_new(table).expect("provider should be
created")
+ }
+
+ fn planned_file_names(scan: &PaimonTableScan) -> Vec<String> {
+ let mut names = scan
+ .planned_partitions()
+ .iter()
+ .flat_map(|partition| partition.iter())
+ .flat_map(|split| split.data_files().iter())
+ .map(|file| file.file_name.clone())
+ .collect::<Vec<_>>();
+ names.sort();
+ names
+ }
+
#[tokio::test]
async fn test_scan_partition_filter_plans_matching_partition_set() {
let provider = create_provider("partitioned_log_table").await;
@@ -475,6 +662,42 @@ mod tests {
assert_eq!(scan.pushed_predicate(), Some(&expected));
}
+ #[tokio::test]
+ async fn test_scan_applies_projection_to_data_evolution_planning() {
+ let provider = data_evolution_projection_pruning_provider().await;
+ let config = SessionConfig::new().with_target_partitions(8);
+ let ctx = SessionContext::new_with_config(config);
+ let state = ctx.state();
+
+ let full_plan = provider
+ .scan(&state, None, &[], None)
+ .await
+ .expect("full scan should succeed");
+ let full_scan = full_plan
+ .as_any()
+ .downcast_ref::<PaimonTableScan>()
+ .expect("Expected PaimonTableScan");
+ assert_eq!(
+ planned_file_names(full_scan),
+ vec!["id.parquet".to_string(), "name.parquet".to_string()]
+ );
+
+ let projection = vec![1];
+ let projected_plan = provider
+ .scan(&state, Some(&projection), &[], None)
+ .await
+ .expect("projected scan should succeed");
+ let projected_scan = projected_plan
+ .as_any()
+ .downcast_ref::<PaimonTableScan>()
+ .expect("Expected PaimonTableScan");
+
+ assert_eq!(
+ planned_file_names(projected_scan),
+ vec!["name.parquet".to_string()]
+ );
+ }
+
#[tokio::test]
async fn test_scan_applies_limit_hint_only_when_safe() {
let provider = create_provider("partitioned_log_table").await;
diff --git a/crates/paimon/src/table/read_builder.rs
b/crates/paimon/src/table/read_builder.rs
index 099a8a0..dbf22ca 100644
--- a/crates/paimon/src/table/read_builder.rs
+++ b/crates/paimon/src/table/read_builder.rs
@@ -223,6 +223,7 @@ impl<'a> ReadBuilder<'a> {
self.limit,
self.row_ranges.clone(),
)
+ .with_projection(self.projected_fields.clone())
}
/// Create a table read for consuming splits (e.g. from a scan plan).
@@ -239,54 +240,90 @@ impl<'a> ReadBuilder<'a> {
))
}
- fn resolve_projected_fields(&self, projected_fields: &[String]) ->
Result<Vec<DataField>> {
- if projected_fields.is_empty() {
- return Ok(Vec::new());
- }
+ pub(super) fn resolve_projected_fields(
+ &self,
+ projected_fields: &[String],
+ ) -> Result<Vec<DataField>> {
+ resolve_projected_fields(
+ self.table.identifier().full_name(),
+ self.table.schema.fields(),
+ projected_fields,
+ )
+ }
+}
- let full_name = self.table.identifier().full_name();
- let field_map: HashMap<&str, &DataField> = self
- .table
- .schema
- .fields()
- .iter()
- .map(|field| (field.name(), field))
- .collect();
+pub(super) fn resolve_projected_fields(
+ full_name: String,
+ fields: &[DataField],
+ projected_fields: &[String],
+) -> Result<Vec<DataField>> {
+ if projected_fields.is_empty() {
+ return Ok(Vec::new());
+ }
- let mut seen = HashSet::with_capacity(projected_fields.len());
- let mut resolved = Vec::with_capacity(projected_fields.len());
+ let field_map: HashMap<&str, &DataField> =
+ fields.iter().map(|field| (field.name(), field)).collect();
- for name in projected_fields {
- if !seen.insert(name.as_str()) {
- return Err(Error::ConfigInvalid {
- message: format!("Duplicate projection column '{name}' for
table {full_name}"),
- });
- }
+ let mut seen = HashSet::with_capacity(projected_fields.len());
+ let mut resolved = Vec::with_capacity(projected_fields.len());
- if name == crate::spec::ROW_ID_FIELD_NAME {
- resolved.push(DataField::new(
- crate::spec::ROW_ID_FIELD_ID,
- crate::spec::ROW_ID_FIELD_NAME.to_string(),
-
crate::spec::DataType::BigInt(crate::spec::BigIntType::with_nullable(true)),
- ));
- continue;
- }
+ for name in projected_fields {
+ if !seen.insert(name.as_str()) {
+ return Err(Error::ConfigInvalid {
+ message: format!("Duplicate projection column '{name}' for
table {full_name}"),
+ });
+ }
- let field = field_map
- .get(name.as_str())
- .ok_or_else(|| Error::ColumnNotExist {
- full_name: full_name.clone(),
- column: name.clone(),
- })?;
- resolved.push((*field).clone());
+ if name == crate::spec::ROW_ID_FIELD_NAME {
+ resolved.push(DataField::new(
+ crate::spec::ROW_ID_FIELD_ID,
+ crate::spec::ROW_ID_FIELD_NAME.to_string(),
+
crate::spec::DataType::BigInt(crate::spec::BigIntType::with_nullable(true)),
+ ));
+ continue;
}
- Ok(resolved)
+ let field = field_map
+ .get(name.as_str())
+ .ok_or_else(|| Error::ColumnNotExist {
+ full_name: full_name.clone(),
+ column: name.clone(),
+ })?;
+ resolved.push((*field).clone());
}
+
+ Ok(resolved)
+}
+
+pub(super) fn projected_read_field_ids(
+ full_name: String,
+ fields: &[DataField],
+ projected_fields: Option<&Vec<String>>,
+) -> Result<Option<HashSet<i32>>> {
+ let Some(projected) = projected_fields else {
+ return Ok(None);
+ };
+ let fields = resolve_projected_fields(full_name, fields, projected)?;
+ let field_ids = fields
+ .into_iter()
+ .filter(|field| !is_system_projection_field(field.id()))
+ .map(|field| field.id())
+ .collect::<HashSet<_>>();
+ Ok(Some(field_ids))
+}
+
+pub(super) fn is_system_projection_field(field_id: i32) -> bool {
+ matches!(
+ field_id,
+ crate::spec::ROW_ID_FIELD_ID
+ | crate::spec::SEQUENCE_NUMBER_FIELD_ID
+ | crate::spec::VALUE_KIND_FIELD_ID
+ )
}
#[cfg(test)]
mod tests {
+ use super::ReadBuilder;
use crate::table::TableRead;
mod test_utils {
include!(concat!(env!("CARGO_MANIFEST_DIR"), "/../test_utils.rs"));
@@ -300,6 +337,7 @@ mod tests {
use crate::table::{DataSplitBuilder, Table};
use arrow_array::{Int32Array, RecordBatch};
use futures::TryStreamExt;
+ use std::collections::HashSet;
use std::fs;
use tempfile::tempdir;
use test_utils::{local_file_path, test_data_file, write_int_parquet_file};
@@ -360,6 +398,70 @@ mod tests {
)
}
+ #[test]
+ fn test_projected_read_field_ids_uses_projection_ids() {
+ let table = simple_table();
+ let projected = vec!["id".to_string()];
+
+ assert_eq!(
+ super::projected_read_field_ids(
+ table.identifier().full_name(),
+ table.schema().fields(),
+ Some(&projected),
+ )
+ .unwrap(),
+ Some(HashSet::from([1]))
+ );
+ }
+
+ #[test]
+ fn test_projected_read_field_ids_ignores_system_only_projection() {
+ let table = simple_table();
+ let projected = vec![crate::spec::ROW_ID_FIELD_NAME.to_string()];
+
+ assert_eq!(
+ super::projected_read_field_ids(
+ table.identifier().full_name(),
+ table.schema().fields(),
+ Some(&projected),
+ )
+ .unwrap(),
+ Some(HashSet::new())
+ );
+ }
+
+ #[tokio::test]
+ async fn test_new_scan_validates_unknown_projection() {
+ let table = simple_table();
+ let mut builder = ReadBuilder::new(&table);
+ builder.with_projection(&["missing"]);
+
+ let err = builder.new_scan().plan().await.unwrap_err();
+
+ assert!(matches!(
+ err,
+ crate::Error::ColumnNotExist {
+ full_name,
+ column,
+ } if full_name == "default.t" && column == "missing"
+ ));
+ }
+
+ #[tokio::test]
+ async fn test_new_scan_validates_duplicate_projection() {
+ let table = simple_table();
+ let mut builder = ReadBuilder::new(&table);
+ builder.with_projection(&["id", "id"]);
+
+ let err = builder.new_scan().plan().await.unwrap_err();
+
+ assert!(matches!(
+ err,
+ crate::Error::ConfigInvalid { message }
+ if message.contains("Duplicate projection column 'id'")
+ ));
+ }
+
#[test]
fn test_exact_filter_pushdown_is_true_for_partition_only_filter() {
let table = simple_table();
diff --git a/crates/paimon/src/table/table_scan.rs
b/crates/paimon/src/table/table_scan.rs
index 27bc8d8..0c9dc95 100644
--- a/crates/paimon/src/table/table_scan.rs
+++ b/crates/paimon/src/table/table_scan.rs
@@ -32,12 +32,14 @@ use crate::io::FileIO;
use crate::spec::{
avro::SharedSchemaCache, bucket_dir_name, BinaryRow, BucketFunctionType,
CoreOptions,
DataField, DataFileMeta, FileKind, IndexManifest, ManifestEntry,
PartitionComputer, Predicate,
- Snapshot,
+ 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,
};
use crate::table::bin_pack::split_for_batch;
use crate::table::merge_tree_split_generator::{
merge_tree_split_for_batch, KeyComparator, SplitGroup,
};
+use crate::table::schema_manager::SchemaManager;
use crate::table::source::{
any_range_overlaps_file, intersect_ranges_with_file, merge_row_ranges,
DataSplit,
DataSplitBuilder, DeletionFile, PartitionBucket, Plan, RowRange,
@@ -300,6 +302,186 @@ fn should_skip_level_zero_for_scan(
deletion_vectors_enabled || merge_engine.is_ok_and(|e| e ==
crate::spec::MergeEngine::FirstRow)
}
+fn is_system_field_id(field_id: i32) -> bool {
+ matches!(
+ field_id,
+ ROW_ID_FIELD_ID | SEQUENCE_NUMBER_FIELD_ID | VALUE_KIND_FIELD_ID
+ )
+}
+
+fn is_system_field_name(name: &str) -> bool {
+ matches!(
+ name,
+ ROW_ID_FIELD_NAME | SEQUENCE_NUMBER_FIELD_NAME | VALUE_KIND_FIELD_NAME
+ )
+}
+
+fn is_vector_store_file_name(file_name: &str) -> bool {
+ file_name.to_ascii_lowercase().contains(".vector.")
+}
+
+fn is_normal_data_file(file: &DataFileMeta) -> bool {
+ !crate::table::blob_file_writer::is_blob_file_name(&file.file_name)
+ && !is_vector_store_file_name(&file.file_name)
+}
+
+type DataFileFieldIdsCache = HashMap<(i64, Option<Vec<String>>), HashSet<i32>>;
+
+fn data_evolution_representative_file(group: &[DataFileMeta]) ->
crate::Result<usize> {
+ let mut representative: Option<usize> = None;
+ for (idx, file) in group.iter().enumerate() {
+ if !is_normal_data_file(file) {
+ continue;
+ }
+ let should_replace = match representative {
+ None => true,
+ Some(current_idx) => {
+ let current = &group[current_idx];
+ (file.max_sequence_number, file.file_name.as_str())
+ < (current.max_sequence_number, current.file_name.as_str())
+ }
+ };
+ if should_replace {
+ representative = Some(idx);
+ }
+ }
+ representative.ok_or_else(|| crate::Error::DataInvalid {
+ message: "Data-evolution row range group requires at least one normal
data file."
+ .to_string(),
+ source: None,
+ })
+}
+
+async fn resolve_data_file_field_ids(
+ table_schema_id: i64,
+ table_fields: &[DataField],
+ schema_manager: &SchemaManager,
+ file: &DataFileMeta,
+) -> crate::Result<HashSet<i32>> {
+ let schema;
+ let fields = if file.schema_id == table_schema_id {
+ table_fields
+ } else {
+ schema = schema_manager.schema(file.schema_id).await?;
+ schema.fields()
+ };
+
+ let field_id_by_name = fields
+ .iter()
+ .map(|field| (field.name(), field.id()))
+ .collect::<HashMap<_, _>>();
+
+ let mut field_ids = HashSet::new();
+ match file.write_cols.as_ref() {
+ None => {
+ field_ids.extend(
+ fields
+ .iter()
+ .filter(|field| !is_system_field_id(field.id()))
+ .map(|field| field.id()),
+ );
+ }
+ Some(write_cols) => {
+ for col in write_cols {
+ if is_system_field_name(col) {
+ continue;
+ }
+ let Some(field_id) = field_id_by_name.get(col.as_str()) else {
+ return Err(crate::Error::DataInvalid {
+ message: format!(
+ "Cannot find write column '{}' in schema {}.",
+ col, file.schema_id
+ ),
+ source: None,
+ });
+ };
+ if !is_system_field_id(*field_id) {
+ field_ids.insert(*field_id);
+ }
+ }
+ }
+ }
+ Ok(field_ids)
+}
+
+async fn data_file_field_ids(
+ table_schema_id: i64,
+ table_fields: &[DataField],
+ schema_manager: &SchemaManager,
+ file: &DataFileMeta,
+ field_ids_cache: &mut DataFileFieldIdsCache,
+) -> crate::Result<HashSet<i32>> {
+ let key = (file.schema_id, file.write_cols.clone());
+ if let Some(field_ids) = field_ids_cache.get(&key) {
+ return Ok(field_ids.clone());
+ }
+
+ let field_ids =
+ resolve_data_file_field_ids(table_schema_id, table_fields,
schema_manager, file).await?;
+ field_ids_cache.insert(key, field_ids.clone());
+ Ok(field_ids)
+}
+
+async fn prune_data_evolution_group_by_read_fields(
+ group: Vec<DataFileMeta>,
+ read_field_ids: &HashSet<i32>,
+ deletion_vectors_enabled: bool,
+ table_schema_id: i64,
+ table_fields: &[DataField],
+ schema_manager: &SchemaManager,
+ field_ids_cache: &mut DataFileFieldIdsCache,
+) -> crate::Result<Vec<DataFileMeta>> {
+ if read_field_ids.is_empty() || group.len() <= 1 {
+ return Ok(group);
+ }
+
+ let anchor_idx = if deletion_vectors_enabled {
+ Some(data_evolution_representative_file(&group)?)
+ } else {
+ None
+ };
+
+ let mut keep = Vec::with_capacity(group.len());
+ for (idx, file) in group.iter().enumerate() {
+ let file_field_ids = data_file_field_ids(
+ table_schema_id,
+ table_fields,
+ schema_manager,
+ file,
+ field_ids_cache,
+ )
+ .await?;
+ if file_field_ids
+ .iter()
+ .any(|field_id| read_field_ids.contains(field_id))
+ {
+ keep.push(idx);
+ }
+ }
+ if let Some(anchor_idx) = anchor_idx {
+ if !keep.contains(&anchor_idx) {
+ keep.push(anchor_idx);
+ }
+ }
+
+ if keep.is_empty() {
+ keep.push(data_evolution_representative_file(&group)?);
+ } else if keep.iter().any(|idx| !is_normal_data_file(&group[*idx]))
+ && !keep.iter().any(|idx| is_normal_data_file(&group[*idx]))
+ {
+ let representative_idx = data_evolution_representative_file(&group)?;
+ if !keep.contains(&representative_idx) {
+ keep.push(representative_idx);
+ }
+ }
+
+ let mut files = group.into_iter().map(Some).collect::<Vec<_>>();
+ Ok(keep
+ .into_iter()
+ .filter_map(|idx| files.get_mut(idx).and_then(Option::take))
+ .collect())
+}
+
/// TableScan for full table scan (no incremental, no predicate).
///
/// Reference:
[pypaimon.read.table_scan.TableScan](https://github.com/apache/paimon/blob/master/paimon-python/pypaimon/read/table_scan.py)
@@ -317,6 +499,7 @@ pub struct TableScan<'a> {
/// Used by non-read paths (overwrite, truncate, writer restore) that need
/// the complete file set. Normal read scans leave this as `false`.
scan_all_files: bool,
+ projected_fields: Option<Vec<String>>,
}
impl<'a> TableScan<'a> {
@@ -336,6 +519,7 @@ impl<'a> TableScan<'a> {
limit,
row_ranges,
scan_all_files: false,
+ projected_fields: None,
}
}
@@ -345,6 +529,7 @@ impl<'a> TableScan<'a> {
/// the complete file set regardless of merge engine or DV settings.
pub fn with_scan_all_files(mut self) -> Self {
self.scan_all_files = true;
+ self.projected_fields = None;
self
}
@@ -361,6 +546,11 @@ impl<'a> TableScan<'a> {
self
}
+ pub(crate) fn with_projection(mut self, projected_fields:
Option<Vec<String>>) -> Self {
+ self.projected_fields = projected_fields;
+ self
+ }
+
/// Plan the full scan: resolve snapshot (via options or latest), then
read manifests and build DataSplits.
///
/// Time travel is resolved from table options:
@@ -371,11 +561,21 @@ impl<'a> TableScan<'a> {
///
/// Reference:
[TimeTravelUtil.tryTravelToSnapshot](https://github.com/apache/paimon/blob/master/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/TimeTravelUtil.java)
pub async fn plan(&self) -> crate::Result<Plan> {
+ let data_evolution_read_field_ids = self.projected_read_field_ids()?;
let snapshot = match self.resolve_snapshot().await? {
Some(snapshot) => snapshot,
None => return Ok(Plan::new(Vec::new())),
};
- self.plan_snapshot(snapshot).await
+ self.plan_snapshot(snapshot, data_evolution_read_field_ids.as_ref())
+ .await
+ }
+
+ fn projected_read_field_ids(&self) -> crate::Result<Option<HashSet<i32>>> {
+ super::read_builder::projected_read_field_ids(
+ self.table.identifier().full_name(),
+ self.table.schema().fields(),
+ self.projected_fields.as_ref(),
+ )
}
async fn resolve_snapshot(&self) -> crate::Result<Option<Snapshot>> {
@@ -544,11 +744,19 @@ impl<'a> TableScan<'a> {
can_push_down_limit_hint_for_scan(&self.data_predicates, row_ranges)
}
- async fn plan_snapshot(&self, snapshot: Snapshot) -> crate::Result<Plan> {
+ async fn plan_snapshot(
+ &self,
+ snapshot: Snapshot,
+ data_evolution_read_field_ids: Option<&HashSet<i32>>,
+ ) -> 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();
+ let schema_manager = self.table.schema_manager();
let core_options = CoreOptions::new(self.table.schema().options());
let data_evolution_enabled = core_options.data_evolution_enabled();
+ let deletion_vectors_enabled = core_options.deletion_vectors_enabled();
let target_split_size = core_options.source_split_target_size();
let open_file_cost = core_options.source_split_open_file_cost();
let partition_keys = self.table.schema().partition_keys();
@@ -669,6 +877,7 @@ impl<'a> TableScan<'a> {
(None, self.row_ranges.clone())
};
+ let mut data_file_field_ids_cache = DataFileFieldIdsCache::new();
for ((partition, bucket), (total_buckets, data_files)) in groups {
let partition_row = BinaryRow::from_serialized_bytes(&partition)?;
@@ -716,6 +925,31 @@ impl<'a> TableScan<'a> {
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
+ } else {
+ let mut pruned =
Vec::with_capacity(row_id_groups.len());
+ for group in row_id_groups {
+ pruned.push(
+ prune_data_evolution_group_by_read_fields(
+ group,
+ read_field_ids,
+ deletion_vectors_enabled,
+ table_schema_id,
+ table_fields,
+ schema_manager,
+ &mut data_file_field_ids_cache,
+ )
+ .await?,
+ );
+ }
+ pruned
+ }
+ } else {
+ row_id_groups
+ };
+
let (singles, multis): (Vec<_>, Vec<_>) = row_id_groups
.into_iter()
.partition(|group| group.len() == 1);
@@ -831,7 +1065,9 @@ impl<'a> TableScan<'a> {
#[cfg(test)]
mod tests {
- use super::{should_skip_level_zero_for_scan, TableScan};
+ use super::{
+ prune_data_evolution_group_by_read_fields,
should_skip_level_zero_for_scan, TableScan,
+ };
use crate::catalog::Identifier;
use crate::io::FileIOBuilder;
use crate::spec::{
@@ -849,8 +1085,9 @@ mod tests {
};
use crate::table::Table;
use crate::Error;
+ use bytes::Bytes;
use chrono::{DateTime, Utc};
- use std::collections::HashSet;
+ use std::collections::{HashMap, HashSet};
/// Helper to build a DataFileMeta with data evolution fields.
fn make_evo_file(
@@ -884,6 +1121,62 @@ mod tests {
}
}
+ fn make_evo_file_with_cols(
+ name: &str,
+ row_count: i64,
+ max_seq: i64,
+ first_row_id: i64,
+ write_cols: &[&str],
+ ) -> DataFileMeta {
+ let mut file = make_evo_file(name, 10, row_count, max_seq,
Some(first_row_id));
+ file.write_cols = Some(write_cols.iter().map(|col|
(*col).to_string()).collect());
+ file
+ }
+
+ 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([(
+ "data-evolution.enabled".to_string(),
+ "true".to_string(),
+ )]));
+ Table::new(
+ file_io,
+ Identifier::new("test_db", "de_table"),
+ table_path.to_string(),
+ schema,
+ None,
+ )
+ }
+
+ fn two_column_schema(id: i64, left: &str, right: &str) -> TableSchema {
+ TableSchema::new(
+ id,
+ &PaimonSchema::builder()
+ .column(left, DataType::Int(IntType::new()))
+ .column(right, DataType::Int(IntType::new()))
+ .build()
+ .unwrap(),
+ )
+ }
+
+ async fn write_schema_file(table: &Table, schema: &TableSchema) {
+ let path = table.schema_manager().schema_path(schema.id());
+ let dir = path.rsplit_once('/').map(|(dir, _)| dir).unwrap();
+ table.file_io().mkdirs(dir).await.unwrap();
+ let json = serde_json::to_vec(schema).unwrap();
+ table
+ .file_io()
+ .new_output(&path)
+ .unwrap()
+ .write(Bytes::from(json))
+ .await
+ .unwrap();
+ }
+
+ fn file_names_from_files(files: &[DataFileMeta]) -> Vec<&str> {
+ files.iter().map(|file| file.file_name.as_str()).collect()
+ }
+
#[test]
fn test_merge_manifest_entries_keeps_in_place_upgraded_file() {
// Reproduces a single-run compaction "upgrade": the SAME file name is
@@ -1226,6 +1519,179 @@ mod tests {
assert_eq!(file_names(&groups), vec![vec!["a", "b"]]);
}
+ #[tokio::test]
+ async fn test_data_evolution_prunes_files_without_projected_columns() {
+ let table =
+ data_evolution_test_table("memory:/de_prune_cols",
two_column_schema(0, "id", "name"));
+ let read_field_ids = HashSet::from([1]);
+ let files = vec![
+ make_evo_file_with_cols("id.parquet", 10, 1, 0, &["id"]),
+ make_evo_file_with_cols("name.parquet", 10, 2, 0, &["name"]),
+ ];
+ let mut field_ids_cache = HashMap::new();
+
+ let pruned = prune_data_evolution_group_by_read_fields(
+ files,
+ &read_field_ids,
+ false,
+ table.schema().id(),
+ table.schema().fields(),
+ table.schema_manager(),
+ &mut field_ids_cache,
+ )
+ .await
+ .unwrap();
+
+ assert_eq!(file_names_from_files(&pruned), vec!["name.parquet"]);
+ }
+
+ #[tokio::test]
+ async fn test_data_evolution_pruning_keeps_dv_anchor() {
+ let table =
+ data_evolution_test_table("memory:/de_prune_dv",
two_column_schema(0, "id", "name"));
+ let read_field_ids = HashSet::from([1]);
+ let files = vec![
+ make_evo_file_with_cols("new-name.parquet", 10, 5, 0, &["name"]),
+ make_evo_file_with_cols("old-id.parquet", 10, 1, 0, &["id"]),
+ ];
+ let mut field_ids_cache = HashMap::new();
+
+ let pruned = prune_data_evolution_group_by_read_fields(
+ files,
+ &read_field_ids,
+ true,
+ table.schema().id(),
+ table.schema().fields(),
+ table.schema_manager(),
+ &mut field_ids_cache,
+ )
+ .await
+ .unwrap();
+
+ assert_eq!(
+ file_names_from_files(&pruned),
+ vec!["new-name.parquet", "old-id.parquet"]
+ );
+ }
+
+ #[tokio::test]
+ async fn test_data_evolution_pruning_keeps_row_count_representative() {
+ let table = data_evolution_test_table(
+ "memory:/de_prune_representative",
+ two_column_schema(0, "id", "name"),
+ );
+ let read_field_ids = HashSet::from([2]);
+ let files = vec![
+ make_evo_file_with_cols("new-name.parquet", 10, 5, 0, &["name"]),
+ make_evo_file_with_cols("old-id.parquet", 10, 1, 0, &["id"]),
+ ];
+ let mut field_ids_cache = HashMap::new();
+
+ let pruned = prune_data_evolution_group_by_read_fields(
+ files,
+ &read_field_ids,
+ false,
+ table.schema().id(),
+ table.schema().fields(),
+ table.schema_manager(),
+ &mut field_ids_cache,
+ )
+ .await
+ .unwrap();
+
+ assert_eq!(file_names_from_files(&pruned), vec!["old-id.parquet"]);
+ }
+
+ #[tokio::test]
+ async fn test_data_evolution_pruning_matches_renamed_columns_by_field_id()
{
+ let schema_v0 = two_column_schema(0, "id", "old_name");
+ let schema_v1 = two_column_schema(1, "id", "new_name");
+ let table = data_evolution_test_table("memory:/de_prune_rename",
schema_v1);
+ write_schema_file(&table, &schema_v0).await;
+
+ let read_field_ids = HashSet::from([1]);
+ let mut file = make_evo_file_with_cols("renamed.parquet", 10, 1, 0,
&["old_name"]);
+ file.schema_id = 0;
+ let pruned = prune_data_evolution_group_by_read_fields(
+ vec![
+ make_evo_file_with_cols("id.parquet", 10, 2, 0, &["id"]),
+ file,
+ ],
+ &read_field_ids,
+ false,
+ table.schema().id(),
+ table.schema().fields(),
+ table.schema_manager(),
+ &mut HashMap::new(),
+ )
+ .await
+ .unwrap();
+
+ assert_eq!(file_names_from_files(&pruned), vec!["renamed.parquet"]);
+ }
+
+ #[tokio::test]
+ async fn
test_data_evolution_pruning_keeps_normal_representative_for_vector_file() {
+ let table =
+ data_evolution_test_table("memory:/de_prune_vector",
two_column_schema(0, "id", "emb"));
+ let read_field_ids = HashSet::from([1]);
+ let files = vec![
+ make_evo_file_with_cols("data.parquet", 10, 1, 0, &["id"]),
+ make_evo_file_with_cols("emb.vector.parquet", 10, 2, 0, &["emb"]),
+ ];
+ let mut field_ids_cache = HashMap::new();
+
+ let pruned = prune_data_evolution_group_by_read_fields(
+ files,
+ &read_field_ids,
+ false,
+ table.schema().id(),
+ table.schema().fields(),
+ table.schema_manager(),
+ &mut field_ids_cache,
+ )
+ .await
+ .unwrap();
+
+ assert_eq!(
+ file_names_from_files(&pruned),
+ vec!["emb.vector.parquet", "data.parquet"]
+ );
+ }
+
+ #[tokio::test]
+ async fn
test_data_evolution_pruning_rejects_group_without_normal_representative() {
+ let table = data_evolution_test_table(
+ "memory:/de_prune_no_normal",
+ two_column_schema(0, "id", "emb"),
+ );
+ let read_field_ids = HashSet::from([1]);
+ let files = vec![
+ make_evo_file_with_cols("emb-1.vector.parquet", 10, 1, 0,
&["emb"]),
+ make_evo_file_with_cols("emb-2.vector.parquet", 10, 2, 0,
&["emb"]),
+ ];
+ let mut field_ids_cache = HashMap::new();
+
+ let err = prune_data_evolution_group_by_read_fields(
+ files,
+ &read_field_ids,
+ false,
+ table.schema().id(),
+ table.schema().fields(),
+ table.schema_manager(),
+ &mut field_ids_cache,
+ )
+ .await
+ .unwrap_err();
+
+ match err {
+ Error::DataInvalid { message, .. } => {
+ assert!(message.contains("requires at least one normal data
file"));
+ }
+ other => panic!("unexpected error: {other:?}"),
+ }
+ }
+
#[test]
fn test_group_by_overlapping_row_id_overlapping_ranges() {
let files = vec![