blob: bd847d1b19a9fb3ace7420ecf33dc39de35a51e8 [file]
// 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.
//! TableScan for full table scan.
//!
//! Reference: [pypaimon.read.table_scan.TableScan](https://github.com/apache/paimon/blob/release-1.3/paimon-python/pypaimon/read/table_scan.py)
//! and [FullStartingScanner](https://github.com/apache/paimon/blob/release-1.3/paimon-python/pypaimon/read/scanner/full_starting_scanner.py).
use super::bucket_filter::compute_target_buckets;
use super::partition_filter::PartitionFilter;
use super::stats_filter::{
data_evolution_group_matches_predicates, data_file_matches_predicates,
data_file_matches_predicates_for_table, group_by_overlapping_row_id, FileStatsRows,
ResolvedStatsSchema,
};
use super::Table;
use crate::io::FileIO;
use crate::spec::{
avro::SharedSchemaCache, bucket_dir_name, BinaryRow, BucketFunctionType, CoreOptions,
DataField, DataFileMeta, FileKind, IndexManifest, ManifestEntry, PartitionComputer, Predicate,
Snapshot,
};
use crate::table::bin_pack::split_for_batch;
use crate::table::merge_tree_split_generator::{
merge_tree_split_for_batch, KeyComparator, SplitGroup,
};
use crate::table::source::{
any_range_overlaps_file, intersect_ranges_with_file, merge_row_ranges, DataSplit,
DataSplitBuilder, DeletionFile, PartitionBucket, Plan, RowRange,
};
use crate::table::SnapshotManager;
use futures::{StreamExt, TryStreamExt};
use std::collections::{HashMap, HashSet};
use std::sync::Arc;
/// Path segment for manifest directory under table.
const MANIFEST_DIR: &str = "manifest";
/// Path segment for index directory under table.
const INDEX_DIR: &str = "index";
/// Reads a manifest list file (Avro) and returns manifest file metas.
async fn read_manifest_list(
file_io: &FileIO,
table_path: &str,
list_name: &str,
) -> crate::Result<Vec<crate::spec::ManifestFileMeta>> {
if list_name.is_empty() {
return Ok(Vec::new());
}
let path = format!(
"{}/{}/{}",
table_path.trim_end_matches('/'),
MANIFEST_DIR,
list_name
);
let input = file_io.new_input(&path)?;
let bytes = input.read().await?;
crate::spec::avro::from_avro_bytes_fast::<crate::spec::ManifestFileMeta>(&bytes)
}
/// Reads all manifest entries for a snapshot (base + delta manifest lists, then each manifest file).
/// Applies filters during concurrent manifest reading to reduce entries early:
/// - Manifest-file-level partition stats pruning (skip entire manifest files)
/// - Level-0 filtering per entry (DV mode or FirstRow engine)
/// - Partition predicate filtering per entry
/// - Data-level stats pruning per entry (current schema only, cross-schema fail-open)
#[allow(clippy::too_many_arguments)]
async fn read_all_manifest_entries(
file_io: &FileIO,
table_path: &str,
snapshot: &Snapshot,
skip_level_zero: bool,
scan_all_files: bool,
has_primary_keys: bool,
partition_filter: Option<&PartitionFilter>,
partition_fields: &[DataField],
data_predicates: &[Predicate],
current_schema_id: i64,
schema_fields: &[DataField],
bucket_predicate: Option<&Predicate>,
bucket_key_fields: &[DataField],
bucket_function_type: BucketFunctionType,
) -> crate::Result<Vec<ManifestEntry>> {
let (mut manifest_files, delta) = futures::try_join!(
read_manifest_list(file_io, table_path, snapshot.base_manifest_list()),
read_manifest_list(file_io, table_path, snapshot.delta_manifest_list()),
)?;
manifest_files.extend(delta);
// Manifest-file-level partition stats pruning: skip entire manifest files
// whose partition range doesn't overlap the partition predicate.
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)
});
}
}
let manifest_path_prefix = format!("{}/{}", table_path.trim_end_matches('/'), MANIFEST_DIR);
let shared_cache = SharedSchemaCache::new();
let all_entries: Vec<ManifestEntry> = futures::stream::iter(manifest_files)
.map(|meta| {
let path = format!("{}/{}", manifest_path_prefix, meta.file_name());
let cache = shared_cache.clone();
async move {
let input_file = file_io.new_input(&path)?;
let content = input_file.read().await?;
// Per-task bucket cache (few distinct total_buckets values per manifest).
let mut bucket_cache: HashMap<i32, Option<HashSet<i32>>> = HashMap::new();
let entries = crate::spec::avro::from_manifest_bytes_filtered_shared(
&content,
&cache,
&mut |_kind, partition_bytes, bucket, total_buckets| {
// Bucket filter (negative bucket = unassigned)
if has_primary_keys && !scan_all_files && bucket < 0 {
return false;
}
if let Some(pred) = bucket_predicate {
let targets = bucket_cache.entry(total_buckets).or_insert_with(|| {
compute_target_buckets(
pred,
bucket_key_fields,
bucket_function_type,
total_buckets,
)
});
if let Some(targets) = targets {
if !targets.contains(&bucket) {
return false;
}
}
}
// Partition filter
if let Some(pf) = partition_filter {
match pf.matches_entry(partition_bytes) {
Ok(false) => return false,
Ok(true) => {}
Err(_) => {}
}
}
true
},
)?;
// Post-filter: level-0 and data predicates (need DataFileMeta)
let filtered: Vec<ManifestEntry> = entries
.into_iter()
.filter(|entry| {
if skip_level_zero && has_primary_keys && entry.file().level == 0 {
return false;
}
if !data_predicates.is_empty()
&& !data_file_matches_predicates(
entry.file(),
data_predicates,
current_schema_id,
schema_fields,
)
{
return false;
}
true
})
.collect();
Ok::<_, crate::Error>(filtered)
}
})
.buffered(64)
.try_collect::<Vec<_>>()
.await?
.into_iter()
.flatten()
.collect();
Ok(all_entries)
}
/// Builds a map from (partition, bucket) to (data_file_name -> DeletionFile) from index manifest entries.
/// Only considers ADD entries with index_type "DELETION_VECTORS" and their deletion_vectors_ranges.
fn build_deletion_files_map(
index_entries: &[crate::spec::IndexManifestEntry],
table_path: &str,
) -> HashMap<PartitionBucket, HashMap<String, DeletionFile>> {
use crate::spec::FileKind;
let table_path = table_path.trim_end_matches('/');
let index_path_prefix = format!("{table_path}/{INDEX_DIR}");
let mut map: HashMap<PartitionBucket, HashMap<String, DeletionFile>> =
HashMap::with_capacity(index_entries.len());
for entry in index_entries {
if entry.kind != FileKind::Add {
continue;
}
if entry.index_file.index_type != "DELETION_VECTORS" {
continue;
}
let ranges = match &entry.index_file.deletion_vectors_ranges {
Some(r) if !r.is_empty() => r,
_ => continue,
};
let key = PartitionBucket::new(entry.partition.clone(), entry.bucket);
let dv_path = format!("{}/{}", index_path_prefix, entry.index_file.file_name);
let per_bucket = map.entry(key).or_default();
for (data_file_name, meta) in ranges {
per_bucket.insert(
data_file_name.clone(),
DeletionFile::new(
dv_path.clone(),
meta.offset as i64,
meta.length as i64,
meta.cardinality,
),
);
}
}
map
}
/// Nets add/delete manifest entries for a scan, returning only the live ADD set.
///
/// Mirrors Java `AbstractFileStoreScan.readAndMergeFileEntries`: first collect
/// the full [`Identifier`] of every DELETE entry, then keep the ADD entries
/// whose identifier is not in that set. The identity is the complete Paimon file
/// identity (`partition, bucket, level, file_name, extra_files, embedded_index,
/// external_path`, matching Java `FileEntry.Identifier`).
///
/// Keying on file name alone is wrong: a single-run compaction upgrades a file
/// *in place* — `DELETE f@oldLevel` plus `ADD f@newLevel` with the same file
/// name (`PojoDataFileMeta.upgrade` reuses the name, only changing `level`). An
/// identity without `level` lets the DELETE cancel the upgraded ADD, dropping
/// the file from the scan and silently losing its rows on read.
///
/// Collecting deletes first (rather than insert/remove while iterating) makes
/// the result independent of ADD/DELETE ordering, matching the Java scan path.
fn merge_manifest_entries(entries: Vec<ManifestEntry>) -> Vec<ManifestEntry> {
use crate::spec::Identifier;
let deleted: HashSet<Identifier> = entries
.iter()
.filter(|e| *e.kind() == FileKind::Delete)
.map(|e| e.identifier())
.collect();
entries
.into_iter()
.filter(|e| *e.kind() == FileKind::Add && !deleted.contains(&e.identifier()))
.collect()
}
/// Whether scan-owned pruning still preserves `merged_row_count()` as a safe
/// row-count hint.
///
/// Data predicates and row ranges can reduce rows within a split after planning,
/// so split-level row counts stop being a conservative bound for final rows.
pub(super) fn can_push_down_limit_hint_for_scan(
data_predicates: &[Predicate],
row_ranges: Option<&[RowRange]>,
) -> bool {
data_predicates.is_empty() && row_ranges.is_none()
}
fn should_skip_level_zero_for_scan(
scan_all_files: bool,
has_primary_keys: bool,
deletion_vectors_enabled: bool,
merge_engine: crate::Result<crate::spec::MergeEngine>,
) -> bool {
if scan_all_files {
return false;
}
if !has_primary_keys {
return false;
}
deletion_vectors_enabled || merge_engine.is_ok_and(|e| e == crate::spec::MergeEngine::FirstRow)
}
/// TableScan for full table scan (no incremental, no predicate).
///
/// Reference: [pypaimon.read.table_scan.TableScan](https://github.com/apache/paimon/blob/master/paimon-python/pypaimon/read/table_scan.py)
#[derive(Debug, Clone)]
pub struct TableScan<'a> {
table: &'a Table,
partition_filter: Option<PartitionFilter>,
data_predicates: Vec<Predicate>,
bucket_predicate: Option<Predicate>,
/// Optional limit on the number of rows to return.
/// When set, the scan will try to return only enough splits to satisfy the limit.
limit: Option<usize>,
row_ranges: Option<Vec<RowRange>>,
/// When true, disables level-0 filtering so all files are visible.
/// Used by non-read paths (overwrite, truncate, writer restore) that need
/// the complete file set. Normal read scans leave this as `false`.
scan_all_files: bool,
}
impl<'a> TableScan<'a> {
pub(crate) fn new(
table: &'a Table,
partition_filter: Option<PartitionFilter>,
data_predicates: Vec<Predicate>,
bucket_predicate: Option<Predicate>,
limit: Option<usize>,
row_ranges: Option<Vec<RowRange>>,
) -> Self {
Self {
table,
partition_filter,
data_predicates,
bucket_predicate,
limit,
row_ranges,
scan_all_files: false,
}
}
/// Disable level-0 filtering so all files are visible.
///
/// Used by non-read paths (overwrite, truncate, writer restore) that need
/// the complete file set regardless of merge engine or DV settings.
pub fn with_scan_all_files(mut self) -> Self {
self.scan_all_files = true;
self
}
/// Set row ranges for scan-time filtering.
///
/// This replaces any existing row_ranges. Typically used to inject
/// results from global index lookups (e.g. full-text search).
pub fn with_row_ranges(mut self, ranges: Vec<RowRange>) -> Self {
self.row_ranges = if ranges.is_empty() {
None
} else {
Some(ranges)
};
self
}
/// Plan the full scan: resolve snapshot (via options or latest), then read manifests and build DataSplits.
///
/// Time travel is resolved from table options:
/// - only one of `scan.version`, `scan.timestamp-millis` may be set
/// - `scan.version` → tag name (if exists) → snapshot id (if parseable) → error
/// - `scan.timestamp-millis` → find the latest snapshot <= that timestamp
/// - otherwise → read the latest snapshot
///
/// Reference: [TimeTravelUtil.tryTravelToSnapshot](https://github.com/apache/paimon/blob/master/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/TimeTravelUtil.java)
pub async fn plan(&self) -> crate::Result<Plan> {
let snapshot = match self.resolve_snapshot().await? {
Some(snapshot) => snapshot,
None => return Ok(Plan::new(Vec::new())),
};
self.plan_snapshot(snapshot).await
}
async fn resolve_snapshot(&self) -> crate::Result<Option<Snapshot>> {
// A table copy produced by `copy_with_time_travel` already resolved
// the selector in its options; reuse it instead of re-reading
// tag/snapshot files on every plan.
if let Some(snapshot) = self.table.travel_snapshot() {
return Ok(Some(snapshot.clone()));
}
// A time-travelled schema without its resolved snapshot means the
// selector was changed after the travel (`copy_with_options`).
// Resolving the new selector here would evolve a different snapshot's
// files to the stale historical schema, so fail instead.
if self.table.is_time_traveled() {
return Err(crate::Error::DataInvalid {
message: "Table options changed after time travel; \
use copy_with_time_travel to re-resolve the snapshot and schema"
.to_string(),
source: None,
});
}
let file_io = self.table.file_io();
let table_path = self.table.location();
match super::time_travel::travel_to_snapshot(
file_io,
table_path,
self.table.schema().options(),
)
.await?
{
Some(snapshot) => Ok(Some(snapshot)),
None => {
let snapshot_manager =
SnapshotManager::new(file_io.clone(), table_path.to_string());
snapshot_manager.get_latest_snapshot().await
}
}
}
/// Apply a limit-pushdown hint to the generated splits.
///
/// Mirrors Java `DataTableBatchScan#applyPushDownLimit`: splits whose
/// `merged_row_count()` is unknown (for example merge-needed PK splits or
/// unknown deletion cardinality) are skipped — they contribute an unknown
/// number of rows, so they cannot help satisfy the limit. Pruning is
/// committed only once the accumulated known row count reaches the limit;
/// if it never does, the original split list is returned unchanged. The
/// caller or query engine must still enforce the final LIMIT.
fn apply_limit_pushdown(&self, splits: Vec<DataSplit>) -> Vec<DataSplit> {
let limit = match self.limit {
Some(l) => l,
None => return splits,
};
if limit == 0 {
return Vec::new();
}
if splits.is_empty() {
return splits;
}
let mut limited_splits = Vec::new();
let mut scanned_row_count: i64 = 0;
for split in &splits {
if let Some(merged_count) = split.merged_row_count() {
limited_splits.push(split.clone());
scanned_row_count += merged_count;
if scanned_row_count >= limit as i64 {
return limited_splits;
}
}
}
splits
}
/// Read all manifest entries from a snapshot, applying filters and merging.
///
/// This is the shared entry point used by both `plan_snapshot` (scan) and
/// `TableCommit` (overwrite). Filters include partition predicate, data
/// predicates, and bucket predicate.
pub(crate) async fn plan_manifest_entries(
&self,
snapshot: &Snapshot,
) -> crate::Result<Vec<ManifestEntry>> {
let file_io = self.table.file_io();
let table_path = self.table.location();
let core_options = CoreOptions::new(self.table.schema().options());
let data_evolution_enabled = core_options.data_evolution_enabled();
let has_primary_keys = !self.table.schema().primary_keys().is_empty();
let deletion_vectors_enabled = core_options.deletion_vectors_enabled();
// Skip level-0 files for PK tables when:
// - DV mode: level-0 files are unmerged, DV handles dedup at higher levels
// - FirstRow engine without DV: reads go through DataFileReader (no merge),
// so only compacted (level > 0) files are safe to read directly
// Deduplicate engine always uses KeyValueFileReader which handles level-0
// via sort-merge, so level-0 files must remain visible.
//
// Non-read paths (overwrite, truncate, writer restore) set scan_all_files=true
// to see all files including level-0, matching Java's CommitScanner behavior.
let skip_level_zero = should_skip_level_zero_for_scan(
self.scan_all_files,
has_primary_keys,
deletion_vectors_enabled,
core_options.merge_engine(),
);
let partition_fields = self.table.schema().partition_fields();
let pushdown_data_predicates = if data_evolution_enabled {
&[][..]
} else {
self.data_predicates.as_slice()
};
let bucket_key_fields: Vec<DataField> = if self.bucket_predicate.is_none() {
Vec::new()
} else {
let bucket_keys = core_options.bucket_key().unwrap_or_else(|| {
if has_primary_keys {
self.table.schema().trimmed_primary_keys()
} else {
Vec::new()
}
});
bucket_keys
.iter()
.filter_map(|key| {
self.table
.schema()
.fields()
.iter()
.find(|f| f.name() == key)
.cloned()
})
.collect::<Vec<_>>()
};
let bucket_function_type = core_options.bucket_function_type()?;
let entries = read_all_manifest_entries(
file_io,
table_path,
snapshot,
skip_level_zero,
self.scan_all_files,
has_primary_keys,
self.partition_filter.as_ref(),
&partition_fields,
pushdown_data_predicates,
self.table.schema().id(),
self.table.schema().fields(),
self.bucket_predicate.as_ref(),
&bucket_key_fields,
bucket_function_type,
)
.await?;
Ok(merge_manifest_entries(entries))
}
fn can_push_down_limit_hint(&self, row_ranges: Option<&[RowRange]>) -> bool {
can_push_down_limit_hint_for_scan(&self.data_predicates, row_ranges)
}
async fn plan_snapshot(&self, snapshot: Snapshot) -> crate::Result<Plan> {
let file_io = self.table.file_io();
let table_path = self.table.location();
let core_options = CoreOptions::new(self.table.schema().options());
let data_evolution_enabled = core_options.data_evolution_enabled();
let target_split_size = core_options.source_split_target_size();
let open_file_cost = core_options.source_split_open_file_cost();
let partition_keys = self.table.schema().partition_keys();
let entries = self.plan_manifest_entries(&snapshot).await?;
if entries.is_empty() {
return Ok(Plan::new(Vec::new()));
}
// For non-data-evolution tables, cross-schema files were kept (fail-open)
// by the pushdown. Apply the full schema-aware filter for those files.
let entries = if self.data_predicates.is_empty() || data_evolution_enabled {
entries
} else {
let current_schema_id = self.table.schema().id();
let has_cross_schema = entries
.iter()
.any(|e| e.file().schema_id != current_schema_id);
if !has_cross_schema {
entries
} else {
let mut kept = Vec::with_capacity(entries.len());
let mut schema_cache: HashMap<i64, Option<Arc<ResolvedStatsSchema>>> =
HashMap::new();
for entry in entries {
if entry.file().schema_id == current_schema_id
|| data_file_matches_predicates_for_table(
self.table,
entry.file(),
&self.data_predicates,
&mut schema_cache,
)
.await
{
kept.push(entry);
}
}
kept
}
};
if entries.is_empty() {
return Ok(Plan::new(Vec::new()));
}
// Group by (partition, bucket), decomposing entries to avoid cloning partition.
let mut groups: HashMap<(Vec<u8>, i32), (i32, Vec<DataFileMeta>)> =
HashMap::with_capacity(entries.len());
for e in entries {
let (partition, bucket, total_buckets, file) = e.into_parts();
let entry = groups
.entry((partition, bucket))
.or_insert_with(|| (total_buckets, Vec::new()));
entry.1.push(file);
}
let snapshot_id = snapshot.id();
let base_path = table_path.trim_end_matches('/');
let mut splits = Vec::with_capacity(groups.len());
let partition_computer = if !partition_keys.is_empty() {
Some(PartitionComputer::new(
partition_keys,
self.table.schema().fields(),
core_options.partition_default_name(),
core_options.legacy_partition_name(),
)?)
} else {
None
};
// Primary-key tables must keep key-overlapping files in one split so the
// sort-merge reader sees every version of a key. The comparator decodes
// the trimmed-PK min/max keys written by the kv writer.
//
// Deletion-vector and first-row tables read without merging (stale rows
// are masked by DVs / level-0 is skipped), so they keep plain size-based
// packing like Java's MergeTreeSplitGenerator fast path.
let read_merges_overlapping_keys = !core_options.deletion_vectors_enabled()
&& !matches!(
core_options.merge_engine(),
Ok(crate::spec::MergeEngine::FirstRow)
);
let pk_comparator = if read_merges_overlapping_keys {
KeyComparator::from_table_schema(self.table.schema())
} else {
None
};
// Read deletion vector index manifest once (like Java generateSplits / scanDvIndex).
let (deletion_files_map, effective_row_ranges) =
if let Some(index_manifest_name) = snapshot.index_manifest() {
let index_manifest_path = format!("{base_path}/{MANIFEST_DIR}");
let path = format!("{index_manifest_path}/{index_manifest_name}");
let index_entries = IndexManifest::read(file_io, &path).await?;
let dv_map = build_deletion_files_map(&index_entries, base_path);
// Use pushed-down row_ranges first; otherwise try global index.
let row_ranges = if self.row_ranges.is_some() {
self.row_ranges.clone()
} else if data_evolution_enabled
&& core_options.global_index_enabled()
&& !self.data_predicates.is_empty()
{
super::global_index_scanner::evaluate_global_index(
file_io,
base_path,
&index_entries,
&self.data_predicates,
self.table.schema().fields(),
)
.await?
} else {
None
};
(Some(dv_map), row_ranges)
} else {
(None, self.row_ranges.clone())
};
for ((partition, bucket), (total_buckets, data_files)) in groups {
let partition_row = BinaryRow::from_serialized_bytes(&partition)?;
let bucket_path = if let Some(ref computer) = partition_computer {
let partition_path = computer.generate_partition_path(&partition_row)?;
format!("{base_path}/{partition_path}{}", bucket_dir_name(bucket))
} else {
format!("{base_path}/{}", bucket_dir_name(bucket))
};
// Original `partition` Vec consumed by PartitionBucket for DV map lookup.
let per_bucket_deletion_map = deletion_files_map
.as_ref()
.and_then(|map| map.get(&PartitionBucket::new(partition, bucket)));
// Data-evolution tables merge overlapping row-id groups column-wise during read.
// Keep that split boundary intact and only bin-pack single-file groups.
// Apply group-level predicate filtering after grouping by row_id range.
let file_groups: Vec<SplitGroup> = if data_evolution_enabled {
let row_id_groups = group_by_overlapping_row_id(data_files);
// Filter groups by merged stats before splitting.
let row_id_groups: Vec<Vec<DataFileMeta>> = if self.data_predicates.is_empty() {
row_id_groups
} else {
row_id_groups
.into_iter()
.filter(|group| {
data_evolution_group_matches_predicates(
group,
&self.data_predicates,
self.table.schema().fields(),
)
})
.collect()
};
// Filter groups by row ID ranges.
let row_id_groups = if let Some(ref ranges) = effective_row_ranges {
row_id_groups
.into_iter()
.filter(|group| group.iter().any(|f| any_range_overlaps_file(ranges, f)))
.collect()
} else {
row_id_groups
};
let (singles, multis): (Vec<_>, Vec<_>) = row_id_groups
.into_iter()
.partition(|group| group.len() == 1);
let mut result = Vec::new();
for group in multis {
// Files sharing a row-id range hold column slices of the
// same logical rows; physical counts overcount them
// (Java DataEvolutionSplitGenerator: not raw convertible).
result.push(SplitGroup {
files: group,
raw_convertible: false,
});
}
let single_files: Vec<DataFileMeta> = singles.into_iter().flatten().collect();
for file_group in split_for_batch(single_files, target_split_size, open_file_cost) {
result.push(SplitGroup {
files: file_group,
raw_convertible: true,
});
}
result
} else if let Some(ref comparator) = pk_comparator {
// Merge-tree path: keep key-overlapping files in one split and
// mark which splits the sort-merge reader can skip (mirrors
// Java MergeTreeSplitGenerator#splitForBatch). Only engines
// whose writer deduplicates at flush guarantee a file never
// holds two rows of one key, so only they may mark groups raw
// convertible; see merge_tree_split_for_batch. (First-row
// tables do not take this path today, but its writer dedups
// too, so keep the gate accurate.)
let file_keys_unique = matches!(
core_options.merge_engine(),
Ok(crate::spec::MergeEngine::Deduplicate)
| Ok(crate::spec::MergeEngine::FirstRow)
);
merge_tree_split_for_batch(
data_files,
comparator,
target_split_size,
open_file_cost,
file_keys_unique,
)
} else {
split_for_batch(data_files, target_split_size, open_file_cost)
.into_iter()
.map(|files| SplitGroup {
files,
raw_convertible: true,
})
.collect()
};
for group in file_groups {
let SplitGroup {
files: file_group,
raw_convertible,
} = group;
let data_deletion_files = per_bucket_deletion_map.map(|per_bucket| {
file_group
.iter()
.map(|f| per_bucket.get(&f.file_name).cloned())
.collect::<Vec<Option<DeletionFile>>>()
});
// Compute row_ranges before moving file_group to avoid clone
let split_row_ranges = if let Some(ref ranges) = effective_row_ranges {
let mut split_ranges = Vec::new();
for file in &file_group {
split_ranges.extend(intersect_ranges_with_file(ranges, file));
}
let split_ranges = merge_row_ranges(split_ranges);
if split_ranges.is_empty() {
None
} else {
Some(split_ranges)
}
} else {
None
};
let mut builder = DataSplitBuilder::new()
.with_snapshot(snapshot_id)
.with_partition(partition_row.clone())
.with_bucket(bucket)
.with_bucket_path(bucket_path.clone())
.with_total_buckets(total_buckets)
.with_data_files(file_group)
.with_raw_convertible(raw_convertible);
if let Some(files) = data_deletion_files {
builder = builder.with_data_deletion_files(files);
}
if let Some(row_ranges) = split_row_ranges {
builder = builder.with_row_ranges(row_ranges);
}
splits.push(builder.build()?);
}
}
// With data predicates or row_ranges, merged_row_count() reflects pre-filter
// row counts, so stopping early could return fewer rows than the limit.
let splits = if self.can_push_down_limit_hint(effective_row_ranges.as_deref()) {
self.apply_limit_pushdown(splits)
} else {
splits
};
Ok(Plan::new(splits))
}
}
#[cfg(test)]
mod tests {
use super::{should_skip_level_zero_for_scan, TableScan};
use crate::catalog::Identifier;
use crate::io::FileIOBuilder;
use crate::spec::{
stats::BinaryTableStats, ArrayType, BinaryRow, BinaryRowBuilder, BucketFunctionType,
DataField, DataFileMeta, DataType, Datum, DeletionVectorMeta, FileKind, IndexFileMeta,
IndexManifestEntry, IntType, Predicate, PredicateBuilder, PredicateOperator,
Schema as PaimonSchema, TableSchema, VarCharType,
};
use crate::table::bucket_filter::{compute_target_buckets, extract_predicate_for_keys};
use crate::table::partition_filter::PartitionFilter;
use crate::table::source::{DataSplit, DataSplitBuilder, DeletionFile};
use crate::table::stats_filter::{data_file_matches_predicates, group_by_overlapping_row_id};
use crate::table::Table;
use crate::Error;
use chrono::{DateTime, Utc};
use std::collections::HashSet;
/// Helper to build a DataFileMeta with data evolution fields.
fn make_evo_file(
name: &str,
file_size: i64,
row_count: i64,
max_seq: i64,
first_row_id: Option<i64>,
) -> DataFileMeta {
DataFileMeta {
file_name: name.to_string(),
file_size,
row_count,
min_key: Vec::new(),
max_key: Vec::new(),
key_stats: BinaryTableStats::new(Vec::new(), Vec::new(), Vec::new()),
value_stats: BinaryTableStats::new(Vec::new(), Vec::new(), Vec::new()),
min_sequence_number: 0,
max_sequence_number: max_seq,
schema_id: 0,
level: 0,
extra_files: Vec::new(),
creation_time: DateTime::<Utc>::from_timestamp(0, 0),
delete_row_count: None,
embedded_index: None,
first_row_id,
write_cols: None,
external_path: None,
file_source: None,
value_stats_cols: None,
}
}
#[test]
fn test_merge_manifest_entries_keeps_in_place_upgraded_file() {
// Reproduces a single-run compaction "upgrade": the SAME file name is
// deleted at level 0 and re-added at a higher level (Paimon promotes a
// lone sorted run in place instead of rewriting it). Netting ADD/DELETE
// by file name alone (ignoring `level`) wrongly drops the upgraded file.
// Matches Java `FileEntry.Identifier`, which includes `level`.
use super::merge_manifest_entries;
use crate::spec::ManifestEntry;
let entry = |kind: FileKind, name: &str, level: i32| -> ManifestEntry {
let mut file = make_evo_file(name, 1, 1, 1, None);
file.level = level;
ManifestEntry::new(kind, Vec::new(), 0, 1, file, 2)
};
let entries = vec![
entry(FileKind::Add, "f.parquet", 0), // original level-0 write
entry(FileKind::Delete, "f.parquet", 0), // compaction removes the L0 version
entry(FileKind::Add, "f.parquet", 5), // same file upgraded to level 5
entry(FileKind::Add, "g.parquet", 0), // unrelated fresh file
];
let mut live: Vec<(String, i32)> = merge_manifest_entries(entries)
.into_iter()
.map(|e| (e.file().file_name.clone(), e.file().level))
.collect();
live.sort();
assert_eq!(
live,
vec![("f.parquet".to_string(), 5), ("g.parquet".to_string(), 0)],
"upgraded file (f@L5) must survive; only f@L0 is cancelled by the DELETE"
);
}
fn file_names(groups: &[Vec<DataFileMeta>]) -> Vec<Vec<&str>> {
groups
.iter()
.map(|g| g.iter().map(|f| f.file_name.as_str()).collect())
.collect()
}
fn int_stats_row(value: Option<i32>) -> Vec<u8> {
let mut builder = BinaryRowBuilder::new(1);
match value {
Some(value) => builder.write_int(0, value),
None => builder.set_null_at(0),
}
builder.build_serialized()
}
fn partition_string_field() -> Vec<DataField> {
vec![DataField::new(
0,
"dt".to_string(),
DataType::VarChar(VarCharType::default()),
)]
}
fn int_field() -> Vec<DataField> {
vec![DataField::new(
0,
"id".to_string(),
DataType::Int(IntType::new()),
)]
}
fn test_data_file_meta(
min_values: Vec<u8>,
max_values: Vec<u8>,
null_counts: Vec<Option<i64>>,
row_count: i64,
) -> DataFileMeta {
test_data_file_meta_with_schema(
min_values,
max_values,
null_counts,
row_count,
0, // default schema_id
)
}
fn test_data_file_meta_with_schema(
min_values: Vec<u8>,
max_values: Vec<u8>,
null_counts: Vec<Option<i64>>,
row_count: i64,
schema_id: i64,
) -> DataFileMeta {
DataFileMeta {
file_name: "test.parquet".into(),
file_size: 128,
row_count,
min_key: Vec::new(),
max_key: Vec::new(),
key_stats: BinaryTableStats::new(Vec::new(), Vec::new(), Vec::new()),
value_stats: BinaryTableStats::new(min_values, max_values, null_counts),
min_sequence_number: 0,
max_sequence_number: 0,
schema_id,
level: 1,
extra_files: Vec::new(),
creation_time: Some(Utc::now()),
delete_row_count: None,
embedded_index: None,
first_row_id: None,
write_cols: None,
external_path: None,
file_source: None,
value_stats_cols: None,
}
}
fn limit_test_table() -> Table {
let file_io = FileIOBuilder::new("file").build().unwrap();
let schema = PaimonSchema::builder().build().unwrap();
let table_schema = TableSchema::new(0, &schema);
Table::new(
file_io,
Identifier::new("test_db", "test_table"),
"/tmp/test-table".to_string(),
table_schema,
None,
)
}
fn limit_test_split(file_name: &str, row_count: i64) -> DataSplit {
let mut file = test_data_file_meta(Vec::new(), Vec::new(), Vec::new(), row_count);
file.file_name = file_name.to_string();
DataSplitBuilder::new()
.with_snapshot(1)
.with_partition(BinaryRow::new(0))
.with_bucket(0)
.with_bucket_path(format!("file:/tmp/{file_name}"))
.with_total_buckets(1)
.with_data_files(vec![file])
.build()
.unwrap()
}
fn limit_test_split_with_unknown_merged_row_count(
file_name: &str,
row_count: i64,
) -> DataSplit {
let mut file = test_data_file_meta(Vec::new(), Vec::new(), Vec::new(), row_count);
file.file_name = file_name.to_string();
DataSplitBuilder::new()
.with_snapshot(1)
.with_partition(BinaryRow::new(0))
.with_bucket(0)
.with_bucket_path(format!("file:/tmp/{file_name}"))
.with_total_buckets(1)
.with_data_files(vec![file])
.with_data_deletion_files(vec![Some(DeletionFile::new(
format!("file:/tmp/{file_name}.dv"),
0,
0,
None,
))])
.build()
.unwrap()
}
fn split_file_names(splits: &[DataSplit]) -> Vec<&str> {
splits
.iter()
.map(|split| split.data_files()[0].file_name.as_str())
.collect()
}
#[test]
fn test_apply_limit_pushdown_zero_returns_empty() {
let table = limit_test_table();
let scan = TableScan::new(&table, None, vec![], None, Some(0), None);
let splits = vec![
limit_test_split("a.parquet", 2),
limit_test_split("b.parquet", 3),
];
let pruned = scan.apply_limit_pushdown(splits);
assert!(pruned.is_empty());
}
/// Java semantics: unknown-count splits are skipped — they cannot prove
/// progress toward the limit — and pruning commits once the counted
/// splits alone cover the limit.
#[test]
fn test_apply_limit_pushdown_skips_unknown_merged_row_count() {
let table = limit_test_table();
let scan = TableScan::new(&table, None, vec![], None, Some(3), None);
let splits = vec![
limit_test_split("a.parquet", 2),
limit_test_split_with_unknown_merged_row_count("b.parquet", 4),
limit_test_split("c.parquet", 3),
];
let pruned = scan.apply_limit_pushdown(splits);
assert_eq!(split_file_names(&pruned), vec!["a.parquet", "c.parquet"]);
}
/// When counted splits never reach the limit, the original split list is
/// returned unchanged (mirrors Java `applyPushDownLimit`).
#[test]
fn test_apply_limit_pushdown_returns_all_when_limit_not_reached() {
let table = limit_test_table();
let scan = TableScan::new(&table, None, vec![], None, Some(100), None);
let splits = vec![
limit_test_split("a.parquet", 2),
limit_test_split_with_unknown_merged_row_count("b.parquet", 4),
limit_test_split("c.parquet", 3),
];
let pruned = scan.apply_limit_pushdown(splits);
assert_eq!(
split_file_names(&pruned),
vec!["a.parquet", "b.parquet", "c.parquet"]
);
}
/// A non-raw-convertible split (merge-needed PK split) has an unknown
/// merged row count: its physical rows overcount merged versions, so it
/// must not satisfy the limit on its own.
#[test]
fn test_apply_limit_pushdown_treats_merge_splits_as_unknown() {
let table = limit_test_table();
let scan = TableScan::new(&table, None, vec![], None, Some(3), None);
let mut merge_file = test_data_file_meta(Vec::new(), Vec::new(), Vec::new(), 10);
merge_file.file_name = "merge.parquet".to_string();
let merge_split = DataSplitBuilder::new()
.with_snapshot(1)
.with_partition(BinaryRow::new(0))
.with_bucket(0)
.with_bucket_path("file:/tmp/merge.parquet".to_string())
.with_total_buckets(1)
.with_data_files(vec![merge_file])
.with_raw_convertible(false)
.build()
.unwrap();
assert_eq!(merge_split.merged_row_count(), None);
let splits = vec![merge_split, limit_test_split("a.parquet", 3)];
let pruned = scan.apply_limit_pushdown(splits);
// The merge split is skipped; only the counted split commits pruning.
assert_eq!(split_file_names(&pruned), vec!["a.parquet"]);
}
#[test]
fn test_first_row_skips_level_zero_by_default() {
assert!(should_skip_level_zero_for_scan(
false,
true,
false,
Ok(crate::spec::MergeEngine::FirstRow),
));
}
#[test]
fn test_scan_all_files_disables_first_row_level_zero_skip() {
assert!(!should_skip_level_zero_for_scan(
true,
true,
false,
Ok(crate::spec::MergeEngine::FirstRow),
));
}
#[test]
fn test_partition_filter_decode_failure_fails_open() {
let fields = partition_string_field();
let predicate = PredicateBuilder::new(&fields)
.equal("dt", Datum::String("2024-01-01".into()))
.unwrap();
// Range predicate to force Predicate variant (fail-open path)
let filter = PartitionFilter::Predicate(predicate);
assert!(filter.matches_entry(&[0xFF, 0x00]).unwrap());
}
#[test]
fn test_partition_filter_eval_error_fails_fast() {
let mut builder = BinaryRowBuilder::new(1);
builder.write_string(0, "2024-01-01");
let serialized = builder.build_serialized();
let predicate = Predicate::Leaf {
column: "dt".into(),
index: 0,
data_type: DataType::Array(ArrayType::new(DataType::Int(IntType::new()))),
op: PredicateOperator::Eq,
literals: vec![Datum::Int(42)],
};
let filter = PartitionFilter::Predicate(predicate);
let err = filter
.matches_entry(&serialized)
.expect_err("eval_row error should propagate");
assert!(
matches!(&err, Error::Unsupported { message } if message.contains("extract_datum")),
"Expected extract_datum unsupported error, got: {err:?}"
);
}
const TEST_SCHEMA_ID: i64 = 0;
fn test_schema_fields() -> Vec<DataField> {
int_field()
}
#[test]
fn test_group_by_overlapping_row_id_empty() {
let result = group_by_overlapping_row_id(vec![]);
assert!(result.is_empty());
}
#[test]
fn test_group_by_overlapping_row_id_no_row_ids() {
let files = vec![
make_evo_file("a", 10, 100, 1, None),
make_evo_file("b", 10, 100, 2, None),
];
let groups = group_by_overlapping_row_id(files);
assert_eq!(file_names(&groups), vec![vec!["b"], vec!["a"]]);
}
#[test]
fn test_group_by_overlapping_row_id_same_range() {
let files = vec![
make_evo_file("a", 10, 100, 2, Some(0)),
make_evo_file("b", 10, 100, 1, Some(0)),
];
let groups = group_by_overlapping_row_id(files);
assert_eq!(groups.len(), 1);
assert_eq!(file_names(&groups), vec![vec!["a", "b"]]);
}
#[test]
fn test_group_by_overlapping_row_id_overlapping_ranges() {
let files = vec![
make_evo_file("a", 10, 100, 1, Some(0)),
make_evo_file("b", 10, 100, 2, Some(50)),
];
let groups = group_by_overlapping_row_id(files);
assert_eq!(groups.len(), 1);
assert_eq!(file_names(&groups), vec![vec!["a", "b"]]);
}
#[test]
fn test_group_by_overlapping_row_id_non_overlapping() {
let files = vec![
make_evo_file("a", 10, 100, 1, Some(0)),
make_evo_file("b", 10, 100, 2, Some(100)),
];
let groups = group_by_overlapping_row_id(files);
assert_eq!(groups.len(), 2);
assert_eq!(file_names(&groups), vec![vec!["a"], vec!["b"]]);
}
#[test]
fn test_group_by_overlapping_row_id_mixed() {
let files = vec![
make_evo_file("a", 10, 100, 1, Some(0)),
make_evo_file("b", 10, 100, 2, Some(0)),
make_evo_file("c", 10, 100, 3, None),
make_evo_file("d", 10, 100, 4, Some(200)),
];
let groups = group_by_overlapping_row_id(files);
assert_eq!(
file_names(&groups),
vec![vec!["c"], vec!["b", "a"], vec!["d"]]
);
}
#[test]
fn test_group_by_overlapping_row_id_sorted_by_seq() {
let files = vec![
make_evo_file("a", 10, 100, 1, Some(0)),
make_evo_file("b", 10, 100, 3, Some(0)),
make_evo_file("c", 10, 100, 2, Some(0)),
];
let groups = group_by_overlapping_row_id(files);
assert_eq!(groups.len(), 1);
assert_eq!(file_names(&groups), vec![vec!["b", "c", "a"]]);
}
#[test]
fn test_data_file_matches_eq_prunes_out_of_range() {
let fields = int_field();
let file = test_data_file_meta(
int_stats_row(Some(10)),
int_stats_row(Some(20)),
vec![Some(0)],
5,
);
let predicate = PredicateBuilder::new(&fields)
.equal("id", Datum::Int(30))
.unwrap();
assert!(!data_file_matches_predicates(
&file,
&[predicate],
TEST_SCHEMA_ID,
&test_schema_fields(),
));
}
#[test]
fn test_data_file_matches_is_null_prunes_when_null_count_is_zero() {
let fields = int_field();
let file = test_data_file_meta(
int_stats_row(Some(10)),
int_stats_row(Some(20)),
vec![Some(0)],
5,
);
let predicate = PredicateBuilder::new(&fields).is_null("id").unwrap();
assert!(!data_file_matches_predicates(
&file,
&[predicate],
TEST_SCHEMA_ID,
&test_schema_fields(),
));
}
#[test]
fn test_data_file_matches_is_not_null_prunes_all_null_file() {
let fields = int_field();
let file = test_data_file_meta(int_stats_row(None), int_stats_row(None), vec![Some(5)], 5);
let predicate = PredicateBuilder::new(&fields).is_not_null("id").unwrap();
assert!(!data_file_matches_predicates(
&file,
&[predicate],
TEST_SCHEMA_ID,
&test_schema_fields(),
));
}
#[test]
fn test_data_file_matches_unsupported_predicate_fails_open() {
let fields = int_field();
let file = test_data_file_meta(
int_stats_row(Some(10)),
int_stats_row(Some(20)),
vec![Some(0)],
5,
);
let pb = PredicateBuilder::new(&fields);
let predicate = Predicate::or(vec![
pb.less_than("id", Datum::Int(5)).unwrap(),
pb.greater_than("id", Datum::Int(25)).unwrap(),
]);
assert!(data_file_matches_predicates(
&file,
&[predicate],
TEST_SCHEMA_ID,
&test_schema_fields(),
));
}
#[test]
fn test_data_file_matches_corrupt_stats_fails_open() {
let fields = int_field();
let file = test_data_file_meta(Vec::new(), Vec::new(), vec![Some(0)], 5);
let predicate = PredicateBuilder::new(&fields)
.equal("id", Datum::Int(30))
.unwrap();
assert!(data_file_matches_predicates(
&file,
&[predicate],
TEST_SCHEMA_ID,
&test_schema_fields(),
));
}
#[test]
fn test_data_file_matches_schema_mismatch_fails_open() {
let fields = int_field();
let file = test_data_file_meta_with_schema(
int_stats_row(Some(10)),
int_stats_row(Some(20)),
vec![Some(0)],
5,
5,
);
let predicate = PredicateBuilder::new(&fields)
.equal("id", Datum::Int(30))
.unwrap();
assert!(data_file_matches_predicates(
&file,
&[predicate],
TEST_SCHEMA_ID,
&test_schema_fields(),
));
}
#[test]
fn test_data_file_matches_always_false_prunes_despite_schema_mismatch() {
let file = test_data_file_meta_with_schema(
int_stats_row(Some(10)),
int_stats_row(Some(20)),
vec![Some(0)],
5,
99,
);
assert!(!data_file_matches_predicates(
&file,
&[Predicate::AlwaysFalse],
TEST_SCHEMA_ID,
&test_schema_fields(),
));
}
#[test]
fn test_data_file_matches_always_true_keeps_file_despite_schema_mismatch() {
let file = test_data_file_meta_with_schema(
int_stats_row(Some(10)),
int_stats_row(Some(20)),
vec![Some(0)],
5,
99,
);
assert!(data_file_matches_predicates(
&file,
&[Predicate::AlwaysTrue],
TEST_SCHEMA_ID,
&test_schema_fields(),
));
}
#[test]
fn test_build_deletion_files_map_preserves_cardinality() {
let entries = vec![IndexManifestEntry {
version: 1,
kind: FileKind::Add,
partition: vec![1, 2, 3],
bucket: 7,
index_file: IndexFileMeta {
index_type: "DELETION_VECTORS".into(),
file_name: "index-file".into(),
file_size: 128,
row_count: 1,
deletion_vectors_ranges: Some(indexmap::IndexMap::from([(
"data-file.parquet".into(),
DeletionVectorMeta {
offset: 11,
length: 22,
cardinality: Some(33),
},
)])),
global_index_meta: None,
},
}];
let map = super::build_deletion_files_map(&entries, "file:/tmp/table");
let by_bucket = map
.get(&super::PartitionBucket::new(vec![1, 2, 3], 7))
.expect("partition bucket should exist");
let deletion_file = by_bucket
.get("data-file.parquet")
.expect("deletion file should exist");
assert_eq!(
deletion_file,
&DeletionFile::new("file:/tmp/table/index/index-file".into(), 11, 22, Some(33))
);
}
// ======================== Bucket predicate filtering ========================
fn bucket_key_fields() -> Vec<DataField> {
vec![DataField::new(
0,
"id".to_string(),
DataType::Int(IntType::new()),
)]
}
#[test]
fn test_extract_predicate_for_keys_eq() {
let fields = vec![
DataField::new(0, "id".to_string(), DataType::Int(IntType::new())),
DataField::new(
1,
"name".to_string(),
DataType::VarChar(VarCharType::default()),
),
];
let pb = PredicateBuilder::new(&fields);
let filter = Predicate::and(vec![
pb.equal("id", Datum::Int(42)).unwrap(),
pb.equal("name", Datum::String("alice".into())).unwrap(),
]);
let keys = vec!["id".to_string()];
let extracted = extract_predicate_for_keys(&filter, &fields, &keys);
assert!(extracted.is_some());
match extracted.unwrap() {
Predicate::Leaf {
column, index, op, ..
} => {
assert_eq!(column, "id");
assert_eq!(index, 0); // remapped to key index
assert_eq!(op, PredicateOperator::Eq);
}
other => panic!("expected Leaf, got {other:?}"),
}
}
#[test]
fn test_extract_predicate_for_keys_no_match() {
let fields = vec![
DataField::new(0, "id".to_string(), DataType::Int(IntType::new())),
DataField::new(
1,
"name".to_string(),
DataType::VarChar(VarCharType::default()),
),
];
let pb = PredicateBuilder::new(&fields);
let filter = pb.equal("name", Datum::String("alice".into())).unwrap();
let keys = vec!["id".to_string()];
let extracted = extract_predicate_for_keys(&filter, &fields, &keys);
assert!(extracted.is_none());
}
#[test]
fn test_compute_target_buckets_single_eq() {
let fields = bucket_key_fields();
// Build a bucket predicate (already projected to bucket key space, index=0)
let pred = Predicate::Leaf {
column: "id".into(),
index: 0,
data_type: DataType::Int(IntType::new()),
op: PredicateOperator::Eq,
literals: vec![Datum::Int(42)],
};
let buckets = compute_target_buckets(&pred, &fields, BucketFunctionType::Default, 4);
assert!(buckets.is_some());
let buckets = buckets.unwrap();
assert_eq!(buckets.len(), 1);
// The bucket should be deterministic
let bucket = *buckets.iter().next().unwrap();
assert!((0..4).contains(&bucket));
}
#[test]
fn test_compute_target_buckets_in_predicate() {
let fields = bucket_key_fields();
let pred = Predicate::Leaf {
column: "id".into(),
index: 0,
data_type: DataType::Int(IntType::new()),
op: PredicateOperator::In,
literals: vec![Datum::Int(1), Datum::Int(2), Datum::Int(3)],
};
let buckets = compute_target_buckets(&pred, &fields, BucketFunctionType::Default, 4);
assert!(buckets.is_some());
let buckets = buckets.unwrap();
// Should have at most 3 buckets (could be fewer if some hash to the same bucket)
assert!(!buckets.is_empty());
assert!(buckets.len() <= 3);
for &b in &buckets {
assert!((0..4).contains(&b));
}
}
#[test]
fn test_compute_target_buckets_range_returns_none() {
let fields = bucket_key_fields();
let pred = Predicate::Leaf {
column: "id".into(),
index: 0,
data_type: DataType::Int(IntType::new()),
op: PredicateOperator::Gt,
literals: vec![Datum::Int(10)],
};
let buckets = compute_target_buckets(&pred, &fields, BucketFunctionType::Default, 4);
assert!(
buckets.is_none(),
"Range predicates cannot determine target buckets"
);
}
#[test]
fn test_compute_target_buckets_composite_key() {
let fields = vec![
DataField::new(0, "a".to_string(), DataType::Int(IntType::new())),
DataField::new(1, "b".to_string(), DataType::Int(IntType::new())),
];
let pred = Predicate::And(vec![
Predicate::Leaf {
column: "a".into(),
index: 0,
data_type: DataType::Int(IntType::new()),
op: PredicateOperator::Eq,
literals: vec![Datum::Int(1)],
},
Predicate::Leaf {
column: "b".into(),
index: 1,
data_type: DataType::Int(IntType::new()),
op: PredicateOperator::Eq,
literals: vec![Datum::Int(2)],
},
]);
let buckets = compute_target_buckets(&pred, &fields, BucketFunctionType::Default, 8);
assert!(buckets.is_some());
let buckets = buckets.unwrap();
assert_eq!(buckets.len(), 1);
let bucket = *buckets.iter().next().unwrap();
assert!((0..8).contains(&bucket));
}
#[test]
fn test_compute_target_buckets_partial_key_returns_none() {
// Only one of two bucket key fields has an eq predicate
let fields = vec![
DataField::new(0, "a".to_string(), DataType::Int(IntType::new())),
DataField::new(1, "b".to_string(), DataType::Int(IntType::new())),
];
let pred = Predicate::Leaf {
column: "a".into(),
index: 0,
data_type: DataType::Int(IntType::new()),
op: PredicateOperator::Eq,
literals: vec![Datum::Int(1)],
};
let buckets = compute_target_buckets(&pred, &fields, BucketFunctionType::Default, 8);
assert!(
buckets.is_none(),
"Partial bucket key should not determine target buckets"
);
}
#[test]
fn test_compute_target_buckets_string_key() {
let fields = vec![DataField::new(
0,
"name".to_string(),
DataType::VarChar(VarCharType::default()),
)];
let pred = Predicate::Leaf {
column: "name".into(),
index: 0,
data_type: DataType::VarChar(VarCharType::default()),
op: PredicateOperator::Eq,
literals: vec![Datum::String("alice".into())],
};
let buckets = compute_target_buckets(&pred, &fields, BucketFunctionType::Default, 4);
assert!(buckets.is_some());
let buckets = buckets.unwrap();
assert_eq!(buckets.len(), 1);
let bucket = *buckets.iter().next().unwrap();
assert!((0..4).contains(&bucket));
}
#[test]
fn test_compute_target_buckets_mod_function() {
let fields = bucket_key_fields();
let pred = Predicate::Leaf {
column: "id".into(),
index: 0,
data_type: DataType::Int(IntType::new()),
op: PredicateOperator::Eq,
literals: vec![Datum::Int(-3)],
};
let buckets = compute_target_buckets(&pred, &fields, BucketFunctionType::Mod, 5);
assert_eq!(buckets, Some(HashSet::from([2])));
}
#[test]
fn test_compute_target_buckets_hive_function() {
let fields = vec![
DataField::new(0, "id".to_string(), DataType::Int(IntType::new())),
DataField::new(
1,
"name".to_string(),
DataType::VarChar(VarCharType::default()),
),
];
let pred = Predicate::And(vec![
Predicate::Leaf {
column: "id".into(),
index: 0,
data_type: DataType::Int(IntType::new()),
op: PredicateOperator::Eq,
literals: vec![Datum::Int(7)],
},
Predicate::Leaf {
column: "name".into(),
index: 1,
data_type: DataType::VarChar(VarCharType::default()),
op: PredicateOperator::Eq,
literals: vec![Datum::String("hello".into())],
},
]);
let buckets = compute_target_buckets(&pred, &fields, BucketFunctionType::Hive, 8);
assert_eq!(buckets, Some(HashSet::from([3])));
}
#[test]
fn test_compute_target_buckets_is_null() {
let fields = bucket_key_fields();
let pred = Predicate::Leaf {
column: "id".into(),
index: 0,
data_type: DataType::Int(IntType::new()),
op: PredicateOperator::IsNull,
literals: vec![],
};
let buckets = compute_target_buckets(&pred, &fields, BucketFunctionType::Default, 4);
assert!(buckets.is_some(), "IsNull should determine a target bucket");
let buckets = buckets.unwrap();
assert_eq!(buckets.len(), 1);
let bucket = *buckets.iter().next().unwrap();
assert!((0..4).contains(&bucket));
// Verify it matches the expected bucket from a null BinaryRow
let mut builder = BinaryRowBuilder::new(1);
builder.set_null_at(0);
let expected = (builder.build().hash_code() % 4).abs();
assert_eq!(bucket, expected);
}
#[test]
fn test_compute_target_buckets_composite_key_with_null() {
let fields = vec![
DataField::new(0, "a".to_string(), DataType::Int(IntType::new())),
DataField::new(1, "b".to_string(), DataType::Int(IntType::new())),
];
// a = 1 AND b IS NULL
let pred = Predicate::And(vec![
Predicate::Leaf {
column: "a".into(),
index: 0,
data_type: DataType::Int(IntType::new()),
op: PredicateOperator::Eq,
literals: vec![Datum::Int(1)],
},
Predicate::Leaf {
column: "b".into(),
index: 1,
data_type: DataType::Int(IntType::new()),
op: PredicateOperator::IsNull,
literals: vec![],
},
]);
let buckets = compute_target_buckets(&pred, &fields, BucketFunctionType::Default, 8);
assert!(
buckets.is_some(),
"Composite key with IsNull should determine a target bucket"
);
let buckets = buckets.unwrap();
assert_eq!(buckets.len(), 1);
let bucket = *buckets.iter().next().unwrap();
assert!((0..8).contains(&bucket));
}
}