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");
+ }
+}