Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 5 additions & 5 deletions crates/iceberg/src/arrow/record_batch_transformer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -59,11 +59,11 @@ fn constants_map(
for (pos, field) in partition_spec.fields().iter().enumerate() {
// Only identity transforms should use constant values from partition metadata
if matches!(field.transform, Transform::Identity) {
// Get the field from schema to extract its type
let iceberg_field = schema.field_by_id(field.source_id).ok_or(Error::new(
ErrorKind::Unexpected,
format!("Field {} not found in schema", field.source_id),
))?;
// The source column may have been dropped from the schema after the spec was
// created. It cannot be projected in that case, so no constant is needed.
let Some(iceberg_field) = schema.field_by_id(field.source_id) else {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Follow-up (non-blocking, coverage): this constants_map branch isn't reached by the native scan flow today — FileScanTask is built with with_partition_spec(None) (scan/context.rs, with a TODO: Pass actual PartitionSpec through context chain), and pipeline.rs only calls with_partition when the spec is Some. So this change is a reasonable forward-looking guard, but it's currently unexercised (the new scan test uses None). A direct unit test here — mirroring the existing with_partition tests in this file — would prevent a future refactor that threads the spec through from silently reintroducing the panic.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the review. I have filed followup to fix these #2869

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We do pass PartitionSpec thru the context chain. It happened last week in this PR: #2695

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

ah you're right, I missed #2695, the spec is threaded through now, so this is a live-path fix, not the forward-looking guard, thanks for the pointer. test's covered in #2869.

continue;
};

// Ensure the field type is primitive
let prim_type = match &*iceberg_field.field_type {
Expand Down
14 changes: 14 additions & 0 deletions crates/iceberg/src/scan/context.rs
Original file line number Diff line number Diff line change
Expand Up @@ -175,9 +175,23 @@ impl PlanContext {
.await
}

/// Returns the partition filter for a manifest, or an always-true filter when the
/// manifest's spec cannot be resolved against the snapshot schema. Files of such
/// manifests are not partition-pruned but still receive the row filter.
fn get_partition_filter(&self, manifest_file: &ManifestFile) -> Result<Arc<BoundPredicate>> {
let partition_spec_id = manifest_file.partition_spec_id;

// Historical specs may reference source columns that were later dropped from the
// schema, in which case the partition type cannot be resolved. A missing spec is
// reported by the partition filter cache instead.
let resolvable = self
.table_metadata
.partition_spec_by_id(partition_spec_id)
.is_none_or(|spec| spec.partition_type(&self.snapshot_schema).is_ok());

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit (non-blocking, efficiency): this resolvability check computes partition_type(&self.snapshot_schema) once per manifest file, but PartitionFilterCache::get already recomputes partition_type internally on a cache miss. For a scan over many manifests sharing one spec, resolvable manifests pay for partition_type twice on the first manifest and re-run the check for every subsequent manifest even though the final filter is cached by spec_id. Cheap enough for a stopgap, but caching resolvability by spec_id (or folding the check into the cache-miss path) would avoid the O(#manifests) extra calls.

if !resolvable {
return Ok(Arc::new(BoundPredicate::AlwaysTrue));

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit (non-blocking): a // TODO(#2844) breadcrumb here would make the stopgap nature discoverable from the code itself — the PR body notes the longer-term fix (deriving partition types from transforms, à la apache/iceberg#17262, which would restore pruning on a historical spec's still-live partition fields), but that context lives only in the PR.

}

let partition_filter = self.partition_filter_cache.get(
partition_spec_id,
&self.table_metadata,
Expand Down
98 changes: 96 additions & 2 deletions crates/iceberg/src/scan/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -646,8 +646,10 @@ pub mod tests {
use crate::scan::FileScanTask;
use crate::spec::{
DEFAULT_SCHEMA_NAME_MAPPING, DataContentType, DataFileBuilder, DataFileFormat, Datum,
Literal, ManifestEntry, ManifestListWriter, ManifestStatus, ManifestWriterBuilder,
NestedField, PartitionSpec, PrimitiveType, Schema, Struct, StructType, TableMetadata, Type,
Literal, MAIN_BRANCH, ManifestEntry, ManifestListWriter, ManifestStatus,
ManifestWriterBuilder, NestedField, Operation, PartitionSpec, PrimitiveType, Schema,
Snapshot, Struct, StructType, Summary, TableMetadata, TableMetadataBuilder, Type,
UnboundPartitionSpec,
};
use crate::table::Table;
use crate::test_utils::test_runtime;
Expand Down Expand Up @@ -1643,6 +1645,98 @@ pub mod tests {
);
}

#[tokio::test]
async fn test_filtered_scan_with_dropped_partition_source_column() {
let mut fixture = TableTestFixture::new();
fixture.setup_manifest_files().await;

// baseline: the same filtered scan against the table before evolution
let baseline = scan_y_gte_5(&fixture.table).await;
assert!(!baseline.is_empty());
assert!(baseline.iter().all(|y| *y >= 5));

// Evolve the table so that the manifests reference a historical spec whose source
// column is no longer in the current schema: make an unpartitioned spec the
// default, then drop the original spec's source column from the schema.
let current_schema = fixture.table.metadata().current_schema();
let evolved_schema = Schema::builder()
.with_fields(
current_schema
.as_struct()
.fields()
.iter()
.filter(|field| field.id != 1)
.cloned(),
)
.with_identifier_field_ids(vec![2])
.build()
.unwrap();
let evolved =
TableMetadataBuilder::new_from_metadata(fixture.table.metadata().clone(), None)
.add_default_partition_spec(UnboundPartitionSpec::builder().build())
.unwrap()
.add_current_schema(evolved_schema)
.unwrap()
.build()
.unwrap()
.metadata;

// a commit after the evolution carries the previous manifests forward: the new
// snapshot uses the evolved schema while its manifests still use historical spec 0
let parent = evolved.current_snapshot().unwrap().clone();
let snapshot = Snapshot::builder()
.with_snapshot_id(parent.snapshot_id() + 1)
.with_parent_snapshot_id(Some(parent.snapshot_id()))
.with_sequence_number(evolved.last_sequence_number() + 1)
.with_timestamp_ms(evolved.last_updated_ms + 1)
.with_schema_id(evolved.current_schema_id())
.with_manifest_list(parent.manifest_list())
.with_summary(Summary {
operation: Operation::Append,
additional_properties: HashMap::new(),
})
.build();
let metadata = TableMetadataBuilder::new_from_metadata(evolved, None)
.set_branch_snapshot(snapshot, MAIN_BRANCH)
.unwrap()
.build()
.unwrap()
.metadata;
let table = fixture.table.clone().with_metadata(Arc::new(metadata));

// planning and reading must succeed, and the results must match the table before
// evolution: no rows wrongly pruned and none returned unfiltered
let evolved = scan_y_gte_5(&table).await;
assert_eq!(evolved, baseline);
}

async fn scan_y_gte_5(table: &Table) -> Vec<i64> {
let table_scan = table
.scan()
.select(["y"])
.with_filter(Reference::new("y").greater_than_or_equal_to(Datum::long(5)))
.build()
.unwrap();
let batches: Vec<_> = table_scan
.to_arrow()
.await
.unwrap()
.try_collect()
.await
.unwrap();

let mut values: Vec<i64> = batches
.iter()
.flat_map(|batch| {
let col = batch.column_by_name("y").unwrap();
let arr = col.as_any().downcast_ref::<Int64Array>().unwrap();
(0..arr.len()).map(|i| arr.value(i)).collect::<Vec<_>>()
})
.collect();
values.sort_unstable();
values
}

#[tokio::test]
async fn test_open_parquet_no_deletions() {
let mut fixture = TableTestFixture::new();
Expand Down
Loading