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 169db37c Align native table writes and merge ordering with Java (#958)
169db37c is described below

commit 169db37c7034d1d7501826ef55833e841ff84e7d
Author: Jingsong Lee <[email protected]>
AuthorDate: Sat Sep 26 08:41:25 2026 +0800

    Align native table writes and merge ordering with Java (#958)
---
 crates/paimon/src/spec/core_options.rs             |  18 +
 crates/paimon/src/spec/partition_utils.rs          | 141 ++++++-
 crates/paimon/src/spec/schema.rs                   |   4 +
 crates/paimon/src/table/cow_writer.rs              |   4 +-
 crates/paimon/src/table/data_file_writer.rs        |  98 ++++-
 .../src/table/dedicated_format_file_writer.rs      |  10 +-
 crates/paimon/src/table/kv_file_writer.rs          | 194 +++++++---
 crates/paimon/src/table/mod.rs                     |   1 +
 crates/paimon/src/table/sort_merge.rs              | 279 ++++++++++----
 crates/paimon/src/table/table_write.rs             |  46 ++-
 crates/paimon/src/table/write_batch_normalize.rs   | 248 ++++++++++++
 .../tests/binary_partition_native_write_test.rs    | 128 +++++++
 .../tests/nested_native_write_compat_test.rs       | 239 ++++++++++++
 ...pk_partial_update_sequence_group_parity_test.rs | 357 +++++++++++++++++
 .../paimon/tests/pk_user_sequence_parity_test.rs   | 426 +++++++++++++++++++++
 crates/paimon/tests/target_file_row_num_test.rs    | 184 +++++++++
 16 files changed, 2211 insertions(+), 166 deletions(-)

diff --git a/crates/paimon/src/spec/core_options.rs 
b/crates/paimon/src/spec/core_options.rs
index 114b89ab..fb517c11 100644
--- a/crates/paimon/src/spec/core_options.rs
+++ b/crates/paimon/src/spec/core_options.rs
@@ -1292,6 +1292,24 @@ impl<'a> CoreOptions<'a> {
             .unwrap_or(DEFAULT_TARGET_FILE_SIZE)
     }
 
+    /// Maximum rows in a newly written data file. Java defaults to 
`Long.MAX_VALUE`.
+    pub fn target_file_row_num(&self) -> crate::Result<i64> {
+        let Some(raw) = self.options.get("target-file-row-num") else {
+            return Ok(i64::MAX);
+        };
+        let rows = raw
+            .parse::<i64>()
+            .map_err(|_| crate::Error::ConfigInvalid {
+                message: format!("target-file-row-num must be a positive 
integer, got '{raw}'"),
+            })?;
+        if rows <= 0 {
+            return Err(crate::Error::ConfigInvalid {
+                message: format!("target-file-row-num must be positive, got 
{rows}"),
+            });
+        }
+        Ok(rows)
+    }
+
     /// Explicit `file.block-size`, in bytes. Formats choose their own default.
     pub(crate) fn file_block_size(&self) -> crate::Result<Option<i64>> {
         self.options
diff --git a/crates/paimon/src/spec/partition_utils.rs 
b/crates/paimon/src/spec/partition_utils.rs
index 60b83eb2..99391ef1 100644
--- a/crates/paimon/src/spec/partition_utils.rs
+++ b/crates/paimon/src/spec/partition_utils.rs
@@ -244,6 +244,16 @@ fn format_partition_value(
             s.to_string()
         }
 
+        DataType::Binary(_) | DataType::VarBinary(_) if !legacy => {
+            // Java's BinaryToStringCastRule wraps the raw bytes in a
+            // BinaryString, whose toString decodes them as UTF-8.
+            let value = decode_java_utf8(row.get_binary(pos)?);
+            if is_java_whitespace_only(&value) {
+                return Ok(default_partition_name.to_string());
+            }
+            value
+        }
+
         DataType::Date(_) => {
             if legacy {
                 // Legacy: field.toString() on the epoch-day Integer → raw int 
value.
@@ -316,6 +326,61 @@ fn format_partition_value(
     Ok(value)
 }
 
+/// Decode as Java's UTF-8 decoder does for `BinaryString.toString()`.
+/// Rust's lossy decoder replaces each byte of a UTF-8 encoded surrogate,
+/// while Java replaces the complete malformed surrogate sequence once.
+fn decode_java_utf8(mut bytes: &[u8]) -> String {
+    let mut decoded = String::with_capacity(bytes.len());
+    loop {
+        match std::str::from_utf8(bytes) {
+            Ok(valid) => {
+                decoded.push_str(valid);
+                break;
+            }
+            Err(error) => {
+                let valid_len = error.valid_up_to();
+                
decoded.push_str(std::str::from_utf8(&bytes[..valid_len]).unwrap());
+                bytes = &bytes[valid_len..];
+
+                let malformed_len =
+                    if bytes.len() >= 2 && bytes[0] == 0xED && 
(0xA0..=0xBF).contains(&bytes[1]) {
+                        // The first two bytes denote a surrogate. Java 
consumes
+                        // its third byte too when it is a continuation byte.
+                        if bytes.get(2).is_some_and(|byte| byte & 0xC0 == 
0x80) {
+                            3
+                        } else {
+                            2
+                        }
+                    } else {
+                        error.error_len().unwrap_or(bytes.len())
+                    };
+                decoded.push('\u{FFFD}');
+                bytes = &bytes[malformed_len..];
+            }
+        }
+    }
+    decoded
+}
+
+/// Java `StringUtils.isNullOrWhitespaceOnly` checks each UTF-16 code unit with
+/// `Character.isWhitespace`; its whitespace set differs from Rust `str::trim`.
+fn is_java_whitespace_only(value: &str) -> bool {
+    value.chars().all(|ch| {
+        matches!(
+            ch,
+            '\u{0009}'..='\u{000D}'
+                | '\u{001C}'..='\u{0020}'
+                | '\u{1680}'
+                | '\u{180E}'
+                | '\u{2000}'..='\u{2006}'
+                | '\u{2008}'..='\u{200A}'
+                | '\u{2028}'..='\u{2029}'
+                | '\u{205F}'
+                | '\u{3000}'
+        )
+    })
+}
+
 /// Format epoch days (since 1970-01-01) to `yyyy-MM-dd`.
 fn format_date(epoch_days: i32) -> String {
     // chrono epoch is 0001-01-01; offset = 719_163 days between 0001-01-01 
and 1970-01-01.
@@ -640,8 +705,12 @@ mod tests {
         }
 
         fn write_string(&mut self, pos: usize, value: &str) {
+            self.write_bytes(pos, value.as_bytes());
+        }
+
+        fn write_bytes(&mut self, pos: usize, value: &[u8]) {
             let var_offset = self.data.len();
-            self.data.extend_from_slice(value.as_bytes());
+            self.data.extend_from_slice(value);
             let len = value.len();
             let encoded = ((var_offset as u64) << 32) | (len as u64);
             let offset = self.field_offset(pos);
@@ -1120,7 +1189,8 @@ mod tests {
 
     #[test]
     fn test_unsupported_types() {
-        // Binary
+        // Legacy binary uses byte[].toString() in Java, which contains an
+        // allocation-specific identity hash and cannot name a stable path.
         assert_single_partition_err(
             "data",
             DataType::Binary(BinaryType::new(10).unwrap()),
@@ -1150,6 +1220,73 @@ mod tests {
         );
     }
 
+    #[test]
+    fn test_binary_partition_uses_utf8_cast_and_path_escaping() {
+        assert_single_partition(
+            "bin",
+            DataType::Binary(BinaryType::new(16).unwrap()),
+            |b| b.write_string(0, "a/b=c"),
+            "bin=a%2Fb%3Dc/",
+            false,
+        );
+        assert_single_partition(
+            "bin",
+            DataType::VarBinary(VarBinaryType::new(16).unwrap()),
+            |b| b.write_string(0, "汉字"),
+            "bin=汉字/",
+            false,
+        );
+        assert_single_partition(
+            "bin",
+            DataType::VarBinary(VarBinaryType::new(16).unwrap()),
+            |b| b.write_string(0, "  "),
+            "bin=__DEFAULT_PARTITION__/",
+            false,
+        );
+    }
+
+    #[test]
+    fn test_binary_partition_matches_java_whitespace() {
+        // JDK 8 Character.isWhitespace includes U+001C and U+180E, but not
+        // U+00A0 or U+2007. Rust str::trim differs for controls and NBSP.
+        for (value, expected) in [
+            (b"\x1c".as_slice(), "bin=__DEFAULT_PARTITION__/"),
+            ("\u{180E}".as_bytes(), "bin=__DEFAULT_PARTITION__/"),
+            ("\u{00A0}".as_bytes(), "bin=\u{00A0}/"),
+            ("\u{2007}".as_bytes(), "bin=\u{2007}/"),
+        ] {
+            assert_single_partition(
+                "bin",
+                DataType::VarBinary(VarBinaryType::new(16).unwrap()),
+                |b| b.write_bytes(0, value),
+                expected,
+                false,
+            );
+        }
+    }
+
+    #[test]
+    fn test_binary_partition_matches_java_malformed_utf8() {
+        // Expected replacements were checked against JDK 8 new String(bytes, 
UTF_8).
+        for (value, expected) in [
+            (b"\xED\xA0\x80".as_slice(), "bin=\u{FFFD}/"),
+            (b"\xED\xA0".as_slice(), "bin=\u{FFFD}/"),
+            (b"\xED\xA0A".as_slice(), "bin=\u{FFFD}A/"),
+            (b"\xED\xA0\xFF".as_slice(), "bin=\u{FFFD}\u{FFFD}/"),
+            (b"\xE0\x80\x80".as_slice(), "bin=\u{FFFD}\u{FFFD}\u{FFFD}/"),
+            (b"\xE2\x82".as_slice(), "bin=\u{FFFD}/"),
+            (b"\xE2(\xA1".as_slice(), "bin=\u{FFFD}(\u{FFFD}/"),
+        ] {
+            assert_single_partition(
+                "bin",
+                DataType::VarBinary(VarBinaryType::new(16).unwrap()),
+                |b| b.write_bytes(0, value),
+                expected,
+                false,
+            );
+        }
+    }
+
     #[test]
     fn test_empty_row_with_partition_keys() {
         let fields = vec![make_field("dt", DataType::Int(IntType::new()))];
diff --git a/crates/paimon/src/spec/schema.rs b/crates/paimon/src/spec/schema.rs
index 6f0e8083..0ca57f5a 100644
--- a/crates/paimon/src/spec/schema.rs
+++ b/crates/paimon/src/spec/schema.rs
@@ -1213,6 +1213,10 @@ impl Schema {
         Self::validate_bucket_keys(options, fields, partition_keys, 
primary_keys)?;
         Self::validate_sequence_field(options, fields, partition_keys, 
primary_keys)?;
         Self::validate_read_batch_size(options)?;
+        let core_options = CoreOptions::new(options);
+        if !core_options.is_format_table() {
+            core_options.target_file_row_num()?;
+        }
         Self::validate_primary_key_vector_index(fields, primary_keys, 
options)?;
         Self::validate_primary_key_full_text_index(fields, primary_keys, 
options)?;
         Ok(())
diff --git a/crates/paimon/src/table/cow_writer.rs 
b/crates/paimon/src/table/cow_writer.rs
index 6a07f49d..b076207e 100644
--- a/crates/paimon/src/table/cow_writer.rs
+++ b/crates/paimon/src/table/cow_writer.rs
@@ -236,6 +236,7 @@ impl CopyOnWriteMergeWriter {
         )?;
 
         let target_file_size = core_options.target_file_size();
+        let target_file_row_num = core_options.target_file_row_num()?;
         let file_compression = core_options.file_compression().to_string();
         let file_compression_zstd_level = 
core_options.file_compression_zstd_level();
         let write_buffer_size = core_options.write_parquet_buffer_size();
@@ -319,7 +320,8 @@ impl CopyOnWriteMergeWriter {
                         Some(0),
                         None,
                         None,
-                    );
+                    )
+                    .with_target_file_row_num(target_file_row_num);
                     writer.write(&rewritten).await?;
                     writer.prepare_commit().await?
                 } else {
diff --git a/crates/paimon/src/table/data_file_writer.rs 
b/crates/paimon/src/table/data_file_writer.rs
index 6039da60..090de04f 100644
--- a/crates/paimon/src/table/data_file_writer.rs
+++ b/crates/paimon/src/table/data_file_writer.rs
@@ -19,7 +19,7 @@
 //! 
[`DataEvolutionPartialWriter`](super::data_evolution_writer::DataEvolutionPartialWriter).
 //!
 //! `DataFileWriter` streams Arrow `RecordBatch`es to Parquet files on storage,
-//! handles file rolling when `target_file_size` is reached, and collects
+//! handles file rolling when the configured file size or row count is 
reached, and collects
 //! [`DataFileMeta`] for the commit path.
 
 use super::data_file_index_writer::{DataFileIndexWriter, FileIndexOptions};
@@ -41,7 +41,7 @@ use tokio::task::JoinSet;
 /// Low-level writer that produces Parquet data files for a single (partition, 
bucket).
 ///
 /// Batches are accumulated into a single `FormatFileWriter` that streams 
directly
-/// to storage. When `target_file_size` is reached the current file is rolled
+/// to storage. When the size or row count target is reached the current file 
is rolled
 /// (closed in the background) and a new one is opened on the next write.
 ///
 /// Call [`prepare_commit`](Self::prepare_commit) to finalize and collect file 
metadata.
@@ -52,6 +52,7 @@ pub(crate) struct DataFileWriter {
     bucket: i32,
     schema_id: i64,
     target_file_size: i64,
+    target_file_row_num: i64,
     file_compression: String,
     file_compression_zstd_level: i32,
     write_buffer_size: i64,
@@ -62,9 +63,10 @@ pub(crate) struct DataFileWriter {
     file_source: Option<i32>,
     first_row_id: Option<i64>,
     write_cols: Option<Vec<String>>,
-    written_files: Vec<DataFileMeta>,
+    written_files: Vec<(usize, DataFileMeta)>,
+    next_file_ordinal: usize,
     /// Background file close tasks spawned during rolling.
-    in_flight_closes: JoinSet<Result<DataFileMeta>>,
+    in_flight_closes: JoinSet<Result<(usize, DataFileMeta)>>,
     /// Current open format writer, lazily created on first write.
     current_writer: Option<Box<dyn FormatFileWriter>>,
     current_file_name: Option<String>,
@@ -105,6 +107,7 @@ impl DataFileWriter {
             bucket,
             schema_id,
             target_file_size,
+            target_file_row_num: i64::MAX,
             file_compression,
             file_compression_zstd_level,
             write_buffer_size,
@@ -116,6 +119,7 @@ impl DataFileWriter {
             first_row_id,
             write_cols,
             written_files: Vec::new(),
+            next_file_ordinal: 0,
             in_flight_closes: JoinSet::new(),
             current_writer: None,
             current_file_name: None,
@@ -132,6 +136,12 @@ impl DataFileWriter {
         self
     }
 
+    pub(crate) fn with_target_file_row_num(mut self, rows: i64) -> Self {
+        debug_assert!(rows > 0);
+        self.target_file_row_num = rows;
+        self
+    }
+
     pub(crate) fn with_resources(mut self, resources: Option<ResourceContext>) 
-> Self {
         self.resources = resources;
         self
@@ -141,7 +151,7 @@ impl DataFileWriter {
         self.resources = resources;
     }
 
-    /// Write a RecordBatch. Rolls to a new file when target size is reached.
+    /// Write a RecordBatch. Rolls when either target size or row count is 
reached.
     pub(crate) async fn write(&mut self, batch: &RecordBatch) -> Result<()> {
         let result = self.write_batch(batch).await;
         if self.index_options.is_some() && result.is_err() {
@@ -165,15 +175,17 @@ impl DataFileWriter {
         }
         self.current_row_count += batch.num_rows() as i64;
 
-        // Roll to a new file if target size is reached — close in background
-        if self.current_writer.as_ref().unwrap().num_bytes() as i64 >= 
self.target_file_size {
+        // Like Java's bundled write, a batch stays intact even if it crosses
+        // the limit. The next batch opens a new file.
+        if self.current_row_count >= self.target_file_row_num
+            || self.current_writer.as_ref().unwrap().num_bytes() as i64 >= 
self.target_file_size
+        {
             self.roll_file();
         }
 
-        // Flush row group if in-progress buffer exceeds write_buffer_size
-        if let Some(w) = self.current_writer.as_mut() {
-            if w.in_progress_size() as i64 >= self.write_buffer_size {
-                w.flush().await?;
+        if let Some(writer) = self.current_writer.as_mut() {
+            if writer.in_progress_size() as i64 >= self.write_buffer_size {
+                writer.flush().await?;
             }
         }
 
@@ -190,7 +202,7 @@ impl DataFileWriter {
             "{}{}-{}.{}",
             self.data_file_prefix,
             uuid::Uuid::new_v4(),
-            self.written_files.len(),
+            self.next_file_ordinal,
             self.file_format,
         );
         let bucket_dir = self.bucket_dir();
@@ -225,7 +237,9 @@ impl DataFileWriter {
     /// Close the current file writer and record the file metadata.
     pub(crate) async fn close_current_file(&mut self) -> Result<()> {
         if let Some(close) = self.take_close() {
-            self.written_files.push(close.await?);
+            let ordinal = self.next_file_ordinal;
+            self.next_file_ordinal += 1;
+            self.written_files.push((ordinal, close.await?));
         }
         Ok(())
     }
@@ -233,7 +247,10 @@ impl DataFileWriter {
     /// Spawn the current writer's close in the background for non-blocking 
rolling.
     fn roll_file(&mut self) {
         if let Some(close) = self.take_close() {
-            self.in_flight_closes.spawn(close);
+            let ordinal = self.next_file_ordinal;
+            self.next_file_ordinal += 1;
+            self.in_flight_closes
+                .spawn(async move { Ok((ordinal, close.await?)) });
         }
     }
 
@@ -297,14 +314,16 @@ impl DataFileWriter {
     async fn finish(&mut self) -> Result<Vec<DataFileMeta>> {
         self.close_current_file().await?;
         while let Some(result) = self.in_flight_closes.join_next().await {
-            let meta = result.map_err(|e| crate::Error::DataInvalid {
+            let file = result.map_err(|e| crate::Error::DataInvalid {
                 message: format!("Background file close task panicked: {e}"),
                 source: None,
             })??;
-            self.written_files.push(meta);
+            self.written_files.push(file);
         }
         self.created_paths.clear();
-        Ok(std::mem::take(&mut self.written_files))
+        let mut files = std::mem::take(&mut self.written_files);
+        files.sort_unstable_by_key(|(ordinal, _)| *ordinal);
+        Ok(files.into_iter().map(|(_, meta)| meta).collect())
     }
 
     pub(super) async fn abort(&mut self) {
@@ -370,6 +389,7 @@ mod tests {
     use super::*;
     use crate::io::{FileIOBuilder, FileIOProvider};
     use crate::spec::{DataType, IntType};
+    use arrow_array::Int32Array;
     use arrow_schema::{DataType as ArrowDataType, Field, Schema};
     use opendal::Operator;
     use std::sync::{Arc, Mutex};
@@ -438,4 +458,48 @@ mod tests {
             .iter()
             .all(|path| !path.contains("//")));
     }
+
+    #[tokio::test]
+    async fn row_limit_rolls_after_whole_batches_and_preserves_file_order() {
+        let mut writer = DataFileWriter::new(
+            FileIOBuilder::new("memory").build().unwrap(),
+            "memory:///row-limit-test".to_string(),
+            String::new(),
+            0,
+            0,
+            i64::MAX,
+            "none".to_string(),
+            0,
+            i64::MAX,
+            "parquet".to_string(),
+            vec![DataField::new(
+                0,
+                "id".to_string(),
+                DataType::Int(IntType::new()),
+            )],
+            HashMap::new(),
+            Some(0),
+            None,
+            None,
+        )
+        .with_target_file_row_num(2);
+        let schema = Arc::new(Schema::new(vec![Field::new(
+            "id",
+            ArrowDataType::Int32,
+            false,
+        )]));
+        for values in [vec![1], vec![2], vec![3, 4, 5], vec![6]] {
+            let batch =
+                RecordBatch::try_new(schema.clone(), 
vec![Arc::new(Int32Array::from(values))])
+                    .unwrap();
+            writer.write(&batch).await.unwrap();
+        }
+
+        let files = writer.prepare_commit().await.unwrap();
+        // Two single-row batches share a file; the three-row batch stays 
intact.
+        assert_eq!(
+            files.iter().map(|file| file.row_count).collect::<Vec<_>>(),
+            vec![2, 3, 1]
+        );
+    }
 }
diff --git a/crates/paimon/src/table/dedicated_format_file_writer.rs 
b/crates/paimon/src/table/dedicated_format_file_writer.rs
index d65b8744..4e6102d2 100644
--- a/crates/paimon/src/table/dedicated_format_file_writer.rs
+++ b/crates/paimon/src/table/dedicated_format_file_writer.rs
@@ -69,6 +69,7 @@ impl AppendDedicatedFormatFileWriter {
         bucket: i32,
         schema_id: i64,
         target_file_size: i64,
+        target_file_row_num: i64,
         blob_target_file_size: i64,
         file_compression: String,
         file_compression_zstd_level: i32,
@@ -120,7 +121,8 @@ impl AppendDedicatedFormatFileWriter {
                         Some(0),
                         None,
                         Some(vec![field.name().to_string()]),
-                    ),
+                    )
+                    .with_target_file_row_num(target_file_row_num),
                     field_name: field.name().to_string(),
                     column_index: idx,
                 });
@@ -158,7 +160,8 @@ impl AppendDedicatedFormatFileWriter {
                         Some(0),
                         None,
                         Some(vector_field_names.clone()),
-                    ),
+                    )
+                    .with_target_file_row_num(target_file_row_num),
                     field_names: vector_field_names,
                     column_indices: vector_column_indices,
                     schema: vector_schema,
@@ -184,7 +187,8 @@ impl AppendDedicatedFormatFileWriter {
             Some(0),
             None,
             Some(normal_field_names),
-        );
+        )
+        .with_target_file_row_num(target_file_row_num);
 
         Self {
             normal_writer,
diff --git a/crates/paimon/src/table/kv_file_writer.rs 
b/crates/paimon/src/table/kv_file_writer.rs
index 721c5b0f..4ddaa63b 100644
--- a/crates/paimon/src/table/kv_file_writer.rs
+++ b/crates/paimon/src/table/kv_file_writer.rs
@@ -61,6 +61,7 @@ use std::sync::Arc;
 pub(crate) struct KeyValueFileWriter {
     file_io: FileIO,
     config: KeyValueWriteConfig,
+    target_file_row_num: usize,
     ignore_delete: bool,
     /// Next sequence number to assign (bucket-local, always auto-incremented).
     next_sequence_number: i64,
@@ -130,8 +131,13 @@ impl KeyValueFileWriter {
         config: KeyValueWriteConfig,
         next_sequence_number: i64,
     ) -> Result<Self> {
-        let ignore_delete = config.merge_engine == MergeEngine::PartialUpdate
-            && CoreOptions::new(&config.table_options).ignore_delete();
+        let core_options = CoreOptions::new(&config.table_options);
+        let target_file_row_num = core_options
+            .target_file_row_num()?
+            .try_into()
+            .unwrap_or(usize::MAX);
+        let ignore_delete =
+            config.merge_engine == MergeEngine::PartialUpdate && 
core_options.ignore_delete();
         if config.merge_engine == MergeEngine::PartialUpdate {
             let partial_update = 
PartialUpdateConfig::new(&config.table_options);
             partial_update.validate_write_mode(true, &config.table_name)?;
@@ -167,6 +173,7 @@ impl KeyValueFileWriter {
         Ok(Self {
             file_io,
             config,
+            target_file_row_num,
             ignore_delete,
             next_sequence_number,
             buffer: Vec::new(),
@@ -268,11 +275,13 @@ impl KeyValueFileWriter {
                 }),
             });
         }
+        let user_sequence_descending =
+            
!CoreOptions::new(&self.config.table_options).sequence_field_sort_order_is_ascending();
         for &idx in &self.config.sequence_field_indices {
             sort_columns.push(SortColumn {
                 values: combined.column(idx).clone(),
                 options: Some(SortOptions {
-                    descending: false,
+                    descending: user_sequence_descending,
                     nulls_first: true,
                 }),
             });
@@ -290,7 +299,7 @@ impl KeyValueFileWriter {
                 source: None,
             })?;
 
-        // After sorting by PK + seq fields + auto-seq (all ascending), merge
+        // After sorting by PK + configured user sequence + auto-seq, merge
         // each key group down to one row, mirroring Java's
         // MergeTreeWriter#flushWriteBuffer (the write buffer runs the merge
         // function before any file is written, so a flushed file never holds
@@ -324,61 +333,77 @@ impl KeyValueFileWriter {
             }
         };
 
-        let data_delete_row_count = 
Self::indexed_delete_row_count(&data_batch, &data_indices)?;
-        let changelog_delete_row_count = if self.config.input_changelog {
-            Some(Self::indexed_delete_row_count(&combined, &sorted_indices)?)
-        } else {
-            None
-        };
-
-        // Java derives file sequence bounds from emitted rows; allocation 
still
-        // advances over all buffered inputs, including rows folded away.
-        let output_sequences = 
data_seq.as_any().downcast_ref::<Int64Array>().unwrap();
-        let (min_output_seq, max_output_seq) = data_indices
-            .values()
-            .iter()
-            .map(|&idx| output_sequences.value(idx as usize))
-            .fold((i64::MAX, i64::MIN), |(min, max), seq| {
-                (min.min(seq), max.max(seq))
-            });
-        let data_file = self
-            .write_indexed_file(
-                &data_batch,
-                data_seq.as_ref(),
-                &data_indices,
-                IndexedFileWrite {
-                    is_changelog: false,
-                    file_prefix: &self.config.data_file_prefix,
-                    file_ordinal: self.written_files.len(),
-                    file_format: &self.config.file_format,
-                    file_compression: &self.config.file_compression,
-                    min_sequence_number: min_output_seq,
-                    max_sequence_number: max_output_seq,
-                    delete_row_count: data_delete_row_count,
-                },
-            )
-            .await?;
-        self.written_files.push(data_file);
-
-        if let Some(delete_row_count) = changelog_delete_row_count {
-            let changelog_file = self
+        // The sorted output is already materialized in FLUSH_CHUNK_ROWS 
batches.
+        // Use the row target as an upper bound for each emitted batch and keep
+        // every file's key bounds, sequence bounds, and index local to its 
rows.
+        let data_sequences = 
data_seq.as_any().downcast_ref::<Int64Array>().unwrap();
+        for offset in 
(0..data_indices.len()).step_by(self.target_file_row_num) {
+            let len = self.target_file_row_num.min(data_indices.len() - 
offset);
+            let file_indices = data_indices.slice(offset, len);
+            let (min_sequence_number, max_sequence_number) = file_indices
+                .values()
+                .iter()
+                .map(|&idx| data_sequences.value(idx as usize))
+                .fold((i64::MAX, i64::MIN), |(min, max), seq| {
+                    (min.min(seq), max.max(seq))
+                });
+            let file = self
                 .write_indexed_file(
-                    &combined,
-                    seq_array.as_ref(),
-                    &sorted_indices,
+                    &data_batch,
+                    data_seq.as_ref(),
+                    &file_indices,
                     IndexedFileWrite {
-                        is_changelog: true,
-                        file_prefix: &self.config.changelog_file_prefix,
-                        file_ordinal: self.written_changelog_files.len(),
-                        file_format: &self.config.changelog_file_format,
-                        file_compression: 
&self.config.changelog_file_compression,
-                        min_sequence_number: start_seq,
-                        max_sequence_number: end_seq,
-                        delete_row_count,
+                        is_changelog: false,
+                        file_prefix: &self.config.data_file_prefix,
+                        file_ordinal: self.written_files.len(),
+                        file_format: &self.config.file_format,
+                        file_compression: &self.config.file_compression,
+                        min_sequence_number,
+                        max_sequence_number,
+                        delete_row_count: Self::indexed_delete_row_count(
+                            &data_batch,
+                            &file_indices,
+                        )?,
                     },
                 )
                 .await?;
-            self.written_changelog_files.push(changelog_file);
+            self.written_files.push(file);
+        }
+
+        if self.config.input_changelog {
+            let input_sequences = 
seq_array.as_any().downcast_ref::<Int64Array>().unwrap();
+            for offset in 
(0..sorted_indices.len()).step_by(self.target_file_row_num) {
+                let len = self.target_file_row_num.min(sorted_indices.len() - 
offset);
+                let file_indices = sorted_indices.slice(offset, len);
+                let (min_sequence_number, max_sequence_number) = file_indices
+                    .values()
+                    .iter()
+                    .map(|&idx| input_sequences.value(idx as usize))
+                    .fold((i64::MAX, i64::MIN), |(min, max), seq| {
+                        (min.min(seq), max.max(seq))
+                    });
+                let file = self
+                    .write_indexed_file(
+                        &combined,
+                        seq_array.as_ref(),
+                        &file_indices,
+                        IndexedFileWrite {
+                            is_changelog: true,
+                            file_prefix: &self.config.changelog_file_prefix,
+                            file_ordinal: self.written_changelog_files.len(),
+                            file_format: &self.config.changelog_file_format,
+                            file_compression: 
&self.config.changelog_file_compression,
+                            min_sequence_number,
+                            max_sequence_number,
+                            delete_row_count: Self::indexed_delete_row_count(
+                                &combined,
+                                &file_indices,
+                            )?,
+                        },
+                    )
+                    .await?;
+                self.written_changelog_files.push(file);
+            }
         }
         Ok(())
     }
@@ -839,7 +864,7 @@ impl KeyValueFileWriter {
                 row_idx: idx as usize,
                 // Ordering is already established by the write buffer's Arrow 
sort.
                 sequence_number: 0,
-                user_sequences: Vec::new(),
+                user_sequence: None,
                 value_kind: value_kinds
                     .filter(|kinds| kinds.is_valid(idx as usize))
                     .map_or(0, |kinds| kinds.value(idx as usize)),
@@ -945,11 +970,15 @@ impl KeyValueFileWriter {
         let rows: Vec<_> = sorted_indices
             .values()
             .iter()
-            .map(|&idx| MergeRow {
+            .enumerate()
+            .map(|(sorted_rank, &idx)| MergeRow {
                 batch_idx: 0,
                 row_idx: idx as usize,
-                sequence_number: 0,
-                user_sequences: Vec::new(),
+                // The write buffer has already sorted by user sequence and
+                // arrival sequence. Preserve that order when the shared merge
+                // function sorts MergeRows again.
+                sequence_number: sorted_rank as i64,
+                user_sequence: None,
                 value_kind: value_kinds
                     .filter(|kinds| kinds.is_valid(idx as usize))
                     .map_or(0, |kinds| kinds.value(idx as usize)),
@@ -1269,6 +1298,49 @@ mod tests {
         .unwrap()
     }
 
+    #[tokio::test]
+    async fn 
target_file_row_num_rolls_data_and_changelog_with_local_metadata() {
+        let schema = Arc::new(ArrowSchema::new(vec![
+            ArrowField::new("id", ArrowDataType::Int32, false),
+            ArrowField::new("seq", ArrowDataType::Int64, false),
+            ArrowField::new("value", ArrowDataType::Int32, true),
+        ]));
+        let batch = RecordBatch::try_new(
+            schema,
+            vec![
+                Arc::new(Int32Array::from(vec![5, 1, 4, 2, 3])),
+                Arc::new(Int64Array::from(vec![50, 10, 40, 20, 30])),
+                Arc::new(Int32Array::from(vec![50, 10, 40, 20, 30])),
+            ],
+        )
+        .unwrap();
+        let mut config = test_write_config(MergeEngine::Deduplicate);
+        config.input_changelog = true;
+        config.write_buffer_size = i64::MAX;
+        config
+            .table_options
+            .insert("target-file-row-num".into(), "2".into());
+        let mut writer =
+            
KeyValueFileWriter::new(FileIOBuilder::new("memory").build().unwrap(), config, 
0)
+                .unwrap();
+        writer.write(&batch).await.unwrap();
+        let prepared = writer.prepare_commit().await.unwrap();
+
+        for files in [&prepared.data_files, &prepared.changelog_files] {
+            assert_eq!(
+                files.iter().map(|file| file.row_count).collect::<Vec<_>>(),
+                vec![2, 2, 1]
+            );
+            assert_eq!(
+                files
+                    .iter()
+                    .map(|file| (file.min_sequence_number, 
file.max_sequence_number))
+                    .collect::<Vec<_>>(),
+                vec![(1, 3), (2, 4), (0, 0)]
+            );
+        }
+    }
+
     #[tokio::test]
     async fn test_pk_value_stats_use_emitted_rows_and_logical_columns() {
         let schema = Arc::new(ArrowSchema::new(vec![
@@ -2088,6 +2160,10 @@ mod tests {
             .as_any()
             .downcast_ref::<Int64Array>()
             .unwrap();
+        let user_converter = 
RowConverter::new(vec![SortField::new(ArrowDataType::Int64)]).unwrap();
+        let user_sequences = user_converter
+            .convert_columns(&[Arc::new(seq_col.clone())])
+            .unwrap();
 
         for (group_idx, group_rows) in [vec![0usize, 2, 3], vec![1usize, 
4]].iter().enumerate() {
             let rows: Vec<MergeRow> = group_rows
@@ -2097,7 +2173,7 @@ mod tests {
                     row_idx,
                     sequence_number: seq_values[row_idx],
                     value_kind: 0,
-                    user_sequences: vec![Some(seq_col.value(row_idx) as i128)],
+                    user_sequence: Some(user_sequences.row(row_idx).owned()),
                 })
                 .collect();
             let result = merge_fn.merge(&rows, &buffer, &identity, 
&schema).unwrap();
diff --git a/crates/paimon/src/table/mod.rs b/crates/paimon/src/table/mod.rs
index 0b343e36..1b918cca 100644
--- a/crates/paimon/src/table/mod.rs
+++ b/crates/paimon/src/table/mod.rs
@@ -137,6 +137,7 @@ pub(crate) mod vector_search_result;
 #[cfg(test)]
 mod vector_search_test_utils;
 mod vindex_index_build_builder;
+mod write_batch_normalize;
 mod write_builder;
 
 use crate::Result;
diff --git a/crates/paimon/src/table/sort_merge.rs 
b/crates/paimon/src/table/sort_merge.rs
index 688bccf1..3c130745 100644
--- a/crates/paimon/src/table/sort_merge.rs
+++ b/crates/paimon/src/table/sort_merge.rs
@@ -32,7 +32,7 @@ use crate::table::ArrowRecordBatchStream;
 use crate::Error;
 use arrow_array::{new_null_array, ArrayRef, Int64Array, Int8Array, 
RecordBatch};
 use arrow_ord::ord::make_comparator;
-use arrow_row::{RowConverter, Rows, SortField};
+use arrow_row::{OwnedRow, RowConverter, Rows, SortField};
 use arrow_schema::{SchemaRef, SortOptions};
 use arrow_select::interleave::interleave;
 use async_stream::try_stream;
@@ -78,8 +78,9 @@ pub(crate) struct MergeRow {
     pub row_idx: usize,
     pub sequence_number: i64,
     pub value_kind: i8,
-    /// User-defined sequence values from `sequence.field` (empty if not 
configured).
-    pub user_sequences: Vec<Option<i128>>,
+    /// Row-encoded user-defined sequence fields, with Java's sort direction
+    /// and nulls-first ordering. `None` when `sequence.field` is not set.
+    pub user_sequence: Option<OwnedRow>,
 }
 
 #[cfg(test)]
@@ -152,21 +153,19 @@ pub(crate) trait MergeFunction: Send + Sync {
 }
 
 /// Deduplicate merge: keeps the row with the highest sequence.
-/// When `sequence.field` is configured (one or more fields), compares user
-/// sequences lexicographically first, then falls back to system
-/// `_SEQUENCE_NUMBER` as tie-breaker.
+/// When `sequence.field` is configured, compares its typed Arrow row encoding
+/// first, then falls back to system `_SEQUENCE_NUMBER` as tie-breaker.
 /// When sequence numbers are equal, keeps the last-added row 
(last-writer-wins).
 /// Filters out DELETE and UPDATE_BEFORE rows.
 pub(crate) struct DeduplicateMergeFunction;
 
 pub(super) fn compare_sequence_order(lhs: &MergeRow, rhs: &MergeRow) -> 
Ordering {
-    match (lhs.user_sequences.is_empty(), rhs.user_sequences.is_empty()) {
-        (false, false) => lhs
-            .user_sequences
-            .cmp(&rhs.user_sequences)
-            .then_with(|| lhs.sequence_number.cmp(&rhs.sequence_number)),
-        _ => lhs.sequence_number.cmp(&rhs.sequence_number),
-    }
+    lhs.user_sequence
+        .cmp(&rhs.user_sequence)
+        .then_with(|| lhs.sequence_number.cmp(&rhs.sequence_number))
+        // Java SortMergeReader compares isAdd() after equal sequence numbers:
+        // retracts come first, adds last.
+        .then_with(|| matches!(lhs.value_kind, 0 | 
2).cmp(&matches!(rhs.value_kind, 0 | 2)))
 }
 
 impl MergeFunction for DeduplicateMergeFunction {
@@ -1026,6 +1025,8 @@ struct SortMergeCursor {
     batch: RecordBatch,
     /// Row-encoded keys for the current batch (via arrow-row).
     rows: Rows,
+    /// Typed user-defined sequence fields for the current batch.
+    user_rows: Option<Rows>,
     offset: usize,
 }
 
@@ -1059,54 +1060,10 @@ impl SortMergeCursor {
         }
     }
 
-    /// Read the user-defined sequence field value (cast to i64 for ordering).
-    /// Returns None if the column is NULL at this row.
-    ///
-    /// Supports the same types as Java Paimon's `UserDefinedSeqComparator`:
-    /// TinyInt, SmallInt, Int, BigInt, Timestamp, Date, Decimal.
-    fn user_sequence(&self, user_seq_index: usize) -> Option<i128> {
-        let col = self.batch.column(user_seq_index);
-        if col.is_null(self.offset) {
-            return None;
-        }
-        use arrow_array::*;
-        let any = col.as_any();
-        if let Some(arr) = any.downcast_ref::<Int64Array>() {
-            return Some(arr.value(self.offset) as i128);
-        }
-        if let Some(arr) = any.downcast_ref::<Int32Array>() {
-            return Some(arr.value(self.offset) as i128);
-        }
-        if let Some(arr) = any.downcast_ref::<Int16Array>() {
-            return Some(arr.value(self.offset) as i128);
-        }
-        if let Some(arr) = any.downcast_ref::<Int8Array>() {
-            return Some(arr.value(self.offset) as i128);
-        }
-        // Timestamps are stored as i64 internally (micros, millis, seconds, 
nanos).
-        if let Some(arr) = any.downcast_ref::<TimestampMicrosecondArray>() {
-            return Some(arr.value(self.offset) as i128);
-        }
-        if let Some(arr) = any.downcast_ref::<TimestampMillisecondArray>() {
-            return Some(arr.value(self.offset) as i128);
-        }
-        if let Some(arr) = any.downcast_ref::<TimestampNanosecondArray>() {
-            return Some(arr.value(self.offset) as i128);
-        }
-        if let Some(arr) = any.downcast_ref::<TimestampSecondArray>() {
-            return Some(arr.value(self.offset) as i128);
-        }
-        if let Some(arr) = any.downcast_ref::<Date32Array>() {
-            return Some(arr.value(self.offset) as i128);
-        }
-        if let Some(arr) = any.downcast_ref::<Date64Array>() {
-            return Some(arr.value(self.offset) as i128);
-        }
-        // Decimal128: use raw i128 value for ordering (same precision/scale 
within a column).
-        if let Some(arr) = any.downcast_ref::<Decimal128Array>() {
-            return Some(arr.value(self.offset));
-        }
-        None
+    fn user_sequence(&self) -> Option<OwnedRow> {
+        self.user_rows
+            .as_ref()
+            .map(|rows| rows.row(self.offset).owned())
     }
 }
 
@@ -1263,14 +1220,36 @@ impl SortMergeReaderBuilder {
             source: Some(Box::new(e)),
         })?;
 
+        let user_sequence_converter = if self.user_sequence_indices.is_empty() 
{
+            None
+        } else {
+            let fields = self
+                .user_sequence_indices
+                .iter()
+                .map(|&idx| {
+                    SortField::new_with_options(
+                        self.input_schema.field(idx).data_type().clone(),
+                        SortOptions {
+                            descending: self.user_sequence_descending,
+                            nulls_first: true,
+                        },
+                    )
+                })
+                .collect();
+            Some(RowConverter::new(fields).map_err(|e| Error::DataInvalid {
+                message: format!("Unsupported user sequence field type: {e}"),
+                source: Some(Box::new(e)),
+            })?)
+        };
+
         sort_merge_stream(
             self.streams,
             row_converter,
+            user_sequence_converter,
             self.key_indices,
             self.seq_index,
             self.value_kind_index,
             self.user_sequence_indices,
-            self.user_sequence_descending,
             self.value_indices,
             self.output_schema,
             self.merge_function,
@@ -1297,6 +1276,27 @@ fn convert_batch_keys(
         })
 }
 
+fn convert_batch_user_sequences(
+    batch: &RecordBatch,
+    indices: &[usize],
+    converter: Option<&mut RowConverter>,
+) -> crate::Result<Option<Rows>> {
+    let Some(converter) = converter else {
+        return Ok(None);
+    };
+    let columns = indices
+        .iter()
+        .map(|&idx| batch.column(idx).clone())
+        .collect::<Vec<_>>();
+    converter
+        .convert_columns(&columns)
+        .map(Some)
+        .map_err(|e| Error::DataInvalid {
+            message: format!("Failed to encode user sequence fields: {e}"),
+            source: Some(Box::new(e)),
+        })
+}
+
 /// Compare two cursors by their current key. `None` cursors are treated as
 /// greater than any value (exhausted streams sink to the bottom).
 fn compare_cursors(cursors: &[Option<SortMergeCursor>], a: usize, b: usize) -> 
Ordering {
@@ -1318,11 +1318,11 @@ fn compare_cursors(cursors: &[Option<SortMergeCursor>], 
a: usize, b: usize) -> O
 fn sort_merge_stream(
     mut streams: Vec<ArrowRecordBatchStream>,
     mut row_converter: RowConverter,
+    mut user_sequence_converter: Option<RowConverter>,
     key_indices: Vec<usize>,
     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>,
@@ -1351,7 +1351,12 @@ fn sort_merge_stream(
                 let batch = batch_result?;
                 if batch.num_rows() > 0 {
                     let rows = convert_batch_keys(&batch, &key_indices, &mut 
row_converter)?;
-                    cursors.push(Some(SortMergeCursor { batch, rows, offset: 0 
}));
+                    let user_rows = convert_batch_user_sequences(
+                        &batch,
+                        &user_sequence_indices,
+                        user_sequence_converter.as_mut(),
+                    )?;
+                    cursors.push(Some(SortMergeCursor { batch, rows, 
user_rows, offset: 0 }));
                     found = true;
                     break;
                 }
@@ -1420,11 +1425,7 @@ 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).map(|value| {
-                                if user_sequence_descending { -value } else { 
value }
-                            })
-                        }).collect(),
+                        user_sequence: cursor.user_sequence(),
                     });
                 }
 
@@ -1440,10 +1441,15 @@ fn sort_merge_stream(
                             let batch = batch_result?;
                             if batch.num_rows() > 0 {
                                 let rows = convert_batch_keys(&batch, 
&key_indices, &mut row_converter)?;
+                                let user_rows = convert_batch_user_sequences(
+                                    &batch,
+                                    &user_sequence_indices,
+                                    user_sequence_converter.as_mut(),
+                                )?;
                                 let buf_idx = batch_buffer.len();
                                 
batch_buffer.push(BufferedBatch::Source(batch.clone()));
                                 stream_batch_idx[current_winner] = 
Some(buf_idx);
-                                cursors[current_winner] = Some(SortMergeCursor 
{ batch, rows, offset: 0 });
+                                cursors[current_winner] = Some(SortMergeCursor 
{ batch, rows, user_rows, offset: 0 });
                                 break;
                             }
                         }
@@ -1613,7 +1619,7 @@ mod tests {
                 row_idx,
                 sequence_number,
                 value_kind: 0,
-                user_sequences: vec![],
+                user_sequence: None,
             })
             .collect();
         let result = FirstRowMergeFunction {
@@ -1639,14 +1645,14 @@ mod tests {
                     row_idx: 0,
                     sequence_number: 1,
                     value_kind: 0,
-                    user_sequences: vec![],
+                    user_sequence: None,
                 },
                 MergeRow {
                     batch_idx: 0,
                     row_idx: 1,
                     sequence_number: 2,
                     value_kind: kind,
-                    user_sequences: vec![],
+                    user_sequence: None,
                 },
             ];
             let merge = FirstRowMergeFunction {
@@ -2142,6 +2148,133 @@ mod tests {
         assert_eq!(values.value(0), "low");
     }
 
+    #[tokio::test]
+    async fn test_string_user_sequence_across_streams_in_both_orders() {
+        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("seq", DataType::Utf8, true),
+            Field::new("value", DataType::Utf8, false),
+        ]));
+        let output_schema = Arc::new(Schema::new(vec![
+            Field::new("pk", DataType::Int32, false),
+            Field::new("seq", DataType::Utf8, true),
+            Field::new("value", DataType::Utf8, false),
+        ]));
+        for (descending, expected) in [(false, "high"), (true, "low")] {
+            let streams = [
+                (Some("z"), 1, "high"),
+                (Some("a"), 2, "low"),
+                (None, 3, "null"),
+            ]
+            .into_iter()
+            .map(|(seq, number, value)| {
+                let batch = RecordBatch::try_new(
+                    schema.clone(),
+                    vec![
+                        Arc::new(Int32Array::from(vec![1])),
+                        Arc::new(Int64Array::from(vec![number])),
+                        Arc::new(Int8Array::from(vec![0])),
+                        Arc::new(StringArray::from(vec![seq])),
+                        Arc::new(StringArray::from(vec![value])),
+                    ],
+                )
+                .unwrap();
+                stream_from_batches(vec![batch])
+            })
+            .collect();
+            let batches = SortMergeReaderBuilder::new(
+                streams,
+                schema.clone(),
+                vec![0],
+                1,
+                2,
+                vec![3],
+                vec![3, 4],
+                output_schema.clone(),
+                Box::new(DeduplicateMergeFunction),
+            )
+            .with_user_sequence_descending(descending)
+            .build()
+            .unwrap()
+            .try_collect::<Vec<_>>()
+            .await
+            .unwrap();
+            assert_eq!(batches.len(), 1);
+            let values = batches[0]
+                .column_by_name("value")
+                .unwrap()
+                .as_any()
+                .downcast_ref::<StringArray>()
+                .unwrap();
+            assert_eq!(values.value(0), expected, "descending={descending}");
+        }
+    }
+
+    #[tokio::test]
+    async fn test_equal_sequence_places_retract_before_add() {
+        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("seq", DataType::Utf8, false),
+            Field::new("value", DataType::Utf8, false),
+        ]));
+        let output_schema = Arc::new(Schema::new(vec![
+            Field::new("pk", DataType::Int32, false),
+            Field::new("seq", DataType::Utf8, false),
+            Field::new("value", DataType::Utf8, false),
+        ]));
+        for kinds in [[3, 0], [0, 3]] {
+            let streams = kinds
+                .into_iter()
+                .map(|kind| {
+                    let batch = RecordBatch::try_new(
+                        schema.clone(),
+                        vec![
+                            Arc::new(Int32Array::from(vec![1])),
+                            Arc::new(Int64Array::from(vec![5])),
+                            Arc::new(Int8Array::from(vec![kind])),
+                            Arc::new(StringArray::from(vec!["same"])),
+                            Arc::new(StringArray::from(vec![if kind == 0 {
+                                "add"
+                            } else {
+                                "delete"
+                            }])),
+                        ],
+                    )
+                    .unwrap();
+                    stream_from_batches(vec![batch])
+                })
+                .collect();
+            let batches = SortMergeReaderBuilder::new(
+                streams,
+                schema.clone(),
+                vec![0],
+                1,
+                2,
+                vec![3],
+                vec![3, 4],
+                output_schema.clone(),
+                Box::new(DeduplicateMergeFunction),
+            )
+            .build()
+            .unwrap()
+            .try_collect::<Vec<_>>()
+            .await
+            .unwrap();
+            assert_eq!(batches.len(), 1);
+            let values = batches[0]
+                .column_by_name("value")
+                .unwrap()
+                .as_any()
+                .downcast_ref::<StringArray>()
+                .unwrap();
+            assert_eq!(values.value(0), "add");
+        }
+    }
+
     #[tokio::test]
     async fn test_delete_row_filtered() {
         let schema = make_schema();
diff --git a/crates/paimon/src/table/table_write.rs 
b/crates/paimon/src/table/table_write.rs
index fb8d128e..b50a7501 100644
--- a/crates/paimon/src/table/table_write.rs
+++ b/crates/paimon/src/table/table_write.rs
@@ -45,9 +45,10 @@ use crate::table::partition_filter::PartitionFilter;
 use crate::table::postpone_file_writer::{PostponeFileWriter, 
PostponeWriteConfig};
 use crate::table::prepared_files::PreparedFiles;
 use crate::table::row_kind_generator::RowKindGenerator;
+use crate::table::write_batch_normalize::normalize_write_array;
 use crate::table::{Snapshot, SnapshotManager, Table, TableScan};
 use crate::Result;
-use arrow_array::RecordBatch;
+use arrow_array::{ArrayRef, RecordBatch};
 use std::collections::{HashMap, HashSet};
 use std::sync::Arc;
 
@@ -133,6 +134,7 @@ pub struct TableWrite {
     partition_keys: Vec<String>,
     schema_id: i64,
     target_file_size: i64,
+    target_file_row_num: i64,
     blob_target_file_size: i64,
     vector_target_file_size: i64,
     file_compression: String,
@@ -201,6 +203,7 @@ impl TableWrite {
             partition_keys: schema.partition_keys().to_vec(),
             schema_id: schema.id(),
             target_file_size: 0,
+            target_file_row_num: i64::MAX,
             blob_target_file_size: 0,
             vector_target_file_size: 0,
             file_compression: String::new(),
@@ -315,6 +318,7 @@ impl TableWrite {
             });
         }
         let target_file_size = core_options.target_file_size();
+        let target_file_row_num = core_options.target_file_row_num()?;
         let blob_target_file_size = core_options.blob_target_file_size();
         let vector_target_file_size = core_options.vector_target_file_size();
         let file_compression = core_options.file_compression().to_string();
@@ -465,6 +469,7 @@ impl TableWrite {
             partition_keys,
             schema_id: schema.id(),
             target_file_size,
+            target_file_row_num,
             blob_target_file_size,
             vector_target_file_size,
             file_compression,
@@ -602,11 +607,11 @@ impl TableWrite {
     }
 
     pub(super) fn normalize_write_batch(&self, batch: &RecordBatch) -> 
Result<Option<RecordBatch>> {
-        self.validate_write_batch_schema(batch)?;
+        let batch = self.validate_write_batch_schema(batch)?;
         if batch.num_rows() == 0 {
             return Ok(None);
         }
-        let batch = self.enrich_rowkind_batch(batch)?;
+        let batch = self.enrich_rowkind_batch(&batch)?;
         Ok((batch.num_rows() != 0).then_some(batch))
     }
 
@@ -619,7 +624,7 @@ impl TableWrite {
         self.write_bucket(partition, bucket, batch).await
     }
 
-    fn validate_write_batch_schema(&self, batch: &RecordBatch) -> Result<()> {
+    fn validate_write_batch_schema(&self, batch: &RecordBatch) -> 
Result<RecordBatch> {
         let expected_schema = &self.write_schema;
         let actual_schema = batch.schema();
         let table_field_count = expected_schema.fields().len();
@@ -689,19 +694,20 @@ impl TableWrite {
             }
         }
 
+        let mut normalized_columns: Vec<ArrayRef> = 
Vec::with_capacity(actual_field_count);
         for (index, expected_field) in 
expected_schema.fields().iter().enumerate() {
             let actual_field = actual_schema.field(index);
-            if actual_field.data_type() != expected_field.data_type() {
-                return Err(crate::Error::DataInvalid {
+            let column = normalize_write_array(batch.column(index), 
expected_field.data_type())
+                .map_err(|error| crate::Error::DataInvalid {
                     message: format!(
-                        "write batch schema data type mismatch for field '{}' 
at index {index}: expected {:?}, actual {:?}",
+                        "write batch schema data type mismatch for field '{}' 
at index {index}: expected {:?}, actual {:?}: {error}",
                         expected_field.name(),
                         expected_field.data_type(),
                         actual_field.data_type()
                     ),
-                    source: None,
-                });
-            }
+                    source: Some(Box::new(error)),
+                })?;
+            normalized_columns.push(column);
         }
         if includes_value_kind {
             let actual_field = actual_schema.field(table_field_count);
@@ -715,9 +721,25 @@ impl TableWrite {
                     source: None,
                 });
             }
+            normalized_columns.push(batch.column(table_field_count).clone());
         }
 
-        Ok(())
+        let schema = if includes_value_kind {
+            let mut fields = 
expected_schema.fields().iter().cloned().collect::<Vec<_>>();
+            fields.push(actual_schema.fields()[table_field_count].clone());
+            Arc::new(arrow_schema::Schema::new_with_metadata(
+                fields,
+                expected_schema.metadata().clone(),
+            ))
+        } else {
+            expected_schema.clone()
+        };
+        RecordBatch::try_new(schema, normalized_columns).map_err(|error| {
+            crate::Error::DataInvalid {
+                message: format!("Failed to normalize write batch schema: 
{error}"),
+                source: Some(Box::new(error)),
+            }
+        })
     }
 
     /// Group rows by (partition_bytes, bucket) and return sub-batches.
@@ -1114,6 +1136,7 @@ impl TableWrite {
                     bucket,
                     self.schema_id,
                     self.target_file_size,
+                    self.target_file_row_num,
                     self.blob_target_file_size,
                     self.file_compression.clone(),
                     self.file_compression_zstd_level,
@@ -1149,6 +1172,7 @@ impl TableWrite {
                     None,
                 )
                 .with_file_index(self.file_index_options.clone())
+                .with_target_file_row_num(self.target_file_row_num)
                 .with_resources(self.resources.clone()),
             ))
         }
diff --git a/crates/paimon/src/table/write_batch_normalize.rs 
b/crates/paimon/src/table/write_batch_normalize.rs
new file mode 100644
index 00000000..01529d9b
--- /dev/null
+++ b/crates/paimon/src/table/write_batch_normalize.rs
@@ -0,0 +1,248 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Convert equivalent Arrow write layouts to the table's canonical schema.
+//!
+//! PyArrow calls list children `item`; Paimon's Arrow schema calls them
+//! `element`. Parquet persists child names, so accepting the input schema
+//! without rebuilding the arrays would produce a file that later projections
+//! cannot reliably read by the table schema. The buffers and values are 
shared;
+//! only nested array wrappers are rebuilt.
+
+use crate::{Error, Result};
+use arrow_array::{
+    Array, ArrayRef, BinaryArray, FixedSizeBinaryArray, ListArray, MapArray, 
StructArray,
+};
+use arrow_schema::DataType;
+use std::sync::Arc;
+
+pub(super) fn normalize_write_array(array: &ArrayRef, expected: &DataType) -> 
Result<ArrayRef> {
+    if array.data_type() == expected {
+        return Ok(array.clone());
+    }
+
+    match (array.data_type(), expected) {
+        (DataType::FixedSizeBinary(_), DataType::Binary) => {
+            let fixed = array
+                .as_any()
+                .downcast_ref::<FixedSizeBinaryArray>()
+                .unwrap();
+            Ok(Arc::new(BinaryArray::from_iter((0..fixed.len()).map(
+                |row| (!fixed.is_null(row)).then(|| fixed.value(row)),
+            ))))
+        }
+        (DataType::List(_), DataType::List(expected_child)) => {
+            let list = array.as_any().downcast_ref::<ListArray>().unwrap();
+            let values = normalize_write_array(list.values(), 
expected_child.data_type())?;
+            let normalized = ListArray::try_new(
+                expected_child.clone(),
+                list.offsets().clone(),
+                values,
+                list.nulls().cloned(),
+            )
+            .map_err(normalization_error)?;
+            Ok(Arc::new(normalized))
+        }
+        (DataType::Struct(actual_fields), DataType::Struct(expected_fields))
+            if actual_fields.len() == expected_fields.len() =>
+        {
+            let row = array.as_any().downcast_ref::<StructArray>().unwrap();
+            let columns = actual_fields
+                .iter()
+                .zip(expected_fields)
+                .enumerate()
+                .map(|(index, (actual, expected))| {
+                    if actual.name() != expected.name() {
+                        return Err(Error::DataInvalid {
+                            message: format!(
+                                "Nested ROW field name mismatch at index 
{index}: expected '{}', actual '{}'",
+                                expected.name(),
+                                actual.name()
+                            ),
+                            source: None,
+                        });
+                    }
+                    normalize_write_array(row.column(index), 
expected.data_type())
+                })
+                .collect::<Result<Vec<_>>>()?;
+            let normalized =
+                StructArray::try_new(expected_fields.clone(), columns, 
row.nulls().cloned())
+                    .map_err(normalization_error)?;
+            Ok(Arc::new(normalized))
+        }
+        (DataType::Map(_, _), DataType::Map(expected_entries, 
expected_ordered)) => {
+            let map = array.as_any().downcast_ref::<MapArray>().unwrap();
+            let entries: ArrayRef = Arc::new(map.entries().clone());
+            let entries = normalize_write_array(&entries, 
expected_entries.data_type())?;
+            let entries = entries
+                .as_any()
+                .downcast_ref::<StructArray>()
+                .unwrap()
+                .clone();
+            let normalized = MapArray::try_new(
+                expected_entries.clone(),
+                map.offsets().clone(),
+                entries,
+                map.nulls().cloned(),
+                *expected_ordered,
+            )
+            .map_err(normalization_error)?;
+            Ok(Arc::new(normalized))
+        }
+        _ => Err(Error::DataInvalid {
+            message: format!(
+                "Arrow write type {:?} is incompatible with table type 
{expected:?}",
+                array.data_type()
+            ),
+            source: None,
+        }),
+    }
+}
+
+fn normalization_error(error: arrow_schema::ArrowError) -> Error {
+    Error::DataInvalid {
+        message: format!("Invalid nested Arrow write array: {error}"),
+        source: Some(Box::new(error)),
+    }
+}
+
+#[cfg(test)]
+mod tests {
+    use super::*;
+    use arrow_array::{Int32Array, StringArray};
+    use arrow_buffer::{OffsetBuffer, ScalarBuffer};
+    use arrow_schema::Field;
+
+    #[test]
+    fn fixed_binary_conversion_preserves_nulls_and_bytes() {
+        let input: ArrayRef = Arc::new(
+            FixedSizeBinaryArray::try_from_sparse_iter_with_size(
+                [Some(b"ab".as_slice()), None, 
Some(b"cd".as_slice())].into_iter(),
+                2,
+            )
+            .unwrap(),
+        );
+        let output = normalize_write_array(&input, &DataType::Binary).unwrap();
+        let output = output.as_any().downcast_ref::<BinaryArray>().unwrap();
+        assert_eq!(output.value(0), b"ab");
+        assert!(output.is_null(1));
+        assert_eq!(output.value(2), b"cd");
+    }
+
+    #[test]
+    fn list_child_alias_changes_only_schema_and_keeps_values() {
+        let input: ArrayRef = Arc::new(ListArray::new(
+            Arc::new(Field::new("item", DataType::Int32, true)),
+            OffsetBuffer::new(ScalarBuffer::from(vec![0, 2, 3])),
+            Arc::new(Int32Array::from(vec![1, 2, 3])),
+            None,
+        ));
+        let expected = DataType::List(Arc::new(Field::new("element", 
DataType::Int32, true)));
+        let output = normalize_write_array(&input, &expected).unwrap();
+        assert_eq!(output.data_type(), &expected);
+        assert_eq!(
+            output
+                .as_any()
+                .downcast_ref::<ListArray>()
+                .unwrap()
+                .value(0)
+                .len(),
+            2
+        );
+
+        let wrong = DataType::List(Arc::new(Field::new("element", 
DataType::Utf8, true)));
+        assert!(normalize_write_array(&input, &wrong).is_err());
+    }
+
+    #[test]
+    fn row_child_names_cannot_be_silently_reordered() {
+        let input: ArrayRef = Arc::new(StructArray::from(vec![
+            (
+                Arc::new(Field::new("right", DataType::Int32, true)),
+                Arc::new(Int32Array::from(vec![1])) as ArrayRef,
+            ),
+            (
+                Arc::new(Field::new("left", DataType::Utf8, true)),
+                Arc::new(StringArray::from(vec!["x"])) as ArrayRef,
+            ),
+        ]));
+        let expected = DataType::Struct(
+            vec![
+                Field::new("left", DataType::Int32, true),
+                Field::new("right", DataType::Utf8, true),
+            ]
+            .into(),
+        );
+        let error = normalize_write_array(&input, &expected).unwrap_err();
+        assert!(error.to_string().contains("Nested ROW field name mismatch"));
+    }
+
+    #[test]
+    fn map_value_list_alias_is_normalized_recursively() {
+        let key_field = Arc::new(Field::new("key", DataType::Int32, false));
+        let value_field = Arc::new(Field::new(
+            "value",
+            DataType::List(Arc::new(Field::new("item", DataType::Int32, 
true))),
+            true,
+        ));
+        let values: ArrayRef = Arc::new(ListArray::new(
+            Arc::new(Field::new("item", DataType::Int32, true)),
+            OffsetBuffer::new(ScalarBuffer::from(vec![0, 1, 3])),
+            Arc::new(Int32Array::from(vec![10, 20, 21])),
+            None,
+        ));
+        let entries = StructArray::try_new(
+            vec![key_field.clone(), value_field].into(),
+            vec![Arc::new(Int32Array::from(vec![1, 2])), values],
+            None,
+        )
+        .unwrap();
+        let input: ArrayRef = Arc::new(
+            MapArray::try_new(
+                Arc::new(Field::new("entries", entries.data_type().clone(), 
false)),
+                OffsetBuffer::new(ScalarBuffer::from(vec![0, 2])),
+                entries,
+                None,
+                false,
+            )
+            .unwrap(),
+        );
+
+        let expected_value = Arc::new(Field::new(
+            "value",
+            DataType::List(Arc::new(Field::new("element", DataType::Int32, 
true))),
+            true,
+        ));
+        let expected_entries = Arc::new(Field::new(
+            "entries",
+            DataType::Struct(vec![key_field, expected_value].into()),
+            false,
+        ));
+        let expected = DataType::Map(expected_entries, false);
+        let output = normalize_write_array(&input, &expected).unwrap();
+        assert_eq!(output.data_type(), &expected);
+        let map = output.as_any().downcast_ref::<MapArray>().unwrap();
+        assert_eq!(map.value_length(0), 2);
+        let lists = map
+            .entries()
+            .column(1)
+            .as_any()
+            .downcast_ref::<ListArray>()
+            .unwrap();
+        assert_eq!(lists.value(1).len(), 2);
+    }
+}
diff --git a/crates/paimon/tests/binary_partition_native_write_test.rs 
b/crates/paimon/tests/binary_partition_native_write_test.rs
new file mode 100644
index 00000000..c04728d0
--- /dev/null
+++ b/crates/paimon/tests/binary_partition_native_write_test.rs
@@ -0,0 +1,128 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+mod common;
+
+use arrow_array::{ArrayRef, BinaryArray, Int32Array, RecordBatch};
+use arrow_schema::{DataType as ArrowDataType, Field as ArrowField, Schema as 
ArrowSchema};
+use common::incremental_helpers::{memory_table, persist_table_schema, 
setup_dirs};
+use futures::TryStreamExt;
+use paimon::spec::{DataType, IntType, Schema, TableSchema, VarBinaryType};
+use std::sync::Arc;
+
+async fn check_binary_partition(primary_key: bool) {
+    let path = if primary_key {
+        "memory:/binary_partition/pk"
+    } else {
+        "memory:/binary_partition/append"
+    };
+    let mut schema = Schema::builder()
+        .column("bin", DataType::VarBinary(VarBinaryType::new(64).unwrap()))
+        .column("id", DataType::Int(IntType::new()))
+        .partition_keys(["bin"])
+        .option("bucket", "1")
+        .option("partition.legacy-name", "false");
+    if primary_key {
+        schema = schema.primary_key(["id"]);
+    } else {
+        schema = schema.option("bucket-key", "id");
+    }
+    let (io, table) = memory_table(path, TableSchema::new(0, 
&schema.build().unwrap()));
+    setup_dirs(&io, path).await;
+    persist_table_schema(&io, path, table.schema()).await;
+
+    let bin: ArrayRef = Arc::new(BinaryArray::from_iter_values([
+        b"a/b".as_slice(),
+        b"a=b".as_slice(),
+        "\u{00A0}".as_bytes(),
+        b"\x1c".as_slice(),
+        b"\xED\xA0\x80".as_slice(),
+    ]));
+    let batch = RecordBatch::try_new(
+        Arc::new(ArrowSchema::new(vec![
+            ArrowField::new("bin", ArrowDataType::Binary, true),
+            ArrowField::new("id", ArrowDataType::Int32, true),
+        ])),
+        vec![bin, Arc::new(Int32Array::from(vec![1, 2, 3, 4, 5]))],
+    )
+    .unwrap();
+    let builder = table.new_write_builder();
+    let mut writer = builder.new_write().unwrap();
+    writer.write_arrow_batch(&batch).await.unwrap();
+    let messages = writer.prepare_commit().await.unwrap();
+    assert_eq!(messages.len(), 5);
+    builder.new_commit().commit(messages).await.unwrap();
+
+    let plan = table.new_read_builder().new_scan().plan().await.unwrap();
+    let paths = plan
+        .splits()
+        .iter()
+        .map(|split| split.bucket_path().to_string())
+        .collect::<Vec<_>>();
+    assert!(
+        paths.iter().any(|path| path.contains("bin=a%2Fb/")),
+        "{paths:?}"
+    );
+    assert!(
+        paths.iter().any(|path| path.contains("bin=a%3Db/")),
+        "{paths:?}"
+    );
+    // These directory names match Java's Character.isWhitespace and UTF-8 
decoder.
+    for expected in [
+        "bin=\u{00A0}/",
+        "bin=__DEFAULT_PARTITION__/",
+        "bin=\u{FFFD}/",
+    ] {
+        assert!(
+            paths.iter().any(|path| path.contains(expected)),
+            "missing {expected} in {paths:?}"
+        );
+    }
+    let batches: Vec<RecordBatch> = table
+        .new_read_builder()
+        .new_read()
+        .unwrap()
+        .to_arrow(plan.splits())
+        .unwrap()
+        .try_collect()
+        .await
+        .unwrap();
+    let mut ids = batches
+        .iter()
+        .flat_map(|batch| {
+            batch
+                .column(1)
+                .as_any()
+                .downcast_ref::<Int32Array>()
+                .unwrap()
+                .values()
+                .to_vec()
+        })
+        .collect::<Vec<_>>();
+    ids.sort_unstable();
+    assert_eq!(ids, vec![1, 2, 3, 4, 5]);
+}
+
+#[tokio::test]
+async fn append_binary_partitions_round_trip() {
+    check_binary_partition(false).await;
+}
+
+#[tokio::test]
+async fn primary_key_binary_partitions_round_trip() {
+    check_binary_partition(true).await;
+}
diff --git a/crates/paimon/tests/nested_native_write_compat_test.rs 
b/crates/paimon/tests/nested_native_write_compat_test.rs
new file mode 100644
index 00000000..1613afa5
--- /dev/null
+++ b/crates/paimon/tests/nested_native_write_compat_test.rs
@@ -0,0 +1,239 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+mod common;
+
+use arrow_array::{
+    Array, ArrayRef, BinaryArray, FixedSizeBinaryArray, Int32Array, ListArray, 
MapArray,
+    RecordBatch, StringArray, StructArray,
+};
+use arrow_buffer::{OffsetBuffer, ScalarBuffer};
+use arrow_schema::{DataType as ArrowDataType, Field as ArrowField, Schema as 
ArrowSchema};
+use common::incremental_helpers::{memory_table, persist_table_schema, 
setup_dirs};
+use futures::TryStreamExt;
+use paimon::spec::{
+    ArrayType, BinaryType, DataField, DataType, IntType, MapType, RowType, 
Schema, TableSchema,
+    VarCharType,
+};
+use std::sync::Arc;
+
+fn input_batch() -> RecordBatch {
+    // PyArrow names list children `item`; the Paimon table schema names them
+    // `element`. The nested ROW exercises the same normalization recursively.
+    let item = Arc::new(ArrowField::new("item", ArrowDataType::Int32, true));
+    let item_type = ArrowDataType::List(item.clone());
+    let ids: ArrayRef = Arc::new(Int32Array::from(vec![2, 1]));
+    let items: ArrayRef = Arc::new(ListArray::new(
+        item.clone(),
+        OffsetBuffer::new(ScalarBuffer::from(vec![0, 2, 3])),
+        Arc::new(Int32Array::from(vec![20, 21, 10])),
+        None,
+    ));
+    let tags: ArrayRef = Arc::new(ListArray::new(
+        item,
+        OffsetBuffer::new(ScalarBuffer::from(vec![0, 1, 3])),
+        Arc::new(Int32Array::from(vec![200, 100, 101])),
+        None,
+    ));
+    let tags_field = Arc::new(ArrowField::new("tags", item_type.clone(), 
true));
+    let payload: ArrayRef = 
Arc::new(StructArray::from(vec![(tags_field.clone(), tags)]));
+    let raw: ArrayRef = Arc::new(
+        FixedSizeBinaryArray::try_from_iter([b"two!".as_slice(), 
b"one!".as_slice()].into_iter())
+            .unwrap(),
+    );
+    let map_values: ArrayRef = Arc::new(ListArray::new(
+        Arc::new(ArrowField::new("item", ArrowDataType::Int32, true)),
+        OffsetBuffer::new(ScalarBuffer::from(vec![0, 1, 3])),
+        Arc::new(Int32Array::from(vec![200, 100, 101])),
+        None,
+    ));
+    let map_entries = StructArray::try_new(
+        vec![
+            Arc::new(ArrowField::new("key", ArrowDataType::Utf8, false)),
+            Arc::new(ArrowField::new("value", item_type.clone(), true)),
+        ]
+        .into(),
+        vec![Arc::new(StringArray::from(vec!["two", "one"])), map_values],
+        None,
+    )
+    .unwrap();
+    let map_entries_field = Arc::new(ArrowField::new(
+        "entries",
+        map_entries.data_type().clone(),
+        false,
+    ));
+    let attributes: ArrayRef = Arc::new(
+        MapArray::try_new(
+            map_entries_field.clone(),
+            OffsetBuffer::new(ScalarBuffer::from(vec![0, 1, 2])),
+            map_entries,
+            None,
+            false,
+        )
+        .unwrap(),
+    );
+
+    RecordBatch::try_new(
+        Arc::new(ArrowSchema::new(vec![
+            ArrowField::new("id", ArrowDataType::Int32, true),
+            ArrowField::new("items", item_type, true),
+            ArrowField::new(
+                "payload",
+                ArrowDataType::Struct(vec![tags_field].into()),
+                true,
+            ),
+            ArrowField::new("raw", ArrowDataType::FixedSizeBinary(4), true),
+            ArrowField::new(
+                "attributes",
+                ArrowDataType::Map(map_entries_field, false),
+                true,
+            ),
+        ])),
+        vec![ids, items, payload, raw, attributes],
+    )
+    .unwrap()
+}
+
+async fn check_nested_write(primary_key: bool) {
+    let path = if primary_key {
+        "memory:/nested_native_compat/pk"
+    } else {
+        "memory:/nested_native_compat/append"
+    };
+    let mut schema = Schema::builder()
+        .column("id", DataType::Int(IntType::new()))
+        .column(
+            "items",
+            DataType::Array(ArrayType::new(DataType::Int(IntType::new()))),
+        )
+        .column(
+            "payload",
+            DataType::Row(RowType::new(vec![DataField::new(
+                3,
+                "tags".into(),
+                DataType::Array(ArrayType::new(DataType::Int(IntType::new()))),
+            )])),
+        )
+        .column("raw", DataType::Binary(BinaryType::new(4).unwrap()))
+        .column(
+            "attributes",
+            DataType::Map(MapType::new(
+                DataType::VarChar(VarCharType::string_type()),
+                DataType::Array(ArrayType::new(DataType::Int(IntType::new()))),
+            )),
+        )
+        .option("bucket", "1");
+    if primary_key {
+        schema = schema.primary_key(["id"]);
+    } else {
+        schema = schema.option("bucket-key", "id");
+    }
+    let (io, table) = memory_table(path, TableSchema::new(0, 
&schema.build().unwrap()));
+    setup_dirs(&io, path).await;
+    persist_table_schema(&io, path, table.schema()).await;
+
+    let builder = table.new_write_builder();
+    let mut writer = builder.new_write().unwrap();
+    writer.write_arrow_batch(&input_batch()).await.unwrap();
+    let messages = writer.prepare_commit().await.unwrap();
+    builder.new_commit().commit(messages).await.unwrap();
+
+    let plan = table.new_read_builder().new_scan().plan().await.unwrap();
+    let batches: Vec<RecordBatch> = table
+        .new_read_builder()
+        .new_read()
+        .unwrap()
+        .to_arrow(plan.splits())
+        .unwrap()
+        .try_collect()
+        .await
+        .unwrap();
+    assert_eq!(batches.iter().map(RecordBatch::num_rows).sum::<usize>(), 2);
+    for batch in batches {
+        let ids = batch
+            .column(0)
+            .as_any()
+            .downcast_ref::<Int32Array>()
+            .unwrap();
+        let items = batch
+            .column(1)
+            .as_any()
+            .downcast_ref::<ListArray>()
+            .unwrap();
+        let payload = batch
+            .column(2)
+            .as_any()
+            .downcast_ref::<StructArray>()
+            .unwrap();
+        let raw = batch
+            .column(3)
+            .as_any()
+            .downcast_ref::<BinaryArray>()
+            .unwrap();
+        let attributes = 
batch.column(4).as_any().downcast_ref::<MapArray>().unwrap();
+        let attribute_values = attributes
+            .entries()
+            .column(1)
+            .as_any()
+            .downcast_ref::<ListArray>()
+            .unwrap();
+        let ArrowDataType::List(attribute_item) = attribute_values.data_type() 
else {
+            unreachable!()
+        };
+        assert_eq!(attribute_item.name(), "element");
+        let ArrowDataType::List(item_field) = items.data_type() else {
+            unreachable!()
+        };
+        assert_eq!(item_field.name(), "element");
+        let tags = payload
+            .column(0)
+            .as_any()
+            .downcast_ref::<ListArray>()
+            .unwrap();
+        let ArrowDataType::List(tag_field) = tags.data_type() else {
+            unreachable!()
+        };
+        assert_eq!(tag_field.name(), "element");
+        for row in 0..batch.num_rows() {
+            match ids.value(row) {
+                1 => {
+                    assert_eq!(items.value(row).len(), 1);
+                    assert_eq!(tags.value(row).len(), 2);
+                    assert_eq!(raw.value(row), b"one!");
+                    assert_eq!(attributes.value_length(row), 1);
+                }
+                2 => {
+                    assert_eq!(items.value(row).len(), 2);
+                    assert_eq!(tags.value(row).len(), 1);
+                    assert_eq!(raw.value(row), b"two!");
+                    assert_eq!(attributes.value_length(row), 1);
+                }
+                id => panic!("unexpected id {id}"),
+            }
+        }
+    }
+}
+
+#[tokio::test]
+async fn append_accepts_pyarrow_list_alias_and_fixed_binary() {
+    check_nested_write(false).await;
+}
+
+#[tokio::test]
+async fn primary_key_accepts_pyarrow_list_alias_and_fixed_binary() {
+    check_nested_write(true).await;
+}
diff --git 
a/crates/paimon/tests/pk_partial_update_sequence_group_parity_test.rs 
b/crates/paimon/tests/pk_partial_update_sequence_group_parity_test.rs
new file mode 100644
index 00000000..e8f90b97
--- /dev/null
+++ b/crates/paimon/tests/pk_partial_update_sequence_group_parity_test.rs
@@ -0,0 +1,357 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Persisted parity for Java `PartialUpdateMergeFunctionTest` sequence-group
+//! cases. These exercise the write buffer and independent files, then read
+//! through both full and projected table readers.
+
+#[path = "common/rowkind_helpers.rs"]
+mod helpers;
+
+use std::sync::Arc;
+
+use arrow_array::{Array, Int32Array, Int8Array, RecordBatch};
+use arrow_schema::{DataType as ArrowDataType, Field as ArrowField, Schema as 
ArrowSchema};
+use futures::StreamExt;
+use helpers::{memory_table, persist_table_schema, setup_dirs, write_batch};
+use paimon::spec::{DataType, IntType, Schema, TableSchema, 
VALUE_KIND_FIELD_NAME};
+use paimon::table::Table;
+
+#[derive(Clone, Copy, Debug)]
+struct Update {
+    seq_a: Option<i32>,
+    seq_b: Option<i32>,
+    amount: Option<i32>,
+    first: Option<i32>,
+    last: Option<i32>,
+    aux_seq: Option<i32>,
+    aux_value: Option<i32>,
+    free: Option<i32>,
+}
+
+impl Update {
+    fn new(
+        seq: (Option<i32>, Option<i32>),
+        grouped: (Option<i32>, Option<i32>, Option<i32>),
+        aux: (Option<i32>, Option<i32>),
+        free: Option<i32>,
+    ) -> Self {
+        Self {
+            seq_a: seq.0,
+            seq_b: seq.1,
+            amount: grouped.0,
+            first: grouped.1,
+            last: grouped.2,
+            aux_seq: aux.0,
+            aux_value: aux.1,
+            free,
+        }
+    }
+}
+
+fn batch(updates: &[Update]) -> RecordBatch {
+    RecordBatch::try_new(
+        Arc::new(ArrowSchema::new(
+            [
+                "id",
+                "seq_a",
+                "seq_b",
+                "amount",
+                "first",
+                "last",
+                "aux_seq",
+                "aux_value",
+                "free",
+            ]
+            .into_iter()
+            .map(|name| ArrowField::new(name, ArrowDataType::Int32, name != 
"id"))
+            .collect::<Vec<_>>(),
+        )),
+        vec![
+            Arc::new(Int32Array::from(vec![1; updates.len()])),
+            Arc::new(Int32Array::from(
+                updates.iter().map(|row| row.seq_a).collect::<Vec<_>>(),
+            )),
+            Arc::new(Int32Array::from(
+                updates.iter().map(|row| row.seq_b).collect::<Vec<_>>(),
+            )),
+            Arc::new(Int32Array::from(
+                updates.iter().map(|row| row.amount).collect::<Vec<_>>(),
+            )),
+            Arc::new(Int32Array::from(
+                updates.iter().map(|row| row.first).collect::<Vec<_>>(),
+            )),
+            Arc::new(Int32Array::from(
+                updates.iter().map(|row| row.last).collect::<Vec<_>>(),
+            )),
+            Arc::new(Int32Array::from(
+                updates.iter().map(|row| row.aux_seq).collect::<Vec<_>>(),
+            )),
+            Arc::new(Int32Array::from(
+                updates.iter().map(|row| row.aux_value).collect::<Vec<_>>(),
+            )),
+            Arc::new(Int32Array::from(
+                updates.iter().map(|row| row.free).collect::<Vec<_>>(),
+            )),
+        ],
+    )
+    .unwrap()
+}
+
+async fn table(path: &str) -> Table {
+    let mut builder = Schema::builder().column("id", 
DataType::Int(IntType::new()));
+    for field in [
+        "seq_a",
+        "seq_b",
+        "amount",
+        "first",
+        "last",
+        "aux_seq",
+        "aux_value",
+        "free",
+    ] {
+        builder = builder.column(field, DataType::Int(IntType::new()));
+    }
+    let schema = builder
+        .primary_key(["id"])
+        .option("bucket", "1")
+        .option("merge-engine", "partial-update")
+        .option("fields.seq_a,seq_b.sequence-group", "amount,first,last")
+        .option("fields.aux_seq.sequence-group", "aux_value")
+        .option("fields.amount.aggregate-function", "sum")
+        .option("fields.first.aggregate-function", "first_value")
+        .option("fields.last.aggregate-function", "last_value")
+        .option("fields.aux_value.aggregate-function", "last_non_null_value")
+        .build()
+        .unwrap();
+    let (file_io, table) = memory_table(path, TableSchema::new(0, &schema));
+    setup_dirs(&file_io, path).await;
+    persist_table_schema(&file_io, path, table.schema()).await;
+    table
+}
+
+async fn scan_optional(
+    table: &Table,
+    projection: Option<&[&str]>,
+) -> Option<Vec<(String, Option<i32>)>> {
+    let plan = table.new_read_builder().new_scan().plan().await.unwrap();
+    let mut builder = table.new_read_builder();
+    if let Some(projection) = projection {
+        builder.with_projection(projection).unwrap();
+    }
+    let mut stream = 
builder.new_read().unwrap().to_arrow(plan.splits()).unwrap();
+    let batch = stream.next().await?.unwrap();
+    assert_eq!(batch.num_rows(), 1);
+    assert!(stream.next().await.is_none());
+    Some(
+        batch
+            .schema()
+            .fields()
+            .iter()
+            .enumerate()
+            .map(|(index, field)| {
+                let column = batch
+                    .column(index)
+                    .as_any()
+                    .downcast_ref::<Int32Array>()
+                    .unwrap();
+                (
+                    field.name().clone(),
+                    (!column.is_null(0)).then(|| column.value(0)),
+                )
+            })
+            .collect(),
+    )
+}
+
+async fn scan_one(table: &Table, projection: Option<&[&str]>) -> Vec<(String, 
Option<i32>)> {
+    scan_optional(table, projection)
+        .await
+        .expect("expected one current-state row")
+}
+
+#[tokio::test]
+async fn 
composite_sequence_group_aggregates_older_inputs_and_independent_groups() {
+    // The three updates mirror Java's multi-sequence first/last test, with an
+    // additional SUM and a second independent sequence group.
+    let inputs = [
+        Update::new(
+            (Some(1), Some(1)),
+            (Some(1), Some(1), Some(1)),
+            (Some(1), Some(1)),
+            Some(10),
+        ),
+        Update::new(
+            (Some(2), Some(2)),
+            (Some(1), Some(2), Some(2)),
+            (Some(0), Some(2)),
+            Some(20),
+        ),
+        Update::new(
+            (Some(0), Some(1)),
+            (Some(3), Some(3), Some(3)),
+            (Some(2), None),
+            Some(30),
+        ),
+    ];
+    for separate_commits in [false, true] {
+        let path = 
format!("memory:/partial_update_parity/group_{separate_commits}");
+        let table = table(&path).await;
+        if separate_commits {
+            for input in inputs {
+                write_batch(&table, &batch(&[input])).await;
+            }
+        } else {
+            write_batch(&table, &batch(&inputs)).await;
+        }
+        let actual = scan_one(&table, None).await;
+        assert_eq!(
+            actual,
+            vec![
+                ("id".into(), Some(1)),
+                ("seq_a".into(), Some(2)),
+                ("seq_b".into(), Some(2)),
+                ("amount".into(), Some(5)),
+                ("first".into(), Some(3)),
+                ("last".into(), Some(2)),
+                ("aux_seq".into(), Some(2)),
+                ("aux_value".into(), Some(1)),
+                ("free".into(), Some(30)),
+            ],
+            "separate_commits={separate_commits}"
+        );
+        let projected = scan_one(
+            &table,
+            Some(&["id", "amount", "first", "last", "aux_value", "free"]),
+        )
+        .await;
+        assert_eq!(
+            projected,
+            vec![
+                ("id".into(), Some(1)),
+                ("amount".into(), Some(5)),
+                ("first".into(), Some(3)),
+                ("last".into(), Some(2)),
+                ("aux_value".into(), Some(1)),
+                ("free".into(), Some(30)),
+            ],
+            "projection, separate_commits={separate_commits}"
+        );
+    }
+}
+
+async fn partial_delete_table(path: &str) -> Table {
+    let mut builder = Schema::builder().column("id", 
DataType::Int(IntType::new()));
+    for field in ["a", "b", "seq_left", "c", "d", "seq_right"] {
+        builder = builder.column(field, DataType::Int(IntType::new()));
+    }
+    let schema = builder
+        .primary_key(["id"])
+        .option("bucket", "1")
+        .option("merge-engine", "partial-update")
+        .option("fields.seq_left.sequence-group", "a,b")
+        .option("fields.seq_right.sequence-group", "c,d")
+        .option(
+            "partial-update.remove-record-on-sequence-group",
+            "seq_right",
+        )
+        .build()
+        .unwrap();
+    let (file_io, table) = memory_table(path, TableSchema::new(0, &schema));
+    setup_dirs(&file_io, path).await;
+    persist_table_schema(&file_io, path, table.schema()).await;
+    table
+}
+
+fn partial_delete_batch(values: [Option<i32>; 6], kind: i8) -> RecordBatch {
+    let mut columns: Vec<Arc<dyn Array>> = 
vec![Arc::new(Int32Array::from(vec![1]))];
+    columns.extend(
+        values
+            .into_iter()
+            .map(|value| Arc::new(Int32Array::from(vec![value])) as Arc<dyn 
Array>),
+    );
+    columns.push(Arc::new(Int8Array::from(vec![kind])));
+    let mut fields = ["id", "a", "b", "seq_left", "c", "d", "seq_right"]
+        .into_iter()
+        .map(|name| ArrowField::new(name, ArrowDataType::Int32, name != "id"))
+        .collect::<Vec<_>>();
+    fields.push(ArrowField::new(
+        VALUE_KIND_FIELD_NAME,
+        ArrowDataType::Int8,
+        false,
+    ));
+    RecordBatch::try_new(Arc::new(ArrowSchema::new(fields)), columns).unwrap()
+}
+
+#[tokio::test]
+async fn sequence_group_delete_matches_java_across_commits_and_projection() {
+    // Java PartialUpdateMergeFunctionTest.testSequenceGroupPartialDelete:
+    // a delete with a sequence for one group clears that group's fields;
+    // another group retains its latest values. The configured right-hand
+    // sequence group can later remove the entire record.
+    let path = "memory:/partial_update_parity/group_delete";
+    let table = partial_delete_table(path).await;
+    let steps = [
+        ([Some(1), Some(1), Some(1), Some(1), Some(1), Some(1)], 0),
+        ([Some(2), Some(2), Some(2), Some(2), Some(2), None], 0),
+        ([Some(3), Some(3), Some(1), Some(3), Some(3), Some(3)], 0),
+        ([Some(1), Some(1), Some(3), Some(1), Some(1), None], 3),
+        ([Some(1), Some(1), Some(3), Some(1), Some(1), Some(4)], 3),
+        ([Some(4), Some(4), Some(4), Some(5), Some(5), Some(5)], 0),
+        ([Some(1), Some(1), Some(6), Some(1), Some(1), Some(6)], 3),
+    ];
+    let expected = [
+        [Some(1), Some(1), Some(1), Some(1), Some(1), Some(1)],
+        [Some(2), Some(2), Some(2), Some(1), Some(1), Some(1)],
+        [Some(2), Some(2), Some(2), Some(3), Some(3), Some(3)],
+        [None, None, Some(3), Some(3), Some(3), Some(3)],
+        [Some(1), Some(1), Some(3), Some(1), Some(1), Some(4)],
+        [Some(4), Some(4), Some(4), Some(5), Some(5), Some(5)],
+        [Some(1), Some(1), Some(6), Some(1), Some(1), Some(6)],
+    ];
+    for (step, ((values, kind), expected)) in 
steps.into_iter().zip(expected).enumerate() {
+        write_batch(&table, &partial_delete_batch(values, kind)).await;
+        if step == 4 || step == 6 {
+            // Java's merge result contains the DELETE payload at this point,
+            // but a current-state table scan must hide the tombstone.
+            assert!(scan_optional(&table, None).await.is_none());
+            assert!(scan_optional(&table, Some(&["id", "b", "d"]))
+                .await
+                .is_none());
+            continue;
+        }
+        let actual = scan_optional(&table, None)
+            .await
+            .unwrap_or_else(|| panic!("expected a current-state row at step 
{step}"));
+        let values = actual
+            .into_iter()
+            .skip(1)
+            .map(|(_, value)| value)
+            .collect::<Vec<_>>();
+        assert_eq!(values, expected, "step {step}");
+        let projected = scan_one(&table, Some(&["id", "b", "d"])).await;
+        assert_eq!(
+            projected,
+            vec![
+                ("id".into(), Some(1)),
+                ("b".into(), expected[1]),
+                ("d".into(), expected[4]),
+            ],
+            "projected step {step}"
+        );
+    }
+}
diff --git a/crates/paimon/tests/pk_user_sequence_parity_test.rs 
b/crates/paimon/tests/pk_user_sequence_parity_test.rs
new file mode 100644
index 00000000..d43c4381
--- /dev/null
+++ b/crates/paimon/tests/pk_user_sequence_parity_test.rs
@@ -0,0 +1,426 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Persisted user-sequence controls for Java's `UserDefinedSeqComparator`.
+
+#[path = "common/rowkind_helpers.rs"]
+mod helpers;
+
+use std::sync::Arc;
+
+use arrow_array::{
+    Array, ArrayRef, BinaryArray, BooleanArray, Date32Array, Decimal128Array, 
Float32Array,
+    Float64Array, Int32Array, RecordBatch, StringArray, Time32MillisecondArray,
+    TimestampMicrosecondArray,
+};
+use arrow_schema::{DataType as ArrowDataType, Field as ArrowField, Schema as 
ArrowSchema};
+use futures::StreamExt;
+use helpers::{memory_table, persist_table_schema, setup_dirs, write_batch};
+use paimon::spec::{
+    BooleanType, DataType, DateType, DecimalType, DoubleType, FloatType, 
IntType,
+    LocalZonedTimestampType, Schema, TableSchema, TimeType, TimestampType, 
VarBinaryType,
+    VarCharType,
+};
+use paimon::table::Table;
+
+fn sequence_batch(rows: &[(Option<&str>, i32)]) -> RecordBatch {
+    RecordBatch::try_new(
+        Arc::new(ArrowSchema::new(vec![
+            ArrowField::new("id", ArrowDataType::Int32, false),
+            ArrowField::new("seq", ArrowDataType::Utf8, true),
+            ArrowField::new("value", ArrowDataType::Int32, true),
+        ])),
+        vec![
+            Arc::new(Int32Array::from(vec![1; rows.len()])),
+            Arc::new(StringArray::from(
+                rows.iter().map(|(seq, _)| *seq).collect::<Vec<_>>(),
+            )),
+            Arc::new(Int32Array::from(
+                rows.iter().map(|(_, value)| *value).collect::<Vec<_>>(),
+            )),
+        ],
+    )
+    .unwrap()
+}
+
+async fn sequence_table(path: &str, engine: &str, descending: bool) -> Table {
+    let schema = Schema::builder()
+        .column("id", DataType::Int(IntType::new()))
+        .column("seq", DataType::VarChar(VarCharType::string_type()))
+        .column("value", DataType::Int(IntType::new()))
+        .primary_key(["id"])
+        .option("bucket", "1")
+        .option("merge-engine", engine)
+        .option("sequence.field", "seq")
+        .option(
+            "sequence.field.sort-order",
+            if descending {
+                "descending"
+            } else {
+                "ascending"
+            },
+        )
+        .build()
+        .unwrap();
+    let (file_io, table) = memory_table(path, TableSchema::new(0, &schema));
+    setup_dirs(&file_io, path).await;
+    persist_table_schema(&file_io, path, table.schema()).await;
+    table
+}
+
+async fn scan_sequence_value(table: &Table) -> (Option<String>, i32) {
+    let plan = table.new_read_builder().new_scan().plan().await.unwrap();
+    let mut stream = table
+        .new_read_builder()
+        .new_read()
+        .unwrap()
+        .to_arrow(plan.splits())
+        .unwrap();
+    let batch = stream.next().await.unwrap().unwrap();
+    assert_eq!(batch.num_rows(), 1);
+    let sequences = batch
+        .column_by_name("seq")
+        .unwrap()
+        .as_any()
+        .downcast_ref::<StringArray>()
+        .unwrap();
+    let values = batch
+        .column_by_name("value")
+        .unwrap()
+        .as_any()
+        .downcast_ref::<Int32Array>()
+        .unwrap();
+    let result = (
+        (!sequences.is_null(0)).then(|| sequences.value(0).to_string()),
+        values.value(0),
+    );
+    assert!(stream.next().await.is_none());
+    result
+}
+
+#[tokio::test]
+async fn string_sequence_respects_order_in_flush_and_across_commits() {
+    // The one-commit case tests the write buffer; separate commits test the
+    // read-side merge of independently persisted files. Java chooses the
+    // lexicographically greatest sequence in ascending mode and the least in
+    // descending mode, regardless of arrival order.
+    for engine in ["deduplicate", "partial-update", "aggregation"] {
+        for descending in [false, true] {
+            for separate_commits in [false, true] {
+                let path =
+                    
format!("memory:/pk_user_sequence/{engine}_{descending}_{separate_commits}");
+                let table = sequence_table(&path, engine, descending).await;
+                if separate_commits {
+                    write_batch(&table, &sequence_batch(&[(Some("z"), 
10)])).await;
+                    write_batch(&table, &sequence_batch(&[(Some("a"), 
20)])).await;
+                } else {
+                    write_batch(&table, &sequence_batch(&[(Some("z"), 10), 
(Some("a"), 20)])).await;
+                }
+                let expected = if descending {
+                    (Some("a".to_string()), 20)
+                } else {
+                    (Some("z".to_string()), 10)
+                };
+                assert_eq!(
+                    scan_sequence_value(&table).await,
+                    expected,
+                    "{engine}, descending={descending}, 
separate_commits={separate_commits}"
+                );
+            }
+        }
+    }
+}
+
+#[tokio::test]
+async fn null_user_sequence_stays_first_in_both_orders() {
+    for descending in [false, true] {
+        let path = format!("memory:/pk_user_sequence/null_{descending}");
+        let table = sequence_table(&path, "deduplicate", descending).await;
+        write_batch(&table, &sequence_batch(&[(Some("z"), 10)])).await;
+        write_batch(&table, &sequence_batch(&[(None, 20)])).await;
+        assert_eq!(
+            scan_sequence_value(&table).await,
+            (Some("z".into()), 10),
+            "descending={descending}"
+        );
+    }
+}
+
+struct TypedSequenceCase {
+    name: &'static str,
+    field_type: DataType,
+    low: ArrayRef,
+    high: ArrayRef,
+}
+
+fn typed_sequence_cases() -> Vec<TypedSequenceCase> {
+    vec![
+        TypedSequenceCase {
+            name: "boolean",
+            field_type: DataType::Boolean(BooleanType::new()),
+            low: Arc::new(BooleanArray::from(vec![false])),
+            high: Arc::new(BooleanArray::from(vec![true])),
+        },
+        TypedSequenceCase {
+            name: "binary",
+            field_type: DataType::VarBinary(VarBinaryType::new(16).unwrap()),
+            low: Arc::new(BinaryArray::from_iter_values([b"\x00".as_slice()])),
+            high: 
Arc::new(BinaryArray::from_iter_values([b"\xff".as_slice()])),
+        },
+        TypedSequenceCase {
+            name: "date",
+            field_type: DataType::Date(DateType::new()),
+            low: Arc::new(Date32Array::from(vec![-1])),
+            high: Arc::new(Date32Array::from(vec![1])),
+        },
+        TypedSequenceCase {
+            name: "decimal",
+            field_type: DataType::Decimal(DecimalType::new(10, 2).unwrap()),
+            low: Arc::new(
+                Decimal128Array::from(vec![-123_i128])
+                    .with_precision_and_scale(10, 2)
+                    .unwrap(),
+            ),
+            high: Arc::new(
+                Decimal128Array::from(vec![456_i128])
+                    .with_precision_and_scale(10, 2)
+                    .unwrap(),
+            ),
+        },
+        TypedSequenceCase {
+            name: "float",
+            field_type: DataType::Float(FloatType::new()),
+            low: Arc::new(Float32Array::from(vec![-0.0])),
+            high: Arc::new(Float32Array::from(vec![0.0])),
+        },
+        TypedSequenceCase {
+            name: "double",
+            field_type: DataType::Double(DoubleType::new()),
+            low: Arc::new(Float64Array::from(vec![-0.0])),
+            high: Arc::new(Float64Array::from(vec![0.0])),
+        },
+        TypedSequenceCase {
+            name: "time",
+            field_type: DataType::Time(TimeType::new(3).unwrap()),
+            low: Arc::new(Time32MillisecondArray::from(vec![0])),
+            high: Arc::new(Time32MillisecondArray::from(vec![86_399_999])),
+        },
+        TypedSequenceCase {
+            name: "timestamp",
+            field_type: DataType::Timestamp(TimestampType::new(6).unwrap()),
+            low: Arc::new(TimestampMicrosecondArray::from(vec![-1])),
+            high: Arc::new(TimestampMicrosecondArray::from(vec![1])),
+        },
+        TypedSequenceCase {
+            name: "timestamp_ltz",
+            field_type: 
DataType::LocalZonedTimestamp(LocalZonedTimestampType::new(6).unwrap()),
+            low: 
Arc::new(TimestampMicrosecondArray::from(vec![-1]).with_timezone("UTC")),
+            high: 
Arc::new(TimestampMicrosecondArray::from(vec![1]).with_timezone("UTC")),
+        },
+    ]
+}
+
+fn typed_sequence_batch(sequence: ArrayRef, value: i32) -> RecordBatch {
+    RecordBatch::try_new(
+        Arc::new(ArrowSchema::new(vec![
+            ArrowField::new("id", ArrowDataType::Int32, false),
+            ArrowField::new("seq", sequence.data_type().clone(), true),
+            ArrowField::new("value", ArrowDataType::Int32, true),
+        ])),
+        vec![
+            Arc::new(Int32Array::from(vec![1])),
+            sequence,
+            Arc::new(Int32Array::from(vec![value])),
+        ],
+    )
+    .unwrap()
+}
+
+async fn typed_sequence_table(
+    path: &str,
+    engine: &str,
+    sequence_type: DataType,
+    descending: bool,
+) -> Table {
+    let schema = Schema::builder()
+        .column("id", DataType::Int(IntType::new()))
+        .column("seq", sequence_type)
+        .column("value", DataType::Int(IntType::new()))
+        .primary_key(["id"])
+        .option("bucket", "1")
+        .option("merge-engine", engine)
+        .option("sequence.field", "seq")
+        .option(
+            "sequence.field.sort-order",
+            if descending {
+                "descending"
+            } else {
+                "ascending"
+            },
+        )
+        .build()
+        .unwrap();
+    let (file_io, table) = memory_table(path, TableSchema::new(0, &schema));
+    setup_dirs(&file_io, path).await;
+    persist_table_schema(&file_io, path, table.schema()).await;
+    table
+}
+
+async fn scan_typed_sequence_value(table: &Table) -> i32 {
+    let plan = table.new_read_builder().new_scan().plan().await.unwrap();
+    let mut stream = table
+        .new_read_builder()
+        .new_read()
+        .unwrap()
+        .to_arrow(plan.splits())
+        .unwrap();
+    let batch = stream.next().await.unwrap().unwrap();
+    assert_eq!(batch.num_rows(), 1);
+    let value = batch
+        .column_by_name("value")
+        .unwrap()
+        .as_any()
+        .downcast_ref::<Int32Array>()
+        .unwrap()
+        .value(0);
+    assert!(stream.next().await.is_none());
+    value
+}
+
+#[tokio::test]
+async fn typed_user_sequences_match_java_in_flush_and_across_files() {
+    // These types previously fell through the reader's integer-only sequence
+    // extraction and were all treated as equal. Both commit layouts exercise
+    // the write-buffer comparator and the cross-file reader comparator.
+    for case in typed_sequence_cases() {
+        for engine in ["deduplicate", "partial-update", "aggregation"] {
+            for descending in [false, true] {
+                for separate_commits in [false, true] {
+                    let path = format!(
+                        "memory:/pk_user_sequence/typed_{}_{}_{}_{}",
+                        case.name, engine, descending, separate_commits
+                    );
+                    let table =
+                        typed_sequence_table(&path, engine, 
case.field_type.clone(), descending)
+                            .await;
+                    let high = typed_sequence_batch(case.high.clone(), 10);
+                    let low = typed_sequence_batch(case.low.clone(), 20);
+                    if separate_commits {
+                        write_batch(&table, &high).await;
+                        write_batch(&table, &low).await;
+                    } else {
+                        let batch =
+                            
arrow_select::concat::concat_batches(&high.schema(), &[high, low])
+                                .unwrap();
+                        write_batch(&table, &batch).await;
+                    }
+                    assert_eq!(
+                        scan_typed_sequence_value(&table).await,
+                        if descending { 20 } else { 10 },
+                        "type={}, engine={engine}, descending={descending}, 
separate_commits={separate_commits}",
+                        case.name
+                    );
+                }
+            }
+        }
+    }
+}
+
+fn composite_sequence_batch(rows: &[(Option<&str>, Option<i32>, i32)]) -> 
RecordBatch {
+    RecordBatch::try_new(
+        Arc::new(ArrowSchema::new(vec![
+            ArrowField::new("id", ArrowDataType::Int32, false),
+            ArrowField::new("seq_text", ArrowDataType::Utf8, true),
+            ArrowField::new("seq_number", ArrowDataType::Int32, true),
+            ArrowField::new("value", ArrowDataType::Int32, true),
+        ])),
+        vec![
+            Arc::new(Int32Array::from(vec![1; rows.len()])),
+            Arc::new(StringArray::from(
+                rows.iter().map(|row| row.0).collect::<Vec<_>>(),
+            )),
+            Arc::new(Int32Array::from(
+                rows.iter().map(|row| row.1).collect::<Vec<_>>(),
+            )),
+            Arc::new(Int32Array::from(
+                rows.iter().map(|row| row.2).collect::<Vec<_>>(),
+            )),
+        ],
+    )
+    .unwrap()
+}
+
+async fn composite_sequence_table(path: &str, engine: &str, descending: bool) 
-> Table {
+    let schema = Schema::builder()
+        .column("id", DataType::Int(IntType::new()))
+        .column("seq_text", DataType::VarChar(VarCharType::string_type()))
+        .column("seq_number", DataType::Int(IntType::new()))
+        .column("value", DataType::Int(IntType::new()))
+        .primary_key(["id"])
+        .option("bucket", "1")
+        .option("merge-engine", engine)
+        .option("sequence.field", "seq_text,seq_number")
+        .option(
+            "sequence.field.sort-order",
+            if descending {
+                "descending"
+            } else {
+                "ascending"
+            },
+        )
+        .build()
+        .unwrap();
+    let (file_io, table) = memory_table(path, TableSchema::new(0, &schema));
+    setup_dirs(&file_io, path).await;
+    persist_table_schema(&file_io, path, table.schema()).await;
+    table
+}
+
+#[tokio::test]
+async fn composite_user_sequence_uses_lexicographic_order_and_nulls_first() {
+    // The second sequence field decides between the two "b" rows. A NULL in
+    // either field is smaller than a non-NULL value in both directions.
+    let rows = [
+        (Some("b"), Some(0), 10),
+        (Some("a"), Some(9), 20),
+        (Some("b"), Some(1), 30),
+        (Some("b"), None, 40),
+        (None, Some(100), 50),
+    ];
+    for engine in ["deduplicate", "partial-update", "aggregation"] {
+        for descending in [false, true] {
+            for separate_commits in [false, true] {
+                let path = format!(
+                    
"memory:/pk_user_sequence/composite_{engine}_{descending}_{separate_commits}"
+                );
+                let table = composite_sequence_table(&path, engine, 
descending).await;
+                if separate_commits {
+                    for row in rows {
+                        write_batch(&table, 
&composite_sequence_batch(&[row])).await;
+                    }
+                } else {
+                    write_batch(&table, 
&composite_sequence_batch(&rows)).await;
+                }
+                assert_eq!(
+                    scan_typed_sequence_value(&table).await,
+                    if descending { 20 } else { 30 },
+                    "engine={engine}, descending={descending}, 
separate_commits={separate_commits}"
+                );
+            }
+        }
+    }
+}
diff --git a/crates/paimon/tests/target_file_row_num_test.rs 
b/crates/paimon/tests/target_file_row_num_test.rs
new file mode 100644
index 00000000..9e6d6735
--- /dev/null
+++ b/crates/paimon/tests/target_file_row_num_test.rs
@@ -0,0 +1,184 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+mod common;
+
+use common::incremental_helpers::{
+    make_batch, memory_table, persist_table_schema, pk_schema, setup_dirs,
+};
+use futures::TryStreamExt;
+use paimon::spec::{DataType, IntType, Schema, TableSchema};
+use paimon::table::Table;
+
+async fn table(path: &str, data_evolution: bool) -> Table {
+    let mut schema = Schema::builder()
+        .column("id", DataType::Int(IntType::new()))
+        .column("value", DataType::Int(IntType::new()))
+        .option("target-file-row-num", "2")
+        .option("target-file-size", "256mb");
+    if data_evolution {
+        schema = schema
+            .option("data-evolution.enabled", "true")
+            .option("row-tracking.enabled", "true");
+    }
+    let (io, table) = memory_table(path, TableSchema::new(0, 
&schema.build().unwrap()));
+    setup_dirs(&io, path).await;
+    persist_table_schema(&io, path, table.schema()).await;
+    table
+}
+
+async fn check_roll(data_evolution: bool) {
+    let path = if data_evolution {
+        "memory:/target_file_row_num/data_evolution"
+    } else {
+        "memory:/target_file_row_num/append"
+    };
+    let table = table(path, data_evolution).await;
+    let builder = table.new_write_builder();
+    let mut writer = builder.new_write().unwrap();
+    for (start, end) in [(0, 1), (1, 2), (2, 5), (5, 6)] {
+        let batch = make_batch(
+            (start..end).collect(),
+            (start..end).map(|id| id * 10).collect(),
+        );
+        writer.write_arrow_batch(&batch).await.unwrap();
+    }
+    let messages = writer.prepare_commit().await.unwrap();
+    assert_eq!(messages.len(), 1);
+    assert_eq!(
+        messages[0]
+            .new_files
+            .iter()
+            .map(|file| file.row_count)
+            .collect::<Vec<_>>(),
+        vec![2, 3, 1]
+    );
+    builder.new_commit().commit(messages).await.unwrap();
+
+    let plan = table
+        .new_read_builder()
+        .new_scan()
+        .with_scan_all_files()
+        .plan()
+        .await
+        .unwrap();
+    if data_evolution {
+        let mut files = plan
+            .splits()
+            .iter()
+            .flat_map(|split| split.data_files())
+            .map(|file| (file.first_row_id.unwrap(), file.row_count))
+            .collect::<Vec<_>>();
+        files.sort_unstable();
+        assert_eq!(files, vec![(0, 2), (2, 3), (5, 1)]);
+    }
+    let batches: Vec<arrow_array::RecordBatch> = table
+        .new_read_builder()
+        .new_read()
+        .unwrap()
+        .to_arrow(plan.splits())
+        .unwrap()
+        .try_collect()
+        .await
+        .unwrap();
+    let mut ids = batches
+        .iter()
+        .flat_map(|batch| {
+            batch
+                .column(0)
+                .as_any()
+                .downcast_ref::<arrow_array::Int32Array>()
+                .unwrap()
+                .values()
+                .to_vec()
+        })
+        .collect::<Vec<_>>();
+    ids.sort_unstable();
+    assert_eq!(ids, vec![0, 1, 2, 3, 4, 5]);
+}
+
+#[tokio::test]
+async fn append_rolls_after_batches() {
+    check_roll(false).await;
+}
+
+#[tokio::test]
+async fn data_evolution_rolls_after_batches_and_assigns_contiguous_row_ids() {
+    check_roll(true).await;
+}
+
+#[tokio::test]
+async fn primary_key_row_limit_files_commit_and_read_all_keys() {
+    let path = "memory:/target_file_row_num/primary_key";
+    let (io, table) = memory_table(path, pk_schema(&[("target-file-row-num", 
"2")]));
+    setup_dirs(&io, path).await;
+    persist_table_schema(&io, path, table.schema()).await;
+
+    let builder = table.new_write_builder();
+    let mut writer = builder.new_write().unwrap();
+    writer
+        .write_arrow_batch(&make_batch(vec![5, 1, 4, 2, 3], vec![50, 10, 40, 
20, 30]))
+        .await
+        .unwrap();
+    let messages = writer.prepare_commit().await.unwrap();
+    assert_eq!(messages.len(), 1);
+    assert_eq!(
+        messages[0]
+            .new_files
+            .iter()
+            .map(|file| file.row_count)
+            .collect::<Vec<_>>(),
+        vec![2, 2, 1]
+    );
+    builder.new_commit().commit(messages).await.unwrap();
+
+    let plan = table.new_read_builder().new_scan().plan().await.unwrap();
+    let batches: Vec<arrow_array::RecordBatch> = table
+        .new_read_builder()
+        .new_read()
+        .unwrap()
+        .to_arrow(plan.splits())
+        .unwrap()
+        .try_collect()
+        .await
+        .unwrap();
+    let mut ids = batches
+        .iter()
+        .flat_map(|batch| {
+            batch
+                .column(0)
+                .as_any()
+                .downcast_ref::<arrow_array::Int32Array>()
+                .unwrap()
+                .values()
+                .to_vec()
+        })
+        .collect::<Vec<_>>();
+    ids.sort_unstable();
+    assert_eq!(ids, vec![1, 2, 3, 4, 5]);
+}
+
+#[test]
+fn reject_invalid_row_targets_at_schema_creation() {
+    for value in ["0", "-1", "bad"] {
+        let result = Schema::builder()
+            .column("id", DataType::Int(IntType::new()))
+            .option("target-file-row-num", value)
+            .build();
+        assert!(result.is_err(), "invalid target {value} was accepted");
+    }
+}

Reply via email to