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()

Reply via email to