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 7715f5c2 Align native read schema evolution, ROW decoding, and merge
semantics (#907)
7715f5c2 is described below
commit 7715f5c216950a0689f422947209cf228cb4bf7d
Author: Jingsong Lee <[email protected]>
AuthorDate: Tue Sep 22 10:23:57 2026 +0800
Align native read schema evolution, ROW decoding, and merge semantics (#907)
---
crates/paimon/src/arrow/format/mod.rs | 87 ++++++--
crates/paimon/src/arrow/format/row.rs | 8 +
crates/paimon/src/arrow/nested_evolution.rs | 273 +++++++++++++++++++++++-
crates/paimon/src/spec/core_options.rs | 10 +
crates/paimon/src/table/audit_log_table/read.rs | 1 +
crates/paimon/src/table/data_file_reader.rs | 238 ++++++++++++++++-----
crates/paimon/src/table/kv_file_reader.rs | 191 ++++++++++++++++-
crates/paimon/src/table/sort_merge.rs | 67 +++++-
crates/paimon/src/table/table_read.rs | 2 +
9 files changed, 804 insertions(+), 73 deletions(-)
diff --git a/crates/paimon/src/arrow/format/mod.rs
b/crates/paimon/src/arrow/format/mod.rs
index b4ede7e6..e9c057e3 100644
--- a/crates/paimon/src/arrow/format/mod.rs
+++ b/crates/paimon/src/arrow/format/mod.rs
@@ -62,6 +62,17 @@ pub(crate) struct FilePredicates {
/// - Row range selection
#[async_trait]
pub(crate) trait FormatFileReader: Send + Sync {
+ /// Choose the fields the decoder must actually read. Most columnar formats
+ /// can read the projection, while positional formats need the complete
+ /// physical data schema to decode each row.
+ fn select_read_fields(
+ &self,
+ _data_schema_fields: &[DataField],
+ projected_fields: &[DataField],
+ ) -> Vec<DataField> {
+ projected_fields.to_vec()
+ }
+
/// Read a single data file, returning a stream of RecordBatches containing
/// at least the projected columns (using names from the file's schema). A
/// reader MAY include extra columns it needed to scan (e.g. predicate
columns
@@ -194,24 +205,38 @@ pub(crate) fn create_format_reader(
create_format_reader_with_budget(
path,
blob_as_descriptor,
- read_fields,
+ FormatReadFields {
+ data_schema: read_fields,
+ projected: read_fields,
+ },
&HashMap::new(),
None,
blob::DEFAULT_BLOB_READ_PARALLELISM,
MosaicPrefetchOptions::default(),
)
+ .map(|configured| configured.reader)
+}
+
+pub(crate) struct ConfiguredFormatReader {
+ pub reader: Box<dyn FormatFileReader>,
+ pub read_fields: Vec<DataField>,
+}
+
+pub(crate) struct FormatReadFields<'a> {
+ pub data_schema: &'a [DataField],
+ pub projected: &'a [DataField],
}
/// Create a format reader with table options and runtime read resources.
pub(crate) fn create_format_reader_with_budget(
path: &str,
blob_as_descriptor: bool,
- read_fields: &[DataField],
+ fields: FormatReadFields<'_>,
table_options: &HashMap<String, String>,
parquet_read_budget: Option<Arc<ReadBudget>>,
blob_parallelism: usize,
mosaic_prefetch: MosaicPrefetchOptions,
-) -> crate::Result<Box<dyn FormatFileReader>> {
+) -> crate::Result<ConfiguredFormatReader> {
let lower = path.to_ascii_lowercase();
let reader: Box<dyn FormatFileReader> = if lower.ends_with(".parquet") {
Box::new(parquet::ParquetFormatReader::with_options(
@@ -229,18 +254,14 @@ pub(crate) fn create_format_reader_with_budget(
Box::new(avro::AvroFormatReader)
} else if lower.ends_with(".row") {
Box::new(row::RowFormatReader)
+ } else if lower.ends_with(".mosaic") {
+ Box::new(mosaic::MosaicFormatReader::with_prefetch(mosaic_prefetch))
} else {
- if lower.ends_with(".mosaic") {
- return Ok(shredding::maybe_wrap_reader(
-
Box::new(mosaic::MosaicFormatReader::with_prefetch(mosaic_prefetch)),
- read_fields,
- ));
- }
#[cfg(feature = "vortex")]
if lower.ends_with(".vortex") {
- return Ok(shredding::maybe_wrap_reader(
+ return Ok(configure_format_reader(
Box::new(vortex::VortexFormatReader),
- read_fields,
+ fields,
));
}
return Err(Error::Unsupported {
@@ -250,7 +271,18 @@ pub(crate) fn create_format_reader_with_budget(
),
});
};
- Ok(shredding::maybe_wrap_reader(reader, read_fields))
+ Ok(configure_format_reader(reader, fields))
+}
+
+fn configure_format_reader(
+ reader: Box<dyn FormatFileReader>,
+ fields: FormatReadFields<'_>,
+) -> ConfiguredFormatReader {
+ let read_fields = reader.select_read_fields(fields.data_schema,
fields.projected);
+ ConfiguredFormatReader {
+ reader: shredding::maybe_wrap_reader(reader, &read_fields),
+ read_fields,
+ }
}
fn supported_read_formats() -> Vec<&'static str> {
@@ -382,6 +414,37 @@ fn timestamp_millis_data_type(data_type:
&arrow_schema::DataType) -> arrow_schem
mod tests {
use super::*;
use crate::io::FileIOBuilder;
+ use crate::spec::{DataType, IntType};
+
+ #[test]
+ fn format_selects_physical_or_projected_fields() {
+ let data_schema_fields = vec![
+ DataField::new(0, "id".to_string(), DataType::Int(IntType::new())),
+ DataField::new(1, "value".to_string(),
DataType::Int(IntType::new())),
+ ];
+ let projected_fields = &data_schema_fields[1..];
+ for (path, expected) in [
+ ("data.row", data_schema_fields.as_slice()),
+ ("data.parquet", projected_fields),
+ ("data.orc", projected_fields),
+ ("data.mosaic", projected_fields),
+ ] {
+ let configured = create_format_reader_with_budget(
+ path,
+ false,
+ FormatReadFields {
+ data_schema: &data_schema_fields,
+ projected: projected_fields,
+ },
+ &HashMap::new(),
+ None,
+ blob::DEFAULT_BLOB_READ_PARALLELISM,
+ MosaicPrefetchOptions::default(),
+ )
+ .unwrap();
+ assert_eq!(configured.read_fields.as_slice(), expected, "{path}");
+ }
+ }
#[tokio::test]
async fn create_format_writer_error_lists_every_supported_format() {
diff --git a/crates/paimon/src/arrow/format/row.rs
b/crates/paimon/src/arrow/format/row.rs
index 56278108..22e37aed 100644
--- a/crates/paimon/src/arrow/format/row.rs
+++ b/crates/paimon/src/arrow/format/row.rs
@@ -225,6 +225,14 @@ pub(super) fn row_type_from_arrow_schema(schema:
&SchemaRef) -> crate::Result<Ve
#[async_trait]
impl FormatFileReader for RowFormatReader {
+ fn select_read_fields(
+ &self,
+ data_schema_fields: &[DataField],
+ _projected_fields: &[DataField],
+ ) -> Vec<DataField> {
+ data_schema_fields.to_vec()
+ }
+
async fn read_batch_stream(
&self,
reader: Box<dyn FileRead>,
diff --git a/crates/paimon/src/arrow/nested_evolution.rs
b/crates/paimon/src/arrow/nested_evolution.rs
index dde19e22..1a65cd38 100644
--- a/crates/paimon/src/arrow/nested_evolution.rs
+++ b/crates/paimon/src/arrow/nested_evolution.rs
@@ -33,7 +33,7 @@
use std::sync::Arc;
-use arrow_array::{new_null_array, Array, ArrayRef, ListArray, MapArray,
StructArray};
+use arrow_array::{new_null_array, Array, ArrayRef, ListArray, MapArray,
StringArray, StructArray};
use arrow_cast::cast;
use arrow_schema::{DataType as ArrowDataType, Field as ArrowField, Fields};
@@ -64,6 +64,18 @@ pub(crate) fn evolve_column(
return Ok(source.clone());
}
+ if matches!(target_type, DataType::VarChar(_) | DataType::Char(_))
+ && matches!(
+ source_type,
+ DataType::Row(_) | DataType::Array(_) | DataType::Map(_)
+ )
+ {
+ return Ok(Arc::new(StringArray::from(render_string_values(
+ source,
+ source_type,
+ )?)));
+ }
+
match (target_type, source_type) {
// A variant-extraction ROW is synthetic: its fields are numbered by
// position, not by schema field id, so pairing them by id would mix
@@ -105,6 +117,122 @@ pub(crate) fn evolve_column(
})
}
+/// Match Paimon's constructed-value string form, including nested values and
+/// NULL containers. Arrow's generic cast cannot convert ROW/ARRAY/MAP to Utf8.
+fn render_string_values(
+ source: &ArrayRef,
+ source_type: &DataType,
+) -> crate::Result<Vec<Option<String>>> {
+ match source_type {
+ DataType::Row(row_type) => {
+ let rows = source
+ .as_any()
+ .downcast_ref::<StructArray>()
+ .ok_or_else(|| crate::Error::DataInvalid {
+ message: format!("expected ROW array, got {:?}",
source.data_type()),
+ source: None,
+ })?;
+ let children = row_type
+ .fields()
+ .iter()
+ .enumerate()
+ .map(|(index, field)| render_string_values(rows.column(index),
field.data_type()))
+ .collect::<crate::Result<Vec<_>>>()?;
+ Ok((0..rows.len())
+ .map(|index| {
+ rows.is_valid(index).then(|| {
+ format!(
+ "{{{}}}",
+ children
+ .iter()
+ .map(|child|
child[index].as_deref().unwrap_or("null"))
+ .collect::<Vec<_>>()
+ .join(", ")
+ )
+ })
+ })
+ .collect())
+ }
+ DataType::Array(array_type) => {
+ let rows =
source.as_any().downcast_ref::<ListArray>().ok_or_else(|| {
+ crate::Error::DataInvalid {
+ message: format!("expected ARRAY array, got {:?}",
source.data_type()),
+ source: None,
+ }
+ })?;
+ let values = render_string_values(rows.values(),
array_type.element_type())?;
+ let offsets = rows.value_offsets();
+ Ok((0..rows.len())
+ .map(|index| {
+ rows.is_valid(index).then(|| {
+ format!(
+ "[{}]",
+ values[offsets[index] as usize..offsets[index + 1]
as usize]
+ .iter()
+ .map(|value|
value.as_deref().unwrap_or("null"))
+ .collect::<Vec<_>>()
+ .join(", ")
+ )
+ })
+ })
+ .collect())
+ }
+ DataType::Map(map_type) => {
+ let rows =
source.as_any().downcast_ref::<MapArray>().ok_or_else(|| {
+ crate::Error::DataInvalid {
+ message: format!("expected MAP array, got {:?}",
source.data_type()),
+ source: None,
+ }
+ })?;
+ let keys = render_string_values(rows.keys(), map_type.key_type())?;
+ let values = render_string_values(rows.values(),
map_type.value_type())?;
+ let offsets = rows.value_offsets();
+ Ok((0..rows.len())
+ .map(|index| {
+ rows.is_valid(index).then(|| {
+ format!(
+ "{{{}}}",
+ (offsets[index] as usize..offsets[index + 1] as
usize)
+ .map(|item| format!(
+ "{} -> {}",
+ keys[item].as_deref().unwrap_or("null"),
+ values[item].as_deref().unwrap_or("null")
+ ))
+ .collect::<Vec<_>>()
+ .join(", ")
+ )
+ })
+ })
+ .collect())
+ }
+ _ => {
+ let casted = cast(source, &ArrowDataType::Utf8).map_err(|error| {
+ crate::Error::UnexpectedError {
+ message: format!(
+ "failed to render {:?} as string during schema
evolution",
+ source.data_type()
+ ),
+ source: Some(Box::new(error)),
+ }
+ })?;
+ let strings = casted
+ .as_any()
+ .downcast_ref::<StringArray>()
+ .ok_or_else(|| crate::Error::UnexpectedError {
+ message: "string rendering did not produce
Utf8".to_string(),
+ source: None,
+ })?;
+ Ok((0..strings.len())
+ .map(|index| {
+ strings
+ .is_valid(index)
+ .then(|| strings.value(index).to_string())
+ })
+ .collect())
+ }
+ }
+}
+
/// Rebuild a struct array to `target_row`: pair children by field id, recurse,
/// NULL-fill target children the source does not have, and preserve the
/// source's row-level validity buffer.
@@ -351,8 +479,10 @@ fn rebuild_map(
#[cfg(test)]
mod tests {
use super::*;
- use crate::spec::{BigIntType, DataField, IntType, VarCharType};
- use arrow_array::{Int32Array, Int64Array, StringArray};
+ use crate::spec::{
+ ArrayType, BigIntType, DataField, DecimalType, IntType, MapType,
VarCharType,
+ };
+ use arrow_array::{Decimal128Array, Int32Array, Int64Array, StringArray};
use arrow_buffer::NullBuffer;
use arrow_schema::{DataType as ArrowDataType, Fields};
@@ -578,6 +708,143 @@ mod tests {
assert_eq!(as_struct(&out).column_names(), vec!["codec", "width"]);
}
+ #[test]
+ fn renders_constructed_values_as_strings_during_schema_evolution() {
+ let source = source_struct();
+ let out = evolve_column(&source, &source_row(),
&string_type()).unwrap();
+ assert_eq!(strings(&out).value(0), "{h264, 1920}");
+ assert_eq!(strings(&out).value(1), "{h265, 3840}");
+
+ let values: ArrayRef = Arc::new(Int32Array::from(vec![Some(1), None,
Some(3)]));
+ let list: ArrayRef = Arc::new(
+ ListArray::try_new(
+ Arc::new(ArrowField::new("element", ArrowDataType::Int32,
true)),
+ arrow_buffer::OffsetBuffer::new(vec![0, 2, 2, 3].into()),
+ values,
+ Some(NullBuffer::from(vec![true, false, true])),
+ )
+ .unwrap(),
+ );
+ let source_type =
DataType::Array(ArrayType::new(DataType::Int(IntType::new())));
+ let out = evolve_column(&list, &source_type, &string_type()).unwrap();
+ let rendered = strings(&out);
+ assert_eq!(rendered.value(0), "[1, null]");
+ assert!(rendered.is_null(1));
+ assert_eq!(rendered.value(2), "[3]");
+
+ let keys: ArrayRef = Arc::new(StringArray::from(vec!["k"]));
+ let values: ArrayRef = Arc::new(Int32Array::from(vec![Some(7)]));
+ let entry_fields = Fields::from(vec![
+ ArrowField::new("key", ArrowDataType::Utf8, false),
+ ArrowField::new("value", ArrowDataType::Int32, true),
+ ]);
+ let entries = StructArray::try_new(entry_fields.clone(), vec![keys,
values], None).unwrap();
+ let map: ArrayRef = Arc::new(
+ MapArray::try_new(
+ Arc::new(ArrowField::new(
+ "entries",
+ ArrowDataType::Struct(entry_fields),
+ false,
+ )),
+ arrow_buffer::OffsetBuffer::new(vec![0, 1].into()),
+ entries,
+ None,
+ false,
+ )
+ .unwrap(),
+ );
+ let source_type = DataType::Map(MapType::new(string_type(),
DataType::Int(IntType::new())));
+ let out = evolve_column(&map, &source_type, &string_type()).unwrap();
+ assert_eq!(strings(&out).value(0), "{k -> 7}");
+ }
+
+ #[test]
+ fn reducing_decimal_scale_rounds_half_up_for_both_signs() {
+ let source: ArrayRef = Arc::new(
+ Decimal128Array::from(vec![
+ Some(4_567),
+ Some(-4_567),
+ Some(4_565),
+ Some(-4_565),
+ Some(4_564),
+ Some(-4_564),
+ None,
+ ])
+ .with_precision_and_scale(6, 3)
+ .unwrap(),
+ );
+ let source_type = DataType::Decimal(DecimalType::new(6, 3).unwrap());
+ let target_type = DataType::Decimal(DecimalType::new(6, 2).unwrap());
+ let out = evolve_column(&source, &source_type, &target_type).unwrap();
+ let decimals = out.as_any().downcast_ref::<Decimal128Array>().unwrap();
+ assert_eq!(decimals.data_type(), &ArrowDataType::Decimal128(6, 2));
+ assert_eq!(decimals.value(0), 457);
+ assert_eq!(decimals.value(1), -457);
+ assert_eq!(decimals.value(2), 457);
+ assert_eq!(decimals.value(3), -457);
+ assert_eq!(decimals.value(4), 456);
+ assert_eq!(decimals.value(5), -456);
+ assert!(decimals.is_null(6));
+ decimals.validate_decimal_precision(6).unwrap();
+ }
+
+ #[test]
+ fn reducing_decimal_scale_nulls_target_precision_overflow() {
+ let source: ArrayRef = Arc::new(
+ Decimal128Array::from(vec![
+ Some(9_994),
+ Some(9_995),
+ Some(-9_995),
+ Some(999_999),
+ None,
+ ])
+ .with_precision_and_scale(6, 3)
+ .unwrap(),
+ );
+ let source_type = DataType::Decimal(DecimalType::new(6, 3).unwrap());
+ let target_type = DataType::Decimal(DecimalType::new(3, 2).unwrap());
+ let out = evolve_column(&source, &source_type, &target_type).unwrap();
+ let decimals = out.as_any().downcast_ref::<Decimal128Array>().unwrap();
+ assert_eq!(decimals.data_type(), &ArrowDataType::Decimal128(3, 2));
+ assert_eq!(decimals.value(0), 999);
+ assert!(decimals.is_null(1));
+ assert!(decimals.is_null(2));
+ assert!(decimals.is_null(3));
+ assert!(decimals.is_null(4));
+ decimals.validate_decimal_precision(3).unwrap();
+ }
+
+ #[test]
+ fn renders_nested_row_child_as_string_without_losing_parent_nulls() {
+ let inner: ArrayRef = Arc::new(StructArray::from(vec![(
+ Arc::new(ArrowField::new("a", ArrowDataType::Int32, true)),
+ Arc::new(Int32Array::from(vec![Some(1), None])) as ArrayRef,
+ )]));
+ let source: ArrayRef = Arc::new(
+ StructArray::try_new(
+ Fields::from(vec![ArrowField::new(
+ "inner",
+ inner.data_type().clone(),
+ true,
+ )]),
+ vec![inner],
+ Some(NullBuffer::from(vec![true, false])),
+ )
+ .unwrap(),
+ );
+ let source_type = row(vec![field(
+ 1,
+ "inner",
+ row(vec![field(2, "a", DataType::Int(IntType::new()))]),
+ )]);
+ let target_type = row(vec![field(1, "inner", string_type())]);
+ let out = evolve_column(&source, &source_type, &target_type).unwrap();
+ let outer = as_struct(&out);
+ let rendered = strings(outer.column(0));
+ assert_eq!(rendered.value(0), "{1}");
+ assert!(outer.is_null(1));
+ }
+
#[test]
fn fills_null_for_a_child_added_inside_an_array_element() {
// source: ARRAY<ROW<codec>>, two elements in one list row.
diff --git a/crates/paimon/src/spec/core_options.rs
b/crates/paimon/src/spec/core_options.rs
index 0073f172..bd20c09b 100644
--- a/crates/paimon/src/spec/core_options.rs
+++ b/crates/paimon/src/spec/core_options.rs
@@ -101,6 +101,7 @@ pub(crate) const MOSAIC_READ_PREFETCH_MAX_BYTES_OPTION:
&str = "mosaic.read.pref
pub(crate) const TABLE_READ_SEQUENCE_NUMBER_ENABLED_OPTION: &str =
"table-read.sequence-number.enabled";
pub(crate) const SEQUENCE_FIELD_OPTION: &str = "sequence.field";
+const SEQUENCE_FIELD_SORT_ORDER_OPTION: &str = "sequence.field.sort-order";
pub(crate) const DISABLE_EXPLICIT_TYPE_CASTING_OPTION: &str =
"disable-explicit-type-casting";
pub(crate) const DISABLE_ALTER_COLUMN_NULL_TO_NOT_NULL_OPTION: &str =
"alter-column-null-to-not-null.disabled";
@@ -622,6 +623,15 @@ impl<'a> CoreOptions<'a> {
.unwrap_or_default()
}
+ /// User sequence fields sort ascending unless explicitly configured
descending.
+ /// Null sequences remain smaller than non-null sequences in either order.
+ pub fn sequence_field_sort_order_is_ascending(&self) -> bool {
+ !self
+ .options
+ .get(SEQUENCE_FIELD_SORT_ORDER_OPTION)
+ .is_some_and(|order| order.eq_ignore_ascii_case("descending"))
+ }
+
/// Merge engine for primary-key tables. Default is `Deduplicate`.
pub fn merge_engine(&self) -> crate::Result<MergeEngine> {
match self.options.get(MERGE_ENGINE_OPTION) {
diff --git a/crates/paimon/src/table/audit_log_table/read.rs
b/crates/paimon/src/table/audit_log_table/read.rs
index 175ce60e..d81f1f3d 100644
--- a/crates/paimon/src/table/audit_log_table/read.rs
+++ b/crates/paimon/src/table/audit_log_table/read.rs
@@ -117,6 +117,7 @@ impl<'a> AuditLogRead<'a> {
read_type,
predicates: self.read.data_predicates.clone(),
primary_keys:
self.read.table.schema.trimmed_primary_keys(),
+ table_primary_keys:
self.read.table.schema.primary_keys().to_vec(),
merge_engine,
sequence_fields: core_options
.sequence_fields()
diff --git a/crates/paimon/src/table/data_file_reader.rs
b/crates/paimon/src/table/data_file_reader.rs
index 95452cfc..a4a18839 100644
--- a/crates/paimon/src/table/data_file_reader.rs
+++ b/crates/paimon/src/table/data_file_reader.rs
@@ -17,7 +17,9 @@
use crate::arrow::build_target_arrow_schema;
use crate::arrow::format::blob::DEFAULT_BLOB_READ_PARALLELISM;
-use crate::arrow::format::{create_format_reader_with_budget,
MosaicPrefetchOptions};
+use crate::arrow::format::{
+ create_format_reader_with_budget, FormatReadFields, MosaicPrefetchOptions,
+};
use crate::arrow::schema_evolution::{create_index_mapping, NULL_FIELD_INDEX};
use crate::arrow::ReadBudget;
use crate::deletion_vector::{DeletionVector, DeletionVectorFactory};
@@ -107,7 +109,7 @@ impl FileRead for TimedFileRead {
}
}
-/// Reads data from Parquet files.
+/// Reads data files through their format-specific readers.
#[derive(Clone)]
pub(crate) struct DataFileReader {
file_io: FileIO,
@@ -335,6 +337,7 @@ impl DataFileReader {
&split,
file_meta,
data_fields,
+ None,
row_selection,
)?;
while let Some(batch) = stream.next().await {
@@ -395,7 +398,7 @@ impl DataFileReader {
}
}
- /// Read a single parquet file from a split, returning a lazy stream of
batches.
+ /// Read a single data file from a split, returning a lazy stream of
batches.
/// Optionally applies a deletion vector.
///
/// Handles schema evolution using field-ID-based index mapping:
@@ -422,7 +425,43 @@ impl DataFileReader {
});
let row_selection =
merge_row_selection(file_meta.row_count, dv.as_deref(),
local_ranges.as_deref());
- self.read_single_file_stream_with_selection(split, file_meta,
data_fields, row_selection)
+ self.read_single_file_stream_with_selection(
+ split,
+ file_meta,
+ data_fields,
+ None,
+ row_selection,
+ )
+ }
+
+ /// Read one file with the complete physical data schema supplied by a
+ /// caller such as the KV reader. The format chooses whether it needs the
+ /// full schema or only the projected fields.
+ pub(super) fn read_single_file_stream_with_schema(
+ &self,
+ split: &DataSplit,
+ file_meta: DataFileMeta,
+ data_fields: Option<Vec<DataField>>,
+ data_schema_fields: Vec<DataField>,
+ dv: Option<Arc<DeletionVector>>,
+ row_ranges: Option<Vec<RowRange>>,
+ ) -> crate::Result<ArrowRecordBatchStream> {
+ let local_ranges = row_ranges.as_ref().map(|ranges| {
+ to_local_row_ranges(
+ ranges,
+ file_meta.first_row_id.unwrap_or(0),
+ file_meta.row_count,
+ )
+ });
+ let row_selection =
+ merge_row_selection(file_meta.row_count, dv.as_deref(),
local_ranges.as_deref());
+ self.read_single_file_stream_with_selection(
+ split,
+ file_meta,
+ data_fields,
+ Some(data_schema_fields),
+ row_selection,
+ )
}
fn read_single_file_stream_with_selection(
@@ -430,6 +469,7 @@ impl DataFileReader {
split: &DataSplit,
file_meta: DataFileMeta,
data_fields: Option<Vec<DataField>>,
+ data_schema_fields: Option<Vec<DataField>>,
row_selection: Option<Vec<RowRange>>,
) -> crate::Result<ArrowRecordBatchStream> {
if row_selection.as_ref().is_some_and(Vec::is_empty) {
@@ -488,8 +528,6 @@ impl DataFileReader {
let target_schema = build_target_arrow_schema(&read_type)?;
let file_fields = data_fields.clone().unwrap_or_else(||
table_fields.clone());
- let is_row_file = is_row_file(&file_meta);
-
// What the reader is asked for.
let projected_read_fields: Vec<DataField> = if let Some(ref df) =
data_fields {
read_data_fields(df, &read_type)?
@@ -500,11 +538,26 @@ impl DataFileReader {
.cloned()
.collect()
};
- let format_read_fields = if is_row_file {
- file_fields.clone()
- } else {
- projected_read_fields
- };
+ let data_schema_fields = data_schema_fields_for_file(
+ &file_fields,
+ file_meta.write_cols.as_deref(),
+ data_schema_fields.as_deref(),
+ )?;
+ let path_to_read = split.data_file_path(&file_meta);
+ let configured_reader = create_format_reader_with_budget(
+ &path_to_read,
+ blob_as_descriptor,
+ FormatReadFields {
+ data_schema: &data_schema_fields,
+ projected: &projected_read_fields,
+ },
+ &table_options,
+ parquet_read_budget,
+ blob_parallelism,
+ mosaic_prefetch,
+ )?;
+ let format_read_fields = configured_reader.read_fields;
+ let format_reader = configured_reader.reader;
// The decoded batch is described by `format_read_fields`, so map
// `read_type` onto *that* list: its entries carry the types the
columns
// actually come back as, which is what reconciling them needs.
@@ -524,7 +577,7 @@ impl DataFileReader {
let remapped = crate::arrow::filtering::remap_predicates_to_file(
&predicates,
&table_fields,
- &file_fields,
+ &data_schema_fields,
);
if remapped.is_empty() && row_filter_factory.is_none() {
None
@@ -532,23 +585,13 @@ impl DataFileReader {
Some(crate::arrow::format::FilePredicates {
predicates: remapped,
row_filter_factory,
- file_fields: file_fields.clone(),
+ file_fields: data_schema_fields.clone(),
})
}
};
Ok(try_stream! {
let schema_open_start = read_timing.as_ref().map(|_|
Instant::now());
- let path_to_read = split.data_file_path(&file_meta);
- let format_reader = create_format_reader_with_budget(
- &path_to_read,
- blob_as_descriptor,
- &format_read_fields,
- &table_options,
- parquet_read_budget,
- blob_parallelism,
- mosaic_prefetch,
- )?;
let input_file = file_io.new_input(&path_to_read)?;
let open_start = read_timing.as_ref().map(|_| Instant::now());
let file_reader = input_file.reader().await?;
@@ -781,8 +824,6 @@ impl DataFileReader {
let target_schema = build_target_arrow_schema(&read_type)?;
let file_fields = data_fields.clone().unwrap_or_else(||
table_fields.clone());
- let is_row_file = is_row_file(&file_meta);
-
// What the reader is asked for.
let projected_read_fields: Vec<DataField> = if let Some(ref df) =
data_fields {
read_data_fields(df, &read_type)?
@@ -793,11 +834,23 @@ impl DataFileReader {
.cloned()
.collect()
};
- let format_read_fields = if is_row_file {
- file_fields.clone()
- } else {
- projected_read_fields
- };
+ let data_schema_fields =
+ data_schema_fields_for_file(&file_fields,
file_meta.write_cols.as_deref(), None)?;
+ let path_to_read = split.data_file_path(&file_meta);
+ let configured_reader = create_format_reader_with_budget(
+ &path_to_read,
+ blob_as_descriptor,
+ FormatReadFields {
+ data_schema: &data_schema_fields,
+ projected: &projected_read_fields,
+ },
+ &table_options,
+ parquet_read_budget,
+ blob_parallelism,
+ mosaic_prefetch,
+ )?;
+ let format_read_fields = configured_reader.read_fields;
+ let format_reader = configured_reader.reader;
// The decoded batch is described by `format_read_fields`, so map
// `read_type` onto *that* list: its entries carry the types the
columns
// actually come back as, which is what reconciling them needs.
@@ -815,7 +868,7 @@ impl DataFileReader {
let remapped = crate::arrow::filtering::remap_predicates_to_file(
&predicates,
&table_fields,
- &file_fields,
+ &data_schema_fields,
);
if remapped.is_empty() {
None
@@ -823,7 +876,7 @@ impl DataFileReader {
Some(crate::arrow::format::FilePredicates {
predicates: remapped,
row_filter_factory: None,
- file_fields: file_fields.clone(),
+ file_fields: data_schema_fields.clone(),
})
}
};
@@ -836,16 +889,6 @@ impl DataFileReader {
merge_row_selection(file_meta.row_count, dv.as_deref(),
Some(&local_ranges));
Ok(try_stream! {
- let path_to_read = split.data_file_path(&file_meta);
- let format_reader = create_format_reader_with_budget(
- &path_to_read,
- blob_as_descriptor,
- &format_read_fields,
- &table_options,
- parquet_read_budget,
- blob_parallelism,
- mosaic_prefetch,
- )?;
let input_file = file_io.new_input(&path_to_read)?;
let file_reader = input_file.reader().await?;
@@ -1041,12 +1084,27 @@ fn data_field_with_type(field: &DataField, data_type:
DataType) -> DataField {
.with_description(field.description().map(ToString::to_string))
}
-fn is_row_file(file_meta: &DataFileMeta) -> bool {
- file_meta.file_name.to_ascii_lowercase().ends_with(".row")
- || file_meta
- .external_path
- .as_deref()
- .is_some_and(|path| path.to_ascii_lowercase().ends_with(".row"))
+fn data_schema_fields_for_file(
+ file_fields: &[DataField],
+ write_cols: Option<&[String]>,
+ data_schema_fields: Option<&[DataField]>,
+) -> crate::Result<Vec<DataField>> {
+ if let Some(write_cols) = write_cols {
+ return write_cols
+ .iter()
+ .map(|name| {
+ file_fields
+ .iter()
+ .find(|field| field.name() == name)
+ .cloned()
+ .ok_or_else(|| Error::DataInvalid {
+ message: format!("write column '{name}' is absent from
the file schema"),
+ source: None,
+ })
+ })
+ .collect();
+ }
+ Ok(data_schema_fields.unwrap_or(file_fields).to_vec())
}
/// Convert ranges from their read-path coordinate system to file-local ranges.
@@ -1378,6 +1436,31 @@ mod row_tests {
DataField::new(id, name.to_string(), data_type)
}
+ #[test]
+ fn complete_data_schema_uses_supplied_physical_schema_and_write_cols() {
+ let fields = vec![
+ field(1, "id", DataType::Int(IntType::new())),
+ field(2, "value", DataType::Int(IntType::new())),
+ ];
+ let physical = vec![field(1_000_001, "_KEY_id",
DataType::Int(IntType::new()))];
+ assert_eq!(
+ data_schema_fields_for_file(&fields, None, None).unwrap(),
+ fields
+ );
+ let selected = data_schema_fields_for_file(&fields, None,
Some(&physical)).unwrap();
+ assert_eq!(
+ selected.iter().map(DataField::name).collect::<Vec<_>>(),
+ vec!["_KEY_id"]
+ );
+ let partial =
+ data_schema_fields_for_file(&fields, Some(&["value".to_string()]),
Some(&physical))
+ .unwrap();
+ assert_eq!(
+ partial.iter().map(DataField::name).collect::<Vec<_>>(),
+ vec!["value"]
+ );
+ }
+
fn data_file(file_name: &str, file_size: i64, row_count: i64, schema_id:
i64) -> DataFileMeta {
DataFileMeta {
file_name: file_name.to_string(),
@@ -1589,6 +1672,65 @@ mod row_tests {
assert_eq!(ages, vec![30, 40, 50]);
}
+ #[tokio::test]
+ async fn
row_partial_write_cols_resolve_predicates_against_physical_schema() {
+ let fields = vec![
+ field(0, "id", DataType::Int(IntType::new())),
+ field(1, "age", DataType::Int(IntType::new())),
+ ];
+ let schema = build_target_arrow_schema(&fields[..1]).unwrap();
+ let batch =
+ RecordBatch::try_new(schema.clone(),
vec![Arc::new(Int32Array::from(vec![1, 2]))])
+ .unwrap();
+ let file_io = FileIOBuilder::new("memory").build().unwrap();
+ let table_path = "memory:/row_partial_predicate";
+ let bucket_path = format!("{table_path}/bucket-0");
+ let file_name = "partial.row";
+ let output = file_io
+ .new_output(&format!("{bucket_path}/{file_name}"))
+ .unwrap();
+ let mut writer = create_format_writer(&output, schema, "zstd", 1,
None, None, None)
+ .await
+ .unwrap();
+ writer.write(&batch).await.unwrap();
+ let file_size = writer.close().await.unwrap().file_size as i64;
+ let mut file = data_file(file_name, file_size, 2, 1);
+ file.write_cols = Some(vec!["id".to_string()]);
+ let split = DataSplitBuilder::new()
+ .with_snapshot(1)
+ .with_partition(BinaryRow::new(0))
+ .with_bucket(0)
+ .with_bucket_path(bucket_path)
+ .with_total_buckets(1)
+ .with_data_files(vec![file])
+ .build()
+ .unwrap();
+ let predicates = PredicateBuilder::new(&fields);
+ for (predicate, expected_rows) in [
+ (predicates.is_null("age").unwrap(), 2),
+ (predicates.greater_than("age", Datum::Int(0)).unwrap(), 0),
+ ] {
+ let reader = DataFileReader::new(
+ file_io.clone(),
+ SchemaManager::new(file_io.clone(), table_path.to_string()),
+ 1,
+ fields.clone(),
+ vec![fields[0].clone()],
+ vec![predicate],
+ );
+ let batches = reader
+ .read(std::slice::from_ref(&split))
+ .unwrap()
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap();
+ assert_eq!(
+ batches.iter().map(RecordBatch::num_rows).sum::<usize>(),
+ expected_rows
+ );
+ }
+ }
+
/// The predicate must not renumber a projected `_ROW_ID`: the surviving
row
/// keeps its original physical position even when the predicate column is
/// absent from the requested projection.
diff --git a/crates/paimon/src/table/kv_file_reader.rs
b/crates/paimon/src/table/kv_file_reader.rs
index cfc265bc..36ca7ed8 100644
--- a/crates/paimon/src/table/kv_file_reader.rs
+++ b/crates/paimon/src/table/kv_file_reader.rs
@@ -72,7 +72,10 @@ pub(crate) struct KeyValueReadConfig {
pub table_fields: Vec<DataField>,
pub read_type: Vec<DataField>,
pub predicates: Vec<Predicate>,
+ /// Physical key fields exclude partition columns for sort-merge.
pub primary_keys: Vec<String>,
+ /// Full table keys also protect partition-PK fields from aggregation.
+ pub table_primary_keys: Vec<String>,
pub merge_engine: MergeEngine,
pub sequence_fields: Vec<String>,
pub read_batch_size: usize,
@@ -172,6 +175,42 @@ fn ensure_merge_input_limit(input_stream_count: usize,
limit: Option<usize>) ->
Ok(())
}
+/// Java's `KeyValueFieldsExtractor` builds the physical file layout from the
+/// schema that wrote the file, including its historical key names and types.
+fn key_value_data_schema_fields(
+ file_fields: &[DataField],
+ trimmed_primary_keys: &[String],
+) -> crate::Result<Vec<DataField>> {
+ let mut physical = Vec::with_capacity(trimmed_primary_keys.len() + 2 +
file_fields.len());
+ for name in trimmed_primary_keys {
+ let field = file_fields
+ .iter()
+ .find(|field| field.name() == name)
+ .ok_or_else(|| Error::DataInvalid {
+ message: format!("KV key field '{name}' is absent from the
file schema"),
+ source: None,
+ })?;
+ physical.push(
+ field
+ .clone()
+ .with_name(format!("_KEY_{name}"))
+ .with_id(field.id() + 1_000_000),
+ );
+ }
+ physical.push(DataField::new(
+ SEQUENCE_NUMBER_FIELD_ID,
+ SEQUENCE_NUMBER_FIELD_NAME.to_string(),
+ PaimonDataType::BigInt(BigIntType::new()),
+ ));
+ physical.push(DataField::new(
+ VALUE_KIND_FIELD_ID,
+ VALUE_KIND_FIELD_NAME.to_string(),
+ PaimonDataType::TinyInt(TinyIntType::new()),
+ ));
+ physical.extend_from_slice(file_fields);
+ Ok(physical)
+}
+
struct MergeRun {
files: Vec<MergeFile>,
}
@@ -307,7 +346,7 @@ impl KeyValueFileReader {
&config.table_options,
&config.table_name,
merge_output_fields,
- &config.primary_keys,
+ &config.table_primary_keys,
&config.sequence_fields,
)?)),
}
@@ -595,18 +634,28 @@ impl KeyValueFileReader {
.with_table_options(config.table_options.clone())
.with_mosaic_prefetch(config.mosaic_prefetch);
let run_schema_manager = config.schema_manager.clone();
+ let run_table_fields = config.table_fields.clone();
+ let run_primary_keys = config.primary_keys.clone();
let run_file_io = file_io.clone();
let deletion_files_by_split =
deletion_files_by_split.clone();
let run_stream: ArrowRecordBatchStream =
Box::pin(try_stream! {
for MergeFile { split, file: file_meta } in files {
- let data_fields: Option<Vec<DataField>> =
+ let data_schema =
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
- };
+
Some(run_schema_manager.schema(file_meta.schema_id).await?)
+ } else {
+ None
+ };
+ let data_fields =
data_schema.as_ref().map(|schema| schema.fields().to_vec());
+ let file_fields = data_schema
+ .as_ref()
+ .map_or(run_table_fields.as_slice(),
|schema| schema.fields());
+ let file_key_names = data_schema
+ .as_ref()
+ .map(|schema|
schema.trimmed_primary_keys());
+ let key_names =
file_key_names.as_deref().unwrap_or(&run_primary_keys);
+ let data_schema_fields =
+ key_value_data_schema_fields(file_fields,
key_names)?;
let deletion_file = deletion_files_by_split
.get(&(Arc::as_ptr(&split) as usize))
.and_then(|files|
files.get(&file_meta.file_name))
@@ -617,10 +666,11 @@ impl KeyValueFileReader {
)),
None => None,
};
- let mut file_stream =
reader.read_single_file_stream(
+ let mut file_stream =
reader.read_single_file_stream_with_schema(
split.as_ref(),
file_meta,
data_fields,
+ data_schema_fields,
deletion_vector,
split.row_ranges().map(|ranges|
ranges.to_vec()),
)?;
@@ -658,6 +708,10 @@ impl KeyValueFileReader {
merge_output_schema.clone(),
merge_function(&config, &merge_output_fields)?,
)
+ .with_user_sequence_descending(
+ !CoreOptions::new(&config.table_options)
+ .sequence_field_sort_order_is_ascending(),
+ )
.build()?;
while let Some(batch) = merge_stream.next().await {
@@ -721,6 +775,8 @@ impl KeyValueFileReader {
#[cfg(test)]
mod tests {
use super::*;
+ use crate::arrow::build_target_arrow_schema;
+ use crate::arrow::format::create_format_writer;
use crate::catalog::Identifier;
use crate::deletion_vector::DeletionVector;
use crate::io::FileIOBuilder;
@@ -742,6 +798,117 @@ mod tests {
use std::collections::HashMap;
use std::sync::Arc;
+ #[test]
+ fn kv_row_layout_preserves_file_schema_key_names_and_fields() {
+ let fields = vec![
+ DataField::new(0, "old_id".to_string(),
DataType::Int(IntType::new()))
+ .with_description(Some("sort key".to_string())),
+ DataField::new(1, "value".to_string(),
DataType::Int(IntType::new())),
+ ];
+ let physical = key_value_data_schema_fields(&fields,
&["old_id".to_string()]).unwrap();
+ assert_eq!(
+ physical.iter().map(DataField::name).collect::<Vec<_>>(),
+ vec![
+ "_KEY_old_id",
+ "_SEQUENCE_NUMBER",
+ "_VALUE_KIND",
+ "old_id",
+ "value"
+ ]
+ );
+ assert_eq!(physical[0].id(), 1_000_000);
+ assert_eq!(physical[0].description(), Some("sort key"));
+ assert_eq!(&physical[3..], fields);
+ assert!(key_value_data_schema_fields(&fields,
&["id".to_string()]).is_err());
+ }
+
+ #[tokio::test]
+ async fn kv_row_read_uses_historical_schema_for_renamed_key() {
+ let file_io = test_file_io();
+ let table_path = "memory:/kv_row_renamed_key";
+ setup_dirs(&file_io, table_path).await;
+ let old_schema = TableSchema::new(
+ 0,
+ &Schema::builder()
+ .column("old_id", DataType::Int(IntType::new()))
+ .column("value", DataType::Int(IntType::new()))
+ .primary_key(["old_id"])
+ .option("bucket", "1")
+ .build()
+ .unwrap(),
+ );
+ let current_schema = TableSchema::new(
+ 1,
+ &Schema::builder()
+ .column("id", DataType::Int(IntType::new()))
+ .column("value", DataType::Int(IntType::new()))
+ .primary_key(["id"])
+ .option("bucket", "1")
+ .build()
+ .unwrap(),
+ );
+ let table = Table::new(
+ file_io.clone(),
+ Identifier::new("default", "kv_row_renamed_key"),
+ table_path.to_string(),
+ current_schema,
+ None,
+ );
+ write_schema_file(&table, &old_schema).await;
+
+ let physical =
+ key_value_data_schema_fields(old_schema.fields(),
&old_schema.trimmed_primary_keys())
+ .unwrap();
+ let schema = build_target_arrow_schema(&physical).unwrap();
+ let batch = RecordBatch::try_new(
+ schema.clone(),
+ vec![
+ Arc::new(Int32Array::from(vec![7])),
+ Arc::new(Int64Array::from(vec![0])),
+ Arc::new(Int8Array::from(vec![0])),
+ Arc::new(Int32Array::from(vec![7])),
+ Arc::new(Int32Array::from(vec![42])),
+ ],
+ )
+ .unwrap();
+ let bucket_path = format!("{table_path}/bucket-0");
+ let file_name = "part-0.row";
+ let output = file_io
+ .new_output(&format!("{bucket_path}/{file_name}"))
+ .unwrap();
+ let mut writer = create_format_writer(&output, schema, "zstd", 1,
None, None, None)
+ .await
+ .unwrap();
+ writer.write(&batch).await.unwrap();
+ let file_size = writer.close().await.unwrap().file_size as i64;
+
+ let mut file = dummy_data_file(file_name.to_string());
+ file.file_size = file_size;
+ file.min_key = int_key(7);
+ file.max_key = int_key(7);
+ let split = DataSplitBuilder::new()
+ .with_snapshot(1)
+ .with_partition(BinaryRow::new(0))
+ .with_bucket(0)
+ .with_bucket_path(bucket_path)
+ .with_total_buckets(1)
+ .with_data_files(vec![file])
+ .with_raw_convertible(false)
+ .build()
+ .unwrap();
+ let batches = table
+ .new_read_builder()
+ .new_read()
+ .unwrap()
+ .to_arrow(&[split])
+ .unwrap()
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap();
+ assert_eq!(int_column(&batches, "id"), vec![7]);
+ assert_eq!(int_column(&batches, "value"), vec![42]);
+ }
+
#[tokio::test]
async fn test_row_id_filter_on_a_primary_key_table_is_rejected() {
let file_io = test_file_io();
@@ -1477,6 +1644,7 @@ mod tests {
read_type: table.schema().fields().to_vec(),
predicates: Vec::new(),
primary_keys: table.schema().trimmed_primary_keys(),
+ table_primary_keys: table.schema().primary_keys().to_vec(),
merge_engine: core_options.merge_engine().unwrap(),
sequence_fields: Vec::new(),
read_batch_size: core_options.read_batch_size().unwrap(),
@@ -1597,6 +1765,7 @@ mod tests {
read_type: table.schema().fields().to_vec(),
predicates: Vec::new(),
primary_keys: table.schema().trimmed_primary_keys(),
+ table_primary_keys: table.schema().primary_keys().to_vec(),
merge_engine: core_options.merge_engine().unwrap(),
sequence_fields: Vec::new(),
read_batch_size: core_options.read_batch_size().unwrap(),
@@ -1801,6 +1970,7 @@ mod tests {
read_type: table.schema().fields().to_vec(),
predicates: Vec::new(),
primary_keys: table.schema().trimmed_primary_keys(),
+ table_primary_keys: table.schema().primary_keys().to_vec(),
merge_engine: core_options.merge_engine().unwrap(),
sequence_fields: core_options
.sequence_fields()
@@ -1877,6 +2047,7 @@ mod tests {
read_type: table.schema().fields().to_vec(),
predicates: Vec::new(),
primary_keys: table.schema().trimmed_primary_keys(),
+ table_primary_keys: table.schema().primary_keys().to_vec(),
merge_engine: core_options.merge_engine().unwrap(),
sequence_fields: Vec::new(),
read_batch_size: core_options.read_batch_size().unwrap(),
@@ -2070,6 +2241,7 @@ mod tests {
read_type: table.schema().fields().to_vec(),
predicates: Vec::new(),
primary_keys: table.schema().trimmed_primary_keys(),
+ table_primary_keys: table.schema().primary_keys().to_vec(),
merge_engine: core_options.merge_engine().unwrap(),
sequence_fields: Vec::new(),
read_batch_size: core_options.read_batch_size().unwrap(),
@@ -2138,6 +2310,7 @@ mod tests {
read_type: table.schema().fields().to_vec(),
predicates: Vec::new(),
primary_keys: table.schema().trimmed_primary_keys(),
+ table_primary_keys: table.schema().primary_keys().to_vec(),
merge_engine: core_options.merge_engine().unwrap(),
sequence_fields: Vec::new(),
read_batch_size: core_options.read_batch_size().unwrap(),
diff --git a/crates/paimon/src/table/sort_merge.rs
b/crates/paimon/src/table/sort_merge.rs
index 1e88005b..fcbbdaf0 100644
--- a/crates/paimon/src/table/sort_merge.rs
+++ b/crates/paimon/src/table/sort_merge.rs
@@ -955,6 +955,7 @@ pub(crate) struct SortMergeReaderBuilder {
value_kind_index: usize,
/// Indices of user-defined sequence field columns in input_schema (if
configured).
user_sequence_indices: Vec<usize>,
+ user_sequence_descending: bool,
/// Indices of user value columns in input_schema (output columns).
value_indices: Vec<usize>,
/// Output schema (key + value columns, no system columns).
@@ -983,6 +984,7 @@ impl SortMergeReaderBuilder {
seq_index,
value_kind_index,
user_sequence_indices,
+ user_sequence_descending: false,
value_indices,
output_schema,
merge_function,
@@ -996,6 +998,11 @@ impl SortMergeReaderBuilder {
self
}
+ pub(crate) fn with_user_sequence_descending(mut self, descending: bool) ->
Self {
+ self.user_sequence_descending = descending;
+ self
+ }
+
/// Build the sort-merge stream.
pub(crate) fn build(self) -> crate::Result<ArrowRecordBatchStream> {
let sort_fields: Vec<SortField> = self
@@ -1016,6 +1023,7 @@ impl SortMergeReaderBuilder {
self.seq_index,
self.value_kind_index,
self.user_sequence_indices,
+ self.user_sequence_descending,
self.value_indices,
self.output_schema,
self.merge_function,
@@ -1067,6 +1075,7 @@ fn sort_merge_stream(
seq_index: usize,
value_kind_index: usize,
user_sequence_indices: Vec<usize>,
+ user_sequence_descending: bool,
value_indices: Vec<usize>,
output_schema: SchemaRef,
merge_function: Box<dyn MergeFunction>,
@@ -1164,7 +1173,11 @@ fn sort_merge_stream(
row_idx: cursor.offset,
sequence_number: cursor.sequence_number(seq_index),
value_kind: cursor.value_kind(value_kind_index),
- user_sequences:
user_sequence_indices.iter().map(|&idx| cursor.user_sequence(idx)).collect(),
+ user_sequences:
user_sequence_indices.iter().map(|&idx| {
+ cursor.user_sequence(idx).map(|value| {
+ if user_sequence_descending { -value } else {
value }
+ })
+ }).collect(),
});
}
@@ -1830,6 +1843,58 @@ mod tests {
assert_eq!(values, vec!["winner_a", "winner_b"]);
}
+ #[tokio::test]
+ async fn test_descending_user_sequence_keeps_nulls_first() {
+ let schema = Arc::new(Schema::new(vec![
+ Field::new("pk", DataType::Int32, false),
+ Field::new("_SEQUENCE_NUMBER", DataType::Int64, false),
+ Field::new("_VALUE_KIND", DataType::Int8, false),
+ Field::new("ts", DataType::Int64, true),
+ Field::new("value", DataType::Utf8, false),
+ ]));
+ let output_schema = Arc::new(Schema::new(vec![
+ Field::new("pk", DataType::Int32, false),
+ Field::new("value", DataType::Utf8, false),
+ ]));
+ let batch = RecordBatch::try_new(
+ schema.clone(),
+ vec![
+ Arc::new(Int32Array::from(vec![1, 1, 1])),
+ Arc::new(Int64Array::from(vec![1, 2, 3])),
+ Arc::new(Int8Array::from(vec![0, 0, 0])),
+ Arc::new(Int64Array::from(vec![Some(100), Some(50), None])),
+ Arc::new(StringArray::from(vec!["high", "low", "null"])),
+ ],
+ )
+ .unwrap();
+
+ let result = SortMergeReaderBuilder::new(
+ vec![stream_from_batches(vec![batch])],
+ schema,
+ vec![0],
+ 1,
+ 2,
+ vec![3],
+ vec![4],
+ output_schema,
+ Box::new(DeduplicateMergeFunction),
+ )
+ .with_user_sequence_descending(true)
+ .build()
+ .unwrap()
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap();
+
+ assert_eq!(result.len(), 1);
+ let values = result[0]
+ .column(1)
+ .as_any()
+ .downcast_ref::<StringArray>()
+ .unwrap();
+ assert_eq!(values.value(0), "low");
+ }
+
#[tokio::test]
async fn test_delete_row_filtered() {
let schema = make_schema();
diff --git a/crates/paimon/src/table/table_read.rs
b/crates/paimon/src/table/table_read.rs
index dde0f512..70576d30 100644
--- a/crates/paimon/src/table/table_read.rs
+++ b/crates/paimon/src/table/table_read.rs
@@ -819,6 +819,7 @@ impl<'a> PaimonTableRead<'a> {
read_type: read_type.to_vec(),
predicates: self.data_predicates.clone(),
primary_keys: self.table.schema.trimmed_primary_keys(),
+ table_primary_keys: self.table.schema.primary_keys().to_vec(),
merge_engine: core_options.merge_engine()?,
sequence_fields: core_options
.sequence_fields()
@@ -946,6 +947,7 @@ impl<'a> PaimonTableRead<'a> {
read_type: self.read_type().to_vec(),
predicates: self.data_predicates.clone(),
primary_keys: self.table.schema.trimmed_primary_keys(),
+ table_primary_keys: self.table.schema.primary_keys().to_vec(),
merge_engine: core_options.merge_engine()?,
sequence_fields: core_options
.sequence_fields()