This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-rust.git
The following commit(s) were added to refs/heads/main by this push:
new 3d445f98 perf(read): feed sorted runs into the PK merge LoserTree
(#668)
3d445f98 is described below
commit 3d445f987e6d65263aa1e1dd7f274b865002b1dc
Author: Junrui Lee <[email protected]>
AuthorDate: Tue Aug 4 21:39:22 2026 +0800
perf(read): feed sorted runs into the PK merge LoserTree (#668)
---
crates/paimon/src/arrow/residual.rs | 96 +--
crates/paimon/src/predicate_stats.rs | 30 +-
crates/paimon/src/spec/mod.rs | 1 -
crates/paimon/src/spec/predicate.rs | 18 +-
crates/paimon/src/table/kv_file_reader.rs | 853 ++++++++++++++++-----
.../paimon/src/table/merge_tree_split_generator.rs | 281 ++++++-
crates/paimon/src/table/sort_merge.rs | 2 +
crates/paimon/src/table/table_read.rs | 4 +-
8 files changed, 1004 insertions(+), 281 deletions(-)
diff --git a/crates/paimon/src/arrow/residual.rs
b/crates/paimon/src/arrow/residual.rs
index 56b77050..fe4e9c43 100644
--- a/crates/paimon/src/arrow/residual.rs
+++ b/crates/paimon/src/arrow/residual.rs
@@ -511,8 +511,7 @@ fn evaluate_between_predicate(
};
// Delegate the two bound comparisons to `evaluate_column_predicate` rather
// than calling `arrow_gt_eq`/`arrow_lt_eq` directly, so Between inherits
the
- // type-faithful comparison paths (e.g. signed-byte order for Binary).
Using
- // Arrow's kernels here directly would reintroduce unsigned binary
ordering.
+ // type-faithful comparison paths (e.g. byte ordering for Binary).
let lo_mask = evaluate_column_predicate(array, &low_scalar,
PredicateOperator::GtEq)?;
let hi_mask = evaluate_column_predicate(array, &high_scalar,
PredicateOperator::LtEq)?;
let between = arrow_arith::boolean::and_kleene(&lo_mask, &hi_mask)?;
@@ -703,22 +702,6 @@ fn evaluate_column_predicate(
) -> Result<BooleanArray, ArrowError> {
let scalar = string_scalar_for_column(scalar, column.data_type())?;
- // Binary ordering must match Paimon's Datum::Bytes semantics (Java
signed-byte
- // order, 0xFF < 0x00), which Arrow's unsigned byte comparison does not.
Route
- // ordering ops on Binary/VarBinary columns through the signed comparator.
- // Eq/NotEq are order-independent, so Arrow's kernels are correct for them.
- if matches!(
- column.data_type(),
- arrow_schema::DataType::Binary | arrow_schema::DataType::LargeBinary
- ) && matches!(
- op,
- PredicateOperator::Lt
- | PredicateOperator::LtEq
- | PredicateOperator::Gt
- | PredicateOperator::GtEq
- ) {
- return evaluate_binary_ordering_predicate(column, &scalar, op);
- }
match op {
PredicateOperator::Eq => arrow_eq(column, &scalar),
PredicateOperator::NotEq => arrow_neq(column, &scalar),
@@ -745,55 +728,6 @@ fn evaluate_column_predicate(
}
}
-/// Evaluate an ordering predicate (`Lt`/`LtEq`/`Gt`/`GtEq`) on a Binary column
-/// using Paimon's Java signed-byte order, matching `Datum::Bytes` semantics.
-/// `scalar` is a single-element Binary array (the literal). NULL column rows
-/// produce NULL in the mask (later collapsed to `false` by
`sanitize_filter_mask`).
-fn evaluate_binary_ordering_predicate(
- column: &ArrayRef,
- scalar: &Scalar<ArrayRef>,
- op: PredicateOperator,
-) -> Result<BooleanArray, ArrowError> {
- use arrow_array::cast::AsArray;
- use std::cmp::Ordering;
-
- // The scalar wraps a length-1 Binary array holding the literal bytes.
- let (scalar_array, _) = scalar.get();
- let literal: &[u8] = if let Some(a) = scalar_array.as_binary_opt::<i32>() {
- a.value(0)
- } else if let Some(a) = scalar_array.as_binary_opt::<i64>() {
- a.value(0)
- } else {
- return Err(ArrowError::ComputeError(
- "binary ordering predicate expects a Binary literal".to_string(),
- ));
- };
-
- // Row-wise comparison via the shared signed-byte comparator (single
source of
- // truth with `Datum` ordering).
- let compare = |bytes: &[u8]| -> bool {
- let ord = crate::spec::java_bytes_cmp(bytes, literal);
- match op {
- PredicateOperator::Lt => ord == Ordering::Less,
- PredicateOperator::LtEq => ord != Ordering::Greater,
- PredicateOperator::Gt => ord == Ordering::Greater,
- PredicateOperator::GtEq => ord != Ordering::Less,
- _ => unreachable!("only ordering ops reach here"),
- }
- };
-
- let mask: BooleanArray = if let Some(a) = column.as_binary_opt::<i32>() {
- a.iter().map(|v| v.map(compare)).collect()
- } else if let Some(a) = column.as_binary_opt::<i64>() {
- a.iter().map(|v| v.map(compare)).collect()
- } else {
- return Err(ArrowError::ComputeError(
- "binary ordering predicate expects a Binary column".to_string(),
- ));
- };
- Ok(mask)
-}
-
/// Arrow comparison and pattern kernels reject mismatched string types. The
/// shared scalar built from Paimon's logical Char/VarChar type is Utf8, while
a
/// decoded file column may be Utf8, LargeUtf8, or Utf8View. Promote the scalar
@@ -1641,10 +1575,9 @@ mod tests {
}
#[test]
- fn test_binary_ordering_uses_java_signed_byte_order() {
- // Paimon Datum::Bytes orders by signed byte (0xFF < 0x00). Arrow's
- // unsigned comparison would order 0xFF as the largest. Verify the
- // residual matches Paimon: filter `col > 0x00` must EXCLUDE 0xFF.
+ fn test_binary_ordering_uses_java_unsigned_byte_order() {
+ // Paimon Datum::Bytes and Arrow both order bytes as unsigned values.
+ // Verify the residual keeps 0xFF for `col > 0x00`.
use crate::spec::BinaryType;
let col = DataField::new(
0,
@@ -1659,7 +1592,7 @@ mod tests {
let values: Vec<Option<&[u8]>> = vec![Some(&[0x00]), Some(&[0x01]),
Some(&[0xFF])];
let batch =
RecordBatch::try_new(schema,
vec![Arc::new(BinaryArray::from(values))]).unwrap();
- // col > 0x00 : signed order -> only 0x01 (0xFF is negative, < 0x00).
+ // col > 0x00 : unsigned order -> 0x01 and 0xFF.
let pred = leaf(
0,
DataType::Binary(BinaryType::new(1).unwrap()),
@@ -1677,18 +1610,15 @@ mod tests {
let got: Vec<&[u8]> = (0..out_col.len()).map(|i|
out_col.value(i)).collect();
assert_eq!(
got,
- vec![&[0x01u8][..]],
- "0xFF must be excluded (signed < 0x00)"
+ vec![&[0x01u8][..], &[0xFFu8][..]],
+ "0xFF must be included (unsigned > 0x00)"
);
}
#[test]
- fn test_binary_between_uses_java_signed_byte_order() {
- // Between must inherit the signed-byte order too (regression: it
called
- // Arrow's unsigned gt_eq/lt_eq directly). `b BETWEEN 0xFF AND 0x01`
is,
- // under signed order, the range [-1, 1] -> keeps 0xFF(-1), 0x00(0),
- // 0x01(1) and excludes 0x7F(127). Under Arrow's unsigned order it
would be
- // [255, 1] = empty, so this distinguishes the two.
+ fn test_binary_between_uses_java_unsigned_byte_order() {
+ // `b BETWEEN 0x01 AND 0xFF` uses unsigned byte ordering, keeping all
+ // values at or above 0x01 through 0xFF and excluding 0x00.
use crate::spec::BinaryType;
let col = DataField::new(
0,
@@ -1708,7 +1638,7 @@ mod tests {
0,
DataType::Binary(BinaryType::new(1).unwrap()),
PredicateOperator::Between,
- vec![Datum::Bytes(vec![0xFF]), Datum::Bytes(vec![0x01])],
+ vec![Datum::Bytes(vec![0x01]), Datum::Bytes(vec![0xFF])],
);
let fp = file_predicates(vec![pred], vec![col.clone()]);
let out =
@@ -1721,8 +1651,8 @@ mod tests {
let got: Vec<&[u8]> = (0..out_col.len()).map(|i|
out_col.value(i)).collect();
assert_eq!(
got,
- vec![&[0xFFu8][..], &[0x00u8][..], &[0x01u8][..]],
- "signed range [-1, 1] keeps 0xFF/0x00/0x01, excludes 0x7F"
+ vec![&[0xFFu8][..], &[0x01u8][..], &[0x7Fu8][..]],
+ "unsigned range [0x01, 0xFF] excludes only 0x00"
);
}
diff --git a/crates/paimon/src/predicate_stats.rs
b/crates/paimon/src/predicate_stats.rs
index 846970b2..71c930dc 100644
--- a/crates/paimon/src/predicate_stats.rs
+++ b/crates/paimon/src/predicate_stats.rs
@@ -526,7 +526,7 @@ fn coerce_stats_datum_for_predicate(datum: Datum,
predicate_data_type: &DataType
#[cfg(test)]
mod tests {
use super::*;
- use crate::spec::{IntType, VarCharType};
+ use crate::spec::{BinaryType, IntType, VarCharType};
struct MockStats {
row_count: i64,
@@ -677,6 +677,34 @@ mod tests {
));
}
+ #[test]
+ fn binary_range_pruning_uses_unsigned_byte_order() {
+ let dt = DataType::Binary(BinaryType::new(1).unwrap());
+ let stats = MockStats {
+ row_count: 10,
+ null_count: Some(0),
+ min: Some(Datum::Bytes(vec![0x80])),
+ max: Some(Datum::Bytes(vec![0xFF])),
+ };
+
+ assert!(data_leaf_may_match(
+ 0,
+ &dt,
+ &dt,
+ PredicateOperator::Gt,
+ &[Datum::Bytes(vec![0x00])],
+ &stats,
+ ));
+ assert!(!data_leaf_may_match(
+ 0,
+ &dt,
+ &dt,
+ PredicateOperator::Lt,
+ &[Datum::Bytes(vec![0x00])],
+ &stats,
+ ));
+ }
+
fn int_stats(min: i32, max: i32) -> MockStats {
MockStats {
row_count: 10,
diff --git a/crates/paimon/src/spec/mod.rs b/crates/paimon/src/spec/mod.rs
index f1d3147f..d210beed 100644
--- a/crates/paimon/src/spec/mod.rs
+++ b/crates/paimon/src/spec/mod.rs
@@ -99,7 +99,6 @@ mod predicate;
pub(crate) use predicate::datum_cmp;
pub(crate) use predicate::eval_row;
pub(crate) use predicate::extract_datum;
-pub(crate) use predicate::java_bytes_cmp;
pub(crate) use predicate::like_match;
pub use predicate::{
field_idx_to_partition_idx, Datum, Predicate, PredicateBuilder,
PredicateOperator, Transform,
diff --git a/crates/paimon/src/spec/predicate.rs
b/crates/paimon/src/spec/predicate.rs
index f91573b7..a6d609a6 100644
--- a/crates/paimon/src/spec/predicate.rs
+++ b/crates/paimon/src/spec/predicate.rs
@@ -210,16 +210,10 @@ fn decimal_cmp(ua: i128, sa: u32, ub: i128, sb: u32) ->
Option<Ordering> {
na.partial_cmp(&nb)
}
-/// Match Java `CompareUtils.compare(byte[], byte[])`, which compares signed
-/// bytes lexicographically.
-pub(crate) fn java_bytes_cmp(a: &[u8], b: &[u8]) -> Ordering {
- for (&lhs, &rhs) in a.iter().zip(b.iter()) {
- let cmp = (lhs as i8).cmp(&(rhs as i8));
- if cmp != Ordering::Equal {
- return cmp;
- }
- }
- a.len().cmp(&b.len())
+/// Match Java `CompareUtils.compare(byte[], byte[])`, which compares bytes as
+/// unsigned values lexicographically.
+fn java_bytes_cmp(a: &[u8], b: &[u8]) -> Ordering {
+ a.cmp(b)
}
/// 10^exp as i128. Returns i128::MAX for exponents that would overflow.
@@ -2286,8 +2280,8 @@ mod tests {
}
#[test]
- fn test_datum_partial_ord_bytes_matches_java_signed_byte_order() {
- assert!(Datum::Bytes(vec![0xFF]) < Datum::Bytes(vec![0x00]));
+ fn test_datum_partial_ord_bytes_matches_java_unsigned_byte_order() {
+ assert!(Datum::Bytes(vec![0x00]) < Datum::Bytes(vec![0xFF]));
}
#[test]
diff --git a/crates/paimon/src/table/kv_file_reader.rs
b/crates/paimon/src/table/kv_file_reader.rs
index b1cae88f..a9612d2a 100644
--- a/crates/paimon/src/table/kv_file_reader.rs
+++ b/crates/paimon/src/table/kv_file_reader.rs
@@ -17,9 +17,9 @@
//! Key-value file reader for primary-key tables using sort-merge with
LoserTree.
//!
-//! Each data file in a split is read as a separate sorted stream. The streams
-//! are merged by primary key using a LoserTree, and rows with the same key are
-//! deduplicated by keeping the one with the highest `_SEQUENCE_NUMBER`.
+//! Data files with disjoint key ranges are concatenated into sorted runs. The
+//! runs are merged by primary key using a LoserTree, and rows with the same
key
+//! are deduplicated by keeping the one with the highest `_SEQUENCE_NUMBER`.
//! Non-primary-key predicate conjuncts are enforced by an exact post-merge
//! residual filter; only primary-key conjuncts are pushed below the merge.
//!
@@ -33,9 +33,9 @@ use super::sort_merge::{
use crate::arrow::{build_target_arrow_schema, ParquetReadBudget};
use crate::io::FileIO;
use crate::spec::{
- BigIntType, DataField, DataType as PaimonDataType, MergeEngine,
PartialUpdateConfig, Predicate,
- TinyIntType, SEQUENCE_NUMBER_FIELD_ID, SEQUENCE_NUMBER_FIELD_NAME,
VALUE_KIND_FIELD_ID,
- VALUE_KIND_FIELD_NAME,
+ BigIntType, DataField, DataFileMeta, DataType as PaimonDataType,
MergeEngine,
+ PartialUpdateConfig, Predicate, TinyIntType, SEQUENCE_NUMBER_FIELD_ID,
+ SEQUENCE_NUMBER_FIELD_NAME, VALUE_KIND_FIELD_ID, VALUE_KIND_FIELD_NAME,
};
use crate::table::schema_manager::SchemaManager;
use crate::table::ArrowRecordBatchStream;
@@ -76,8 +76,9 @@ pub(crate) struct KeyValueReadConfig {
pub read_batch_size: usize,
/// Merge files from all supplied splits into one globally key-sorted
stream.
pub merge_splits: bool,
- /// Optional cap on file streams opened by a single sort-merge group.
- pub max_merge_file_streams: Option<usize>,
+ /// Optional cap on sorted-run inputs merged concurrently by one LoserTree.
+ /// This limits merge fan-in, not files: files within a run are opened
serially.
+ pub max_merge_input_streams: Option<usize>,
/// Scan-shared Parquet concurrency and projected-byte budget.
pub parquet_read_budget: Option<Arc<ParquetReadBudget>>,
}
@@ -153,20 +154,108 @@ fn widen_partial_update_sequence_group_fields(
Ok(user_fields)
}
-fn ensure_merge_fan_in_limit(stream_count: usize, limit: Option<usize>) ->
crate::Result<()> {
+fn ensure_merge_input_limit(input_stream_count: usize, limit: Option<usize>)
-> crate::Result<()> {
if let Some(limit) = limit {
- if stream_count <= limit {
+ if input_stream_count <= limit {
return Ok(());
}
return Err(Error::Unsupported {
message: format!(
- "KeyValueFileReader refuses to merge {stream_count} file
streams in one sort-merge group; maximum is {limit}. Compact the table before
reading this highly fragmented group"
+ "KeyValueFileReader refuses to merge {input_stream_count}
overlapping sorted-run input streams in one sort-merge group; maximum is
{limit}. Compact the table before reading this highly fragmented group"
),
});
}
Ok(())
}
+struct MergeRun {
+ files: Vec<MergeFile>,
+}
+
+struct MergeFile {
+ split: Arc<DataSplit>,
+ file: DataFileMeta,
+}
+
+fn plan_merge_groups(
+ split_group: &[Arc<DataSplit>],
+ comparator: Option<&super::merge_tree_split_generator::KeyComparator>,
+ merge_splits: bool,
+) -> Vec<Vec<MergeRun>> {
+ let Some(comparator) = comparator else {
+ let runs = split_group
+ .iter()
+ .flat_map(|split| {
+ let files = split.data_files().to_vec();
+ let split = Arc::clone(split);
+ files.into_iter().map(move |file| MergeRun {
+ files: vec![MergeFile {
+ split: Arc::clone(&split),
+ file,
+ }],
+ })
+ })
+ .collect::<Vec<_>>();
+ return if runs.is_empty() {
+ Vec::new()
+ } else {
+ vec![runs]
+ };
+ };
+
+ if merge_splits {
+ let files = split_group
+ .iter()
+ .flat_map(|split| {
+ let files = split.data_files().to_vec();
+ let split = Arc::clone(split);
+ files.into_iter().map(move |file| MergeFile {
+ split: Arc::clone(&split),
+ file,
+ })
+ })
+ .collect::<Vec<_>>();
+ let runs = super::merge_tree_split_generator::pack_sorted_runs_by(
+ files,
+ comparator,
+ |merge_file| &merge_file.file,
+ )
+ .into_iter()
+ .map(|files| MergeRun { files })
+ .collect::<Vec<_>>();
+ return if runs.is_empty() {
+ Vec::new()
+ } else {
+ vec![runs]
+ };
+ }
+
+ let mut groups = Vec::new();
+ for split in split_group {
+ for section in super::merge_tree_split_generator::interval_partition(
+ split.data_files().to_vec(),
+ comparator,
+ ) {
+ let runs =
super::merge_tree_split_generator::pack_sorted_runs(section, comparator)
+ .into_iter()
+ .map(|files| MergeRun {
+ files: files
+ .into_iter()
+ .map(|file| MergeFile {
+ split: Arc::clone(split),
+ file,
+ })
+ .collect(),
+ })
+ .collect::<Vec<_>>();
+ if !runs.is_empty() {
+ groups.push(runs);
+ }
+ }
+ }
+ groups
+}
+
impl KeyValueFileReader {
pub(crate) fn new(file_io: FileIO, config: KeyValueReadConfig) -> Self {
let pushdown_predicates = retain_primary_key_conjuncts(
@@ -264,7 +353,16 @@ impl KeyValueFileReader {
})
})
.collect::<crate::Result<Vec<_>>>()?;
-
+ let key_comparator = if key_fields.is_empty() {
+ None
+ } else {
+ Some(super::merge_tree_split_generator::KeyComparator::new(
+ key_fields
+ .iter()
+ .map(|field| field.data_type().clone())
+ .collect(),
+ ))
+ };
// User columns = read_type fields + any key fields not already in
read_type
// + any sequence fields not already included.
let read_type_names: std::collections::HashSet<&str> =
@@ -389,14 +487,16 @@ impl KeyValueFileReader {
}
}
- let split_groups: Vec<Vec<DataSplit>> = if self.config.merge_splits {
- vec![data_splits.to_vec()]
+ let merge_splits = self.config.merge_splits;
+ let data_splits = data_splits
+ .iter()
+ .cloned()
+ .map(Arc::new)
+ .collect::<Vec<_>>();
+ let split_groups: Vec<Vec<Arc<DataSplit>>> = if merge_splits {
+ vec![data_splits]
} else {
- data_splits
- .iter()
- .cloned()
- .map(|split| vec![split])
- .collect()
+ data_splits.into_iter().map(|split| vec![split]).collect()
};
let file_io = self.file_io;
let merge_engine = self.config.merge_engine;
@@ -410,7 +510,7 @@ impl KeyValueFileReader {
let primary_keys = self.config.primary_keys;
let sequence_fields = self.config.sequence_fields;
let read_batch_size = self.config.read_batch_size;
- let max_merge_file_streams = self.config.max_merge_file_streams;
+ let max_merge_input_streams = self.config.max_merge_input_streams;
let parquet_read_budget = self.config.parquet_read_budget;
#[cfg(test)]
let input_batch_sizes = self.input_batch_sizes;
@@ -434,147 +534,148 @@ impl KeyValueFileReader {
})?;
}
}
- let file_count = split_group
- .iter()
- .map(|split| split.data_files().len())
- .sum::<usize>();
- if file_count == 0 {
- continue;
- }
- ensure_merge_fan_in_limit(file_count, max_merge_file_streams)?;
- // Sort-merge must first obtain one batch from every input
stream.
- // A concurrent Parquet reader keeps its row-group permits
until
- // the complete row group has been consumed, so enabling it on
- // several lockstep inputs can let the first file occupy the
- // entire scan budget while the merge waits for the second
file.
- // Keep multi-file merge inputs on the sequential Parquet path.
- let group_parquet_read_budget = if file_count == 1 {
- parquet_read_budget.clone()
- } else {
- None
- };
- // Create one stream per data file.
- let mut file_streams: Vec<ArrowRecordBatchStream> = Vec::new();
-
- for split in split_group {
- for file_meta in split.data_files().to_vec() {
- let data_fields: Option<Vec<DataField>> = if
file_meta.schema_id != table_schema_id {
- let data_schema =
schema_manager.schema(file_meta.schema_id).await?;
- Some(data_schema.fields().to_vec())
+ for merge_group in plan_merge_groups(
+ split_group,
+ key_comparator.as_ref(),
+ merge_splits,
+ ) {
+ let input_stream_count = merge_group.len();
+ ensure_merge_input_limit(input_stream_count,
max_merge_input_streams)?;
+ // Sort-merge must first obtain one batch from every input
+ // stream. Keep concurrent row-group reads disabled
whenever
+ // multiple runs advance in lockstep; one run may still use
+ // the shared budget because its files are opened serially.
+ let group_parquet_read_budget = if input_stream_count == 1
{
+ parquet_read_budget.clone()
} else {
None
};
-
- let reader = DataFileReader::new(
- file_io.clone(),
- schema_manager.clone(),
- table_schema_id,
- table_fields.clone(),
- internal_read_type.clone(),
- pushdown_predicates.clone(),
- )
- .with_batch_size(Some(read_batch_size))
-
.with_parquet_read_budget(group_parquet_read_budget.clone());
-
- let stream = reader.read_single_file_stream(
- split,
- file_meta,
- data_fields,
- None,
- split.row_ranges().map(|ranges| ranges.to_vec()),
- )?;
- #[cfg(test)]
- let stream = if let Some(batch_sizes) =
input_batch_sizes.clone() {
- stream
- .inspect(move |batch| {
- if let Ok(batch) = batch {
-
batch_sizes.lock().unwrap().push(batch.num_rows());
+ let mut file_streams: Vec<ArrowRecordBatchStream> =
Vec::new();
+
+ for MergeRun { files } in merge_group {
+ let reader = DataFileReader::new(
+ file_io.clone(),
+ schema_manager.clone(),
+ table_schema_id,
+ table_fields.clone(),
+ internal_read_type.clone(),
+ pushdown_predicates.clone(),
+ )
+ .with_batch_size(Some(read_batch_size))
+
.with_parquet_read_budget(group_parquet_read_budget.clone());
+ let run_schema_manager = schema_manager.clone();
+ let run_stream: ArrowRecordBatchStream =
Box::pin(try_stream! {
+ for MergeFile { split, file: file_meta } in files {
+ let data_fields: Option<Vec<DataField>> =
+ if file_meta.schema_id != table_schema_id {
+ let data_schema =
+
run_schema_manager.schema(file_meta.schema_id).await?;
+ Some(data_schema.fields().to_vec())
+ } else {
+ None
+ };
+ let mut file_stream =
reader.read_single_file_stream(
+ split.as_ref(),
+ file_meta,
+ data_fields,
+ None,
+ split.row_ranges().map(|ranges|
ranges.to_vec()),
+ )?;
+ while let Some(batch) =
file_stream.next().await {
+ yield batch?;
}
- })
- .boxed()
- } else {
- stream
- };
- file_streams.push(stream);
+ }
+ });
+ #[cfg(test)]
+ let run_stream = if let Some(batch_sizes) =
input_batch_sizes.clone() {
+ run_stream
+ .inspect(move |batch| {
+ if let Ok(batch) = batch {
+
batch_sizes.lock().unwrap().push(batch.num_rows());
+ }
+ })
+ .boxed()
+ } else {
+ run_stream
+ };
+ file_streams.push(run_stream);
}
- }
- if file_streams.is_empty() {
- continue;
- }
-
- // Always go through sort-merge even for a single file: files
- // written before the writer merged key groups at flush may
- // still contain duplicate keys.
- let mut merge_stream = SortMergeReaderBuilder::new(
- file_streams,
- internal_schema.clone(),
- key_indices.clone(),
- seq_index,
- value_kind_index,
- user_sequence_indices.clone(),
- value_indices.clone(),
- merge_output_schema.clone(),
- Self::new_merge_function(
- merge_engine,
- &table_options,
- &table_name,
- &table_fields,
- &merge_output_fields,
- &primary_keys,
- &sequence_fields,
- )?,
- )
- .build()?;
-
- while let Some(batch) = merge_stream.next().await {
- let batch = batch?;
- // The post-merge residual enforces the FULL data predicate
- // on merged rows. PK
- // conjuncts are also in this set (they were already pushed
- // down pre-merge); re-evaluating them on already-matching
- // rows is a no-op and keeps one shared evaluator instead
of
- // deriving a non-PK subset. Runs on the merge-output batch
- // (keys + values, including widened predicate columns);
the
- // reorder below projects the output back to read_type.
- let batch = if residual_predicates.is_empty() {
- batch
- } else {
- match crate::arrow::residual::evaluate_predicates_mask(
- &batch,
- &residual_predicates,
+ // Always go through sort-merge even for a single file:
files
+ // written before the writer merged key groups at flush may
+ // still contain duplicate keys.
+ let mut merge_stream = SortMergeReaderBuilder::new(
+ file_streams,
+ internal_schema.clone(),
+ key_indices.clone(),
+ seq_index,
+ value_kind_index,
+ user_sequence_indices.clone(),
+ value_indices.clone(),
+ merge_output_schema.clone(),
+ Self::new_merge_function(
+ merge_engine,
+ &table_options,
+ &table_name,
&table_fields,
&merge_output_fields,
- )? {
- Some(mask) => {
-
arrow_select::filter::filter_record_batch(&batch, &mask).map_err(
- |e| Error::DataInvalid {
- message: format!(
- "Failed to filter merged batch by
predicates: {e}"
- ),
- source: Some(Box::new(e)),
- },
- )?
+ &primary_keys,
+ &sequence_fields,
+ )?,
+ )
+ .build()?;
+
+ while let Some(batch) = merge_stream.next().await {
+ let batch = batch?;
+ // The post-merge residual enforces the FULL data
predicate
+ // on merged rows. PK conjuncts are also in this set
(they
+ // were already pushed down pre-merge); re-evaluating
them
+ // on already-matching rows is a no-op and keeps one
shared
+ // evaluator instead of deriving a non-PK subset. Runs
on
+ // the merge-output batch (keys + values, including
widened
+ // predicate columns); the reorder below projects the
output
+ // back to read_type.
+ let batch = if residual_predicates.is_empty() {
+ batch
+ } else {
+ match
crate::arrow::residual::evaluate_predicates_mask(
+ &batch,
+ &residual_predicates,
+ &table_fields,
+ &merge_output_fields,
+ )? {
+ Some(mask) =>
arrow_select::filter::filter_record_batch(
+ &batch, &mask,
+ )
+ .map_err(|e| Error::DataInvalid {
+ message: format!(
+ "Failed to filter merged batch by
predicates: {e}"
+ ),
+ source: Some(Box::new(e)),
+ })?,
+ None => batch,
}
- None => batch,
- }
- };
- // Reorder columns from [keys..., values...] to read_type
order.
- let columns: Vec<_> = reorder_map
- .iter()
- .map(|&src| batch.column(src).clone())
- .collect();
- // An explicit row count keeps empty projections working
- // (e.g. COUNT(*) reads no columns).
- let options =
-
RecordBatchOptions::new().with_row_count(Some(batch.num_rows()));
- let reordered =
-
RecordBatch::try_new_with_options(output_schema.clone(), columns, &options)
- .map_err(|e| Error::UnexpectedError {
- message: format!("Failed to reorder merged
RecordBatch: {e}"),
- source: Some(Box::new(e)),
- })?;
- yield reordered;
+ };
+ // Reorder columns from [keys..., values...] to
read_type order.
+ let columns: Vec<_> = reorder_map
+ .iter()
+ .map(|&src| batch.column(src).clone())
+ .collect();
+ // An explicit row count keeps empty projections
working
+ // (e.g. COUNT(*) reads no columns).
+ let options =
+
RecordBatchOptions::new().with_row_count(Some(batch.num_rows()));
+ let reordered = RecordBatch::try_new_with_options(
+ output_schema.clone(),
+ columns,
+ &options,
+ )
+ .map_err(|e| Error::UnexpectedError {
+ message: format!("Failed to reorder merged
RecordBatch: {e}"),
+ source: Some(Box::new(e)),
+ })?;
+ yield reordered;
+ }
}
}
}
@@ -768,10 +869,17 @@ mod tests {
}
}
+ fn int_key(value: i32) -> Vec<u8> {
+ let mut builder = crate::spec::BinaryRowBuilder::new(1);
+ builder.write_int(0, value);
+ builder.build_serialized()
+ }
+
async fn write_multi_row_group_kv_file(
file_io: &FileIO,
table_path: &str,
file_name: &str,
+ start_id: i32,
sequence: i64,
value: i32,
) -> DataFileMeta {
@@ -795,7 +903,7 @@ mod tests {
vec![
Arc::new(Int64Array::from_value(sequence, 128)),
Arc::new(Int8Array::from_value(0, 128)),
- Arc::new(Int32Array::from_iter_values(0..128)),
+ Arc::new(Int32Array::from_iter_values(start_id..start_id +
128)),
Arc::new(Int32Array::from_value(value, 128)),
],
)
@@ -828,9 +936,34 @@ mod tests {
let mut file = dummy_data_file(file_name.to_string());
file.file_size = parquet_bytes.len() as i64;
file.row_count = 128;
+ file.min_key = int_key(start_id);
+ file.max_key = int_key(start_id + 127);
file
}
+ fn kv_reader_with_budget(table: &Table, budget: Arc<ParquetReadBudget>) ->
KeyValueFileReader {
+ let core_options = table.schema().core_options();
+ KeyValueFileReader::new(
+ table.file_io().clone(),
+ KeyValueReadConfig {
+ table_name: table.identifier().full_name(),
+ table_options: table.schema().options().clone(),
+ schema_manager: table.schema_manager().clone(),
+ table_schema_id: table.schema().id(),
+ table_fields: table.schema().fields().to_vec(),
+ read_type: table.schema().fields().to_vec(),
+ predicates: Vec::new(),
+ primary_keys: table.schema().trimmed_primary_keys(),
+ merge_engine: core_options.merge_engine().unwrap(),
+ sequence_fields: Vec::new(),
+ read_batch_size: core_options.read_batch_size().unwrap(),
+ merge_splits: true,
+ max_merge_input_streams: None,
+ parquet_read_budget: Some(budget),
+ },
+ )
+ }
+
#[test]
fn retain_primary_key_conjuncts_semantics() {
let fields = vec![
@@ -911,7 +1044,7 @@ mod tests {
}
#[tokio::test]
- async fn kv_merge_rejects_too_many_file_streams_on_read_path() {
+ async fn kv_merge_rejects_too_many_sorted_runs_on_read_path() {
let file_io = test_file_io();
let table_path = "memory:/kv_merge_fan_in_limit";
let table = pk_table(&file_io, table_path, &[]);
@@ -944,7 +1077,7 @@ mod tests {
sequence_fields: Vec::new(),
read_batch_size: core_options.read_batch_size().unwrap(),
merge_splits: true,
- max_merge_file_streams: Some(256),
+ max_merge_input_streams: Some(256),
parquet_read_budget: None,
},
);
@@ -956,11 +1089,113 @@ mod tests {
.await
.unwrap_err();
assert!(
- matches!(err, Error::Unsupported { message } if
message.contains("file streams")),
- "KV merge must fail before opening an unbounded number of file
streams"
+ matches!(err, Error::Unsupported { message } if
message.contains("sorted-run input streams")),
+ "KV merge must fail before opening an unbounded number of
sorted-run inputs"
);
}
+ #[test]
+ fn sorted_run_planning_limits_each_section_to_overlap_depth() {
+ let file = |name: &str, min: i32, max: i32| {
+ let mut file = dummy_data_file(name.to_string());
+ file.min_key = int_key(min);
+ file.max_key = int_key(max);
+ file
+ };
+ let split = Arc::new(
+ DataSplitBuilder::new()
+ .with_snapshot(1)
+ .with_partition(BinaryRow::new(0))
+ .with_bucket(0)
+
.with_bucket_path("memory:/sorted-run-plan/bucket-0".to_string())
+ .with_total_buckets(1)
+ .with_data_files(vec![
+ file("a", 1, 10),
+ file("b", 5, 15),
+ file("c", 20, 30),
+ file("d", 25, 35),
+ file("e", 40, 50),
+ file("f", 45, 55),
+ ])
+ .build()
+ .unwrap(),
+ );
+ let comparator =
+
super::super::merge_tree_split_generator::KeyComparator::new(vec![DataType::Int(
+ IntType::new(),
+ )]);
+
+ let grouped = plan_merge_groups(std::slice::from_ref(&split),
Some(&comparator), false);
+ assert_eq!(grouped.len(), 3);
+ assert!(grouped.iter().all(|section| section.len() == 2));
+
+ let fallback = plan_merge_groups(std::slice::from_ref(&split), None,
false);
+ assert_eq!(fallback.len(), 1);
+ assert_eq!(fallback[0].len(), 6);
+ }
+
+ #[test]
+ fn sorted_run_planning_merges_disjoint_sections_across_splits() {
+ let file = |name: String, key: i32| {
+ let mut file = dummy_data_file(name);
+ file.min_key = int_key(key);
+ file.max_key = int_key(key);
+ file
+ };
+ let split = |path: &str, files| {
+ Arc::new(
+ DataSplitBuilder::new()
+ .with_snapshot(1)
+ .with_partition(BinaryRow::new(0))
+ .with_bucket(0)
+ .with_bucket_path(path.to_string())
+ .with_total_buckets(1)
+ .with_data_files(files)
+ .build()
+ .unwrap(),
+ )
+ };
+ let first = split(
+ "memory:/sorted-run-plan/first",
+ (0..129)
+ .map(|index| file(format!("first-{index}"), index * 4))
+ .collect(),
+ );
+ let second = split(
+ "memory:/sorted-run-plan/second",
+ (0..128)
+ .map(|index| file(format!("second-{index}"), index * 4 + 2))
+ .collect(),
+ );
+ let comparator =
+
super::super::merge_tree_split_generator::KeyComparator::new(vec![DataType::Int(
+ IntType::new(),
+ )]);
+
+ let grouped = plan_merge_groups(&[first, second], Some(&comparator),
true);
+ assert_eq!(grouped.len(), 1);
+ assert_eq!(grouped[0].len(), 1, "global overlap depth is one");
+ ensure_merge_input_limit(grouped[0].len(), Some(256)).unwrap();
+ let files = &grouped[0][0].files;
+ assert_eq!(files.len(), 257);
+ assert_eq!(
+ files
+ .iter()
+ .take(4)
+ .map(|file| file.file.file_name.as_str())
+ .collect::<Vec<_>>(),
+ vec!["first-0", "second-0", "first-1", "second-1"]
+ );
+ for file in files {
+ let expected_path = if file.file.file_name.starts_with("first-") {
+ "memory:/sorted-run-plan/first"
+ } else {
+ "memory:/sorted-run-plan/second"
+ };
+ assert_eq!(file.split.bucket_path(), expected_path);
+ }
+ }
+
#[tokio::test]
async fn
kv_input_decode_honors_read_batch_size_without_changing_merge_batching() {
let file_io = test_file_io();
@@ -1008,7 +1243,7 @@ mod tests {
.collect(),
read_batch_size: core_options.read_batch_size().unwrap(),
merge_splits: false,
- max_merge_file_streams: None,
+ max_merge_input_streams: None,
parquet_read_budget: None,
},
)
@@ -1029,7 +1264,7 @@ mod tests {
}
#[tokio::test]
- async fn kv_merge_with_shared_budget_does_not_deadlock_between_files() {
+ async fn kv_merge_with_multiple_runs_does_not_deadlock() {
let file_io = test_file_io();
let table_path = "memory:/kv_shared_parquet_budget";
setup_dirs(&file_io, table_path).await;
@@ -1042,18 +1277,27 @@ mod tests {
],
);
let first =
- write_multi_row_group_kv_file(&file_io, table_path,
"first.parquet", 0, 10).await;
+ write_multi_row_group_kv_file(&file_io, table_path,
"first.parquet", 0, 0, 10).await;
let second =
- write_multi_row_group_kv_file(&file_io, table_path,
"second.parquet", 1, 11).await;
- let split = DataSplitBuilder::new()
- .with_snapshot(1)
- .with_partition(BinaryRow::new(0))
- .with_bucket(0)
- .with_bucket_path(format!("{table_path}/bucket-0"))
- .with_total_buckets(1)
- .with_data_files(vec![first, second])
- .build()
- .unwrap();
+ write_multi_row_group_kv_file(&file_io, table_path,
"second.parquet", 0, 1, 11).await;
+ let split = Arc::new(
+ DataSplitBuilder::new()
+ .with_snapshot(1)
+ .with_partition(BinaryRow::new(0))
+ .with_bucket(0)
+ .with_bucket_path(format!("{table_path}/bucket-0"))
+ .with_total_buckets(1)
+ .with_data_files(vec![first, second])
+ .build()
+ .unwrap(),
+ );
+ let comparator =
+
super::super::merge_tree_split_generator::KeyComparator::new(vec![DataType::Int(
+ IntType::new(),
+ )]);
+ let planned = plan_merge_groups(std::slice::from_ref(&split),
Some(&comparator), false);
+ assert_eq!(planned.len(), 1);
+ assert_eq!(planned[0].len(), 2);
let core_options = table.schema().core_options();
let reader = KeyValueFileReader::new(
table.file_io().clone(),
@@ -1070,16 +1314,19 @@ mod tests {
sequence_fields: Vec::new(),
read_batch_size: core_options.read_batch_size().unwrap(),
merge_splits: false,
- max_merge_file_streams: None,
+ max_merge_input_streams: None,
parquet_read_budget: Some(Arc::new(ParquetReadBudget::new(2,
256 << 20).unwrap())),
},
);
let batches = tokio::time::timeout(
std::time::Duration::from_secs(5),
- reader.read(&[split]).unwrap().try_collect::<Vec<_>>(),
+ reader
+ .read(std::slice::from_ref(split.as_ref()))
+ .unwrap()
+ .try_collect::<Vec<_>>(),
)
.await
- .expect("multi-file sort-merge must not wait forever for a shared
Parquet permit")
+ .expect("multiple sorted-run inputs must not deadlock on shared
Parquet permits")
.unwrap();
assert_eq!(
@@ -1088,6 +1335,266 @@ mod tests {
);
}
+ #[tokio::test]
+ async fn single_sorted_run_uses_shared_budget_across_files() {
+ let file_io = test_file_io();
+ let table_path = "memory:/kv_single_run_shared_budget";
+ setup_dirs(&file_io, table_path).await;
+ let table = pk_table(
+ &file_io,
+ table_path,
+ &[
+ ("read.batch-size", "1"),
+ ("read.parquet.row-group.parallelism", "1"),
+ ],
+ );
+ let low =
+ write_multi_row_group_kv_file(&file_io, table_path, "low.parquet",
0, 0, 10).await;
+ let high =
+ write_multi_row_group_kv_file(&file_io, table_path,
"high.parquet", 200, 0, 20).await;
+ let split = Arc::new(
+ DataSplitBuilder::new()
+ .with_snapshot(1)
+ .with_partition(BinaryRow::new(0))
+ .with_bucket(0)
+ .with_bucket_path(format!("{table_path}/bucket-0"))
+ .with_total_buckets(1)
+ .with_data_files(vec![high, low])
+ .build()
+ .unwrap(),
+ );
+ let comparator =
+
super::super::merge_tree_split_generator::KeyComparator::new(vec![DataType::Int(
+ IntType::new(),
+ )]);
+ let planned = plan_merge_groups(std::slice::from_ref(&split),
Some(&comparator), true);
+ assert_eq!(planned.len(), 1);
+ assert_eq!(planned[0].len(), 1);
+ assert_eq!(planned[0][0].files.len(), 2);
+
+ let reader = kv_reader_with_budget(
+ &table,
+ Arc::new(ParquetReadBudget::new(1, 256 << 20).unwrap()),
+ );
+ let batches = tokio::time::timeout(
+ std::time::Duration::from_secs(5),
+ reader
+ .read(std::slice::from_ref(split.as_ref()))
+ .unwrap()
+ .try_collect::<Vec<_>>(),
+ )
+ .await
+ .expect("one sorted run must reuse a single shared Parquet permit
across files")
+ .unwrap();
+
+ assert_eq!(
+ batches.iter().map(RecordBatch::num_rows).sum::<usize>(),
+ 256
+ );
+ }
+
+ #[tokio::test]
+ async fn concurrent_single_runs_share_one_parquet_budget() {
+ let file_io = test_file_io();
+ let table_path = "memory:/kv_concurrent_single_run_budget";
+ setup_dirs(&file_io, table_path).await;
+ let table = pk_table(
+ &file_io,
+ table_path,
+ &[
+ ("read.batch-size", "1"),
+ ("read.parquet.row-group.parallelism", "1"),
+ ],
+ );
+ let first_low =
+ write_multi_row_group_kv_file(&file_io, table_path,
"first-low.parquet", 0, 0, 10)
+ .await;
+ let first_high =
+ write_multi_row_group_kv_file(&file_io, table_path,
"first-high.parquet", 200, 0, 20)
+ .await;
+ let second_low =
+ write_multi_row_group_kv_file(&file_io, table_path,
"second-low.parquet", 400, 0, 30)
+ .await;
+ let second_high =
+ write_multi_row_group_kv_file(&file_io, table_path,
"second-high.parquet", 600, 0, 40)
+ .await;
+ let split = |files| {
+ DataSplitBuilder::new()
+ .with_snapshot(1)
+ .with_partition(BinaryRow::new(0))
+ .with_bucket(0)
+ .with_bucket_path(format!("{table_path}/bucket-0"))
+ .with_total_buckets(1)
+ .with_data_files(files)
+ .build()
+ .unwrap()
+ };
+ let first_split = split(vec![first_high, first_low]);
+ let second_split = split(vec![second_high, second_low]);
+ let budget = Arc::new(ParquetReadBudget::new(1, 256 << 20).unwrap());
+ let first_reader = kv_reader_with_budget(&table, budget.clone());
+ let second_reader = kv_reader_with_budget(&table, budget);
+
+ let (first_batches, second_batches) =
+ tokio::time::timeout(std::time::Duration::from_secs(5), async {
+ tokio::try_join!(
+ first_reader
+ .read(&[first_split])
+ .unwrap()
+ .try_collect::<Vec<_>>(),
+ second_reader
+ .read(&[second_split])
+ .unwrap()
+ .try_collect::<Vec<_>>()
+ )
+ })
+ .await
+ .expect("concurrent single-run readers must make progress with one
shared permit")
+ .unwrap();
+
+ assert_eq!(
+ first_batches
+ .iter()
+ .map(RecordBatch::num_rows)
+ .sum::<usize>(),
+ 256
+ );
+ assert_eq!(
+ second_batches
+ .iter()
+ .map(RecordBatch::num_rows)
+ .sum::<usize>(),
+ 256
+ );
+ }
+
+ #[tokio::test]
+ async fn sorted_run_read_matches_per_file_fan_out() {
+ let file_io = test_file_io();
+ let table_path = "memory:/kv_sorted_run";
+ let table = pk_table(&file_io, table_path, &[]);
+ let low =
+ write_multi_row_group_kv_file(&file_io, table_path, "low.parquet",
0, 0, 10).await;
+ let high =
+ write_multi_row_group_kv_file(&file_io, table_path,
"high.parquet", 200, 0, 20).await;
+ let split = |files| {
+ DataSplitBuilder::new()
+ .with_snapshot(1)
+ .with_partition(BinaryRow::new(0))
+ .with_bucket(0)
+ .with_bucket_path(format!("{table_path}/bucket-0"))
+ .with_total_buckets(1)
+ .with_data_files(files)
+ .build()
+ .unwrap()
+ };
+ let grouped_split = split(vec![high.clone(), low.clone()]);
+ let per_file_splits = vec![split(vec![high]), split(vec![low])];
+ let core_options = table.schema().core_options();
+
+ let read = |splits: &[DataSplit], merge_splits| {
+ KeyValueFileReader::new(
+ table.file_io().clone(),
+ KeyValueReadConfig {
+ table_name: table.identifier().full_name(),
+ table_options: table.schema().options().clone(),
+ schema_manager: table.schema_manager().clone(),
+ table_schema_id: table.schema().id(),
+ table_fields: table.schema().fields().to_vec(),
+ read_type: table.schema().fields().to_vec(),
+ predicates: Vec::new(),
+ primary_keys: table.schema().trimmed_primary_keys(),
+ merge_engine: core_options.merge_engine().unwrap(),
+ sequence_fields: Vec::new(),
+ read_batch_size: core_options.read_batch_size().unwrap(),
+ merge_splits,
+ max_merge_input_streams: None,
+ parquet_read_budget: None,
+ },
+ )
+ .read(splits)
+ .unwrap()
+ };
+
+ let grouped = read(std::slice::from_ref(&grouped_split), false)
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap();
+ let per_file = read(&per_file_splits, true)
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap();
+
+ let expected = (0..128).chain(200..328).collect::<Vec<_>>();
+ assert_eq!(int_column(&grouped, "id"), expected);
+ assert_eq!(int_column(&grouped, "id"), int_column(&per_file, "id"));
+ }
+
+ #[tokio::test]
+ async fn sorted_runs_preserve_global_merge_across_splits() {
+ let file_io = test_file_io();
+ let table_path = "memory:/kv_sorted_run_merge_splits";
+ let table = pk_table(&file_io, table_path, &[]);
+ let low =
+ write_multi_row_group_kv_file(&file_io, table_path, "low.parquet",
0, 0, 10).await;
+ let high =
+ write_multi_row_group_kv_file(&file_io, table_path,
"high.parquet", 300, 0, 20).await;
+ let middle =
+ write_multi_row_group_kv_file(&file_io, table_path,
"middle.parquet", 100, 1, 30).await;
+ let split_with_run = DataSplitBuilder::new()
+ .with_snapshot(1)
+ .with_partition(BinaryRow::new(0))
+ .with_bucket(0)
+ .with_bucket_path(format!("{table_path}/bucket-0"))
+ .with_total_buckets(1)
+ .with_data_files(vec![high, low])
+ .build()
+ .unwrap();
+ let overlapping_split = DataSplitBuilder::new()
+ .with_snapshot(1)
+ .with_partition(BinaryRow::new(0))
+ .with_bucket(0)
+ .with_bucket_path(format!("{table_path}/bucket-0"))
+ .with_total_buckets(1)
+ .with_data_files(vec![middle])
+ .build()
+ .unwrap();
+ let core_options = table.schema().core_options();
+ let reader = KeyValueFileReader::new(
+ table.file_io().clone(),
+ KeyValueReadConfig {
+ table_name: table.identifier().full_name(),
+ table_options: table.schema().options().clone(),
+ schema_manager: table.schema_manager().clone(),
+ table_schema_id: table.schema().id(),
+ table_fields: table.schema().fields().to_vec(),
+ read_type: table.schema().fields().to_vec(),
+ predicates: Vec::new(),
+ primary_keys: table.schema().trimmed_primary_keys(),
+ merge_engine: core_options.merge_engine().unwrap(),
+ sequence_fields: Vec::new(),
+ read_batch_size: core_options.read_batch_size().unwrap(),
+ merge_splits: true,
+ max_merge_input_streams: Some(256),
+ parquet_read_budget: None,
+ },
+ );
+ let batches = reader
+ .read(&[split_with_run, overlapping_split])
+ .unwrap()
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap();
+
+ let expected_ids = (0..228).chain(300..428).collect::<Vec<_>>();
+ let expected_values = std::iter::repeat_n(10, 100)
+ .chain(std::iter::repeat_n(30, 128))
+ .chain(std::iter::repeat_n(20, 128))
+ .collect::<Vec<_>>();
+ assert_eq!(int_column(&batches, "id"), expected_ids);
+ assert_eq!(int_column(&batches, "value"), expected_values);
+ }
+
/// Non-PK equality filter on a dedup PK table read through the sort-merge
/// path must return only matching rows. Before the post-merge residual,
/// the non-PK conjunct was silently dropped and all rows came back.
diff --git a/crates/paimon/src/table/merge_tree_split_generator.rs
b/crates/paimon/src/table/merge_tree_split_generator.rs
index 276aea54..cceee451 100644
--- a/crates/paimon/src/table/merge_tree_split_generator.rs
+++ b/crates/paimon/src/table/merge_tree_split_generator.rs
@@ -94,6 +94,8 @@ impl KeyComparator {
/// Compare decoded keys field-by-field. NULL sorts first; fields that
/// `datum_cmp` cannot order (e.g. float NaN) compare as equal, which forces
/// the files into the same section — conservative but never incorrect.
+/// Binary keys use unsigned lexicographic order, matching the generated Java
+/// key comparator and the on-disk row order.
fn compare_decoded(a: &DecodedKey, b: &DecodedKey) -> Ordering {
for (fa, fb) in a.iter().zip(b.iter()) {
let ord = match (fa, fb) {
@@ -116,8 +118,9 @@ struct KeyedFile {
max: DecodedKey,
}
-/// Decode every file's key range up front. Returns `None` if any file lacks
-/// a usable key range, in which case callers must assume full overlap.
+/// Decode every file's key range up front. Returns the original files as `Err`
+/// if any range is missing, undecodable, or inverted, in which case callers
+/// must assume full overlap.
fn decode_all(
files: Vec<DataFileMeta>,
comparator: &KeyComparator,
@@ -129,11 +132,15 @@ fn decode_all(
comparator.decode(&file.min_key),
comparator.decode(&file.max_key),
) {
- (Some(min), Some(max)) if !undecodable => keyed.push(KeyedFile {
- file: file.clone(),
- min,
- max,
- }),
+ (Some(min), Some(max))
+ if !undecodable && compare_decoded(&min, &max) !=
Ordering::Greater =>
+ {
+ keyed.push(KeyedFile {
+ file: file.clone(),
+ min,
+ max,
+ })
+ }
_ => undecodable = true,
}
}
@@ -151,8 +158,8 @@ fn decode_all(
/// bound starts a new section. Sections never overlap each other, while files
/// inside one section all transitively overlap and must be merged together.
///
-/// Files with empty or undecodable key ranges collapse everything into one
-/// section: no parallelism, but never a missed merge.
+/// Files with empty, undecodable, or inverted key ranges collapse everything
+/// into one section: no parallelism, but never a missed merge.
pub(crate) fn interval_partition(
files: Vec<DataFileMeta>,
comparator: &KeyComparator,
@@ -196,6 +203,112 @@ pub(crate) fn interval_partition(
sections
}
+/// Pack files into sorted runs. Files within a run have strictly disjoint key
+/// ranges and can therefore be read by concatenation. The number of runs
equals
+/// the maximum key-range overlap depth, even when the input spans multiple
+/// non-overlapping sections. Undecodable or inconsistent manifest keys safely
+/// degrade to one run per file, preserving the previous merge fan-in and
+/// correctness.
+pub(crate) fn pack_sorted_runs(
+ files: Vec<DataFileMeta>,
+ comparator: &KeyComparator,
+) -> Vec<Vec<DataFileMeta>> {
+ pack_sorted_runs_by(files, comparator, |file| file)
+}
+
+/// Pack arbitrary payloads into key-sorted runs using `file_meta` to select
the
+/// [`DataFileMeta`] that defines each payload's key range.
+///
+/// Files are appended to a run only when the previous maximum key is strictly
+/// less than the next minimum key, so concatenating that run remains
monotonic.
+/// Missing, undecodable, or inverted ranges degrade to one item per run. A
final
+/// independent range check verifies the concatenation precondition and applies
+/// the same fallback if the constructed runs are not sound.
+pub(crate) fn pack_sorted_runs_by<T, F>(
+ items: Vec<T>,
+ comparator: &KeyComparator,
+ file_meta: F,
+) -> Vec<Vec<T>>
+where
+ F: Fn(&T) -> &DataFileMeta + Copy,
+{
+ if items.len() <= 1 {
+ return if items.is_empty() {
+ Vec::new()
+ } else {
+ vec![items]
+ };
+ }
+
+ let mut decoded_ranges = Vec::with_capacity(items.len());
+ for item in &items {
+ let file = file_meta(item);
+ match (
+ comparator.decode(&file.min_key),
+ comparator.decode(&file.max_key),
+ ) {
+ (Some(min), Some(max)) if compare_decoded(&min, &max) !=
Ordering::Greater => {
+ decoded_ranges.push((min, max));
+ }
+ _ => return items.into_iter().map(|item| vec![item]).collect(),
+ }
+ }
+
+ let mut keyed = items
+ .into_iter()
+ .zip(decoded_ranges)
+ .map(|(item, (min, max))| (item, min, max))
+ .collect::<Vec<_>>();
+ keyed.sort_by(|a, b| compare_decoded(&a.1, &b.1).then_with(||
compare_decoded(&a.2, &b.2)));
+
+ let mut runs: Vec<Vec<T>> = Vec::new();
+ let mut run_ends: Vec<DecodedKey> = Vec::new();
+ for (item, min, max) in keyed {
+ let mut best_run = None;
+ for (index, end) in run_ends.iter().enumerate() {
+ if compare_decoded(end, &min) != Ordering::Less {
+ continue;
+ }
+ match best_run {
+ Some(best) if compare_decoded(&run_ends[best], end) !=
Ordering::Less => {}
+ _ => best_run = Some(index),
+ }
+ }
+
+ match best_run {
+ Some(index) => {
+ runs[index].push(item);
+ run_ends[index] = max;
+ }
+ None => {
+ runs.push(vec![item]);
+ run_ends.push(max);
+ }
+ }
+ }
+
+ let sound = runs.iter().all(|run| {
+ run.windows(2).all(|pair| {
+ let previous = file_meta(&pair[0]);
+ let next = file_meta(&pair[1]);
+ match (
+ comparator.decode(&previous.max_key),
+ comparator.decode(&next.min_key),
+ ) {
+ (Some(previous_max), Some(next_min)) => {
+ compare_decoded(&previous_max, &next_min) == Ordering::Less
+ }
+ _ => false,
+ }
+ })
+ });
+ if sound {
+ runs
+ } else {
+ runs.into_iter().flatten().map(|item| vec![item]).collect()
+ }
+}
+
/// Bin-pack whole sections into splits. A section is atomic: its files
/// overlap on primary key and must never be separated, even when the section
/// alone exceeds `target_split_size`.
@@ -374,6 +487,146 @@ mod tests {
assert_eq!(compare_decoded(&big, &one), Ordering::Greater);
}
+ #[test]
+ fn pack_sorted_runs_collapses_shallow_overlap() {
+ let files = vec![
+ keyed_file("a", 1, 10, 100, 0),
+ keyed_file("b", 5, 15, 100, 0),
+ keyed_file("c", 20, 30, 100, 0),
+ keyed_file("d", 25, 35, 100, 0),
+ keyed_file("e", 40, 50, 100, 0),
+ keyed_file("f", 45, 55, 100, 0),
+ ];
+ let runs = pack_sorted_runs(files, &int_comparator());
+ assert_eq!(runs.len(), 2);
+ assert_eq!(runs.iter().map(Vec::len).sum::<usize>(), 6);
+ }
+
+ #[test]
+ fn pack_sorted_runs_chains_disjoint_files_in_key_order() {
+ let files = vec![
+ keyed_file("c", 21, 30, 100, 0),
+ keyed_file("a", 1, 10, 100, 0),
+ keyed_file("b", 11, 20, 100, 0),
+ ];
+ let runs = pack_sorted_runs(files, &int_comparator());
+ assert_eq!(runs.len(), 1);
+ assert_eq!(
+ runs[0]
+ .iter()
+ .map(|file| file.file_name.as_str())
+ .collect::<Vec<_>>(),
+ vec!["a", "b", "c"]
+ );
+ }
+
+ #[test]
+ fn pack_sorted_runs_treats_touching_ranges_as_overlapping() {
+ let files = vec![
+ keyed_file("a", 1, 10, 100, 0),
+ keyed_file("b", 10, 20, 100, 0),
+ ];
+ assert_eq!(pack_sorted_runs(files, &int_comparator()).len(), 2);
+ }
+
+ #[test]
+ fn pack_sorted_runs_degrades_when_keys_are_undecodable() {
+ let mut undecodable = keyed_file("a", 1, 10, 100, 0);
+ undecodable.min_key.clear();
+ undecodable.max_key.clear();
+ let files = vec![undecodable, keyed_file("b", 5, 15, 100, 0)];
+ assert_eq!(pack_sorted_runs(files, &int_comparator()).len(), 2);
+ }
+
+ #[test]
+ fn pack_sorted_runs_degrades_when_a_file_range_is_inverted() {
+ let files = vec![
+ keyed_file("invalid", 10, 5, 100, 0),
+ keyed_file("valid", 6, 9, 100, 0),
+ ];
+ assert_eq!(pack_sorted_runs(files, &int_comparator()).len(), 2);
+ }
+
+ #[test]
+ fn pack_sorted_runs_orders_binary_keys_unsigned() {
+ fn bytes_key(value: u8) -> Vec<u8> {
+ let mut builder = BinaryRowBuilder::new(1);
+ builder.write_binary(0, &[value]);
+ builder.build_serialized()
+ }
+
+ fn binary_file(name: &str, min: u8, max: u8) -> DataFileMeta {
+ let mut file = keyed_file(name, 0, 0, 100, 0);
+ file.min_key = bytes_key(min);
+ file.max_key = bytes_key(max);
+ file
+ }
+
+ let comparator = KeyComparator::new(vec![DataType::VarBinary(
+ crate::spec::VarBinaryType::new(16).unwrap(),
+ )]);
+ let runs = pack_sorted_runs(
+ vec![
+ binary_file("high", 0x80, 0xFE),
+ binary_file("low", 0x01, 0x7F),
+ ],
+ &comparator,
+ );
+ assert_eq!(runs.len(), 1);
+ assert_eq!(
+ runs[0]
+ .iter()
+ .map(|file| file.file_name.as_str())
+ .collect::<Vec<_>>(),
+ vec!["low", "high"]
+ );
+ }
+
+ #[test]
+ fn pack_sorted_runs_handles_multi_column_keys() {
+ fn key(first: i32, second: &str) -> Vec<u8> {
+ let mut builder = BinaryRowBuilder::new(2);
+ builder.write_int(0, first);
+ builder.write_string(1, second);
+ builder.build_serialized()
+ }
+
+ fn file(name: &str, min: (i32, &str), max: (i32, &str)) ->
DataFileMeta {
+ let mut file = keyed_file(name, 0, 0, 100, 0);
+ file.min_key = key(min.0, min.1);
+ file.max_key = key(max.0, max.1);
+ file
+ }
+
+ let comparator = KeyComparator::new(vec![
+ DataType::Int(IntType::new()),
+ DataType::VarChar(crate::spec::VarCharType::new(16).unwrap()),
+ ]);
+
+ let overlapping = pack_sorted_runs(
+ vec![file("a", (1, "a"), (2, "a")), file("b", (1, "b"), (2, "b"))],
+ &comparator,
+ );
+ assert_eq!(
+ overlapping.len(),
+ 2,
+ "second-key overlap must keep files in separate runs"
+ );
+
+ let disjoint = pack_sorted_runs(
+ vec![file("d", (3, "a"), (4, "a")), file("c", (1, "a"), (2, "a"))],
+ &comparator,
+ );
+ assert_eq!(disjoint.len(), 1);
+ assert_eq!(
+ disjoint[0]
+ .iter()
+ .map(|file| file.file_name.as_str())
+ .collect::<Vec<_>>(),
+ vec!["c", "d"]
+ );
+ }
+
#[test]
fn interval_partition_groups_overlapping_files() {
let files = vec![
@@ -431,6 +684,16 @@ mod tests {
assert_eq!(section_names(§ions), vec![vec!["a", "b"]]);
}
+ #[test]
+ fn interval_partition_inverted_range_degrades_to_single_section() {
+ let files = vec![
+ keyed_file("invalid", 10, 5, 100, 0),
+ keyed_file("valid", 6, 9, 100, 0),
+ ];
+ let sections = interval_partition(files, &int_comparator());
+ assert_eq!(section_names(§ions), vec![vec!["invalid", "valid"]]);
+ }
+
#[test]
fn pack_sections_respects_target_size() {
let sections = vec![
diff --git a/crates/paimon/src/table/sort_merge.rs
b/crates/paimon/src/table/sort_merge.rs
index 39aba7c5..a2619700 100644
--- a/crates/paimon/src/table/sort_merge.rs
+++ b/crates/paimon/src/table/sort_merge.rs
@@ -892,6 +892,8 @@ impl LoserTree {
/// Configuration for building a [`SortMergeReader`].
pub(crate) struct SortMergeReaderBuilder {
+ /// Inputs must each be monotonic by primary key. Concatenated files are
valid
+ /// only when their key ranges are strictly disjoint and ordered.
streams: Vec<ArrowRecordBatchStream>,
/// Full schema of the input streams (key + seq + value_kind + value
columns).
input_schema: SchemaRef,
diff --git a/crates/paimon/src/table/table_read.rs
b/crates/paimon/src/table/table_read.rs
index 7e7d3c3a..c96ffab7 100644
--- a/crates/paimon/src/table/table_read.rs
+++ b/crates/paimon/src/table/table_read.rs
@@ -684,7 +684,7 @@ impl<'a> PaimonTableRead<'a> {
.collect(),
read_batch_size: core_options.read_batch_size()?,
merge_splits: true,
- max_merge_file_streams: Some(256),
+ max_merge_input_streams: Some(256),
// Diff primes the before and after streams in sequence.
Keeping
// a row-group permit across yielded batches can otherwise let
// the first side block the second side indefinitely.
@@ -824,7 +824,7 @@ impl<'a> PaimonTableRead<'a> {
.collect(),
read_batch_size: core_options.read_batch_size()?,
merge_splits: false,
- max_merge_file_streams: None,
+ max_merge_input_streams: None,
parquet_read_budget: Some(self.parquet_read_budget()?),
},
);