diff --git a/crates/integrations/datafusion/src/lateral_vector_search.rs b/crates/integrations/datafusion/src/lateral_vector_search.rs index 75a6250b0..2f5bea5dd 100644 --- a/crates/integrations/datafusion/src/lateral_vector_search.rs +++ b/crates/integrations/datafusion/src/lateral_vector_search.rs @@ -107,6 +107,10 @@ pub(crate) fn optimizer_rules() -> Vec> { Arc::new(RewriteLateralVectorSearch::new()), ]; rules.extend(Optimizer::default().rules); + // After the default rules, so filters have already been pushed into the scan. + rules.push(Arc::new( + crate::partition_count_pushdown::PushDownPartitionCount, + )); rules } diff --git a/crates/integrations/datafusion/src/lib.rs b/crates/integrations/datafusion/src/lib.rs index 95d1efc95..5d581137f 100644 --- a/crates/integrations/datafusion/src/lib.rs +++ b/crates/integrations/datafusion/src/lib.rs @@ -52,6 +52,7 @@ mod full_text_search; mod hybrid_search; mod lateral_vector_search; mod merge_into; +mod partition_count_pushdown; mod physical_plan; mod procedures; mod relation_planner; diff --git a/crates/integrations/datafusion/src/partition_count_pushdown.rs b/crates/integrations/datafusion/src/partition_count_pushdown.rs new file mode 100644 index 000000000..0956c29af --- /dev/null +++ b/crates/integrations/datafusion/src/partition_count_pushdown.rs @@ -0,0 +1,550 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! Answers `SELECT , COUNT(*) ... GROUP BY ` from +//! manifests. +//! +//! DataFusion's `aggregate_statistics` only folds an ungrouped `COUNT(*)`, and it +//! still needs the scan planned first — every live file's metadata and column +//! statistics held as splits, which is what runs out of memory on very large +//! tables. A grouped count additionally opens every data file. +//! +//! Whenever the grouping keys are partition columns and the filter is decided by +//! partition values alone, the answer is a function of the manifests. This rule +//! rewrites +//! +//! ```text +//! Aggregate: groupBy=[[t.dt]], aggr=[[count(1)]] +//! TableScan: t, full_filters=[t.region = 'eu'] +//! ``` +//! +//! into a `SUM` over [`Table::partition_row_counts_with_filter`], which streams +//! manifests in bounded memory and counts data-evolution row ranges once: +//! +//! ```text +//! Projection: t.dt, coalesce(sum(row_count), 0) AS count(1) +//! Aggregate: groupBy=[[t.dt]], aggr=[[sum(row_count)]] +//! TableScan: t (PartitionRowCountProvider) +//! ``` +//! +//! Physical planning pins the selected snapshot, then returns a lazy +//! [`StreamingTableExec`]. `EXPLAIN` does not read manifests; manifest I/O starts +//! when execution polls the plan. + +use std::collections::HashMap; +use std::sync::Arc; + +use async_trait::async_trait; +use datafusion::arrow::array::{ArrayRef, Int64Array, RecordBatch}; +use datafusion::arrow::datatypes::{DataType, Field, Schema, SchemaRef}; +use datafusion::catalog::Session; +use datafusion::common::tree_node::Transformed; +use datafusion::common::{ + internal_err, project_schema, Column, DataFusionError, Result as DFResult, ScalarValue, +}; +use datafusion::datasource::{provider_as_source, source_as_provider, TableProvider, TableType}; +use datafusion::execution::{SendableRecordBatchStream, SessionState, TaskContext}; +use datafusion::functions::core::expr_fn::coalesce; +use datafusion::functions_aggregate::count::count_udaf; +use datafusion::functions_aggregate::expr_fn::{count, sum}; +use datafusion::logical_expr::expr::AggregateFunction; +use datafusion::logical_expr::{ + col, lit, Aggregate, Expr, LogicalPlan, LogicalPlanBuilder, TableScan, TableSource, +}; +use datafusion::optimizer::{ApplyOrder, OptimizerConfig, OptimizerRule}; +use datafusion::physical_plan::stream::RecordBatchStreamAdapter; +use datafusion::physical_plan::streaming::{PartitionStream, StreamingTableExec}; +use datafusion::physical_plan::ExecutionPlan; +use datafusion::sql::TableReference; +use futures::{stream, TryStreamExt}; +use paimon::spec::{CoreOptions, DataField, Datum, Predicate}; +use paimon::table::Table; + +use crate::error::to_datafusion_error; +use crate::filter_pushdown::analyze_filters; +use crate::physical_plan::scan::datum_to_scalar; +use crate::table::PaimonTableProvider; + +const ROW_COUNT_COLUMN: &str = "__paimon_partition_row_count"; + +#[derive(Debug)] +pub(crate) struct PushDownPartitionCount; + +impl OptimizerRule for PushDownPartitionCount { + fn name(&self) -> &str { + "paimon_push_down_partition_count" + } + + fn apply_order(&self) -> Option { + Some(ApplyOrder::TopDown) + } + + fn rewrite( + &self, + plan: LogicalPlan, + _config: &dyn OptimizerConfig, + ) -> DFResult> { + let LogicalPlan::Aggregate(aggregate) = &plan else { + return Ok(Transformed::no(plan)); + }; + match rewrite_aggregate(aggregate)? { + Some(rewritten) => Ok(Transformed::yes(rewritten)), + None => Ok(Transformed::no(plan)), + } + } +} + +fn rewrite_aggregate(aggregate: &Aggregate) -> DFResult> { + let (scan, qualifier) = match aggregate.input.as_ref() { + LogicalPlan::TableScan(scan) => (scan, &scan.table_name), + LogicalPlan::SubqueryAlias(alias) => { + let LogicalPlan::TableScan(scan) = alias.input.as_ref() else { + return Ok(None); + }; + (scan, &alias.alias) + } + _ => return Ok(None), + }; + if scan.fetch.is_some() + || aggregate.aggr_expr.is_empty() + || !aggregate.aggr_expr.iter().all(is_count_star) + { + return Ok(None); + } + let Ok(provider) = source_as_provider(&scan.source) else { + return Ok(None); + }; + let Some(paimon) = provider.downcast_ref::() else { + return Ok(None); + }; + let table = paimon.table(); + let table_schema = table.schema(); + // Manifest row counts of a primary-key table are physical: several versions + // of a key count separately until they are merged at read time. + if CoreOptions::new(table_schema.options()).is_format_table() + || !table_schema.primary_keys().is_empty() + { + return Ok(None); + } + + let partition_keys = table_schema.partition_keys(); + // The synthetic scan must not contain duplicate qualified field names. + if partition_keys.iter().any(|name| name == ROW_COUNT_COLUMN) { + return Ok(None); + } + let groups_by_partition_columns = aggregate + .group_expr + .iter() + .all(|expr| matches!(expr, Expr::Column(column) if partition_keys.contains(&column.name))); + if !groups_by_partition_columns { + return Ok(None); + } + + // Every filter must be decided by partition values alone, with nothing left + // for DataFusion to re-check on rows. + let analysis = analyze_filters(&scan.filters, table_schema.fields(), true); + if analysis.requires_residual { + return Ok(None); + } + match &analysis.pushed_predicate { + Some(predicate) => { + if !table.new_read_builder().is_exact_filter_pushdown(predicate) { + return Ok(None); + } + } + None if !scan.filters.is_empty() => return Ok(None), + None => {} + } + + let partition_fields = table_schema.partition_fields(); + let arrow_schema = paimon.schema(); + let mut fields = Vec::with_capacity(partition_fields.len() + 1); + for field in &partition_fields { + let Ok(arrow_field) = arrow_schema.field_with_name(field.name()) else { + return Ok(None); + }; + fields.push(arrow_field.clone()); + } + fields.push(Field::new(ROW_COUNT_COLUMN, DataType::Int64, false)); + + let counts = PartitionRowCountProvider { + table: table.clone(), + partition_fields, + predicate: analysis.pushed_predicate, + schema: Arc::new(Schema::new(fields)), + fallback_provider: paimon.clone(), + table_name: scan.table_name.clone(), + filters: scan.filters.clone(), + }; + // Only the synthetic scan changes qualifier; fallback filters still refer + // to the original scan's table name. + let counts_scan = LogicalPlan::TableScan(TableScan::try_new( + qualifier.clone(), + provider_as_source(Arc::new(counts)), + None, + vec![], + None, + )?); + + let row_count = Expr::Column(Column::new(Some(qualifier.clone()), ROW_COUNT_COLUMN)); + let summed = LogicalPlan::Aggregate(Aggregate::try_new( + Arc::new(counts_scan), + aggregate.group_expr.clone(), + vec![sum(row_count)], + )?); + + // Reproduce the original output columns exactly: grouping keys keep their + // names, and each COUNT(*) reads the one SUM. An ungrouped aggregate over no + // partitions sums to NULL where COUNT(*) is 0. + let group_len = aggregate.group_expr.len(); + let summed_column = Expr::Column(Column::from(summed.schema().qualified_field(group_len))); + let mut projection = Vec::with_capacity(aggregate.schema.fields().len()); + for index in 0..aggregate.schema.fields().len() { + let (qualifier, field) = aggregate.schema.qualified_field(index); + if index < group_len { + projection.push(Expr::Column(Column::new(qualifier.cloned(), field.name()))); + } else { + projection.push( + coalesce(vec![summed_column.clone(), lit(0i64)]) + .alias_qualified(qualifier.cloned(), field.name()), + ); + } + } + LogicalPlanBuilder::from(summed) + .project(projection)? + .build() + .map(Some) +} + +/// `COUNT(*)` / `COUNT()` with no DISTINCT, FILTER or ORDER BY. +fn is_count_star(expr: &Expr) -> bool { + let expr = match expr { + Expr::Alias(alias) => alias.expr.as_ref(), + other => other, + }; + let Expr::AggregateFunction(AggregateFunction { func, params }) = expr else { + return false; + }; + func == &count_udaf() + && !params.distinct + && params.filter.is_none() + && params.order_by.is_empty() + && matches!(params.args.as_slice(), [Expr::Literal(value, _)] if !value.is_null()) +} + +/// One row per live partition: its typed partition values and its real row count. +/// Planning pins the snapshot and constructs a lazy [`StreamingTableExec`]. +struct PartitionRowCountProvider { + table: Table, + partition_fields: Vec, + predicate: Option, + schema: SchemaRef, + // The provider this scan replaced, kept for partitions whose count the + // manifests cannot give exactly. + fallback_provider: PaimonTableProvider, + table_name: TableReference, + filters: Vec, +} + +impl std::fmt::Debug for PartitionRowCountProvider { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("PartitionRowCountProvider") + .field("table", &self.table.identifier()) + .field("predicate", &self.predicate) + .finish() + } +} + +#[async_trait] +impl TableProvider for PartitionRowCountProvider { + fn schema(&self) -> SchemaRef { + Arc::clone(&self.schema) + } + + fn table_type(&self) -> TableType { + TableType::View + } + + async fn scan( + &self, + state: &dyn Session, + projection: Option<&Vec>, + _filters: &[Expr], + _limit: Option, + ) -> DFResult> { + let state = state + .as_any() + .downcast_ref::() + .ok_or_else(|| { + DataFusionError::Internal( + "partition count execution requires a SessionState".to_string(), + ) + })? + .clone(); + let table = self.table.clone(); + let table = crate::runtime::await_with_runtime(async move { + CoreOptions::new(table.schema().options()).ensure_read_authorized()?; + if table.travel_snapshot().is_some() { + return Ok(Some(table)); + } + let selected = table.copy_with_time_travel_strict(HashMap::new()).await?; + let snapshot = match selected.travel_snapshot() { + Some(snapshot) => snapshot.clone(), + None => { + let Some(snapshot) = table.snapshot_manager().get_latest_snapshot().await? + else { + return Ok(None); + }; + snapshot + } + }; + // Pin data, not the snapshot's schema: DDL may have changed field + // positions since the latest data commit or since logical planning. + Ok(Some(table.copy_with_pinned_snapshot(&snapshot))) + }) + .await + .map_err(to_datafusion_error)?; + let fallback_provider = match table.as_ref().and_then(Table::travel_snapshot) { + Some(snapshot) => self + .fallback_provider + .clone() + .with_pinned_snapshot(snapshot), + None => self.fallback_provider.clone(), + }; + let partition: Arc = Arc::new(PartitionRowCountStream::new( + self, + table, + provider_as_source(Arc::new(fallback_provider)), + projection.cloned(), + state, + )?); + Ok(Arc::new(StreamingTableExec::try_new( + Arc::clone(partition.schema()), + vec![partition], + None, + std::iter::empty(), + false, + None, + )?)) + } +} + +#[derive(Clone)] +struct PartitionRowCountStream { + /// The snapshot selected during physical planning, or `None` for an empty table. + table: Option, + partition_fields: Vec, + predicate: Option, + unprojected_schema: SchemaRef, + projection: Option>, + output_schema: SchemaRef, + source: Arc, + table_name: TableReference, + filters: Vec, + state: Arc, +} + +impl PartitionRowCountStream { + fn new( + provider: &PartitionRowCountProvider, + table: Option
, + source: Arc, + projection: Option>, + state: SessionState, + ) -> DFResult { + let output_schema = project_schema(&provider.schema, projection.as_ref())?; + Ok(Self { + table, + partition_fields: provider.partition_fields.clone(), + predicate: provider.predicate.clone(), + unprojected_schema: Arc::clone(&provider.schema), + projection, + output_schema, + source, + table_name: provider.table_name.clone(), + filters: provider.filters.clone(), + state: Arc::new(state), + }) + } + + async fn execute_stream( + &self, + context: Arc, + ) -> DFResult { + let counts = match self.table.clone() { + Some(table) => { + let predicate = self.predicate.clone(); + crate::runtime::await_with_runtime(async move { + table + .exact_partition_row_counts_with_filter(predicate) + .await + }) + .await + .map_err(to_datafusion_error)? + } + None => Some(Vec::new()), + }; + + // Unknown DV cardinalities stop the metadata path before data manifests + // are aggregated. Keep the fallback pinned to the same snapshot. + let Some(mut counts) = counts else { + log::warn!( + "Partition count metadata cannot determine exact counts; falling back to a data scan \ + (table={}, snapshot_id={:?})", + self.table_name, + self.table.as_ref().and_then(Table::travel_snapshot).map(|snapshot| snapshot.id()), + ); + let plan = crate::runtime::await_with_runtime(self.scan_by_reading()).await?; + if plan.schema() != self.output_schema { + return internal_err!( + "partition count fallback schema mismatch: expected {:?}, got {:?}", + self.output_schema, + plan.schema() + ); + } + return datafusion::physical_plan::execute_stream(plan, context); + }; + + // A partition with no surviving rows must not create a GROUP BY key. + counts.retain(|count| count.record_count != Some(0)); + let batch = self.counts_to_batch(&counts)?; + Ok(Box::pin(RecordBatchStreamAdapter::new( + Arc::clone(&self.output_schema), + Box::pin(stream::iter([Ok(batch)])), + ))) + } + + fn counts_to_batch( + &self, + counts: &[paimon::table::PartitionRowCount], + ) -> DFResult { + let all_columns = (0..self.unprojected_schema.fields().len()).collect::>(); + let projection = self.projection.as_deref().unwrap_or(&all_columns); + let mut columns: Vec = Vec::with_capacity(projection.len()); + for &index in projection { + if index == self.partition_fields.len() { + columns.push(Arc::new(Int64Array::from_iter_values( + counts.iter().filter_map(|count| count.record_count), + ))); + continue; + } + + let field = &self.partition_fields[index]; + let arrow_type = self.unprojected_schema.field(index).data_type(); + let mut values = Vec::with_capacity(counts.len()); + for count in counts { + let datum = count + .partition_row + .get_datum(index, field.data_type()) + .map_err(to_datafusion_error)?; + values.push(match datum { + None => ScalarValue::try_from(arrow_type)?, + Some(datum) => { + partition_datum_to_scalar(datum, arrow_type).ok_or_else(|| { + DataFusionError::Internal(format!( + "cannot represent partition column '{}' as {arrow_type}", + field.name() + )) + })? + } + }); + } + columns.push(if values.is_empty() { + datafusion::arrow::array::new_empty_array(arrow_type) + } else { + ScalarValue::iter_to_array(values)? + }); + } + + Ok(RecordBatch::try_new( + Arc::clone(&self.output_schema), + columns, + )?) + } + + /// The same rows, computed the ordinary way: count the original scan per partition. + async fn scan_by_reading(&self) -> DFResult> { + let source_schema = self.source.schema(); + let partition_indices = self + .partition_fields + .iter() + .map(|field| source_schema.index_of(field.name())) + .collect::, _>>()?; + let scan = LogicalPlan::TableScan(TableScan::try_new( + self.table_name.clone(), + Arc::clone(&self.source), + Some(partition_indices), + self.filters.clone(), + None, + )?); + + let group_by: Vec = self + .partition_fields + .iter() + .map(|field| col(Column::new(Some(self.table_name.clone()), field.name()))) + .collect(); + let mut output = group_by.clone(); + output.push(col(Column::from_name(ROW_COUNT_COLUMN))); + let selected = match &self.projection { + Some(indices) => indices.iter().map(|&index| output[index].clone()).collect(), + None => output.clone(), + }; + let counted = LogicalPlanBuilder::from(scan) + .aggregate(group_by, vec![count(lit(1i64)).alias(ROW_COUNT_COLUMN)])? + .project(selected)? + .build()?; + + // Planned directly, not through the optimizer: this rule must not + // rewrite its own fallback. + self.state + .query_planner() + .create_physical_plan(&counted, &self.state) + .await + } +} + +impl std::fmt::Debug for PartitionRowCountStream { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("PartitionRowCountStream") + .field("table", &self.table_name) + .field("predicate", &self.predicate) + .finish() + } +} + +impl PartitionStream for PartitionRowCountStream { + fn schema(&self) -> &SchemaRef { + &self.output_schema + } + + fn execute(&self, context: Arc) -> SendableRecordBatchStream { + let partition = self.clone(); + let stream = + stream::once(async move { partition.execute_stream(context).await }).try_flatten(); + Box::pin(RecordBatchStreamAdapter::new( + Arc::clone(&self.output_schema), + stream, + )) + } +} + +fn partition_datum_to_scalar(value: Datum, data_type: &DataType) -> Option { + match (value, data_type) { + (Datum::Bytes(value), DataType::Binary) => Some(ScalarValue::Binary(Some(value))), + (value, data_type) => datum_to_scalar(value, data_type), + } +} diff --git a/crates/integrations/datafusion/src/physical_plan/scan.rs b/crates/integrations/datafusion/src/physical_plan/scan.rs index 01a9909f3..f75f0f693 100644 --- a/crates/integrations/datafusion/src/physical_plan/scan.rs +++ b/crates/integrations/datafusion/src/physical_plan/scan.rs @@ -692,7 +692,7 @@ impl ColumnStatsAccumulator { } } -fn datum_to_scalar(value: Datum, data_type: &ArrowDataType) -> Option { +pub(crate) fn datum_to_scalar(value: Datum, data_type: &ArrowDataType) -> Option { match (value, data_type) { (Datum::Bool(value), ArrowDataType::Boolean) => Some(ScalarValue::Boolean(Some(value))), (Datum::TinyInt(value), ArrowDataType::Int8) => Some(ScalarValue::Int8(Some(value))), diff --git a/crates/integrations/datafusion/src/table/mod.rs b/crates/integrations/datafusion/src/table/mod.rs index f7db6c3e2..5d04761e4 100644 --- a/crates/integrations/datafusion/src/table/mod.rs +++ b/crates/integrations/datafusion/src/table/mod.rs @@ -32,7 +32,7 @@ use datafusion::logical_expr::dml::InsertOp; use datafusion::logical_expr::{Expr, TableProviderFilterPushDown}; use datafusion::physical_plan::ExecutionPlan; use paimon::spec::{ - BigIntType, CoreOptions, DataField, DataType, ROW_ID_FIELD_ID, ROW_ID_FIELD_NAME, + BigIntType, CoreOptions, DataField, DataType, Snapshot, ROW_ID_FIELD_ID, ROW_ID_FIELD_NAME, }; use paimon::table::Table; @@ -173,6 +173,11 @@ impl PaimonTableProvider { pub fn table(&self) -> &Table { &self.table } + + pub(crate) fn with_pinned_snapshot(mut self, snapshot: &Snapshot) -> Self { + self.table = self.table.copy_with_pinned_snapshot(snapshot); + self + } } /// Build a `CREATE TABLE` DDL string for a Paimon table. diff --git a/crates/integrations/datafusion/tests/partition_count_pushdown.rs b/crates/integrations/datafusion/tests/partition_count_pushdown.rs new file mode 100644 index 000000000..357d72820 --- /dev/null +++ b/crates/integrations/datafusion/tests/partition_count_pushdown.rs @@ -0,0 +1,935 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! `SELECT , COUNT(*) ... GROUP BY ` must be +//! answered from manifests — no table scan in the plan — and still agree with +//! counting rows, including on data-evolution tables. + +use std::sync::Arc; + +use datafusion::arrow::array::{Array, Int64Array}; +use datafusion::arrow::util::display::array_value_to_string; +use datafusion::physical_plan::displayable; +use paimon::catalog::Identifier; +use paimon::spec::IndexManifest; +use paimon::table::SnapshotManager; +use paimon::{Catalog, CatalogOptions, FileSystemCatalog, Options}; +use paimon_datafusion::SQLContext; +use tempfile::TempDir; + +async fn exec(ctx: &SQLContext, sql: &str) { + ctx.sql(sql) + .await + .unwrap_or_else(|e| panic!("Failed to plan `{sql}`: {e}")) + .collect() + .await + .unwrap_or_else(|e| panic!("Failed to execute `{sql}`: {e}")); +} + +/// Rows as `(leading columns joined by '|', trailing Int64 count)`, sorted. +async fn rows(ctx: &SQLContext, sql: &str) -> Vec<(String, i64)> { + let batches = ctx.sql(sql).await.unwrap().collect().await.unwrap(); + let mut out = Vec::new(); + for batch in &batches { + let last = batch.num_columns() - 1; + let counts = batch + .column(last) + .as_any() + .downcast_ref::() + .unwrap_or_else(|| panic!("last column of `{sql}` must be Int64")); + for row in 0..batch.num_rows() { + let key = (0..last) + .map(|c| array_value_to_string(batch.column(c), row).unwrap()) + .collect::>() + .join("|"); + assert!(!counts.is_null(row)); + out.push((key, counts.value(row))); + } + } + out.sort(); + out +} + +async fn scans_table(ctx: &SQLContext, sql: &str) -> bool { + let plan = ctx + .sql(sql) + .await + .unwrap() + .create_physical_plan() + .await + .unwrap(); + let rendered = displayable(plan.as_ref()).indent(true).to_string(); + rendered.contains("PaimonTableScan") +} + +fn row(key: &str, count: i64) -> (String, i64) { + (key.to_string(), count) +} + +/// A data-evolution table partitioned by `(dt, content_key)` whose `name` column +/// was rewritten by MERGE INTO, so several files cover the same rows. +async fn setup() -> (TempDir, Arc, SQLContext) { + let temp_dir = TempDir::new().unwrap(); + let mut options = Options::new(); + options.set( + CatalogOptions::WAREHOUSE, + format!("file://{}", temp_dir.path().display()), + ); + let catalog = Arc::new(FileSystemCatalog::new(options).unwrap()); + let mut ctx = SQLContext::new(); + ctx.register_catalog("paimon", catalog.clone()) + .await + .unwrap(); + + exec(&ctx, "CREATE SCHEMA paimon.test_db").await; + exec( + &ctx, + "CREATE TABLE paimon.test_db.t (\ + id INT NOT NULL, name STRING, dt STRING, content_key STRING\ + ) PARTITIONED BY (dt, content_key) WITH (\ + 'row-tracking.enabled' = 'true',\ + 'data-evolution.enabled' = 'true'\ + )", + ) + .await; + exec( + &ctx, + "INSERT INTO paimon.test_db.t (id, name, dt, content_key) VALUES \ + (1, 'a', '2024-01-01', 'head'), (2, 'b', '2024-01-01', 'head'), \ + (3, 'c', '2024-01-01', 'tail'), (4, 'd', '2024-01-02', 'head')", + ) + .await; + exec( + &ctx, + "INSERT INTO paimon.test_db.t (id, name, dt, content_key) VALUES \ + (5, 'e', '2024-01-02', 'head'), (6, 'f', '2024-01-03', 'tail')", + ) + .await; + exec( + &ctx, + "CREATE TEMPORARY TABLE paimon.test_db.src AS \ + SELECT * FROM (VALUES (1, 'a2'), (2, 'b2'), (4, 'd2')) AS s(id, name)", + ) + .await; + exec( + &ctx, + "MERGE INTO paimon.test_db.t t USING paimon.test_db.src s ON t.id = s.id \ + WHEN MATCHED THEN UPDATE SET name = s.name", + ) + .await; + (temp_dir, catalog, ctx) +} + +#[tokio::test] +async fn test_grouped_count_with_partition_filter_is_answered_from_manifests() { + let (_tmp, _catalog, ctx) = setup().await; + let sql = "SELECT dt, COUNT(*) AS row_count FROM paimon.test_db.t \ + WHERE content_key = 'head' GROUP BY dt"; + + assert!( + !scans_table(&ctx, sql).await, + "count must not scan the table" + ); + assert_eq!( + rows(&ctx, sql).await, + vec![row("2024-01-01", 2), row("2024-01-02", 2)] + ); + + // COUNT(id) is not rewritten and really reads the rows: an independent oracle. + let oracle = "SELECT dt, COUNT(id) FROM paimon.test_db.t \ + WHERE content_key = 'head' GROUP BY dt"; + assert!(scans_table(&ctx, oracle).await); + assert_eq!(rows(&ctx, oracle).await, rows(&ctx, sql).await); +} + +#[tokio::test] +async fn test_mixed_partition_predicates_preserve_count_exactness() { + use paimon::spec::{Datum, Predicate, PredicateBuilder}; + + let (_tmp, catalog, ctx) = setup().await; + for condition in [ + "dt IN ('2024-01-01', '2024-01-02') AND dt > '2024-01-01'", + "dt IN ('2024-01-01', '2024-01-02', '2024-01-03', '2024-01-04', '2024-01-05') AND dt > '2024-01-01'", + ] { + let sql = format!("SELECT dt, COUNT(*) FROM paimon.test_db.t WHERE content_key = 'head' AND {condition} GROUP BY dt"); + assert!(!scans_table(&ctx, &sql).await, "{sql}"); + assert_eq!(rows(&ctx, &sql).await, vec![row("2024-01-02", 2)], "{sql}"); + assert_eq!( + rows(&ctx, &sql.replace("COUNT(*)", "COUNT(id)")).await, + vec![row("2024-01-02", 2)], + "ordinary scan: {sql}" + ); + } + let sql = "SELECT COUNT(*) FROM paimon.test_db.t WHERE content_key = 'head' AND dt = '2024-01-01' AND dt > '2024-01-01'"; + assert_eq!(rows(&ctx, sql).await, vec![row("", 0)]); + + // Native predicates bypass SQL simplification. Both count APIs must enforce + // the full predicate, independently of the ordinary scan's shared helper. + let table = catalog + .get_table(&Identifier::new("test_db", "t")) + .await + .unwrap(); + let pb = PredicateBuilder::new(table.schema().fields()); + let greater = pb + .greater_than("dt", Datum::String("2024-01-01".into())) + .unwrap(); + let head = pb + .equal("content_key", Datum::String("head".into())) + .unwrap(); + let dates = pb + .is_in( + "dt", + vec![ + Datum::String("2024-01-01".into()), + Datum::String("2024-01-02".into()), + ], + ) + .unwrap(); + for (date_filter, expected) in [ + (dates, vec![("2024-01-02", "head", Some(2))]), + ( + pb.equal("dt", Datum::String("2024-01-01".into())).unwrap(), + vec![], + ), + ] { + let predicate = Predicate::and(vec![date_filter, greater.clone(), head.clone()]); + let partial = table + .partition_row_counts_with_filter(Some(predicate.clone())) + .await + .unwrap(); + let exact = table + .exact_partition_row_counts_with_filter(Some(predicate)) + .await + .unwrap() + .unwrap(); + assert_eq!(partial, exact); + let actual: Vec<_> = exact + .iter() + .map(|count| { + ( + count.partition_row.get_string(0).unwrap(), + count.partition_row.get_string(1).unwrap(), + count.record_count, + ) + }) + .collect(); + assert_eq!(actual, expected); + } +} + +#[tokio::test] +async fn test_table_alias_preserves_count_pushdown_and_eligibility() { + let (_tmp, _catalog, ctx) = setup().await; + for sql in [ + "SELECT p.dt, COUNT(*) FROM paimon.test_db.t AS p GROUP BY p.dt", + "SELECT p.dt AS day, COUNT(*) AS n, COUNT(*) AS n2 FROM paimon.test_db.t AS p \ + WHERE p.content_key = 'head' GROUP BY p.dt ORDER BY day LIMIT 1", + "SELECT COUNT(*) FROM paimon.test_db.t AS p WHERE p.content_key = 'tail'", + "SELECT p.dt, COUNT(*) FROM paimon.test_db.t AS p WHERE p.dt = 'missing' GROUP BY p.dt", + "SELECT p.dt, COUNT(*) FROM paimon.test_db.t VERSION AS OF 1 AS p GROUP BY p.dt", + ] { + assert!(!scans_table(&ctx, sql).await, "{sql}"); + assert_eq!( + rows(&ctx, sql).await, + rows(&ctx, &sql.replace("COUNT(*)", "COUNT(p.id)")).await, + "{sql}" + ); + } + for sql in [ + "SELECT p.dt, COUNT(*) FROM paimon.test_db.t AS p WHERE p.id > 2 GROUP BY p.dt", + "SELECT p.name, COUNT(*) FROM paimon.test_db.t AS p GROUP BY p.name", + "SELECT p.dt, COUNT(*) FROM (SELECT * FROM paimon.test_db.t ORDER BY id LIMIT 2) AS p GROUP BY p.dt", + ] { + assert!(scans_table(&ctx, sql).await, "{sql}"); + assert_eq!(rows(&ctx, sql).await, rows(&ctx, &sql.replace("COUNT(*)", "COUNT(p.id)")).await, "{sql}"); + } +} + +#[tokio::test] +async fn test_distinct_and_grouping_without_count_skip_rewrite() { + let (_tmp, _catalog, ctx) = setup().await; + for sql in [ + "SELECT DISTINCT dt FROM paimon.test_db.t", + "SELECT dt FROM paimon.test_db.t GROUP BY dt", + ] { + assert!(scans_table(&ctx, sql).await, "{sql}"); + let batches = ctx.sql(sql).await.unwrap().collect().await.unwrap(); + let mut values = Vec::new(); + for batch in batches { + for row in 0..batch.num_rows() { + values.push(array_value_to_string(batch.column(0), row).unwrap()); + } + } + values.sort(); + assert_eq!(values, ["2024-01-01", "2024-01-02", "2024-01-03"]); + } + // The outer COUNT needs no projected columns from the DISTINCT result. + let sql = "SELECT COUNT(*) FROM (SELECT DISTINCT dt FROM paimon.test_db.t) d"; + assert!(scans_table(&ctx, sql).await); + assert_eq!(rows(&ctx, sql).await, vec![row("", 3)]); +} + +#[tokio::test] +async fn test_physical_plan_pins_snapshot() { + let (_tmp, _catalog, ctx) = setup().await; + let plan = ctx + .sql("SELECT dt, COUNT(*) FROM paimon.test_db.t GROUP BY dt") + .await + .unwrap() + .create_physical_plan() + .await + .unwrap(); + + exec( + &ctx, + "INSERT INTO paimon.test_db.t (id, name, dt, content_key) \ + VALUES (7, 'g', '2099-01-01', 'head')", + ) + .await; + let batches = datafusion::physical_plan::collect(plan, ctx.ctx().task_ctx()) + .await + .unwrap(); + assert_eq!( + batches.iter().map(|batch| batch.num_rows()).sum::(), + 3 + ); + assert_eq!( + rows( + &ctx, + "SELECT dt, COUNT(*) FROM paimon.test_db.t GROUP BY dt" + ) + .await + .len(), + 4 + ); +} + +#[tokio::test] +async fn test_manifest_reads_are_deferred_until_execution() { + let (_tmp, catalog, ctx) = setup().await; + let sql = "SELECT dt, COUNT(*) FROM paimon.test_db.t GROUP BY dt"; + let plan = ctx + .sql(sql) + .await + .unwrap() + .create_physical_plan() + .await + .unwrap(); + let rendered = displayable(plan.as_ref()).indent(true).to_string(); + assert!(rendered.contains("StreamingTableExec"), "{rendered}"); + assert!(!rendered.contains("PaimonTableScan"), "{rendered}"); + + // If planning had already read the manifests, deleting the list now would + // not affect this plan. Lazy execution must observe the missing file. + let table = catalog + .get_table(&Identifier::new("test_db", "t")) + .await + .unwrap(); + let snapshots = SnapshotManager::new(table.file_io().clone(), table.location().to_string()); + let snapshot = snapshots.get_latest_snapshot().await.unwrap().unwrap(); + table + .file_io() + .delete_file(&snapshots.manifest_path(snapshot.delta_manifest_list())) + .await + .unwrap(); + + ctx.sql(&format!("EXPLAIN {sql}")) + .await + .unwrap() + .collect() + .await + .expect("EXPLAIN must not read manifests"); + assert!( + datafusion::physical_plan::collect(plan, ctx.ctx().task_ctx()) + .await + .is_err(), + "manifest reads must happen during execution" + ); +} + +#[tokio::test] +async fn test_grouped_count_shapes() { + let (_tmp, _catalog, ctx) = setup().await; + + let all_keys = + "SELECT dt, content_key, COUNT(*) FROM paimon.test_db.t GROUP BY dt, content_key"; + assert!(!scans_table(&ctx, all_keys).await); + assert_eq!( + rows(&ctx, all_keys).await, + vec![ + row("2024-01-01|head", 2), + row("2024-01-01|tail", 1), + row("2024-01-02|head", 2), + row("2024-01-03|tail", 1), + ] + ); + + let filtered = "SELECT content_key, COUNT(*) AS c FROM paimon.test_db.t \ + WHERE dt IN ('2024-01-01', '2024-01-03') \ + GROUP BY content_key HAVING COUNT(*) > 1 ORDER BY c DESC"; + assert!(!scans_table(&ctx, filtered).await); + assert_eq!( + rows(&ctx, filtered).await, + vec![row("head", 2), row("tail", 2)] + ); + + // A fetch-capable source must not truncate partition counts below aggregation. + let limited = "SELECT dt, COUNT(*) FROM paimon.test_db.t \ + GROUP BY dt ORDER BY dt LIMIT 1"; + assert_eq!(rows(&ctx, limited).await, vec![row("2024-01-01", 3)]); + + let ungrouped = "SELECT COUNT(*) FROM paimon.test_db.t WHERE content_key = 'tail'"; + assert!(!scans_table(&ctx, ungrouped).await); + assert_eq!(rows(&ctx, ungrouped).await, vec![row("", 2)]); + + let no_match = "SELECT COUNT(*) FROM paimon.test_db.t WHERE content_key = 'nope'"; + assert_eq!(rows(&ctx, no_match).await, vec![row("", 0)]); + let no_match_grouped = "SELECT dt, COUNT(*) FROM paimon.test_db.t \ + WHERE content_key = 'nope' GROUP BY dt"; + assert!(rows(&ctx, no_match_grouped).await.is_empty()); + + let historical = "SELECT dt, COUNT(*) FROM paimon.test_db.t VERSION AS OF 1 GROUP BY dt"; + assert!(!scans_table(&ctx, historical).await); + assert_eq!( + rows(&ctx, historical).await, + vec![row("2024-01-01", 3), row("2024-01-02", 1)] + ); +} + +#[tokio::test] +async fn test_count_preserves_query_schema_after_ddl() { + let (_tmp, _catalog, ctx) = setup().await; + // DDL advances the table schema but leaves the latest data snapshot on the + // old schema. Dropping a column moves the partition column's field index. + exec(&ctx, "ALTER TABLE paimon.test_db.t DROP COLUMN name").await; + let sql = "SELECT dt, COUNT(*) FROM paimon.test_db.t \ + WHERE dt = '2024-01-01' GROUP BY dt"; + let oracle = "SELECT dt, COUNT(id) FROM paimon.test_db.t \ + WHERE dt = '2024-01-01' GROUP BY dt"; + assert_eq!(rows(&ctx, oracle).await, vec![row("2024-01-01", 3)]); + assert!(!scans_table(&ctx, sql).await); + assert_eq!(rows(&ctx, sql).await, rows(&ctx, oracle).await); +} + +#[tokio::test] +async fn test_ineligible_counts_still_scan() { + let (_tmp, _catalog, ctx) = setup().await; + + // A data-column filter cannot be decided from manifests. + let data_filter = "SELECT dt, COUNT(*) FROM paimon.test_db.t WHERE id > 2 GROUP BY dt"; + assert!(scans_table(&ctx, data_filter).await); + assert_eq!( + rows(&ctx, data_filter).await, + vec![ + row("2024-01-01", 1), + row("2024-01-02", 2), + row("2024-01-03", 1) + ] + ); + + // Nor can grouping by a data column. + let data_group = "SELECT name, COUNT(*) FROM paimon.test_db.t GROUP BY name"; + assert!(scans_table(&ctx, data_group).await); + assert_eq!(rows(&ctx, data_group).await.len(), 6); + + exec( + &ctx, + "CREATE TABLE paimon.test_db.pk (id INT NOT NULL, dt STRING NOT NULL, v INT, \ + PRIMARY KEY (id, dt)) PARTITIONED BY (dt) WITH ('bucket' = '1')", + ) + .await; + exec( + &ctx, + "INSERT INTO paimon.test_db.pk VALUES (1, 'a', 1), (2, 'a', 1)", + ) + .await; + // A second version of key 1: two physical rows, one logical row. + exec(&ctx, "INSERT INTO paimon.test_db.pk VALUES (1, 'a', 2)").await; + + let primary_key = "SELECT dt, COUNT(*) FROM paimon.test_db.pk GROUP BY dt"; + assert!(scans_table(&ctx, primary_key).await); + assert_eq!(rows(&ctx, primary_key).await, vec![row("a", 2)]); +} + +/// A deletion vector that does not record its cardinality leaves the manifests +/// unable to give an exact count; the rewritten plan must then count by reading. +#[tokio::test] +async fn test_deletion_vector_zero_groups_and_unknown_cardinality_fallback() { + let temp_dir = TempDir::new().unwrap(); + let mut options = Options::new(); + options.set( + CatalogOptions::WAREHOUSE, + format!("file://{}", temp_dir.path().display()), + ); + let catalog = Arc::new(FileSystemCatalog::new(options).unwrap()); + let mut ctx = SQLContext::new(); + ctx.register_catalog("paimon", catalog.clone()) + .await + .unwrap(); + + exec(&ctx, "CREATE SCHEMA paimon.test_db").await; + exec( + &ctx, + "CREATE TABLE paimon.test_db.t (id INT NOT NULL, name STRING, dt STRING) \ + PARTITIONED BY (dt) WITH (\ + 'row-tracking.enabled' = 'true',\ + 'data-evolution.enabled' = 'true',\ + 'deletion-vectors.enabled' = 'true'\ + )", + ) + .await; + exec( + &ctx, + "INSERT INTO paimon.test_db.t (id, name, dt) VALUES \ + (1, 'a', '2024-01-01'), (2, 'b', '2024-01-01'), (3, 'c', '2024-01-02')", + ) + .await; + exec( + &ctx, + "CREATE TEMPORARY TABLE paimon.test_db.del AS SELECT * FROM (VALUES (1)) AS s(id)", + ) + .await; + exec( + &ctx, + "MERGE INTO paimon.test_db.t t USING paimon.test_db.del s ON t.id = s.id \ + WHEN MATCHED THEN DELETE", + ) + .await; + + let sql = "SELECT dt, COUNT(*) FROM paimon.test_db.t GROUP BY dt"; + let expected = vec![row("2024-01-01", 1), row("2024-01-02", 1)]; + assert!(!scans_table(&ctx, sql).await); + assert_eq!(rows(&ctx, sql).await, expected); + + // Removing every row from a partition must remove its GROUP BY key instead + // of producing a synthetic `(partition, 0)` row. + exec( + &ctx, + "CREATE TEMPORARY TABLE paimon.test_db.del2 AS SELECT * FROM (VALUES (3)) AS s(id)", + ) + .await; + exec( + &ctx, + "MERGE INTO paimon.test_db.t t USING paimon.test_db.del2 s ON t.id = s.id \ + WHEN MATCHED THEN DELETE", + ) + .await; + let expected = vec![row("2024-01-01", 1)]; + assert_eq!(rows(&ctx, sql).await, expected); + assert_eq!( + rows( + &ctx, + "SELECT dt, COUNT(id) FROM paimon.test_db.t GROUP BY dt" + ) + .await, + expected + ); + + // Erase the cardinalities, as written by producers that never recorded them. + let table = catalog + .get_table(&Identifier::new("test_db", "t")) + .await + .unwrap(); + let snapshots = SnapshotManager::new(table.file_io().clone(), table.location().to_string()); + let snapshot = snapshots.get_latest_snapshot().await.unwrap().unwrap(); + let path = snapshots.manifest_path(snapshot.index_manifest().unwrap()); + let mut entries = IndexManifest::read(table.file_io(), &path).await.unwrap(); + let mut erased = 0; + for entry in &mut entries { + for vector in entry + .index_file + .deletion_vectors_ranges + .iter_mut() + .flat_map(|ranges| ranges.values_mut()) + { + vector.cardinality = None; + erased += 1; + } + } + assert!( + erased > 0, + "the delete must have produced a deletion vector" + ); + table.file_io().delete_file(&path).await.unwrap(); + IndexManifest::write(table.file_io(), &path, &entries) + .await + .unwrap(); + + // The fallback is selected lazily during execution, so it is intentionally + // absent from the physical plan produced above. + assert!(!scans_table(&ctx, sql).await); + assert_eq!(rows(&ctx, sql).await, expected); + + // Both fields are strings: reading the old snapshot's `name` column at the + // new schema's `dt` index would silently produce a wrong grouping key. + exec(&ctx, "ALTER TABLE paimon.test_db.t DROP COLUMN name").await; + assert_eq!( + rows( + &ctx, + "SELECT dt, COUNT(id) FROM paimon.test_db.t GROUP BY dt" + ) + .await, + expected + ); + assert_eq!(rows(&ctx, sql).await, expected); + assert_eq!( + rows( + &ctx, + "SELECT dt, COUNT(*) FROM paimon.test_db.t \ + WHERE dt = '2024-01-01' GROUP BY dt" + ) + .await, + expected + ); + + // The lazily planned fallback must use the same pinned snapshot. + let plan = ctx + .sql(sql) + .await + .unwrap() + .create_physical_plan() + .await + .unwrap(); + exec( + &ctx, + "INSERT INTO paimon.test_db.t (id, dt) VALUES (4, '2099-01-01')", + ) + .await; + let batches = datafusion::physical_plan::collect(plan, ctx.ctx().task_ctx()) + .await + .unwrap(); + assert_eq!( + batches.iter().map(|batch| batch.num_rows()).sum::(), + 1 + ); + assert_eq!( + rows(&ctx, sql).await, + vec![row("2024-01-01", 1), row("2099-01-01", 1)] + ); +} + +#[tokio::test] +async fn test_ungrouped_count_on_unpartitioned_and_empty_tables() { + let (_tmp, _catalog, ctx) = setup().await; + exec(&ctx, "CREATE TABLE paimon.test_db.flat (id INT NOT NULL)").await; + + let sql = "SELECT COUNT(*) FROM paimon.test_db.flat"; + assert!(!scans_table(&ctx, sql).await); + assert_eq!(rows(&ctx, sql).await, vec![row("", 0)]); + + exec(&ctx, "INSERT INTO paimon.test_db.flat VALUES (1), (2)").await; + exec(&ctx, "INSERT INTO paimon.test_db.flat VALUES (3)").await; + assert_eq!(rows(&ctx, sql).await, vec![row("", 3)]); +} + +async fn setup_deletion_vectors( + second_partition: &str, +) -> (TempDir, Arc, SQLContext) { + let fixture = setup().await; + let ctx = &fixture.2; + exec( + ctx, + "CREATE TABLE paimon.test_db.dv (id INT NOT NULL, dt STRING) \ + PARTITIONED BY (dt) WITH ('row-tracking.enabled'='true', \ + 'data-evolution.enabled'='true', 'deletion-vectors.enabled'='true')", + ) + .await; + exec( + ctx, + "INSERT INTO paimon.test_db.dv (id, dt) VALUES (1, 'known'), (2, 'known')", + ) + .await; + exec(ctx, &format!("INSERT INTO paimon.test_db.dv (id, dt) VALUES (3, '{second_partition}'), (4, '{second_partition}')")).await; + exec( + ctx, + "CREATE TEMPORARY TABLE paimon.test_db.del AS SELECT 3 AS id", + ) + .await; + exec( + ctx, + "MERGE INTO paimon.test_db.dv t USING paimon.test_db.del s \ + ON t.id = s.id WHEN MATCHED THEN DELETE", + ) + .await; + fixture +} + +#[tokio::test] +async fn test_removed_file_deletion_vector_does_not_reduce_count() { + let (_tmp, catalog, ctx) = setup_deletion_vectors("known").await; + let sql = "SELECT dt, COUNT(*) FROM paimon.test_db.dv GROUP BY dt"; + let oracle = "SELECT dt, COUNT(id) FROM paimon.test_db.dv GROUP BY dt"; + assert_eq!(rows(&ctx, sql).await, vec![row("known", 3)]); + assert_eq!(rows(&ctx, oracle).await, vec![row("known", 3)]); + + let table = catalog + .get_table(&Identifier::new("test_db", "dv")) + .await + .unwrap(); + let snapshots = SnapshotManager::new(table.file_io().clone(), table.location().to_owned()); + let snapshot = snapshots.get_latest_snapshot().await.unwrap().unwrap(); + let entries = IndexManifest::read( + table.file_io(), + &snapshots.manifest_path(snapshot.index_manifest().unwrap()), + ) + .await + .unwrap(); + let entry = entries + .iter() + .find(|entry| entry.index_file.deletion_vectors_ranges.is_some()) + .unwrap(); + let ranges = entry.index_file.deletion_vectors_ranges.as_ref().unwrap(); + assert_eq!(ranges.len(), 1); + let (removed_name, vector) = ranges.iter().next().unwrap(); + assert_eq!(vector.cardinality, Some(1)); + let plan = table.new_read_builder().new_scan().plan().await.unwrap(); + let removed = plan + .splits() + .iter() + .flat_map(|s| s.data_files()) + .find(|f| &f.file_name == removed_name) + .unwrap() + .clone(); + assert_eq!(removed.row_count, 2); + + // The public commit API can remove a data file while retaining its DV index. + let mut message = + paimon::table::CommitMessage::new(entry.partition.clone(), entry.bucket, vec![]); + message.deleted_files.push(removed); + table + .new_write_builder() + .new_commit() + .commit(vec![message]) + .await + .unwrap(); + let after = snapshots.get_latest_snapshot().await.unwrap().unwrap(); + assert_eq!(after.index_manifest(), snapshot.index_manifest()); + let plan = table.new_read_builder().new_scan().plan().await.unwrap(); + assert!(plan + .splits() + .iter() + .flat_map(|s| s.data_files()) + .all(|f| &f.file_name != removed_name)); + assert!(scans_table(&ctx, oracle).await); + assert!(!scans_table(&ctx, sql).await); + let expected = rows(&ctx, oracle).await; + assert_eq!(expected, vec![row("known", 2)]); + assert_eq!(rows(&ctx, sql).await, expected); + let counts = table.partition_row_counts().await.unwrap(); + assert_eq!(counts.len(), 1); + assert_eq!(counts[0].record_count, Some(2)); +} + +#[tokio::test] +async fn test_internal_count_column_name_collision_skips_rewrite() { + let (_tmp, _catalog, ctx) = setup().await; + exec(&ctx, "CREATE TABLE paimon.test_db.collision (id INT NOT NULL, __paimon_partition_row_count STRING) PARTITIONED BY (__paimon_partition_row_count)").await; + exec(&ctx, "INSERT INTO paimon.test_db.collision VALUES (1, 'p')").await; + // The existing ungrouped-count optimization is still allowed to run. + let sql = "SELECT COUNT(*) FROM paimon.test_db.collision"; + assert_eq!(rows(&ctx, sql).await, vec![row("", 1)]); + let grouped = "SELECT __paimon_partition_row_count, COUNT(*) \ + FROM paimon.test_db.collision GROUP BY __paimon_partition_row_count"; + assert!(scans_table(&ctx, grouped).await); + assert_eq!(rows(&ctx, grouped).await, vec![row("p", 1)]); +} + +#[derive(Debug, Default)] +struct ReadTrace(std::sync::Mutex>); + +#[async_trait::async_trait] +impl paimon::io::FileBlockCache for ReadTrace { + async fn get(&self, _: &str, _: std::ops::Range) -> Option { + None + } + async fn put(&self, path: &str, _: u64, _: bytes::Bytes) { + // Always miss: a put records one block actually fetched from the backend. + self.0.lock().unwrap().push(path.to_owned()); + } + async fn invalidate_path(&self, _: &str) {} + async fn invalidate_prefix(&self, _: &str) {} +} + +struct FallbackWarnings(std::sync::Mutex>); + +impl log::Log for FallbackWarnings { + fn enabled(&self, metadata: &log::Metadata<'_>) -> bool { + metadata.level() == log::Level::Warn + && metadata.target() == "paimon_datafusion::partition_count_pushdown" + } + + fn log(&self, record: &log::Record<'_>) { + if self.enabled(record.metadata()) { + self.0.lock().unwrap().push(record.args().to_string()); + } + } + + fn flush(&self) {} +} + +static FALLBACK_WARNINGS: FallbackWarnings = FallbackWarnings(std::sync::Mutex::new(Vec::new())); + +#[tokio::test] +async fn test_unknown_cardinality_falls_back_before_data_manifests() { + log::set_logger(&FALLBACK_WARNINGS).unwrap(); + log::set_max_level(log::LevelFilter::Warn); + // Other tests may log concurrently, but only this test registers "observed". + let warnings = || { + FALLBACK_WARNINGS + .0 + .lock() + .unwrap() + .iter() + .filter(|message| message.contains("table=observed,")) + .cloned() + .collect::>() + }; + let (_tmp, catalog, _fixture_ctx) = setup_deletion_vectors("unknown").await; + let table = catalog + .get_table(&Identifier::new("test_db", "dv")) + .await + .unwrap(); + let snapshots = SnapshotManager::new(table.file_io().clone(), table.location().to_owned()); + let snapshot = snapshots.get_latest_snapshot().await.unwrap().unwrap(); + let index_path = snapshots.manifest_path(snapshot.index_manifest().unwrap()); + let mut entries = IndexManifest::read(table.file_io(), &index_path) + .await + .unwrap(); + let mut erased = 0; + for vector in entries + .iter_mut() + .flat_map(|entry| entry.index_file.deletion_vectors_ranges.iter_mut()) + .flat_map(|ranges| ranges.values_mut()) + { + vector.cardinality = None; + erased += 1; + } + assert_eq!(erased, 1); + table.file_io().delete_file(&index_path).await.unwrap(); + IndexManifest::write(table.file_io(), &index_path, &entries) + .await + .unwrap(); + // The existing API still returns all partitions, including partial counts. + let counts = table.partition_row_counts().await.unwrap(); + assert_eq!(counts.len(), 2); + assert_eq!( + counts.iter().filter(|c| c.record_count == Some(2)).count(), + 1 + ); + assert_eq!( + counts.iter().filter(|c| c.record_count.is_none()).count(), + 1 + ); + + let mut manifests = Vec::new(); + for list in [ + snapshot.base_manifest_list(), + snapshot.delta_manifest_list(), + ] { + manifests.extend( + paimon::spec::ManifestList::read(table.file_io(), &snapshots.manifest_path(list)) + .await + .unwrap(), + ); + } + let trace = Arc::new(ReadTrace::default()); + let observed = paimon::Table::new( + table + .file_io() + .clone() + .with_file_block_cache(trace.clone(), 1024 * 1024, "meta,data") + .unwrap(), + table.identifier().clone(), + table.location().to_owned(), + table.schema().clone(), + None, + ); + assert!(observed + .exact_partition_row_counts_with_filter(None) + .await + .unwrap() + .is_none()); + assert!(trace.0.lock().unwrap().iter().all(|path| !manifests + .iter() + .any(|meta| path.ends_with(meta.file_name())))); + + let ctx = SQLContext::new(); + ctx.ctx() + .register_table( + "observed", + Arc::new(paimon_datafusion::PaimonTableProvider::try_new(observed).unwrap()), + ) + .unwrap(); + let oracle = "SELECT dt, COUNT(id) FROM observed GROUP BY dt"; + let sql = "SELECT dt, COUNT(*) FROM observed GROUP BY dt"; + assert!(scans_table(&ctx, oracle).await); + assert!(!scans_table(&ctx, sql).await); + assert!( + warnings().is_empty(), + "planning must not warn about a fallback" + ); + trace.0.lock().unwrap().clear(); + let expected = rows(&ctx, oracle).await; + assert_eq!(expected, vec![row("known", 2), row("unknown", 1)]); + let ordinary_reads = std::mem::take(&mut *trace.0.lock().unwrap()); + assert_eq!(rows(&ctx, sql).await, expected); + let messages = warnings(); + assert_eq!(messages.len(), 1, "one warning per fallback execution"); + assert!(messages[0].contains("metadata cannot determine exact counts")); + assert!(messages[0].contains(&format!("snapshot_id=Some({})", snapshot.id()))); + let optimized_reads = std::mem::take(&mut *trace.0.lock().unwrap()); + for manifest in manifests { + let ordinary = ordinary_reads + .iter() + .filter(|path| path.ends_with(manifest.file_name())) + .count(); + let optimized = optimized_reads + .iter() + .filter(|path| path.ends_with(manifest.file_name())) + .count(); + assert_eq!(ordinary, 1); + assert_eq!( + optimized, ordinary, + "fallback must not aggregate data manifests first" + ); + } + + // An unknown DV outside the selected partitions must not force a fallback. + assert_eq!( + rows( + &ctx, + "SELECT dt, COUNT(*) FROM observed WHERE dt = 'known' GROUP BY dt" + ) + .await, + vec![row("known", 2)] + ); + assert!(trace + .0 + .lock() + .unwrap() + .iter() + .all(|path| !path.ends_with(".parquet"))); + assert_eq!(warnings().len(), 1, "exact metadata must not warn"); + + // Alias qualification must also preserve filters in the deferred fallback. + let aliased = "SELECT src.dt, COUNT(*) FROM observed AS src \ + WHERE src.dt = 'unknown' GROUP BY src.dt"; + assert!(!scans_table(&ctx, aliased).await); + assert_eq!(warnings().len(), 1); + assert_eq!(rows(&ctx, aliased).await, vec![row("unknown", 1)]); + assert_eq!(warnings().len(), 2); +} diff --git a/crates/paimon/src/spec/avro/decode_helpers.rs b/crates/paimon/src/spec/avro/decode_helpers.rs index 38165329b..cc5620424 100644 --- a/crates/paimon/src/spec/avro/decode_helpers.rs +++ b/crates/paimon/src/spec/avro/decode_helpers.rs @@ -77,7 +77,7 @@ pub(crate) fn read_nullable_string_field( Ok(Some(cursor.read_string()?.to_string())) } -const EMPTY_PARTITION: [u8; 4] = [0, 0, 0, 0]; +pub(super) const EMPTY_PARTITION: &[u8] = &[0, 0, 0, 0]; /// Null/missing/empty partition → valid empty BinaryRow (arity=0). pub(crate) fn normalize_partition(partition: Option>) -> Vec { diff --git a/crates/paimon/src/spec/avro/index_manifest_entry_decode.rs b/crates/paimon/src/spec/avro/index_manifest_entry_decode.rs index 7c94ddf6d..4547b7341 100644 --- a/crates/paimon/src/spec/avro/index_manifest_entry_decode.rs +++ b/crates/paimon/src/spec/avro/index_manifest_entry_decode.rs @@ -19,13 +19,14 @@ use super::cursor::AvroCursor; use super::decode::{neg_count_to_usize, AvroRecordDecode}; use super::decode_helpers::{ extract_record_schema, normalize_partition, read_bytes_field, read_int_field, read_long_field, - read_nullable_string_field, read_string_field, + read_nullable_string_field, read_string_field, EMPTY_PARTITION, }; use super::schema::{skip_nullable_field, FieldSchema, WriterSchema}; use crate::spec::index_manifest::IndexManifestEntry; use crate::spec::manifest_common::FileKind; use crate::spec::{DeletionVectorMeta, GlobalIndexMeta, IndexFileMeta}; use indexmap::IndexMap; +use std::collections::HashMap; impl AvroRecordDecode for IndexManifestEntry { fn decode(cursor: &mut AvroCursor, writer_schema: &WriterSchema) -> crate::Result { @@ -96,16 +97,107 @@ impl AvroRecordDecode for IndexManifestEntry { } } +#[derive(Debug)] +pub(crate) struct SlimIndexManifestEntry<'a> { + pub kind: FileKind, + pub partition: &'a [u8], + pub bucket: i32, + pub index_type: &'a str, + /// File-name mappings, with `None` for an unknown cardinality. Later mappings + /// replace earlier ones, just as in the full index-manifest decoder. + pub deletion_vector_cardinalities: HashMap<&'a str, Option>, +} + +pub(crate) fn decode_slim_index_manifest_entry<'a>( + cursor: &mut AvroCursor<'a>, + writer_schema: &WriterSchema, + is_union_wrapped: bool, +) -> crate::Result> { + if is_union_wrapped && cursor.read_union_index()? == 0 { + return Err(crate::Error::UnexpectedError { + message: "avro decode: unexpected null in top-level union".into(), + source: None, + }); + } + + let mut entry = SlimIndexManifestEntry { + kind: FileKind::Add, + partition: EMPTY_PARTITION, + bucket: 0, + index_type: "", + deletion_vector_cardinalities: HashMap::new(), + }; + for field in &writer_schema.fields { + match field.name.as_str() { + "_KIND" => { + entry.kind = match read_int_field(cursor, field.nullable)? { + 0 => FileKind::Add, + 1 => FileKind::Delete, + v => { + return Err(crate::Error::UnexpectedError { + message: format!("unknown FileKind: {v}"), + source: None, + }) + } + } + } + "_PARTITION" => { + if !field.nullable || cursor.read_union_index()? != 0 { + let partition = cursor.read_bytes()?; + if partition.len() >= 4 { + entry.partition = partition; + } + } + } + "_BUCKET" => entry.bucket = read_int_field(cursor, field.nullable)?, + "_INDEX_TYPE" => { + if !field.nullable || cursor.read_union_index()? != 0 { + entry.index_type = cursor.read_string()?; + } + } + "_DELETIONS_VECTORS_RANGES" | "_DELETION_VECTORS_RANGES" => { + entry.deletion_vector_cardinalities = + decode_nullable_dv_cardinalities(cursor, field.nullable, &field.schema)?; + } + _ => skip_nullable_field(cursor, &field.schema, field.nullable)?, + } + } + Ok(entry) +} + +fn decode_nullable_dv_cardinalities<'a>( + cursor: &mut AvroCursor<'a>, + nullable: bool, + schema: &FieldSchema, +) -> crate::Result>> { + let mut cardinalities = HashMap::new(); + visit_nullable_dv_ranges(cursor, nullable, schema, |name, meta| { + cardinalities.insert(name, meta.cardinality.filter(|value| *value >= 0)); + })?; + Ok(cardinalities) +} + fn decode_nullable_dv_ranges( cursor: &mut AvroCursor, nullable: bool, schema: &FieldSchema, ) -> crate::Result>> { - if nullable { - let idx = cursor.read_union_index()?; - if idx == 0 { - return Ok(None); - } + let mut map = IndexMap::new(); + let present = visit_nullable_dv_ranges(cursor, nullable, schema, |name, meta| { + map.insert(name.to_owned(), meta); + })?; + Ok(present.then_some(map)) +} + +/// Visit DV records with borrowed names; `false` means the outer array was null. +fn visit_nullable_dv_ranges<'a>( + cursor: &mut AvroCursor<'a>, + nullable: bool, + schema: &FieldSchema, + mut visit: impl FnMut(&'a str, DeletionVectorMeta), +) -> crate::Result { + if nullable && cursor.read_union_index()? == 0 { + return Ok(false); } let FieldSchema::Array(item_schema) = schema else { return Err(crate::Error::UnexpectedError { @@ -113,7 +205,6 @@ fn decode_nullable_dv_ranges( source: None, }); }; - let mut map = IndexMap::new(); loop { let count = cursor.read_long()?; if count == 0 { @@ -150,13 +241,17 @@ fn decode_nullable_dv_ranges( source: None, }); }; - let mut file_name = String::new(); + let mut file_name = ""; let mut offset = 0; let mut length = 0; let mut cardinality = None; for field in &record.fields { match field.name.as_str() { - "f0" => file_name = read_string_field(cursor, field.nullable)?, + "f0" => { + if !field.nullable || cursor.read_union_index()? != 0 { + file_name = cursor.read_string()?; + } + } "f1" => offset = read_int_field(cursor, field.nullable)?, "f2" => length = read_int_field(cursor, field.nullable)?, "_CARDINALITY" => { @@ -167,7 +262,7 @@ fn decode_nullable_dv_ranges( _ => skip_nullable_field(cursor, &field.schema, field.nullable)?, } } - map.insert( + visit( file_name, DeletionVectorMeta { offset, @@ -177,7 +272,7 @@ fn decode_nullable_dv_ranges( ); } } - Ok(Some(map)) + Ok(true) } fn decode_nullable_global_index( diff --git a/crates/paimon/src/spec/avro/manifest_entry_decode.rs b/crates/paimon/src/spec/avro/manifest_entry_decode.rs index b4af8a36f..5dc4e60c1 100644 --- a/crates/paimon/src/spec/avro/manifest_entry_decode.rs +++ b/crates/paimon/src/spec/avro/manifest_entry_decode.rs @@ -18,11 +18,11 @@ use super::cursor::AvroCursor; use super::decode::{neg_count_to_usize, AvroRecordDecode}; use super::decode_helpers::{ - extract_record_schema, normalize_partition, read_bytes_field, read_int_field, read_long_field, - read_string_field, + extract_record_schema, read_bytes_field, read_int_field, read_long_field, read_string_field, + EMPTY_PARTITION, }; use super::manifest_file_meta_decode::decode_nullable_binary_table_stats; -use super::schema::{skip_nullable_field, FieldSchema, WriterSchema}; +use super::schema::{skip_nullable_field, WriterSchema}; use crate::spec::manifest_common::FileKind; use crate::spec::stats::BinaryTableStats; use crate::spec::DataFileMeta; @@ -31,55 +31,16 @@ use chrono::{DateTime, Utc}; impl AvroRecordDecode for ManifestEntry { fn decode(cursor: &mut AvroCursor, writer_schema: &WriterSchema) -> crate::Result { - let mut kind: Option = None; - let mut partition: Option> = None; - let mut bucket: Option = None; - let mut total_buckets: Option = None; - let mut file: Option = None; - let mut version: Option = None; - - for field in &writer_schema.fields { - match field.name.as_str() { - "_KIND" => { - let v = read_int_field(cursor, field.nullable)?; - kind = Some(match v { - 0 => FileKind::Add, - 1 => FileKind::Delete, - _ => { - return Err(crate::Error::UnexpectedError { - message: format!("unknown FileKind: {v}"), - source: None, - }) - } - }); - } - "_PARTITION" => partition = Some(read_bytes_field(cursor, field.nullable)?), - "_BUCKET" => bucket = Some(read_int_field(cursor, field.nullable)?), - "_TOTAL_BUCKETS" => total_buckets = Some(read_int_field(cursor, field.nullable)?), - "_FILE" => { - file = decode_nullable_data_file_meta(cursor, &field.schema, field.nullable)?; - } - "_VERSION" => version = Some(read_int_field(cursor, field.nullable)?), - _ => skip_nullable_field(cursor, &field.schema, field.nullable)?, - } - } - - Ok(ManifestEntry::new( - kind.unwrap_or(FileKind::Add), - normalize_partition(partition), - bucket.unwrap_or(0), - total_buckets.unwrap_or(0), - file.unwrap_or_else(default_data_file_meta), - version.unwrap_or(0), - )) + // The generic OCF decoder already consumed the top-level union. + decode_manifest_entries_filtered(cursor, writer_schema, false, &mut |_, _, _, _| true)? + .ok_or_else(missing_file_metadata) } } /// Decode ManifestEntry records with a filter applied on lightweight fields. /// -/// Decodes only _KIND, _PARTITION, _BUCKET, _TOTAL_BUCKETS, _VERSION first. -/// If `filter` returns false, skips the expensive _FILE (DataFileMeta) decoding. -/// Returns only entries that pass the filter. +/// When the writer places the lightweight fields before _FILE, rejected entries +/// skip DataFileMeta decoding entirely. Otherwise, retain the entry conservatively. pub(crate) fn decode_manifest_entries_filtered( cursor: &mut AvroCursor, writer_schema: &WriterSchema, @@ -89,34 +50,66 @@ pub(crate) fn decode_manifest_entries_filtered( where F: FnMut(FileKind, &[u8], i32, i32) -> bool, { - if is_union_wrapped { - let idx = cursor.read_union_index()?; - if idx == 0 { - return Err(crate::Error::UnexpectedError { - message: "avro decode: unexpected null in top-level union".into(), - source: None, - }); - } + Ok(decode_manifest_entry_with( + cursor, + writer_schema, + is_union_wrapped, + filter, + decode_data_file_meta, + )? + .map(|(fields, file)| { + ManifestEntry::new( + fields.kind.unwrap_or(FileKind::Add), + fields.partition().to_vec(), + fields.bucket.unwrap_or(0), + fields.total_buckets.unwrap_or(0), + file, + fields.version, + ) + })) +} + +#[derive(Default)] +struct ManifestEntryFields<'a> { + kind: Option, + partition: Option<&'a [u8]>, + bucket: Option, + total_buckets: Option, + version: i32, +} + +impl<'a> ManifestEntryFields<'a> { + fn partition(&self) -> &'a [u8] { + self.partition + .filter(|bytes| bytes.len() >= 4) + .unwrap_or(EMPTY_PARTITION) } +} - // Two-pass decode: first collect lightweight fields and record _FILE position, - // then conditionally decode _FILE. - let mut kind: Option = None; - let mut partition: Option> = None; - let mut bucket: Option = None; - let mut total_buckets: Option = None; - let mut version: Option = None; - let mut file: Option = None; +/// Shared wire walk; collectors choose full or borrowed _FILE decoding. +fn decode_manifest_entry_with<'a, T>( + cursor: &mut AvroCursor<'a>, + writer_schema: &WriterSchema, + is_union_wrapped: bool, + filter: &mut impl FnMut(FileKind, &[u8], i32, i32) -> bool, + mut decode_file: impl FnMut(&mut AvroCursor<'a>, &WriterSchema) -> crate::Result, +) -> crate::Result, T)>> { + if is_union_wrapped && cursor.read_union_index()? == 0 { + return Err(crate::Error::UnexpectedError { + message: "avro decode: unexpected null in top-level union".into(), + source: None, + }); + } + let mut fields = ManifestEntryFields::default(); + let mut file = None; let mut file_skipped = false; - for field in &writer_schema.fields { match field.name.as_str() { "_KIND" => { - let v = read_int_field(cursor, field.nullable)?; - kind = Some(match v { + fields.kind = Some(match read_int_field(cursor, field.nullable)? { 0 => FileKind::Add, 1 => FileKind::Delete, - _ => { + v => { return Err(crate::Error::UnexpectedError { message: format!("unknown FileKind: {v}"), source: None, @@ -124,66 +117,178 @@ where } }); } - "_PARTITION" => partition = Some(read_bytes_field(cursor, field.nullable)?), - "_BUCKET" => bucket = Some(read_int_field(cursor, field.nullable)?), - "_TOTAL_BUCKETS" => total_buckets = Some(read_int_field(cursor, field.nullable)?), + "_PARTITION" => { + fields.partition = + Some(decode_nullable_bytes_ref(cursor, field.nullable)?.unwrap_or(&[])) + } + "_BUCKET" => fields.bucket = Some(read_int_field(cursor, field.nullable)?), + "_TOTAL_BUCKETS" => { + fields.total_buckets = Some(read_int_field(cursor, field.nullable)?) + } "_FILE" => { - let can_filter = kind.is_some() - && partition.is_some() - && bucket.is_some() - && total_buckets.is_some(); - if can_filter { - let k = kind.unwrap_or(FileKind::Add); - let p = partition.as_deref().unwrap_or(&[]); - let b = bucket.unwrap_or(0); - let tb = total_buckets.unwrap_or(0); - if filter(k, p, b, tb) { - file = - decode_nullable_data_file_meta(cursor, &field.schema, field.nullable)?; - } else { + if let (Some(kind), Some(partition), Some(bucket), Some(total_buckets)) = ( + fields.kind, + fields.partition, + fields.bucket, + fields.total_buckets, + ) { + if !filter(kind, partition, bucket, total_buckets) { skip_nullable_field(cursor, &field.schema, field.nullable)?; file_skipped = true; + continue; } - } else { - file = decode_nullable_data_file_meta(cursor, &field.schema, field.nullable)?; } + file = if field.nullable && cursor.read_union_index()? == 0 { + None + } else { + let schema = extract_record_schema(&field.schema).ok_or_else(|| { + crate::Error::UnexpectedError { + message: "avro decode: _FILE field is not a record".into(), + source: None, + } + })?; + Some(decode_file(cursor, schema)?) + }; } - "_VERSION" => version = Some(read_int_field(cursor, field.nullable)?), + "_VERSION" => fields.version = read_int_field(cursor, field.nullable)?, _ => skip_nullable_field(cursor, &field.schema, field.nullable)?, } } - if file_skipped { return Ok(None); } + Ok(Some((fields, file.ok_or_else(missing_file_metadata)?))) +} - Ok(Some(ManifestEntry::new( - kind.unwrap_or(FileKind::Add), - normalize_partition(partition), - bucket.unwrap_or(0), - total_buckets.unwrap_or(0), - file.unwrap_or_else(default_data_file_meta), - version.unwrap_or(0), - ))) +/// Borrowed view of the manifest-entry fields needed to count rows per partition. +/// Everything else (key/value stats, min/max keys, ...) is skipped in place, so +/// decoding never materializes the per-file statistics that dominate manifest size. +/// Ordinary files allocate nothing; files with `extra_files` allocate only that list. +#[derive(Debug)] +pub(crate) struct SlimManifestEntry<'a> { + pub kind: FileKind, + pub partition: &'a [u8], + pub bucket: i32, + pub level: i32, + pub file_name: &'a str, + pub row_count: i64, + pub first_row_id: Option, + pub extra_files: Vec<&'a str>, + pub embedded_index: Option<&'a [u8]>, + pub external_path: Option<&'a str>, } -fn decode_nullable_data_file_meta( - cursor: &mut AvroCursor, - field_schema: &FieldSchema, +/// Decode one manifest entry as a [`SlimManifestEntry`] borrowing from the block. +pub(crate) fn decode_slim_manifest_entry<'a>( + cursor: &mut AvroCursor<'a>, + writer_schema: &WriterSchema, + is_union_wrapped: bool, +) -> crate::Result> { + let (fields, mut entry) = decode_manifest_entry_with( + cursor, + writer_schema, + is_union_wrapped, + &mut |_, _, _, _| true, + decode_slim_data_file, + )? + .ok_or_else(missing_file_metadata)?; + entry.kind = fields.kind.unwrap_or(FileKind::Add); + entry.partition = fields.partition(); + entry.bucket = fields.bucket.unwrap_or(0); + Ok(entry) +} + +fn decode_slim_data_file<'a>( + cursor: &mut AvroCursor<'a>, + writer_schema: &WriterSchema, +) -> crate::Result> { + let mut entry = SlimManifestEntry { + kind: FileKind::Add, + partition: EMPTY_PARTITION, + bucket: 0, + level: 0, + file_name: "", + row_count: DataFileMeta::ROW_COUNT_UNKNOWN, + first_row_id: None, + extra_files: Vec::new(), + embedded_index: None, + external_path: None, + }; + for file_field in &writer_schema.fields { + match file_field.name.as_str() { + "_FILE_NAME" => { + if !file_field.nullable || cursor.read_union_index()? != 0 { + entry.file_name = cursor.read_string()?; + } + } + "_ROW_COUNT" => { + entry.row_count = decode_nullable_long(cursor, file_field.nullable)? + .unwrap_or(DataFileMeta::ROW_COUNT_UNKNOWN) + } + "_LEVEL" => entry.level = read_int_field(cursor, file_field.nullable)?, + "_EXTRA_FILES" => { + entry.extra_files = decode_borrowed_string_array(cursor, file_field.nullable)? + } + "_EMBEDDED_FILE_INDEX" => { + entry.embedded_index = decode_nullable_bytes_ref(cursor, file_field.nullable)? + } + "_EXTERNAL_PATH" => { + entry.external_path = decode_nullable_string_ref(cursor, file_field.nullable)? + } + "_FIRST_ROW_ID" => { + entry.first_row_id = decode_nullable_long(cursor, file_field.nullable)? + } + _ => skip_nullable_field(cursor, &file_field.schema, file_field.nullable)?, + } + } + Ok(entry) +} + +fn decode_borrowed_string_array<'a>( + cursor: &mut AvroCursor<'a>, nullable: bool, -) -> crate::Result> { - if nullable { - let idx = cursor.read_union_index()?; - if idx == 0 { - return Ok(None); +) -> crate::Result> { + if nullable && cursor.read_union_index()? == 0 { + return Ok(Vec::new()); + } + let mut values = Vec::new(); + loop { + let count = cursor.read_long()?; + if count == 0 { + break; + } + let count = if count < 0 { + cursor.skip_long()?; + neg_count_to_usize(count)? + } else { + count as usize + }; + values.reserve(count.min(cursor.remaining())); + for _ in 0..count { + values.push(cursor.read_string()?); } } - let record_schema = - extract_record_schema(field_schema).ok_or_else(|| crate::Error::UnexpectedError { - message: "avro decode: _FILE field is not a record".into(), - source: None, - })?; - decode_data_file_meta(cursor, record_schema).map(Some) + Ok(values) +} + +fn decode_nullable_bytes_ref<'a>( + cursor: &mut AvroCursor<'a>, + nullable: bool, +) -> crate::Result> { + if nullable && cursor.read_union_index()? == 0 { + return Ok(None); + } + cursor.read_bytes().map(Some) +} + +fn decode_nullable_string_ref<'a>( + cursor: &mut AvroCursor<'a>, + nullable: bool, +) -> crate::Result> { + if nullable && cursor.read_union_index()? == 0 { + return Ok(None); + } + cursor.read_string().map(Some) } /// Read string array, handling both `{"type":"array",...}` and `["null", {"type":"array",...}]`. @@ -227,7 +332,7 @@ fn decode_data_file_meta( match field.name.as_str() { "_FILE_NAME" => file_name = Some(read_string_field(cursor, field.nullable)?), "_FILE_SIZE" => file_size = Some(read_long_field(cursor, field.nullable)?), - "_ROW_COUNT" => row_count = Some(read_long_field(cursor, field.nullable)?), + "_ROW_COUNT" => row_count = decode_nullable_long(cursor, field.nullable)?, "_MIN_KEY" => min_key = Some(read_bytes_field(cursor, field.nullable)?), "_MAX_KEY" => max_key = Some(read_bytes_field(cursor, field.nullable)?), "_KEY_STATS" => { @@ -271,7 +376,7 @@ fn decode_data_file_meta( Ok(DataFileMeta { file_name: file_name.unwrap_or_default(), file_size: file_size.unwrap_or(0), - row_count: row_count.unwrap_or(0), + row_count: row_count.unwrap_or(DataFileMeta::ROW_COUNT_UNKNOWN), min_key: min_key.unwrap_or_default(), max_key: max_key.unwrap_or_default(), key_stats: key_stats.unwrap_or_else(BinaryTableStats::empty), @@ -438,28 +543,9 @@ fn decode_nullable_timestamp_millis( Ok(DateTime::from_timestamp(secs, nanos)) } -fn default_data_file_meta() -> DataFileMeta { - DataFileMeta { - file_name: String::new(), - file_size: 0, - row_count: 0, - min_key: vec![], - max_key: vec![], - key_stats: BinaryTableStats::empty(), - value_stats: BinaryTableStats::empty(), - min_sequence_number: 0, - max_sequence_number: 0, - schema_id: 0, - level: 0, - extra_files: vec![], - creation_time: None, - delete_row_count: None, - embedded_index: None, - file_source: None, - value_stats_cols: None, - external_path: None, - first_row_id: None, - write_cols: None, - column_max_sequence_numbers: None, +fn missing_file_metadata() -> crate::Error { + crate::Error::DataInvalid { + message: "manifest entry is missing non-null _FILE metadata".into(), + source: None, } } diff --git a/crates/paimon/src/spec/avro/mod.rs b/crates/paimon/src/spec/avro/mod.rs index f6cbb720c..490c234bf 100644 --- a/crates/paimon/src/spec/avro/mod.rs +++ b/crates/paimon/src/spec/avro/mod.rs @@ -159,6 +159,67 @@ where decode_manifest_streaming(&mut block_iter, &writer_schema, filter) } +pub(crate) use index_manifest_entry_decode::SlimIndexManifestEntry; +pub(crate) use manifest_entry_decode::SlimManifestEntry; + +/// Visit every entry of a manifest file as a borrowed [`SlimManifestEntry`]. +/// +/// Nothing is collected: each entry is handed to `visit` and dropped, and +/// per-file statistics are skipped rather than decoded, so peak memory is the +/// manifest's decompressed block regardless of how many files it lists. +pub(crate) fn visit_slim_manifest_entries( + bytes: &[u8], + shared_cache: &SharedSchemaCache, + visit: &mut F, +) -> crate::Result<()> +where + F: FnMut(SlimManifestEntry<'_>) -> crate::Result<()>, +{ + visit_ocf_records(bytes, shared_cache, |cursor, schema| { + visit(manifest_entry_decode::decode_slim_manifest_entry( + cursor, + schema, + schema.is_union_wrapped, + )?) + }) +} + +/// Visit index-manifest entries without materializing deletion-vector maps. +pub(crate) fn visit_slim_index_manifest_entries( + bytes: &[u8], + shared_cache: &SharedSchemaCache, + visit: &mut F, +) -> crate::Result<()> +where + F: FnMut(SlimIndexManifestEntry<'_>) -> crate::Result<()>, +{ + visit_ocf_records(bytes, shared_cache, |cursor, schema| { + visit( + index_manifest_entry_decode::decode_slim_index_manifest_entry( + cursor, + schema, + schema.is_union_wrapped, + )?, + ) + }) +} + +fn visit_ocf_records( + bytes: &[u8], + shared_cache: &SharedSchemaCache, + mut visit: impl FnMut(&mut AvroCursor<'_>, &WriterSchema) -> crate::Result<()>, +) -> crate::Result<()> { + let (header, mut block_iter) = parse_ocf_streaming(bytes)?; + let writer_schema = shared_cache.get_or_parse(&header.schema_json)?; + while let Some(block) = block_iter.next_block()? { + let mut cursor = AvroCursor::new(&block.data); + for _ in 0..block.object_count { + visit(&mut cursor, &writer_schema)?; + } + } + Ok(()) +} + fn decode_manifest_streaming( block_iter: &mut ocf::OcfBlockIter<'_>, writer_schema: &WriterSchema, diff --git a/crates/paimon/src/spec/index_manifest.rs b/crates/paimon/src/spec/index_manifest.rs index 967bab8c4..830b2ae74 100644 --- a/crates/paimon/src/spec/index_manifest.rs +++ b/crates/paimon/src/spec/index_manifest.rs @@ -216,7 +216,7 @@ mod tests { DeletionVectorMeta { offset: 31, length: 22, - cardinality: None, + cardinality: has_cardinality.then_some(4), }, ), ])), @@ -224,14 +224,53 @@ mod tests { global_index_meta: None, }, }; - // Two entries also catch a cursor shifted past the final DV. - let entries = vec![entry.clone(), entry]; + // The second entry also catches cursor shifts and verifies that + // unknown cardinalities retain their file-name mapping. + let mut partially_unknown = entry.clone(); + partially_unknown + .index_file + .deletion_vectors_ranges + .as_mut() + .unwrap() + .insert( + "data-unknown.parquet".into(), + DeletionVectorMeta { + offset: 53, + length: 18, + cardinality: None, + }, + ); + let entries = vec![entry, partially_unknown]; let bytes = crate::spec::to_avro_bytes(&schema.to_string(), &entries).unwrap(); assert_eq!( IndexManifest::read_from_bytes(&bytes).unwrap(), entries, "nullable_items={nullable_items}, has_cardinality={has_cardinality}" ); + let mut seen = 0; + crate::spec::avro::visit_slim_index_manifest_entries( + &bytes, + &crate::spec::avro::SharedSchemaCache::new(), + &mut |entry| { + let expected = &entries[seen]; + assert_eq!(entry.bucket, expected.bucket); + assert_eq!( + entry.deletion_vector_cardinalities, + expected + .index_file + .deletion_vectors_ranges + .as_ref() + .unwrap() + .iter() + .map(|(name, meta)| (name.as_str(), meta.cardinality)) + .collect() + ); + seen += 1; + Ok(()) + }, + ) + .unwrap(); + assert_eq!(seen, entries.len()); } } } @@ -264,7 +303,7 @@ mod tests { DeletionVectorMeta { offset: 17, length: 31, - cardinality: Some(2), + cardinality: Some(-1), }, )])), external_path: Some("memory:/external/index".into()), @@ -303,10 +342,36 @@ mod tests { "future".into(), Value::Array(vec![Value::String("ignored".into())]), )); + // Last duplicate wins: full decoding preserves the negative cardinality, + // while slim decoding must report it as unknown, not retain the earlier 2. + let last = Value::Union(1, Box::new(Value::Record(fields.clone()))); + fields + .iter_mut() + .find(|(name, _)| name == "_CARDINALITY") + .unwrap() + .1 = Value::Union(1, Box::new(Value::Long(2))); + items.push(last); let mut writer = apache_avro::Writer::new(&schema, Vec::new()); writer.append(value.resolve(&schema).unwrap()).unwrap(); let bytes = writer.into_inner().unwrap(); assert_eq!(IndexManifest::read_from_bytes(&bytes).unwrap(), vec![entry]); + + let mut seen = 0; + crate::spec::avro::visit_slim_index_manifest_entries( + &bytes, + &crate::spec::avro::SharedSchemaCache::new(), + &mut |entry| { + assert_eq!(entry.bucket, 7); + assert_eq!( + entry.deletion_vector_cardinalities, + std::collections::HashMap::from([("data.parquet", None)]) + ); + seen += 1; + Ok(()) + }, + ) + .unwrap(); + assert_eq!(seen, 1); } #[test] diff --git a/crates/paimon/src/table/mod.rs b/crates/paimon/src/table/mod.rs index 74581c019..5f8923c15 100644 --- a/crates/paimon/src/table/mod.rs +++ b/crates/paimon/src/table/mod.rs @@ -69,6 +69,7 @@ mod lumina_index_build_builder; pub(crate) mod merge_tree_split_generator; mod object_table; mod partition_filter; +mod partition_row_count; mod partition_stat; #[cfg(feature = "fulltext")] mod pk_full_text_bucket_search; @@ -156,6 +157,7 @@ pub use incremental_scan::{ }; pub use lumina_index_build_builder::LuminaIndexBuildBuilder; pub use object_table::{ObjectEntry, ObjectTable}; +pub use partition_row_count::PartitionRowCount; pub use partition_stat::PartitionStat; pub use pk_vector_bucket_split::{BucketVectorPayload, BucketVectorSearchSplit}; pub use postpone_bucket_plan::{PostponeBucketPlan, POSTPONE_BUCKET_PLAN_TOTAL_BUCKETS_FIELD}; @@ -511,6 +513,23 @@ impl Table { /// same snapshot. The snapshot's schema is loaded when it differs from the /// current table schema. pub(crate) async fn copy_with_resolved_snapshot(&self, snapshot: &Snapshot) -> Result { + let mut table = self.copy_with_pinned_snapshot(snapshot); + if snapshot.schema_id() != self.schema.id() { + table.schema = self + .schema_manager + .schema(snapshot.schema_id()) + .await? + .copy_with_replaced_options(table.schema.options().clone()); + } + Ok(table) + } + + /// Create a read-only copy pinned to a snapshot from this table and branch, + /// preserving the current read schema and options other than scan selectors. + /// + /// Unlike time travel, pinning must not change the fields used by an already + /// planned query. The resolved snapshot is cached without additional I/O. + pub fn copy_with_pinned_snapshot(&self, snapshot: &Snapshot) -> Self { let mut options = self.schema.options().clone(); for selector in [ SCAN_TIMESTAMP_MILLIS_OPTION, @@ -526,26 +545,12 @@ impl Table { snapshot.id().to_string(), ); - let schema = if snapshot.schema_id() == self.schema.id() { - self.schema.copy_with_replaced_options(options) - } else { - self.schema_manager - .schema(snapshot.schema_id()) - .await? - .copy_with_replaced_options(options) - }; - Ok(Self { - file_io: self.file_io.clone(), - identifier: self.identifier.clone(), - location: self.location.clone(), - schema, - schema_manager: self.schema_manager.clone(), - branch: self.branch.clone(), - branch_reference: self.branch_reference, - rest_env: self.rest_env.clone(), + Self { + schema: self.schema.copy_with_replaced_options(options), time_traveled: true, travel_snapshot: Some(snapshot.clone()), - }) + ..self.clone() + } } /// Create a copy of this table with extra options merged in, switching to diff --git a/crates/paimon/src/table/partition_filter.rs b/crates/paimon/src/table/partition_filter.rs index bbea7104b..17ad4635b 100644 --- a/crates/paimon/src/table/partition_filter.rs +++ b/crates/paimon/src/table/partition_filter.rs @@ -21,8 +21,8 @@ use crate::predicate_stats::data_leaf_may_match; use crate::spec::{ - eval_row, extract_datum, BinaryRow, BinaryRowBuilder, DataField, Datum, Predicate, - PredicateBuilder, PredicateOperator, + eval_row, extract_datum, BinaryRow, BinaryRowBuilder, DataField, Datum, ManifestFileMeta, + Predicate, PredicateBuilder, PredicateOperator, }; use crate::table::stats_filter::FileStatsRows; use std::collections::HashSet; @@ -34,6 +34,11 @@ pub(crate) struct FieldBounds { max: Predicate, } +struct FieldCandidates<'a> { + predicate: &'a Predicate, + values: Vec>, +} + #[derive(Debug, Clone)] pub(crate) enum PartitionFilter { /// Multiple known partitions: O(1) entry matching via HashSet, @@ -53,10 +58,9 @@ impl PartitionFilter { } let num_fields = partition_fields.len(); - let mut field_candidates: Vec>>> = vec![None; num_fields]; - if !collect_eq_candidates(&predicate, &mut field_candidates) { - return PartitionFilter::Predicate(predicate); - } + let mut field_candidates = (0..num_fields).map(|_| None).collect::>(); + let mut residuals = Vec::new(); + collect_eq_candidates(&predicate, &mut field_candidates, &mut residuals); if field_candidates.iter().any(|c| c.is_none()) { return PartitionFilter::Predicate(predicate); @@ -67,7 +71,7 @@ impl PartitionFilter { loop { let mut builder = BinaryRowBuilder::new(num_fields as i32); for i in 0..num_fields { - let vals = field_candidates[i].as_ref().unwrap(); + let vals = &field_candidates[i].as_ref().unwrap().values; match vals[combo[i]] { Some(datum) => { builder.write_datum(i, datum, partition_fields[i].data_type()); @@ -75,13 +79,44 @@ impl PartitionFilter { None => builder.set_null_at(i), } } - partitions.insert(builder.build_serialized()); + let row = builder.build(); + // Prove the selected constraint with one comparison, not an IN rescan. + // Fall back if encoding changes a literal (e.g. timestamp precision) + // or it is not equal to itself (NaN), preserving full evaluation. + for (i, candidate) in field_candidates.iter().enumerate() { + let FieldCandidates { + predicate: Predicate::Leaf { data_type, .. }, + values, + } = candidate.as_ref().unwrap() + else { + return PartitionFilter::Predicate(predicate); + }; + match extract_datum(&row, i, data_type) { + Ok(value) if value.as_ref() == values[combo[i]] => {} + _ => return PartitionFilter::Predicate(predicate), + } + } + let matches = residuals.iter().try_fold(true, |matched, residual| { + if matched { + eval_row(residual, &row) + } else { + Ok(false) + } + }); + match matches { + Ok(true) => { + partitions.insert(row.to_serialized_bytes()); + } + Ok(false) => {} + // Keep construction infallible; entry matching reports evaluation errors. + Err(_) => return PartitionFilter::Predicate(predicate), + } let mut carry = true; for i in (0..num_fields).rev() { if carry { combo[i] += 1; - if combo[i] < field_candidates[i].as_ref().unwrap().len() { + if combo[i] < field_candidates[i].as_ref().unwrap().values.len() { carry = false; } else { combo[i] = 0; @@ -93,6 +128,7 @@ impl PartitionFilter { } } + // These bounds may be wider than the retained set, but remain safe for pruning. let bounds = match build_bounds_from_candidates(&field_candidates, partition_fields) { Some(b) => b, None => return PartitionFilter::Predicate(predicate), @@ -125,21 +161,31 @@ impl PartitionFilter { pub(super) fn matches_manifest( &self, - stats: &FileStatsRows, + meta: &ManifestFileMeta, partition_fields: &[DataField], ) -> bool { + if partition_fields.is_empty() { + return true; + } + let stats = meta.partition_stats(); + let stats = FileStatsRows::for_manifest_partition( + meta.num_added_files() + meta.num_deleted_files(), + BinaryRow::from_serialized_bytes(stats.min_values()).ok(), + BinaryRow::from_serialized_bytes(stats.max_values()).ok(), + stats.null_counts().clone(), + ); match self { PartitionFilter::PartitionSet { bounds, .. } => { for b in bounds { - if !predicate_may_match(&b.min, stats, partition_fields) - || !predicate_may_match(&b.max, stats, partition_fields) + if !predicate_may_match(&b.min, &stats, partition_fields) + || !predicate_may_match(&b.max, &stats, partition_fields) { return false; } } true } - PartitionFilter::Predicate(pred) => predicate_may_match(pred, stats, partition_fields), + PartitionFilter::Predicate(pred) => predicate_may_match(pred, &stats, partition_fields), } } } @@ -177,7 +223,7 @@ fn predicate_may_match( /// Build per-field min/max bounds from candidate values (from `collect_eq_candidates`). fn build_bounds_from_candidates( - field_candidates: &[Option>>], + field_candidates: &[Option>], partition_fields: &[DataField], ) -> Option> { let pb = PredicateBuilder::new(partition_fields); @@ -185,7 +231,7 @@ fn build_bounds_from_candidates( .iter() .enumerate() .map(|(i, candidates)| { - let vals = candidates.as_ref().unwrap(); + let vals = &candidates.as_ref().unwrap().values; build_field_bounds(&pb, partition_fields[i].name(), vals) }) .collect() @@ -275,47 +321,47 @@ fn build_field_bounds( /// Collect `Eq`/`In`/`IsNull` candidate values per partition field. /// -/// Returns `false` as soon as any node of the tree is not fully represented by -/// the collected candidates. Callers must then keep the original predicate: a -/// `PartitionSet` is the sole authority in `matches_entry`, which never looks at -/// the predicate again, and `ReadBuilder::is_exact_filter_pushdown` lets -/// DataFusion drop its residual filter for a partition-only predicate. +/// Preserve every other condition, including earlier constraints on the same +/// field, as a residual to evaluate before inserting a candidate partition. +/// A `PartitionSet` is the sole authority in `matches_entry`, and exact partition +/// filter pushdown lets DataFusion drop its residual filter, so the set must +/// enforce the complete predicate. fn collect_eq_candidates<'a>( predicate: &'a Predicate, - field_candidates: &mut Vec>>>, -) -> bool { + field_candidates: &mut [Option>], + residuals: &mut Vec<&'a Predicate>, +) { match predicate { - Predicate::And(children) => children - .iter() - .all(|child| collect_eq_candidates(child, field_candidates)), + Predicate::And(children) => { + for child in children { + collect_eq_candidates(child, field_candidates, residuals); + } + } Predicate::Leaf { index, op, literals, .. } if *index < field_candidates.len() => { - // A second conjunct on the same field used to overwrite the first, - // keeping only whichever came last. - if field_candidates[*index].is_some() { - return false; - } - match op { - PredicateOperator::Eq if !literals.is_empty() => { - field_candidates[*index] = Some(vec![Some(&literals[0])]); - true - } + let values = match op { + PredicateOperator::Eq if !literals.is_empty() => vec![Some(&literals[0])], PredicateOperator::In if !literals.is_empty() => { - field_candidates[*index] = Some(literals.iter().map(Some).collect()); - true + literals.iter().map(Some).collect() } - PredicateOperator::IsNull => { - field_candidates[*index] = Some(vec![None]); - true + PredicateOperator::IsNull => vec![None], + _ => { + residuals.push(predicate); + return; } - _ => false, + }; + if let Some(previous) = + field_candidates[*index].replace(FieldCandidates { predicate, values }) + { + // A later constraint on the same field does not supersede the earlier one. + residuals.push(previous.predicate); } } - _ => false, + _ => residuals.push(predicate), } } @@ -343,6 +389,64 @@ mod tests { )] } + #[test] + fn test_manifest_pruning_uses_partition_statistics() { + use crate::spec::stats::BinaryTableStats; + + let fields = partition_fields_dt(); + let pb = PredicateBuilder::new(&fields); + let mut min = BinaryRowBuilder::new(1); + min.write_string(0, "2024-01-01"); + let mut max = BinaryRowBuilder::new(1); + max.write_string(0, "2024-01-03"); + let meta = ManifestFileMeta::new( + "manifest".into(), + 1, + 10, + 5, + BinaryTableStats::new( + min.build_serialized(), + max.build_serialized(), + vec![Some(0)], + ), + 0, + ); + for (predicate, expected) in [ + ( + pb.equal("dt", Datum::String("2024-01-02".into())).unwrap(), + true, + ), + ( + pb.equal("dt", Datum::String("2024-01-04".into())).unwrap(), + false, + ), + ( + pb.greater_than("dt", Datum::String("2024-01-01".into())) + .unwrap(), + true, + ), + ( + pb.greater_than("dt", Datum::String("2024-01-03".into())) + .unwrap(), + false, + ), + (pb.is_null("dt").unwrap(), false), + ] { + let filter = PartitionFilter::from_predicate(predicate, &fields); + assert_eq!(filter.matches_manifest(&meta, &fields), expected); + assert!(filter.matches_manifest(&meta, &[])); + let unknown = ManifestFileMeta::new( + "unknown-stats".into(), + 1, + 10, + 5, + BinaryTableStats::new(vec![0xFF], vec![0xFF], vec![]), + 0, + ); + assert!(filter.matches_manifest(&unknown, &fields)); + } + } + #[test] fn test_eq_builds_partition_set() { let fields = partition_fields_dt(); @@ -447,6 +551,171 @@ mod tests { assert!(!filter.matches_entry(&builder.build_serialized()).unwrap()); } + #[test] + fn test_partition_candidates_preserve_all_conjuncts() { + let fields = partition_fields_dt(); + let pb = PredicateBuilder::new(&fields); + let a = pb.equal("dt", Datum::String("a".into())).unwrap(); + let b = pb.equal("dt", Datum::String("b".into())).unwrap(); + let candidates = pb + .is_in( + "dt", + vec![Datum::String("a".into()), Datum::String("b".into())], + ) + .unwrap(); + for (predicate, expected) in [ + ( + Predicate::and(vec![ + candidates.clone(), + pb.greater_than("dt", Datum::String("a".into())).unwrap(), + ]), + vec![Some("b")], + ), + ( + Predicate::and(vec![ + candidates.clone(), + pb.is_in( + "dt", + vec![Datum::String("b".into()), Datum::String("c".into())], + ) + .unwrap(), + ]), + vec![Some("b")], + ), + (Predicate::and(vec![a.clone(), b.clone()]), vec![]), + ( + Predicate::and(vec![ + candidates.clone(), + Predicate::or(vec![b, pb.equal("dt", Datum::String("c".into())).unwrap()]), + ]), + vec![Some("b")], + ), + ( + Predicate::and(vec![candidates, Predicate::negate(a)]), + vec![Some("b")], + ), + ( + Predicate::and(vec![ + pb.is_null("dt").unwrap(), + pb.is_not_null("dt").unwrap(), + ]), + vec![], + ), + ] { + let filter = PartitionFilter::from_predicate(predicate.clone(), &fields); + assert!(matches!(filter, PartitionFilter::PartitionSet { .. })); + for value in [None, Some("a"), Some("b"), Some("c")] { + let mut row = BinaryRowBuilder::new(1); + match value { + Some(value) => { + row.write_datum(0, &Datum::String(value.into()), fields[0].data_type()) + } + None => row.set_null_at(0), + } + assert_eq!( + filter.matches_entry(&row.build_serialized()).unwrap(), + expected.contains(&value), + "{predicate:?}, {value:?}" + ); + } + } + } + + #[test] + fn test_large_in_candidates_with_residual_filters() { + let fields = vec![DataField::new( + 0, + "id".into(), + DataType::Int(IntType::new()), + )]; + let pb = PredicateBuilder::new(&fields); + for size in [2_000, 4_000] { + for with_range in [false, true] { + let list = pb.is_in("id", (0..size).map(Datum::Int).collect()).unwrap(); + let predicate = if with_range { + Predicate::and(vec![ + list, + pb.greater_or_equal("id", Datum::Int(size / 2)).unwrap(), + ]) + } else { + list + }; + { + let mut candidates = vec![None]; + let mut residuals = Vec::new(); + collect_eq_candidates(&predicate, &mut candidates, &mut residuals); + assert_eq!(candidates[0].as_ref().unwrap().values.len(), size as usize); + assert_eq!(residuals.len(), usize::from(with_range)); + assert!(residuals.iter().all(|p| matches!( + p, + Predicate::Leaf { + op: PredicateOperator::GtEq, + .. + } + ))); + } + let started = std::time::Instant::now(); + let filter = PartitionFilter::from_predicate(predicate, &fields); + println!( + "IN size={size}, range={with_range}: {:?}", + started.elapsed() + ); + let PartitionFilter::PartitionSet { partitions, .. } = &filter else { + panic!("large IN must retain constant-time partition lookup"); + }; + assert_eq!( + partitions.len(), + if with_range { size / 2 } else { size } as usize + ); + for value in [-1, 0, size / 2 - 1, size / 2, size - 1, size] { + let mut row = BinaryRowBuilder::new(1); + row.write_int(0, value); + assert_eq!( + filter.matches_entry(&row.build_serialized()).unwrap(), + (if with_range { size / 2 } else { 0 }..size).contains(&value) + ); + } + } + } + } + + #[test] + fn test_non_roundtripping_candidates_keep_full_predicate() { + let timestamp = |millis, nanos| Datum::Timestamp { millis, nanos }; + for (data_type, literals, probes) in [ + ( + DataType::Double(crate::spec::DoubleType::new()), + vec![Datum::Double(f64::NAN), Datum::Double(1.0)], + vec![ + (Datum::Double(f64::NAN), false), + (Datum::Double(1.0), true), + (Datum::Double(2.0), false), + ], + ), + ( + DataType::Timestamp(crate::spec::TimestampType::new(3).unwrap()), + vec![timestamp(5, 1), timestamp(5, 0)], + vec![(timestamp(5, 0), true), (timestamp(6, 0), false)], + ), + ] { + let fields = vec![DataField::new(0, "key".into(), data_type.clone())]; + let predicate = PredicateBuilder::new(&fields) + .is_in("key", literals) + .unwrap(); + let filter = PartitionFilter::from_predicate(predicate, &fields); + assert!(matches!(filter, PartitionFilter::Predicate(_))); + for (value, expected) in probes { + let mut row = BinaryRowBuilder::new(1); + row.write_datum(0, &value, &data_type); + assert_eq!( + filter.matches_entry(&row.build_serialized()).unwrap(), + expected, + "{value:?}" + ); + } + } + } + #[test] fn test_range_predicate_falls_back() { let fields = partition_fields_dt(); @@ -481,9 +750,9 @@ mod tests { builder.build_serialized() } - /// Coverage is complete, but `>=` is not expressible as a set of values. + /// A residual range must reject a candidate that contradicts it. #[test] - fn test_unexpressible_conjunct_on_covered_field_falls_back() { + fn test_unexpressible_conjunct_on_covered_field_is_preserved() { let fields = partition_fields_dt(); let pb = PredicateBuilder::new(&fields); let pred = Predicate::and(vec![ @@ -492,7 +761,6 @@ mod tests { .unwrap(), ]); let filter = PartitionFilter::from_predicate(pred, &fields); - assert!(matches!(filter, PartitionFilter::Predicate(_))); assert!(!filter .matches_entry(&serialized_dt(&fields, "2024-01-01")) .unwrap()); @@ -501,7 +769,7 @@ mod tests { /// Two expressible conjuncts on one field: the second assignment used to /// overwrite the first, keeping whichever came last — here the wider `In`. #[test] - fn test_second_conjunct_on_same_field_falls_back() { + fn test_second_conjunct_on_same_field_is_preserved() { let fields = partition_fields_dt(); let pb = PredicateBuilder::new(&fields); let pred = Predicate::and(vec![ @@ -516,7 +784,6 @@ mod tests { .unwrap(), ]); let filter = PartitionFilter::from_predicate(pred, &fields); - assert!(matches!(filter, PartitionFilter::Predicate(_))); assert!(!filter .matches_entry(&serialized_dt(&fields, "2024-01-01")) .unwrap()); @@ -527,7 +794,7 @@ mod tests { /// An `Or` over the partition field narrows the `In` beside it. #[test] - fn test_or_conjunct_beside_covering_in_falls_back() { + fn test_or_conjunct_beside_covering_in_is_preserved() { let fields = partition_fields_dt(); let pb = PredicateBuilder::new(&fields); let pred = Predicate::and(vec![ @@ -546,7 +813,6 @@ mod tests { .unwrap(), ]); let filter = PartitionFilter::from_predicate(pred, &fields); - assert!(matches!(filter, PartitionFilter::Predicate(_))); assert!(!filter .matches_entry(&serialized_dt(&fields, "2024-01-03")) .unwrap()); diff --git a/crates/paimon/src/table/partition_row_count.rs b/crates/paimon/src/table/partition_row_count.rs new file mode 100644 index 000000000..2888f9345 --- /dev/null +++ b/crates/paimon/src/table/partition_row_count.rs @@ -0,0 +1,1573 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! Real per-partition row counts, computed from manifests alone in bounded memory. +//! +//! [`Table::partition_stats`] and the `$files` system table both materialize every +//! manifest entry together with its per-column statistics, which does not fit in +//! memory for very large tables, and summing `row_count` over files over-counts +//! data-evolution tables, where several column-group files cover the same rows. +//! +//! This module instead: +//! - decodes manifests through [`SlimManifestEntry`], which borrows the handful of +//! fields it needs and skips statistics in place, so nothing per-file is retained; +//! - for data-evolution tables, counts files carrying a `first_row_id` by the +//! *union* of their row-id ranges, collapsing overlapping column-group/blob files; +//! - nets ADD/DELETE entries with a delete set built only from the manifests that +//! contain deletes, the same semantics as the scan's manifest merge. +//! +//! Peak memory is bounded by partitions, live DELETE entries, up to one million +//! retained ADD identities, deletion-vector mappings, disjoint row-id ranges, and +//! in-flight manifest buffers rather than all live file metadata and column +//! statistics. Highly fragmented row-id space or large embedded indexes may still +//! increase retained state. + +use std::collections::{BTreeMap, HashMap}; +use std::sync::atomic::{AtomicI64, Ordering}; +use std::sync::Arc; + +use futures::{StreamExt, TryStreamExt}; + +use crate::io::FileIO; +use crate::spec::avro::{ + visit_slim_index_manifest_entries, visit_slim_manifest_entries, SharedSchemaCache, + SlimManifestEntry, +}; +use crate::spec::{BinaryRow, CoreOptions, FileKind, ManifestFileMeta, ManifestList, Predicate}; +use crate::table::partition_filter::PartitionFilter; +use crate::table::read_builder::split_scan_predicates; +use crate::table::Table; + +/// Independent I/O and blocking decode limits for data manifests. +const MANIFEST_READ_CONCURRENCY: usize = 32; + +/// ADD identities retained from delete-bearing manifests to avoid fetching them twice. +/// This bounds entry count, not embedded-index payload bytes. +const RETAINED_ADD_BUDGET: i64 = 1_000_000; + +fn manifest_decode_concurrency() -> usize { + std::thread::available_parallelism() + .map_or(2, |parallelism| parallelism.get()) + .clamp(2, MANIFEST_READ_CONCURRENCY) +} + +const DELETION_VECTORS_INDEX_TYPE: &str = "DELETION_VECTORS"; + +/// Real row count of one partition. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct PartitionRowCount { + /// The partition's typed values, one field per partition key. + pub partition_row: BinaryRow, + /// Rows in the partition: overlapping data-evolution files are counted once + /// and deletion-vector rows are subtracted. `None` when it cannot be known + /// exactly (a file without a row count, or a deletion vector without a + /// cardinality) — never a guess. + pub record_count: Option, +} + +/// Disjoint set of inclusive row-id ranges, coalescing overlapping and adjacent ones. +#[derive(Debug, Default)] +struct RowRangeSet { + ranges: BTreeMap, +} + +impl RowRangeSet { + fn insert(&mut self, mut start: i64, mut end: i64) { + if start > end { + return; + } + if let Some((&prev_start, &prev_end)) = self.ranges.range(..=start).next_back() { + if prev_end >= end { + return; + } + if prev_end.saturating_add(1) >= start { + start = prev_start; + self.ranges.remove(&prev_start); + } + } + while let Some((&next_start, &next_end)) = + self.ranges.range(start..=end.saturating_add(1)).next() + { + end = end.max(next_end); + self.ranges.remove(&next_start); + } + self.ranges.insert(start, end); + } + + fn merge(&mut self, other: RowRangeSet) { + if self.ranges.is_empty() { + self.ranges = other.ranges; + return; + } + for (start, end) in other.ranges { + self.insert(start, end); + } + } + + fn total(&self) -> i128 { + self.ranges + .iter() + .map(|(start, end)| i128::from(*end) - i128::from(*start) + 1) + .sum() + } +} + +#[derive(Debug, Default)] +struct PartitionAccum { + /// Rows of files without a `first_row_id`. + plain_rows: i128, + /// Row-id ranges of files with a `first_row_id`. + row_ranges: RowRangeSet, + deleted_rows: i128, + row_count_unknown: bool, +} + +impl PartitionAccum { + fn add_file( + &mut self, + row_count: i64, + first_row_id: Option, + data_evolution_enabled: bool, + ) { + if row_count < 0 { + self.row_count_unknown = true; + return; + } + if data_evolution_enabled { + if let Some(first) = first_row_id { + if row_count > 0 { + let Some(last) = first.checked_add(row_count - 1) else { + self.row_count_unknown = true; + return; + }; + self.row_ranges.insert(first, last); + } + return; + } + } + self.plain_rows += i128::from(row_count); + } + + fn merge(&mut self, other: PartitionAccum) { + self.plain_rows += other.plain_rows; + self.row_ranges.merge(other.row_ranges); + self.deleted_rows += other.deleted_rows; + self.row_count_unknown |= other.row_count_unknown; + } + + fn record_count(&self) -> Option { + if self.row_count_unknown { + return None; + } + let rows = self.plain_rows + self.row_ranges.total() - self.deleted_rows; + (rows >= 0).then(|| i64::try_from(rows).ok()).flatten() + } +} + +type PartitionAccums = HashMap, PartitionAccum>; + +/// Only DV mappings are retained, not metadata for every live data file. +/// Matches the scan's (partition, bucket, file name) lookup and last-write wins. +#[derive(Debug, Default)] +struct DeletionVectors { + by_partition: HashMap, DeletionVectorFiles>, +} + +type DeletionVectorFiles = HashMap, Option>>; + +impl DeletionVectors { + fn cardinality(&self, partition: &[u8], bucket: i32, file_name: &str) -> Option { + self.by_partition + .get(partition) + .and_then(|buckets| buckets.get(&bucket)) + .and_then(|files| files.get(file_name)) + .copied() + .unwrap_or(Some(0)) + } + + fn has_unknown(&self) -> bool { + self.by_partition + .values() + .flat_map(HashMap::values) + .flat_map(HashMap::values) + .any(Option::is_none) + } +} + +fn accumulate( + accums: &mut PartitionAccums, + partition: &[u8], + row_count: i64, + first_row_id: Option, + data_evolution_enabled: bool, + deleted_rows: Option, +) { + let accum = match accums.get_mut(partition) { + Some(accum) => accum, + None => accums.entry(partition.to_vec()).or_default(), + }; + accum.add_file(row_count, first_row_id, data_evolution_enabled); + match deleted_rows { + Some(rows) => accum.deleted_rows += i128::from(rows), + None => accum.row_count_unknown = true, + } +} + +/// Identifiers of deleted files, matching the full Paimon `Identifier` semantics. +/// +/// Nested so ADD lookups borrow partition/file-name bytes from the decode buffer +/// and each partition is stored only once. +#[derive(Debug, Default)] +struct DeleteSet { + by_partition: HashMap, DeletedFiles>, +} + +type DeletedFiles = HashMap, Vec>; + +#[derive(Debug, PartialEq, Eq)] +struct DeletedFile { + bucket: i32, + level: i32, + extra_files: Vec>, + // Exact identity matching requires the payload. A spill-backed delete set can + // replace this copy if embedded-index memory becomes a measured bottleneck. + embedded_index: Option>, + external_path: Option>, +} + +impl DeletedFile { + fn from_entry(entry: &SlimManifestEntry<'_>) -> Self { + Self { + bucket: entry.bucket, + level: entry.level, + extra_files: entry + .extra_files + .iter() + .map(|value| Box::from(*value)) + .collect(), + embedded_index: entry.embedded_index.map(Box::from), + external_path: entry.external_path.map(Box::from), + } + } + + fn matches(&self, entry: &SlimManifestEntry<'_>) -> bool { + self.bucket == entry.bucket + && self.level == entry.level + && self.embedded_index.as_deref() == entry.embedded_index + && self.external_path.as_deref() == entry.external_path + && self.extra_files.len() == entry.extra_files.len() + && self + .extra_files + .iter() + .zip(&entry.extra_files) + .all(|(left, right)| &**left == *right) + } +} + +struct RetainedAdd { + partition: Arc<[u8]>, + file_name: Box, + identity: DeletedFile, + row_count: i64, + first_row_id: Option, +} + +#[derive(Default)] +struct PartitionInterner { + last: Option>, +} + +impl PartitionInterner { + fn intern(&mut self, partition: &[u8]) -> Arc<[u8]> { + match &self.last { + Some(last) if &**last == partition => Arc::clone(last), + _ => { + let partition = Arc::from(partition); + self.last = Some(Arc::clone(&partition)); + partition + } + } + } +} + +impl DeleteSet { + fn is_empty(&self) -> bool { + self.by_partition.is_empty() + } + + fn insert(&mut self, entry: &SlimManifestEntry<'_>) { + let files = match self.by_partition.get_mut(entry.partition) { + Some(files) => files, + None => self + .by_partition + .entry(Box::from(entry.partition)) + .or_default(), + }; + let deleted = DeletedFile::from_entry(entry); + let slots = files.entry(Box::from(entry.file_name)).or_default(); + if !slots.contains(&deleted) { + slots.push(deleted); + } + } + + fn contains(&self, entry: &SlimManifestEntry<'_>) -> bool { + self.by_partition + .get(entry.partition) + .and_then(|files| files.get(entry.file_name)) + .is_some_and(|slots| slots.iter().any(|deleted| deleted.matches(entry))) + } + + fn contains_retained(&self, add: &RetainedAdd) -> bool { + self.by_partition + .get(&*add.partition) + .and_then(|files| files.get(&*add.file_name)) + .is_some_and(|slots| slots.contains(&add.identity)) + } + + fn merge(&mut self, other: DeleteSet) { + for (partition, files) in other.by_partition { + let target = self.by_partition.entry(partition).or_default(); + for (file_name, deleted) in files { + let slots = target.entry(file_name).or_default(); + for entry in deleted { + if !slots.contains(&entry) { + slots.push(entry); + } + } + } + } + } +} + +/// Entry-level partition filter, remembering the last verdict because manifests +/// list long runs of files from the same partition. +struct PartitionMatcher<'a> { + filter: Option<&'a PartitionFilter>, + last_partition: Vec, + last_verdict: Option, +} + +impl<'a> PartitionMatcher<'a> { + fn new(filter: Option<&'a PartitionFilter>) -> Self { + Self { + filter, + last_partition: Vec::new(), + last_verdict: None, + } + } + + fn matches(&mut self, partition: &[u8]) -> crate::Result { + let Some(filter) = self.filter else { + return Ok(true); + }; + if let Some(verdict) = self.last_verdict { + if self.last_partition == partition { + return Ok(verdict); + } + } + let verdict = filter.matches_entry(partition)?; + self.last_partition.clear(); + self.last_partition.extend_from_slice(partition); + self.last_verdict = Some(verdict); + Ok(verdict) + } +} + +struct DeleteManifestSummary { + deletes: DeleteSet, + /// `None` when the retention budget ran out and the manifest must be re-read. + adds: Option>, +} + +/// Reserve a whole manifest, so concurrent decoders cannot each retain a prefix +/// then all abandon it when the shared budget runs out. +fn reserve_retained_adds(budget: &AtomicI64, count: i64) -> Option { + let capacity = usize::try_from(count).ok()?; + budget + .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |available| { + available + .checked_sub(count) + .filter(|remaining| *remaining >= 0) + }) + .ok() + .map(|_| capacity) +} + +fn summarize_delete_manifest( + bytes: &[u8], + cache: &SharedSchemaCache, + budget: &AtomicI64, + num_added_files: i64, + filter: Option<&PartitionFilter>, +) -> crate::Result { + // ponytail: the whole-manifest count can over-reserve under a partition + // filter and cause extra rereads; selected-entry reservations need profiling. + let reservation = reserve_retained_adds(budget, num_added_files); + let mut matcher = PartitionMatcher::new(filter); + let mut interner = PartitionInterner::default(); + let mut deletes = DeleteSet::default(); + let mut adds = reservation.map(|_| Vec::new()); + visit_slim_manifest_entries(bytes, cache, &mut |entry| { + if !matcher.matches(entry.partition)? { + return Ok(()); + } + match entry.kind { + FileKind::Delete => deletes.insert(&entry), + FileKind::Add => { + if let Some(retained) = adds.as_mut() { + if retained.len() < reservation.unwrap_or(0) { + retained.push(RetainedAdd { + partition: interner.intern(entry.partition), + file_name: Box::from(entry.file_name), + identity: DeletedFile::from_entry(&entry), + row_count: entry.row_count, + first_row_id: entry.first_row_id, + }); + } else { + // An understated manifest count must not exceed the + // reservation or make us lose ADDs: re-read instead. + adds = None; + } + } + } + } + Ok(()) + })?; + if let Some(reserved) = reservation { + let unused = reserved - adds.as_ref().map_or(0, Vec::len); + if unused > 0 { + budget.fetch_add(unused as i64, Ordering::Relaxed); + } + } + Ok(DeleteManifestSummary { deletes, adds }) +} + +/// Second-pass result of one manifest: its live ADD files, already aggregated. +fn aggregate_manifest( + bytes: &[u8], + cache: &SharedSchemaCache, + deletes: Option<&DeleteSet>, + filter: Option<&PartitionFilter>, + data_evolution_enabled: bool, + deletion_vectors: &DeletionVectors, +) -> crate::Result { + let mut matcher = PartitionMatcher::new(filter); + let mut accums = PartitionAccums::new(); + visit_slim_manifest_entries(bytes, cache, &mut |entry| { + if entry.kind == FileKind::Add + && matcher.matches(entry.partition)? + && !deletes.is_some_and(|deletes| deletes.contains(&entry)) + { + accumulate( + &mut accums, + entry.partition, + entry.row_count, + entry.first_row_id, + data_evolution_enabled, + deletion_vectors.cardinality(entry.partition, entry.bucket, entry.file_name), + ); + } + Ok(()) + })?; + Ok(accums) +} + +/// Fetch and decode at most 32 manifests in one pipeline. A semaphore limits +/// CPU-heavy Avro work on Tokio's blocking pool without adding another buffer. +fn read_manifests( + file_io: &FileIO, + manifest_dir: &str, + manifests: Vec, + decode: F, +) -> impl futures::Stream> +where + T: Send + 'static, + F: Fn(&[u8], &ManifestFileMeta) -> crate::Result + Send + Sync + 'static, +{ + let file_io = file_io.clone(); + let manifest_dir = manifest_dir.to_string(); + let decode = Arc::new(decode); + let decode_permits = Arc::new(tokio::sync::Semaphore::new(manifest_decode_concurrency())); + futures::stream::iter(manifests) + .map(move |meta| { + let file_io = file_io.clone(); + let path = format!("{}/{}", manifest_dir, meta.file_name()); + let decode = Arc::clone(&decode); + let decode_permits = Arc::clone(&decode_permits); + async move { + let bytes = file_io.new_input(&path)?.read().await?; + let permit = decode_permits.acquire_owned().await.map_err(|error| { + crate::Error::UnexpectedError { + message: format!("manifest decode semaphore closed: {error}"), + source: Some(Box::new(error)), + } + })?; + tokio::task::spawn_blocking(move || { + let _permit = permit; + decode(&bytes, &meta).map(|decoded| (meta, decoded)) + }) + .await + .map_err(|error| crate::Error::UnexpectedError { + message: format!("manifest decode task failed: {error}"), + source: Some(Box::new(error)), + })? + } + }) + .buffer_unordered(MANIFEST_READ_CONCURRENCY) +} + +async fn aggregate_manifests( + file_io: &FileIO, + manifest_dir: &str, + manifests: Vec, + filter: Option>, + data_evolution_enabled: bool, + retained_add_budget: i64, + deletion_vectors: Arc, +) -> crate::Result, PartitionAccum>> { + let cache = SharedSchemaCache::new(); + let (with_deletes, mut second_pass): (Vec<_>, Vec<_>) = manifests + .into_iter() + .partition(|meta| meta.num_deleted_files() > 0); + + // Keep a bounded number of ADDs while collecting the global delete set. + // Manifests without a complete reservation are fetched again. + let mut deletes = DeleteSet::default(); + let mut retained_adds = Vec::new(); + { + let cache = cache.clone(); + let filter = filter.clone(); + let budget = AtomicI64::new(retained_add_budget); + let mut summaries = std::pin::pin!(read_manifests( + file_io, + manifest_dir, + with_deletes, + move |bytes, meta| { + summarize_delete_manifest( + bytes, + &cache, + &budget, + meta.num_added_files(), + filter.as_deref(), + ) + }, + )); + while let Some((meta, summary)) = summaries.try_next().await? { + deletes.merge(summary.deletes); + match summary.adds { + Some(adds) => retained_adds.extend(adds), + None => second_pass.push(meta), + } + } + } + + let mut totals: BTreeMap, PartitionAccum> = BTreeMap::new(); + let mut retained_accums = PartitionAccums::new(); + for add in retained_adds { + if !deletes.contains_retained(&add) { + accumulate( + &mut retained_accums, + &add.partition, + add.row_count, + add.first_row_id, + data_evolution_enabled, + deletion_vectors.cardinality(&add.partition, add.identity.bucket, &add.file_name), + ); + } + } + for (partition, accum) in retained_accums { + totals.entry(partition).or_default().merge(accum); + } + + // Aggregate manifests without DELETEs and those that exceeded the ADD budget. + let deletes = (!deletes.is_empty()).then(|| Arc::new(deletes)); + let mut aggregated = std::pin::pin!(read_manifests( + file_io, + manifest_dir, + second_pass, + move |bytes, _| { + aggregate_manifest( + bytes, + &cache, + deletes.as_deref(), + filter.as_deref(), + data_evolution_enabled, + &deletion_vectors, + ) + }, + )); + while let Some((_, accums)) = aggregated.try_next().await? { + for (partition, accum) in accums { + totals.entry(partition).or_default().merge(accum); + } + } + + Ok(totals) +} + +/// Load matching DV mappings before visiting data files, so only live files +/// contribute deletions and exact-only callers can stop on unknown cardinalities. +async fn read_deletion_vectors( + file_io: &FileIO, + index_manifest_path: &str, + filter: Option>, +) -> crate::Result { + let bytes = file_io.new_input(index_manifest_path)?.read().await?; + tokio::task::spawn_blocking(move || { + let mut matcher = PartitionMatcher::new(filter.as_deref()); + let mut deleted = DeletionVectors::default(); + visit_slim_index_manifest_entries(&bytes, &SharedSchemaCache::new(), &mut |entry| { + if entry.kind != FileKind::Add + || entry.index_type != DELETION_VECTORS_INDEX_TYPE + || !matcher.matches(entry.partition)? + || entry.deletion_vector_cardinalities.is_empty() + { + return Ok(()); + } + let files = deleted + .by_partition + .entry(Box::from(entry.partition)) + .or_default() + .entry(entry.bucket) + .or_default(); + files.extend( + entry + .deletion_vector_cardinalities + .into_iter() + .map(|(name, cardinality)| (Box::from(name), cardinality)), + ); + Ok(()) + })?; + Ok(deleted) + }) + .await + .map_err(|e| crate::Error::UnexpectedError { + message: format!("index manifest decode task failed: {e}"), + source: Some(Box::new(e)), + })? +} + +impl Table { + /// Real row count of every partition in the latest (or time-travelled) snapshot. + /// + /// Reads manifests only — never data files — without retaining every live + /// file or its column statistics. Memory is bounded by partitions, live + /// DELETE entries, up to one million retained ADD identities, deletion-vector + /// mappings, disjoint data-evolution row-id ranges, and in-flight manifest + /// buffers. + /// Data-evolution files sharing a row-id range contribute once. Results are + /// ordered by serialized partition bytes. + /// + /// Primary-key and format tables return [`crate::Error::Unsupported`]. + /// Returns an empty Vec when a supported table has no snapshots yet. + pub async fn partition_row_counts(&self) -> crate::Result> { + self.partition_row_counts_with_filter(None).await + } + + /// [`Table::partition_row_counts`] restricted to the partitions matching `filter`. + /// + /// `filter` may only reference partition columns: anything else cannot be + /// decided from manifests, and is rejected rather than ignored. Manifests + /// whose partition range cannot match are never fetched. + pub async fn partition_row_counts_with_filter( + &self, + filter: Option, + ) -> crate::Result> { + Ok(self + .read_partition_row_counts(filter, false) + .await? + .expect("partial partition counts never stop early")) + } + + /// Exact-only variant of [`Table::partition_row_counts_with_filter`]. + /// Returns `None` for primary-key or format tables, or when metadata cannot + /// establish all counts. Unknown DV cardinalities are checked before fetching + /// data manifests, allowing callers to fall back without first doing a full + /// metadata aggregation. This may conservatively return `None` for an unknown + /// DV on a no-longer-live file. + pub async fn exact_partition_row_counts_with_filter( + &self, + filter: Option, + ) -> crate::Result>> { + self.read_partition_row_counts(filter, true).await + } + + async fn read_partition_row_counts( + &self, + filter: Option, + require_exact: bool, + ) -> crate::Result>> { + let schema = self.schema(); + let core = CoreOptions::new(schema.options()); + // Manifests carry partition values. + core.ensure_read_authorized()?; + // Primary-key counts need merging; format tables do not use Paimon snapshots. + if core.is_format_table() || !schema.primary_keys().is_empty() { + return if require_exact { + Ok(None) + } else { + Err(crate::Error::Unsupported { + message: + "partition row counts are not supported for primary-key or format tables" + .to_string(), + }) + }; + } + + let file_io = self.file_io(); + let Some(snapshot) = super::time_travel::resolve_snapshot(self).await? else { + return Ok(Some(Vec::new())); + }; + + let manifest_sm = self.snapshot_manager(); + let partition_fields = schema.partition_fields(); + let partition_filter = match filter { + None => None, + Some(filter) => { + let (partition_predicate, data_predicates) = split_scan_predicates(self, filter); + if !data_predicates.is_empty() { + return Err(crate::Error::Unsupported { + message: "partition row counts can only be filtered by partition columns" + .to_string(), + }); + } + partition_predicate + .map(|predicate| PartitionFilter::from_predicate(predicate, &partition_fields)) + } + }; + let partition_filter = partition_filter.map(Arc::new); + let deletion_vectors = match snapshot.index_manifest() { + Some(index_manifest) if core.deletion_vectors_enabled() => { + read_deletion_vectors( + file_io, + &manifest_sm.manifest_path(index_manifest), + partition_filter.clone(), + ) + .await? + } + _ => DeletionVectors::default(), + }; + // ponytail: unknown stale DVs can also fall back; check liveness first + // only if these conservative fallbacks become a measured bottleneck. + if require_exact && deletion_vectors.has_unknown() { + return Ok(None); + } + + let base_path = manifest_sm.manifest_path(snapshot.base_manifest_list()); + let delta_path = manifest_sm.manifest_path(snapshot.delta_manifest_list()); + let (mut manifests, delta) = futures::try_join!( + ManifestList::read(file_io, &base_path), + ManifestList::read(file_io, &delta_path), + )?; + manifests.extend(delta); + if let Some(filter) = &partition_filter { + manifests.retain(|meta| filter.matches_manifest(meta, &partition_fields)); + } + + let totals = aggregate_manifests( + file_io, + &manifest_sm.manifest_dir(), + manifests, + partition_filter, + core.data_evolution_enabled(), + RETAINED_ADD_BUDGET, + Arc::new(deletion_vectors), + ) + .await?; + + let mut out = Vec::with_capacity(totals.len()); + for (partition_bytes, accum) in totals { + let record_count = accum.record_count(); + if require_exact && record_count.is_none() { + return Ok(None); + } + out.push(PartitionRowCount { + partition_row: BinaryRow::from_serialized_bytes(&partition_bytes)?, + record_count, + }); + } + Ok(Some(out)) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::io::FileIOBuilder; + use crate::spec::stats::BinaryTableStats; + use crate::spec::{ + DataFileMeta, DeletionVectorMeta, IndexFileMeta, IndexManifest, IndexManifestEntry, + Manifest, ManifestEntry, + }; + + const MANIFEST_DIR: &str = "memory:/partition_row_count/manifest"; + + fn partition(value: u8) -> Vec { + vec![0, 0, 0, 1, value, 0, 0, 0] + } + + fn file(name: &str, level: i32, row_count: i64, first_row_id: Option) -> DataFileMeta { + let stats = BinaryTableStats::empty(); + DataFileMeta { + file_name: name.to_string(), + file_size: 100, + row_count, + min_key: vec![], + max_key: vec![], + key_stats: stats.clone(), + value_stats: stats, + min_sequence_number: 0, + max_sequence_number: 0, + schema_id: 0, + level, + extra_files: vec![], + creation_time: None, + delete_row_count: None, + embedded_index: None, + file_source: None, + value_stats_cols: None, + external_path: None, + first_row_id, + write_cols: None, + column_max_sequence_numbers: None, + } + } + + fn add(partition: Vec, file: DataFileMeta) -> ManifestEntry { + ManifestEntry::new(FileKind::Add, partition, 0, 1, file, 2) + } + + fn delete(partition: Vec, file: DataFileMeta) -> ManifestEntry { + ManifestEntry::new(FileKind::Delete, partition, 0, 1, file, 2) + } + + async fn write_manifest( + file_io: &FileIO, + name: &str, + entries: &[ManifestEntry], + ) -> ManifestFileMeta { + Manifest::write(file_io, &format!("{MANIFEST_DIR}/{name}"), entries) + .await + .unwrap(); + let deleted = entries + .iter() + .filter(|e| *e.kind() == FileKind::Delete) + .count() as i64; + ManifestFileMeta::new( + name.to_string(), + 1, + entries.len() as i64 - deleted, + deleted, + BinaryTableStats::empty(), + 0, + ) + } + + #[test] + fn test_retained_add_budget_reserves_whole_manifests() { + let budget = AtomicI64::new(2); + let start = std::sync::Barrier::new(2); + std::thread::scope(|scope| { + let handles: Vec<_> = (0..2) + .map(|_| { + scope.spawn(|| { + start.wait(); + reserve_retained_adds(&budget, 2) + }) + }) + .collect(); + assert_eq!( + handles + .into_iter() + .filter_map(|handle| handle.join().unwrap()) + .collect::>(), + vec![2] + ); + }); + assert_eq!(budget.load(Ordering::Relaxed), 0); + assert_eq!(reserve_retained_adds(&budget, 0), Some(0)); + assert_eq!(reserve_retained_adds(&budget, -1), None); + assert_eq!(reserve_retained_adds(&budget, i64::MAX), None); + assert_eq!(budget.load(Ordering::Relaxed), 0); + } + + #[tokio::test] + async fn test_retained_add_reservation_refunds_filtered_and_understated_counts() { + let file_io = FileIOBuilder::new("memory").build().unwrap(); + file_io.mkdirs(&format!("{MANIFEST_DIR}/")).await.unwrap(); + let meta = write_manifest( + &file_io, + "reservation", + &[ + add(partition(1), file("a", 0, 2, None)), + add(partition(2), file("b", 0, 3, None)), + delete(partition(1), file("old", 0, 1, None)), + ], + ) + .await; + let bytes = file_io + .new_input(&format!("{MANIFEST_DIR}/{}", meta.file_name())) + .unwrap() + .read() + .await + .unwrap(); + let cache = SharedSchemaCache::new(); + let filter = PartitionFilter::PartitionSet { + partitions: std::collections::HashSet::from([partition(1)]), + bounds: vec![], + }; + let budget = AtomicI64::new(2); + let summary = summarize_delete_manifest(&bytes, &cache, &budget, 2, Some(&filter)).unwrap(); + assert!(!summary.deletes.is_empty()); + assert_eq!(summary.adds.unwrap().len(), 1); + assert_eq!(budget.load(Ordering::Relaxed), 1); + + // The next whole-manifest reservation can use the refunded slot. + assert_eq!(reserve_retained_adds(&budget, 1), Some(1)); + assert_eq!(budget.load(Ordering::Relaxed), 0); + for understated in [0, 1] { + let budget = AtomicI64::new(2); + let summary = + summarize_delete_manifest(&bytes, &cache, &budget, understated, None).unwrap(); + assert!(summary.adds.is_none()); + assert!(!summary.deletes.is_empty()); + assert_eq!(budget.load(Ordering::Relaxed), 2); + } + // Document the conservative tradeoff: the selected ADD fits, but an + // oversized whole-manifest estimate cannot reserve even a partial slot. + let budget = AtomicI64::new(1); + let summary = summarize_delete_manifest(&bytes, &cache, &budget, 2, Some(&filter)).unwrap(); + assert!(summary.adds.is_none()); + assert!(!summary.deletes.is_empty()); + assert_eq!(budget.load(Ordering::Relaxed), 1); + } + + #[tokio::test] + async fn test_partition_row_counts_table_support_and_authorization() { + use crate::catalog::Identifier; + use crate::spec::{DataType, IntType, Schema, TableSchema}; + + for kind in ["append", "primary-key", "format-table"] { + for query_auth in [false, true] { + let mut schema = Schema::builder() + .column("id", DataType::Int(IntType::new())) + .option("query-auth.enabled", query_auth.to_string()); + if kind == "primary-key" { + schema = schema.primary_key(["id"]); + } else if kind == "format-table" { + schema = schema.option("type", "format-table"); + } + // Unsupported tables must not be mistaken for empty tables, + // even when no snapshot exists. Authorization must still fail closed. + let table = Table::new( + FileIOBuilder::new("memory").build().unwrap(), + Identifier::new("default", kind), + format!("memory:/partition-count-support/{kind}/{query_auth}"), + TableSchema::new(0, &schema.build().unwrap()), + None, + ); + let partial = [ + table.partition_row_counts().await, + table.partition_row_counts_with_filter(None).await, + ]; + let exact = table.exact_partition_row_counts_with_filter(None).await; + if query_auth { + for result in partial + .into_iter() + .map(|result| result.map(|_| ())) + .chain([exact.map(|_| ())]) + { + assert!( + matches!(&result, Err(crate::Error::Unsupported { message }) + if message.contains("query-auth.enabled")), + "{kind}: {result:?}" + ); + } + } else if kind == "append" { + assert_eq!(exact.unwrap(), Some(Vec::new())); + for result in partial { + assert!(result.unwrap().is_empty()); + } + } else { + assert_eq!(exact.unwrap(), None, "{kind}"); + for result in partial { + assert!( + matches!(&result, Err(crate::Error::Unsupported { message }) + if message.contains("partition row counts")), + "{kind}: {result:?}" + ); + } + } + } + } + } + + #[tokio::test] + async fn test_deletion_vector_counts_respect_partition_filter() { + let file_io = FileIOBuilder::new("memory").build().unwrap(); + file_io.mkdirs(&format!("{MANIFEST_DIR}/")).await.unwrap(); + let path = format!("{MANIFEST_DIR}/filtered-index-manifest"); + let entry = |value, cardinality| IndexManifestEntry { + version: 1, + kind: FileKind::Add, + partition: partition(value), + bucket: 0, + index_file: IndexFileMeta { + index_type: DELETION_VECTORS_INDEX_TYPE.to_string(), + file_name: format!("index-{value}"), + file_size: 1, + row_count: 1, + deletion_vectors_ranges: Some(indexmap::IndexMap::from([( + format!("data-{value}"), + DeletionVectorMeta { + offset: 0, + length: 1, + cardinality: Some(cardinality), + }, + )])), + external_path: None, + global_index_meta: None, + }, + }; + IndexManifest::write(&file_io, &path, &[entry(1, 3), entry(2, 7)]) + .await + .unwrap(); + let filter = PartitionFilter::PartitionSet { + partitions: std::collections::HashSet::from([partition(1)]), + bounds: Vec::new(), + }; + + let deleted = read_deletion_vectors(&file_io, &path, Some(Arc::new(filter))) + .await + .unwrap(); + + assert_eq!(deleted.by_partition.len(), 1); + assert_eq!(deleted.cardinality(&partition(1), 0, "data-1"), Some(3)); + assert_eq!(deleted.cardinality(&partition(2), 0, "data-2"), Some(0)); + } + + #[tokio::test] + async fn test_deletion_vectors_match_live_partition_bucket_and_file() { + let file_io = FileIOBuilder::new("memory").build().unwrap(); + file_io.mkdirs(&format!("{MANIFEST_DIR}/")).await.unwrap(); + let path = format!("{MANIFEST_DIR}/live-index-manifest"); + let entry = |value, bucket, name: &str, cardinality| IndexManifestEntry { + version: 1, + kind: FileKind::Add, + partition: partition(value), + bucket, + index_file: IndexFileMeta { + index_type: DELETION_VECTORS_INDEX_TYPE.to_string(), + file_name: "index".to_string(), + file_size: 1, + row_count: 1, + deletion_vectors_ranges: Some(indexmap::IndexMap::from([( + name.to_string(), + DeletionVectorMeta { + offset: 0, + length: 1, + cardinality, + }, + )])), + external_path: None, + global_index_meta: None, + }, + }; + IndexManifest::write( + &file_io, + &path, + &[ + entry(1, 0, "a", None), + entry(1, 0, "a", Some(2)), // Last mapping wins, including unknown -> known. + entry(1, 1, "a", Some(4)), + entry(2, 0, "a", Some(5)), + entry(1, 0, "upgraded", Some(1)), + entry(1, 0, "removed", None), // An unknown DV must not poison live counts. + ], + ) + .await + .unwrap(); + let vectors = Arc::new(read_deletion_vectors(&file_io, &path, None).await.unwrap()); + assert_eq!(vectors.cardinality(&partition(1), 0, "a"), Some(2)); + let metas = vec![ + write_manifest( + &file_io, + "live-0", + &[ + add(partition(1), file("a", 0, 10, None)), + ManifestEntry::new( + FileKind::Add, + partition(1), + 1, + 2, + file("a", 0, 12, None), + 2, + ), + add(partition(2), file("a", 0, 8, None)), + add(partition(1), file("upgraded", 0, 3, None)), + add(partition(1), file("removed", 0, 5, None)), + ], + ) + .await, + write_manifest( + &file_io, + "live-1", + &[ + delete(partition(1), file("removed", 0, 5, None)), + delete(partition(1), file("upgraded", 0, 3, None)), + add(partition(1), file("upgraded", 1, 3, None)), + ], + ) + .await, + ]; + for budget in [RETAINED_ADD_BUDGET, 0, 1] { + let totals = aggregate_manifests( + &file_io, + MANIFEST_DIR, + metas.clone(), + None, + false, + budget, + Arc::clone(&vectors), + ) + .await + .unwrap(); + assert_eq!( + totals[&partition(1)].record_count(), + Some(18), + "budget {budget}" + ); + assert_eq!( + totals[&partition(2)].record_count(), + Some(3), + "budget {budget}" + ); + } + } + + async fn counts( + test: &str, + manifests: Vec>, + data_evolution_enabled: bool, + retained_add_budget: i64, + ) -> BTreeMap, i128> { + let file_io = FileIOBuilder::new("memory").build().unwrap(); + file_io.mkdirs(&format!("{MANIFEST_DIR}/")).await.unwrap(); + let mut metas = Vec::new(); + for (i, entries) in manifests.iter().enumerate() { + metas.push(write_manifest(&file_io, &format!("{test}-{i}"), entries).await); + } + aggregate_manifests( + &file_io, + MANIFEST_DIR, + metas, + None, + data_evolution_enabled, + retained_add_budget, + Arc::new(DeletionVectors::default()), + ) + .await + .unwrap() + .into_iter() + .map(|(partition, accum)| { + assert!(!accum.row_count_unknown); + (partition, accum.plain_rows + accum.row_ranges.total()) + }) + .collect() + } + + #[test] + fn test_row_range_set_coalesces_overlapping_and_adjacent_ranges() { + let mut set = RowRangeSet::default(); + set.insert(0, 99); + set.insert(0, 99); // column-group file over the same rows + set.insert(0, 39); // blob file inside the data file's range + set.insert(200, 299); + assert_eq!(set.total(), 200); + assert_eq!(set.ranges.len(), 2); + + set.insert(100, 199); // adjacent on both sides + assert_eq!(set.total(), 300); + assert_eq!(set.ranges.len(), 1); + + set.insert(250, 349); // partial overlap + set.insert(5, 4); // empty + assert_eq!(set.total(), 350); + assert_eq!(set.ranges.len(), 1); + + let mut other = RowRangeSet::default(); + other.insert(340, 399); + other.insert(1000, 1009); + set.merge(other); + assert_eq!(set.total(), 410); + assert_eq!(set.ranges.len(), 2); + } + + #[tokio::test] + async fn test_data_evolution_files_sharing_rows_count_once() { + let result = counts( + "de", + vec![ + vec![ + add(partition(1), file("base-0", 0, 100, Some(0))), + add(partition(1), file("base-1", 0, 50, Some(100))), + add(partition(2), file("other", 0, 7, Some(150))), + ], + // Column-group files written later over the same row-id ranges. + vec![ + add(partition(1), file("cols-0", 0, 100, Some(0))), + add(partition(1), file("cols-1", 0, 50, Some(100))), + ], + ], + true, + RETAINED_ADD_BUDGET, + ) + .await; + + assert_eq!(result[&partition(1)], 150); + assert_eq!(result[&partition(2)], 7); + } + + #[tokio::test] + async fn test_row_ranges_are_not_merged_without_data_evolution() { + let result = counts( + "row-tracking", + vec![vec![ + add(partition(1), file("a", 0, 100, Some(0))), + add(partition(1), file("b", 0, 100, Some(0))), + ]], + false, + RETAINED_ADD_BUDGET, + ) + .await; + + assert_eq!(result[&partition(1)], 200); + } + + #[tokio::test] + async fn test_deleted_files_are_netted_by_complete_identity() { + for budget in [RETAINED_ADD_BUDGET, 0, 1] { + let mut external_add = file("external", 0, 1, None); + external_add.external_path = Some("file:///a".to_string()); + let mut external_delete = external_add.clone(); + external_delete.external_path = Some("file:///b".to_string()); + + let mut extra_add = file("extra", 0, 1, None); + extra_add.extra_files = vec!["a.idx".to_string()]; + let mut extra_delete = extra_add.clone(); + extra_delete.extra_files = vec!["b.idx".to_string()]; + + let mut embedded_add = file("embedded", 0, 1, None); + embedded_add.embedded_index = Some(vec![1]); + let mut embedded_delete = embedded_add.clone(); + embedded_delete.embedded_index = Some(vec![2]); + let mut embedded_match = file("embedded-match", 0, 1, None); + embedded_match.embedded_index = Some(vec![3]); + let embedded_match_delete = embedded_match.clone(); + + let result = counts( + &format!("net-{budget}"), + vec![ + vec![ + add(partition(1), file("a", 0, 10, None)), + add(partition(1), file("b", 0, 20, None)), + add(partition(2), file("gone", 0, 5, None)), + ], + vec![ + delete(partition(1), file("a", 0, 10, None)), + delete(partition(1), file("b", 0, 20, None)), + delete(partition(1), external_delete), + delete(partition(1), extra_delete), + delete(partition(1), embedded_delete), + delete(partition(1), embedded_match_delete), + add(partition(1), external_add), + add(partition(1), extra_add), + add(partition(1), embedded_add), + add(partition(1), embedded_match), + add(partition(1), file("c", 5, 30, None)), + add(partition(1), file("up", 0, 4, None)), + delete(partition(2), file("gone", 0, 5, None)), + ], + vec![ + delete(partition(1), file("up", 0, 4, None)), + add(partition(1), file("up", 5, 4, None)), + ], + ], + false, + budget, + ) + .await; + + assert_eq!(result[&partition(1)], 37, "budget {budget}"); + assert!(!result.contains_key(&partition(2)), "budget {budget}"); + } + } + + #[tokio::test] + async fn test_unknown_row_count_is_reported_as_unknown() { + let file_io = FileIOBuilder::new("memory").build().unwrap(); + file_io.mkdirs(&format!("{MANIFEST_DIR}/")).await.unwrap(); + let meta = write_manifest( + &file_io, + "unknown-0", + &[ + add(partition(1), file("known", 0, 10, None)), + add( + partition(1), + file("unknown", 0, DataFileMeta::ROW_COUNT_UNKNOWN, None), + ), + ], + ) + .await; + let totals = aggregate_manifests( + &file_io, + MANIFEST_DIR, + vec![meta], + None, + false, + RETAINED_ADD_BUDGET, + Arc::new(DeletionVectors::default()), + ) + .await + .unwrap(); + assert!(totals[&partition(1)].row_count_unknown); + } + + #[test] + fn test_missing_manifest_row_counts_are_unknown() { + use crate::spec::avro::{from_avro_bytes_fast, from_manifest_bytes_filtered, SchemaCache}; + use apache_avro::types::Value; + use serde_json::json; + + // Missing counts stay unknown, and absent file identities are rejected. + // Full and slim decoders must agree, including for nonstandard schemas. + for ((case, expected), reordered) in [ + ("missing_count", None), + ("null_count", None), + ("null_file", None), + ("null_deleted_file", None), + ("missing_file", None), + ("zero", Some(0)), + ("known", Some(2)), + ] + .into_iter() + .flat_map(|case| [(case, false), (case, true)]) + { + let mut file_fields = vec![json!({"name": "_FILE_NAME", "type": "string"})]; + let mut file_values = vec![("_FILE_NAME".into(), Value::String("data.parquet".into()))]; + if case != "missing_count" { + file_fields.push(json!({"name": "_ROW_COUNT", "type": ["null", "long"]})); + file_values.push(( + "_ROW_COUNT".into(), + match expected { + Some(count) => Value::Union(1, Box::new(Value::Long(count))), + None => Value::Union(0, Box::new(Value::Null)), + }, + )); + } + file_fields.push(json!({"name": "_FUTURE_FILE_FIELD", "type": ["null", {"type": "array", "items": "long"}]})); + file_values.push(( + "_FUTURE_FILE_FIELD".into(), + Value::Union(1, Box::new(Value::Array(vec![Value::Long(7)]))), + )); + if reordered { + file_fields.reverse(); + file_values.reverse(); + } + let partition = crate::spec::EMPTY_SERIALIZED_ROW.clone(); + let mut fields = vec![ + json!({"name": "_PARTITION", "type": "bytes"}), + json!({"name": "_KIND", "type": "int"}), + json!({"name": "_BUCKET", "type": "int"}), + json!({"name": "_TOTAL_BUCKETS", "type": "int"}), + json!({"name": "_VERSION", "type": ["null", "int"]}), + ]; + let mut values = vec![ + ("_PARTITION".into(), Value::Bytes(partition.clone())), + ( + "_KIND".into(), + Value::Int(i32::from(case == "null_deleted_file")), + ), + ("_BUCKET".into(), Value::Int(3)), + ("_TOTAL_BUCKETS".into(), Value::Int(4)), + ("_VERSION".into(), Value::Union(1, Box::new(Value::Int(2)))), + ]; + let null_file = matches!(case, "null_file" | "null_deleted_file"); + if case != "missing_file" { + fields.push(json!({"name": "_FILE", "type": ["null", { + "type": "record", "name": "file", "fields": file_fields + }]})); + values.push(( + "_FILE".into(), + if null_file { + Value::Union(0, Box::new(Value::Null)) + } else { + Value::Union(1, Box::new(Value::Record(file_values))) + }, + )); + } + fields.push(json!({"name": "_FUTURE_ENTRY_FIELD", "type": ["null", "string"]})); + values.push(( + "_FUTURE_ENTRY_FIELD".into(), + Value::Union(1, Box::new(Value::String("ignored".into()))), + )); + if reordered { + fields.reverse(); + values.reverse(); + } + let record_schema = json!({"type": "record", "name": "manifest", "fields": fields}); + let schema = apache_avro::Schema::parse_str( + &if reordered { + json!(["null", record_schema]) + } else { + record_schema + } + .to_string(), + ) + .unwrap(); + let mut writer = apache_avro::Writer::new(&schema, Vec::new()); + let record = Value::Record(values); + writer + .append(if reordered { + Value::Union(1, Box::new(record)) + } else { + record + }) + .unwrap(); + let bytes = writer.into_inner().unwrap(); + if !reordered && case != "missing_file" { + // The early filter must skip _FILE, even if it is null. + assert!(from_manifest_bytes_filtered( + &bytes, + &mut SchemaCache::new(), + &mut |_, _, _, _| false + ) + .unwrap() + .is_empty()); + } + let full = from_avro_bytes_fast::(&bytes); + let filtered = + from_manifest_bytes_filtered(&bytes, &mut SchemaCache::new(), &mut |_, _, _, _| { + true + }); + let totals = aggregate_manifest( + &bytes, + &SharedSchemaCache::new(), + None, + None, + false, + &DeletionVectors::default(), + ); + if null_file || case == "missing_file" { + for result in [full.map(|_| ()), filtered.map(|_| ()), totals.map(|_| ())] { + assert!( + matches!(result, Err(crate::Error::DataInvalid { .. })), + "{case}" + ); + } + continue; + } + let full = full.unwrap(); + assert_eq!(full, filtered.unwrap(), "{case}, reordered={reordered}"); + assert_eq!( + full[0], + ManifestEntry::new( + FileKind::Add, + partition.clone(), + 3, + 4, + full[0].file().clone(), + 2 + ) + ); + assert_eq!( + full[0].file().row_count, + expected.unwrap_or(DataFileMeta::ROW_COUNT_UNKNOWN), + "{case}" + ); + assert_eq!( + totals.unwrap()[&partition].record_count(), + expected, + "{case}" + ); + } + } + + #[test] + fn test_invalid_row_id_range_is_unknown() { + let mut invalid = PartitionAccum::default(); + invalid.add_file(2, Some(i64::MAX), true); + assert_eq!(invalid.record_count(), None); + + let mut valid = PartitionAccum::default(); + valid.add_file(1, Some(i64::MAX), true); + assert_eq!(valid.record_count(), Some(1)); + } + + /// The slim decoder must agree with the full decoder on every field it reads. + #[test] + fn test_slim_decode_matches_full_decode() { + let path = std::env::current_dir() + .unwrap() + .join("tests/fixtures/manifest/manifest-8ded1f09-fcda-489e-9167-582ac0f9f846-0"); + let bytes = std::fs::read(path).unwrap(); + let full = crate::spec::avro::from_avro_bytes_fast::(&bytes).unwrap(); + + let mut slim = Vec::new(); + visit_slim_manifest_entries(&bytes, &SharedSchemaCache::new(), &mut |e| { + let file = full[slim.len()].file(); + assert_eq!( + e.extra_files, + file.extra_files + .iter() + .map(String::as_str) + .collect::>() + ); + assert_eq!(e.embedded_index, file.embedded_index.as_deref()); + assert_eq!(e.external_path, file.external_path.as_deref()); + slim.push(( + e.kind, + e.partition.to_vec(), + e.bucket, + e.level, + e.file_name.to_string(), + e.row_count, + e.first_row_id, + )); + Ok(()) + }) + .unwrap(); + + let expected: Vec<_> = full + .iter() + .map(|e| { + let f = e.file(); + ( + *e.kind(), + e.partition().to_vec(), + e.bucket(), + f.level, + f.file_name.clone(), + f.row_count, + f.first_row_id, + ) + }) + .collect(); + assert!(!expected.is_empty()); + assert_eq!(slim, expected); + } +} diff --git a/crates/paimon/src/table/table_scan.rs b/crates/paimon/src/table/table_scan.rs index 0decc94e9..2a84357af 100644 --- a/crates/paimon/src/table/table_scan.rs +++ b/crates/paimon/src/table/table_scan.rs @@ -30,8 +30,7 @@ use super::partition_filter::PartitionFilter; use super::row_position_selection::RowPositionSelection; use super::stats_filter::{ data_evolution_group_matches_predicates_for_table, data_file_matches_predicates_for_table, - data_file_matches_predicates_with_key_stats, group_by_overlapping_row_id, FileStatsRows, - ResolvedStatsSchema, + data_file_matches_predicates_with_key_stats, group_by_overlapping_row_id, ResolvedStatsSchema, }; use super::{find_field_id_by_name, Table}; use crate::io::FileIO; @@ -276,21 +275,7 @@ async fn read_all_manifest_entries( // whose partition range doesn't overlap the partition predicate. let manifest_files_before_partition_pruning = manifest_files.len(); if let Some(pf) = partition_filter { - if !partition_fields.is_empty() { - manifest_files.retain(|meta| { - let stats = meta.partition_stats(); - let min_values = BinaryRow::from_serialized_bytes(stats.min_values()).ok(); - let max_values = BinaryRow::from_serialized_bytes(stats.max_values()).ok(); - let null_counts = stats.null_counts().clone(); - let file_stats = FileStatsRows::for_manifest_partition( - meta.num_added_files() + meta.num_deleted_files(), - min_values, - max_values, - null_counts, - ); - pf.matches_manifest(&file_stats, partition_fields) - }); - } + manifest_files.retain(|meta| pf.matches_manifest(meta, partition_fields)); } if let Some(trace) = trace.as_deref_mut() { trace.manifest_files_before_partition_pruning = manifest_files_before_partition_pruning; @@ -2112,21 +2097,7 @@ impl<'a> PaimonTableScan<'a> { read_manifest_list(file_io, table_path, manifest_list_name).await?; if let Some(pf) = self.partition_filter.as_ref() { - if !partition_fields.is_empty() { - manifest_metas.retain(|meta| { - let stats = meta.partition_stats(); - let min_values = BinaryRow::from_serialized_bytes(stats.min_values()).ok(); - let max_values = BinaryRow::from_serialized_bytes(stats.max_values()).ok(); - let null_counts = stats.null_counts().clone(); - let file_stats = FileStatsRows::for_manifest_partition( - meta.num_added_files() + meta.num_deleted_files(), - min_values, - max_values, - null_counts, - ); - pf.matches_manifest(&file_stats, &partition_fields) - }); - } + manifest_metas.retain(|meta| pf.matches_manifest(meta, &partition_fields)); } if let Some(index) = row_range_index { retain_manifest_row_ranges(&mut manifest_metas, index); diff --git a/crates/paimon/src/table/time_travel.rs b/crates/paimon/src/table/time_travel.rs index f5676bce7..838f0d3cc 100644 --- a/crates/paimon/src/table/time_travel.rs +++ b/crates/paimon/src/table/time_travel.rs @@ -405,6 +405,36 @@ mod tests { assert!(recopied.travel_snapshot().is_none()); } + #[tokio::test] + async fn test_pinned_snapshot_preserves_read_schema_without_rereading() { + let (file_io, table_path) = setup_evolved_table().await; + let table = latest_table(&file_io, &table_path) + .copy_with_options(options(&[("scan.version", "2"), ("custom", "value")])); + let manager = table.snapshot_manager(); + let snapshot = manager.get_snapshot(1).await.unwrap(); + // Pinning and later resolution must reuse the supplied snapshot, not + // read it again or switch the table's schema to schema 0. + file_io + .delete_file(&manager.snapshot_path(1)) + .await + .unwrap(); + let pinned = table.copy_with_pinned_snapshot(&snapshot); + assert_eq!(pinned.schema().id(), table.schema().id()); + assert_eq!(pinned.schema().fields(), table.schema().fields()); + assert_eq!(pinned.schema().options().get("custom").unwrap(), "value"); + assert!(!pinned.schema().options().contains_key("scan.version")); + assert_eq!( + pinned.schema().options().get("scan.snapshot-id").unwrap(), + "1" + ); + assert_eq!( + super::resolve_snapshot(&pinned).await.unwrap().unwrap(), + snapshot + ); + assert!(pinned.new_write_builder().new_write().is_err()); + assert!(table.travel_snapshot().is_none()); + } + #[tokio::test] async fn test_copy_with_time_travel_same_schema_still_rejects_write() { let (file_io, table_path) = setup_evolved_table().await;