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 5c87ec53 fix(core): close gaps for PyPaimon native Parquet writers 
(#970)
5c87ec53 is described below

commit 5c87ec53e8bbaa415047e9a5a2beba27fe1b0acf
Author: Jingsong Lee <[email protected]>
AuthorDate: Sun Sep 27 08:32:57 2026 +0800

    fix(core): close gaps for PyPaimon native Parquet writers (#970)
---
 bindings/python/src/write.rs                       |  37 ----
 bindings/python/tests/test_write.py                |  73 +++++++
 crates/paimon/src/arrow/format/parquet.rs          |  38 +++-
 crates/paimon/src/arrow/shredding/variant.rs       |  50 ++++-
 crates/paimon/src/spec/aggregation.rs              |  33 +++-
 crates/paimon/src/spec/core_options.rs             |  19 +-
 crates/paimon/src/spec/schema.rs                   |  37 ++--
 crates/paimon/src/spec/types.rs                    |  41 +++-
 .../paimon/src/table/format_table_write_tests.rs   |  23 ++-
 crates/paimon/src/table/format_table_writer.rs     |   3 +-
 crates/paimon/src/table/kv_file_writer.rs          | 211 ++++++++++++++++++---
 crates/paimon/src/table/table_write.rs             |  46 +++--
 .../paimon/tests/parquet_write_page_index_test.rs  | 122 ++++++++++++
 crates/paimon/tests/rowkind_field_test.rs          |  90 ++++++++-
 14 files changed, 699 insertions(+), 124 deletions(-)

diff --git a/bindings/python/src/write.rs b/bindings/python/src/write.rs
index 5aacc16a..bb02e891 100644
--- a/bindings/python/src/write.rs
+++ b/bindings/python/src/write.rs
@@ -34,39 +34,6 @@ use pyo3::types::{PyBytes, PyDict, PyString};
 use crate::error::to_py_err;
 use crate::predicate::{dict_to_table_predicate, py_to_datum};
 
-/// Validate an incoming batch schema against the table's target Arrow schema:
-/// field count, order, and names must match, and types must match exactly. The
-/// nullable flag is intentionally NOT compared, since 
`build_target_arrow_schema`
-/// derives nullability from the Paimon field while pyarrow-constructed batches
-/// infer nullable=true. No cast — callers supply correctly-typed batches.
-///
-/// Type matching is strict (no binary-family interchange): the lower write 
path
-/// downcasts to the exact Arrow array for each Paimon type (e.g. a `Binary` /
-/// `VarBinary` field requires `arrow_array::BinaryArray`, not `LargeBinary` /
-/// `FixedSizeBinary`). Accepting a near-equivalent type here would pass
-/// validation but then fail deeper with a type-mismatch (or write files whose
-/// Arrow schema differs from the table), so it is rejected up front.
-fn validate_batch_schema(input: &ArrowSchema, target: &ArrowSchema) -> 
PyResult<()> {
-    let mismatch = || {
-        PyValueError::new_err(format!(
-            "Input schema is not consistent with the table schema. \
-             input: {input:?}, table: {target:?}"
-        ))
-    };
-    if input.fields().len() != target.fields().len() {
-        return Err(mismatch());
-    }
-    for (i, t) in input.fields().iter().zip(target.fields().iter()) {
-        if i.name() != t.name() {
-            return Err(mismatch());
-        }
-        if i.data_type() != t.data_type() {
-            return Err(mismatch());
-        }
-    }
-    Ok(())
-}
-
 type PartitionSpec = HashMap<String, Option<Datum>>;
 type PythonPartitionSpec = HashMap<String, Py<PyAny>>;
 
@@ -95,8 +62,6 @@ impl WriteContext {
         };
         Ok(WriteState {
             inner: Some(builder.new_write().map_err(to_py_err)?),
-            target_schema: 
paimon::arrow::build_target_arrow_schema(self.table.schema().fields())
-                .map_err(to_py_err)?,
             table_location: self.table.location().to_string(),
             commit_user: self.commit_user.clone(),
         })
@@ -307,7 +272,6 @@ impl PyStreamWriteBuilder {
 
 struct WriteState {
     inner: Option<TableWrite>,
-    target_schema: Arc<ArrowSchema>,
     table_location: String,
     commit_user: String,
 }
@@ -503,7 +467,6 @@ impl UpdateContext {
 impl WriteState {
     fn write_arrow(&mut self, py: Python<'_>, batch: &Bound<'_, PyAny>) -> 
PyResult<()> {
         let batch = RecordBatch::from_pyarrow_bound(batch)?;
-        validate_batch_schema(&batch.schema(), &self.target_schema)?;
         let inner = self
             .inner
             .as_mut()
diff --git a/bindings/python/tests/test_write.py 
b/bindings/python/tests/test_write.py
index 052a2c11..657965fc 100644
--- a/bindings/python/tests/test_write.py
+++ b/bindings/python/tests/test_write.py
@@ -340,3 +340,76 @@ def test_ltz_schema_alias_write_roundtrip(tmp_path, 
precision, unit, micros, fra
     assert sum(split.row_count() for split in 
table.new_read_builder().new_scan().plan().splits()) == 1
     rows = pa.Table.from_batches(ctx.sql("SELECT id, ts FROM paimon.wdb.t"))
     assert rows.to_pydict() == {"id": [1], "ts": [value]}
+
+
[email protected]("primary_key", [False, True])
[email protected]("stream", [False, True])
+def test_write_normalizes_nested_and_fixed_binary(tmp_path, primary_key, 
stream):
+    ctx = SQLContext()
+    ctx.register_catalog("paimon", {"warehouse": str(tmp_path)})
+    ctx.sql("CREATE SCHEMA paimon.wdb")
+    key = ", PRIMARY KEY (id)" if primary_key else ""
+    options = " WITH ('bucket' = '1')" if primary_key else ""
+    ctx.sql("CREATE TABLE paimon.wdb.t (id INT, data BINARY, items 
INT[]{}){}".format(
+        key, options))
+    table = _get_table(str(tmp_path))
+    builder = table.new_stream_write_builder() if stream else 
table.new_batch_write_builder()
+    writer = builder.new_write()
+    schema = pa.schema([("id", pa.int32()), ("data", pa.binary(2)),
+                        ("items", pa.list_(pa.int32()))])
+    data = pa.record_batch([[0, 1, 2, 3], [b"00", b"ab", None, b"cd"],
+                           [[], [1, None], None, []]], schema=schema)
+    # Rejected batches must not poison the writer or cause partial writes.
+    wrong = data.rename_columns(["items", "data", "id"])
+    with pytest.raises(ValueError):
+        writer.write_arrow(wrong)
+    writer.write_arrow(data.slice(1))
+    messages = writer.prepare_commit(True, 1) if stream else 
writer.prepare_commit()
+    commit = builder.new_commit()
+    if stream:
+        commit.commit(1, messages)
+    else:
+        commit.commit(messages)
+    result = pa.Table.from_batches(ctx.sql("SELECT * FROM paimon.wdb.t ORDER 
BY id"))
+    assert result.to_pylist() == data.slice(1).to_pylist()
+    writer.close()
+
+
+def test_format_write_schema_error_preserves_accepted_rows(tmp_path):
+    ctx = SQLContext()
+    ctx.register_catalog("paimon", {"warehouse": str(tmp_path)})
+    ctx.sql("CREATE SCHEMA paimon.wdb")
+    ctx.sql("CREATE TABLE paimon.wdb.t (id INT, name STRING) "
+            "WITH ('type' = 'format-table', 'file.format' = 'parquet')")
+    builder = _get_table(str(tmp_path)).new_batch_write_builder()
+    writer = builder.new_write()
+    try:
+        writer.write_arrow(_batch([1], ["first"]))
+        with pytest.raises(ValueError):
+            writer.write_arrow(_batch([9], 
["invalid"]).rename_columns(["wrong", "name"]))
+        writer.write_arrow(_batch([2], ["second"]))
+        builder.new_commit().commit(writer.prepare_commit())
+        result = pa.Table.from_batches(ctx.sql("SELECT * FROM paimon.wdb.t 
ORDER BY id"))
+        assert result.to_pydict() == {"id": [1, 2], "name": ["first", 
"second"]}
+    finally:
+        writer.close()
+
+
+def test_write_explicit_row_kinds_use_core_ignore_delete_filter(tmp_path):
+    ctx = SQLContext()
+    ctx.register_catalog("paimon", {"warehouse": str(tmp_path)})
+    ctx.sql("CREATE SCHEMA paimon.wdb")
+    ctx.sql("CREATE TABLE paimon.wdb.t (id INT, value INT, PRIMARY KEY (id)) "
+            "WITH ('bucket' = '1', 'merge-engine' = 'aggregation', "
+            "'fields.value.aggregate-function' = 'sum', 'ignore-delete' = 
'true')")
+    builder = _get_table(str(tmp_path)).new_batch_write_builder()
+    writer = builder.new_write()
+    try:
+        schema = pa.schema([("id", pa.int32()), ("value", pa.int32()),
+                            ("_VALUE_KIND", pa.int8())])
+        writer.write_arrow(pa.record_batch([[1, 1, 2], [10, 3, 7], [0, 3, 1]], 
schema=schema))
+        builder.new_commit().commit(writer.prepare_commit())
+        result = pa.Table.from_batches(ctx.sql("SELECT * FROM paimon.wdb.t"))
+        assert result.to_pydict() == {"id": [1], "value": [10]}
+    finally:
+        writer.close()
diff --git a/crates/paimon/src/arrow/format/parquet.rs 
b/crates/paimon/src/arrow/format/parquet.rs
index b3b54358..441fe049 100644
--- a/crates/paimon/src/arrow/format/parquet.rs
+++ b/crates/paimon/src/arrow/format/parquet.rs
@@ -49,7 +49,7 @@ use parquet::file::metadata::{
     KeyValue, PageIndexPolicy, ParquetMetaData, ParquetMetaDataReader, 
RowGroupMetaData,
 };
 use parquet::file::page_index::column_index::ColumnIndexMetaData;
-use parquet::file::properties::WriterProperties;
+use parquet::file::properties::{EnabledStatistics, WriterProperties};
 use parquet::file::statistics::Statistics as ParquetStatistics;
 use std::cmp::Ordering;
 use std::collections::HashMap;
@@ -339,11 +339,13 @@ impl ParquetFormatWriter {
     ) -> crate::Result<Self> {
         // Reject a bad codec before allocating the writer.
         let codec = parse_compression(compression, zstd_level)?;
+        let core_options = CoreOptions::new(format_options);
+        let page_index_enabled = 
core_options.parquet_write_page_index_enabled()?;
         let async_write = output.async_writer().await?;
         let input_schema = schema;
         let schema = timestamp_millis_schema(&input_schema);
-        let inner = create_parquet_arrow_writer(async_write, schema.clone(), 
codec)?;
-        let core_options = CoreOptions::new(format_options);
+        let inner =
+            create_parquet_arrow_writer(async_write, schema.clone(), codec, 
page_index_enabled)?;
         let stats_modes = write_fields
             .map(|fields| 
core_options.metadata_stats_modes(fields.iter().map(DataField::name)))
             .transpose()?;
@@ -362,8 +364,17 @@ fn create_parquet_arrow_writer(
     async_write: Box<dyn crate::io::AsyncFileWrite>,
     schema: arrow_schema::SchemaRef,
     codec: Compression,
+    page_index_enabled: bool,
 ) -> crate::Result<AsyncArrowWriter<Box<dyn crate::io::AsyncFileWrite>>> {
-    let props = WriterProperties::builder().set_compression(codec).build();
+    let props = WriterProperties::builder()
+        .set_compression(codec)
+        .set_statistics_enabled(if page_index_enabled {
+            EnabledStatistics::Page
+        } else {
+            EnabledStatistics::Chunk
+        })
+        .set_offset_index_disabled(!page_index_enabled)
+        .build();
     AsyncArrowWriter::try_new(async_write, schema, Some(props)).map_err(|e| {
         crate::Error::DataInvalid {
             message: format!("Failed to create parquet writer: {e}"),
@@ -3800,6 +3811,25 @@ mod tests {
         assert_eq!(total_rows, 3);
     }
 
+    #[tokio::test]
+    async fn test_parquet_writer_invalid_page_index_option_creates_no_file() {
+        let file_io = FileIOBuilder::new("memory").build().unwrap();
+        let path = "memory:/invalid_page_index_option.parquet";
+        let output = file_io.new_output(path).unwrap();
+        let options = HashMap::from([(
+            "parquet.write-page-index.enabled".to_string(),
+            "invalid".to_string(),
+        )]);
+        let result =
+            ParquetFormatWriter::new(&output, writer_arrow_schema(), "zstd", 
1, None, &options)
+                .await;
+        assert!(
+            matches!(result, Err(crate::Error::ConfigInvalid { message })
+            if message.contains("parquet.write-page-index.enabled"))
+        );
+        assert!(!file_io.exists(path).await.unwrap());
+    }
+
     #[tokio::test]
     async fn test_parquet_writer_multiple_batches() {
         let file_io = FileIOBuilder::new("memory").build().unwrap();
diff --git a/crates/paimon/src/arrow/shredding/variant.rs 
b/crates/paimon/src/arrow/shredding/variant.rs
index 6246fc80..aec6b18a 100644
--- a/crates/paimon/src/arrow/shredding/variant.rs
+++ b/crates/paimon/src/arrow/shredding/variant.rs
@@ -1390,7 +1390,14 @@ mod tests {
         ];
         let options = HashMap::from([(
             "parquet.variant.shreddingSchema".to_string(),
-            
r#"{"type":"ROW","fields":[{"name":"v","type":{"type":"ROW","fields":[{"name":"age","type":"BIGINT"},{"name":"city","type":"STRING"}]}}]}"#.to_string(),
+            
r#"{"type":"ROW","fields":[{"name":"v","type":{"type":"ROW","fields":[
+                {"name":"age","type":"BIGINT"},
+                {"name":"city","type":"VARCHAR"},
+                {"name":"fixed","type":"BINARY"},
+                {"name":"raw","type":"VARBINARY"},
+                {"name":"tags","type":{"type":"ARRAY","element":"VARCHAR"}},
+                
{"name":"address","type":{"type":"ROW","fields":[{"name":"street","type":"VARCHAR"}]}}
+            ]}}]}"#.to_string(),
         )]);
         let physical_fields = 
configured_variant_shredding_fields(&logical_fields, &options)
             .unwrap()
@@ -1398,14 +1405,24 @@ mod tests {
         assert_ne!(logical_fields, physical_fields);
 
         let variants = vec![
-            
GenericVariant::parse_json(r#"{"age":27,"city":"Beijing"}"#).unwrap(),
+            GenericVariant::parse_json(
+                
r#"{"age":27,"city":"Beijing","tags":["first","second"],"address":{"street":"Paimon
 Road"}}"#,
+            )
+            .unwrap(),
             
GenericVariant::parse_json(r#"{"city":"Hangzhou","other":"x"}"#).unwrap(),
             GenericVariant::parse_json(r#"{"age":"old"}"#).unwrap(),
+            // PyPaimon GenericVariant.from_python({'fixed': b'\x00\x01\x02',
+            //                                    'raw': b'\x03\x04\x05\x06'}).
+            GenericVariant::from_parts(
+                
hex::decode("020200010008113c030000000001023c0400000003040506").unwrap(),
+                hex::decode("01020005086669786564726177").unwrap(),
+            )
+            .unwrap(),
         ];
         let batch = RecordBatch::try_new(
             build_target_arrow_schema(&logical_fields).unwrap(),
             vec![
-                Arc::new(Int32Array::from(vec![1, 2, 3])),
+                Arc::new(Int32Array::from(vec![1, 2, 3, 4])),
                 variant_array_for_test(&variants),
             ],
         )
@@ -1415,6 +1432,33 @@ mod tests {
             batch_to_shredded_physical(&batch, &logical_fields, 
&physical_fields).unwrap();
         assert!(is_shredded_variant_array(physical.column(1).as_ref()));
 
+        let shredded = physical
+            .column(1)
+            .as_any()
+            .downcast_ref::<StructArray>()
+            .unwrap();
+        let typed = shredded
+            .column_by_name("typed_value")
+            .unwrap()
+            .as_any()
+            .downcast_ref::<StructArray>()
+            .unwrap();
+        for (name, expected) in [("fixed", &[0, 1, 2][..]), ("raw", &[3, 4, 5, 
6][..])] {
+            let field = typed
+                .column_by_name(name)
+                .unwrap()
+                .as_any()
+                .downcast_ref::<StructArray>()
+                .unwrap();
+            let binary = field
+                .column_by_name("typed_value")
+                .unwrap()
+                .as_any()
+                .downcast_ref::<BinaryArray>()
+                .unwrap();
+            assert_eq!(binary.value(3), expected);
+        }
+
         let assembled = 
assemble_shredded_variant_array(physical.column(1).as_ref()).unwrap();
         let assembled = 
assembled.as_any().downcast_ref::<StructArray>().unwrap();
         let values = assembled
diff --git a/crates/paimon/src/spec/aggregation.rs 
b/crates/paimon/src/spec/aggregation.rs
index 4742482d..9a05e68f 100644
--- a/crates/paimon/src/spec/aggregation.rs
+++ b/crates/paimon/src/spec/aggregation.rs
@@ -17,11 +17,11 @@
 
 use std::collections::HashMap;
 
+use crate::spec::core_options::IGNORE_DELETE_FALLBACK_KEYS;
 use crate::spec::{CoreOptions, DataField, DataType, VarCharType};
 
 const MERGE_ENGINE_OPTION: &str = "merge-engine";
 const AGGREGATION_ENGINE: &str = "aggregation";
-const IGNORE_DELETE_OPTION: &str = "ignore-delete";
 const IGNORE_DELETE_SUFFIX: &str = ".ignore-delete";
 const AGGREGATION_REMOVE_RECORD_ON_DELETE_OPTION: &str = 
"aggregation.remove-record-on-delete";
 const FIELDS_DEFAULT_AGG_FUNCTION_OPTION: &str = 
"fields.default-aggregate-function";
@@ -547,8 +547,9 @@ pub(crate) fn validate_no_aggregation_on_sequence_field(
 }
 
 fn is_unsupported_aggregation_option(key: &str) -> bool {
-    key == IGNORE_DELETE_OPTION
-        || key.ends_with(IGNORE_DELETE_SUFFIX)
+    // Global ignore-delete and its Java aliases filter rows before 
aggregation.
+    // Field-level aggregation uses ignore-retract instead.
+    (key.ends_with(IGNORE_DELETE_SUFFIX) && 
!IGNORE_DELETE_FALLBACK_KEYS.contains(&key))
         || is_fields_option_with_suffix(key, SEQUENCE_GROUP_SUFFIX)
 }
 
@@ -621,6 +622,30 @@ mod tests {
         );
     }
 
+    #[test]
+    fn test_aggregation_accepts_global_ignore_delete_and_java_aliases() {
+        for key in
+            
std::iter::once("ignore-delete").chain(IGNORE_DELETE_FALLBACK_KEYS.iter().copied())
+        {
+            for value in ["true", "false"] {
+                let options = aggregation_options(&[(key, value)]);
+                let config = AggregationConfig::new(&options);
+                assert_eq!(
+                    config
+                        .validate_create_mode(&pk(), &sample_fields())
+                        .unwrap(),
+                    Some(AggregationMode::Basic),
+                    "{key}={value}"
+                );
+                assert_eq!(
+                    config.validate_runtime_mode(true, "default.t").unwrap(),
+                    Some(AggregationMode::Basic),
+                    "{key}={value}"
+                );
+            }
+        }
+    }
+
     #[test]
     fn test_validate_create_mode_ignores_non_pk_tables() {
         let options = aggregation_options(&[("fields.x.ignore-retract", 
"true")]);
@@ -648,7 +673,7 @@ mod tests {
     #[test]
     fn test_validate_create_mode_rejects_unsupported_options() {
         for key in [
-            IGNORE_DELETE_OPTION,
+            "aggregation.ignore-delete",
             "fields.price.ignore-delete",
             "fields.price.sequence-group",
         ] {
diff --git a/crates/paimon/src/spec/core_options.rs 
b/crates/paimon/src/spec/core_options.rs
index cb69c888..984e1380 100644
--- a/crates/paimon/src/spec/core_options.rs
+++ b/crates/paimon/src/spec/core_options.rs
@@ -95,6 +95,7 @@ const MANIFEST_SORT_ENABLED_OPTION: &str = 
"manifest-sort.enabled";
 const WRITE_PARQUET_BUFFER_SIZE_OPTION: &str = "write.parquet-buffer-size";
 const READ_BATCH_SIZE_OPTION: &str = "read.batch-size";
 const PARQUET_FILTER_COLUMN_INDEX_ENABLED_OPTION: &str = 
"parquet.filter.columnindex.enabled";
+const PARQUET_WRITE_PAGE_INDEX_ENABLED_OPTION: &str = 
"parquet.write-page-index.enabled";
 const PARQUET_ROW_GROUP_PARALLELISM_OPTION: &str = 
"read.parquet.row-group.parallelism";
 pub(crate) const PARQUET_ROW_GROUP_MAX_INFLIGHT_BYTES_OPTION: &str =
     "read.parquet.row-group.max-inflight-bytes";
@@ -112,7 +113,7 @@ pub(crate) const CHANGELOG_PRODUCER_OPTION: &str = 
"changelog-producer";
 const ROWKIND_FIELD_OPTION: &str = "rowkind.field";
 const IGNORE_DELETE_OPTION: &str = "ignore-delete";
 const IGNORE_UPDATE_BEFORE_OPTION: &str = "ignore-update-before";
-const IGNORE_DELETE_FALLBACK_KEYS: &[&str] = &[
+pub(super) const IGNORE_DELETE_FALLBACK_KEYS: &[&str] = &[
     "first-row.ignore-delete",
     "deduplicate.ignore-delete",
     "partial-update.ignore-delete",
@@ -770,6 +771,22 @@ impl<'a> CoreOptions<'a> {
         }
     }
 
+    /// Whether newly written Parquet files include page indexes. Default is 
true.
+    pub fn parquet_write_page_index_enabled(&self) -> crate::Result<bool> {
+        let Some(raw) = 
self.options.get(PARQUET_WRITE_PAGE_INDEX_ENABLED_OPTION) else {
+            return Ok(true);
+        };
+        match raw.trim().to_ascii_lowercase().as_str() {
+            "true" => Ok(true),
+            "false" => Ok(false),
+            _ => Err(crate::Error::ConfigInvalid {
+                message: format!(
+                    "Option '{PARQUET_WRITE_PAGE_INDEX_ENABLED_OPTION}' must 
be true or false, got: {raw}"
+                ),
+            }),
+        }
+    }
+
     /// The declared [`TableType`], defaulting to [`TableType::Table`].
     /// Fails on a value this client does not know.
     pub fn table_type(&self) -> crate::Result<TableType> {
diff --git a/crates/paimon/src/spec/schema.rs b/crates/paimon/src/spec/schema.rs
index 0ca57f5a..c7d15736 100644
--- a/crates/paimon/src/spec/schema.rs
+++ b/crates/paimon/src/spec/schema.rs
@@ -1823,15 +1823,6 @@ impl Schema {
             });
         }
 
-        let merge_engine = core
-            .merge_engine()
-            .map_err(Self::options_error_to_config_invalid)?;
-        if merge_engine != MergeEngine::Deduplicate {
-            return Err(crate::Error::ConfigInvalid {
-                message: "rowkind.field only supports 
merge-engine=deduplicate".to_string(),
-            });
-        }
-
         let producer = core
             .try_changelog_producer()
             .map_err(Self::options_error_to_config_invalid)?;
@@ -3353,7 +3344,8 @@ mod tests {
     #[test]
     fn test_aggregation_schema_validation_rejects_unsupported_options() {
         for (key, value) in [
-            ("ignore-delete", "true"),
+            ("aggregation.ignore-delete", "true"),
+            ("fields.value.ignore-delete", "true"),
             ("fields.value.sequence-group", "g1"),
         ] {
             let err = Schema::builder()
@@ -5812,20 +5804,17 @@ mod tests {
     }
 
     #[test]
-    fn rowkind_field_rejects_non_deduplicate_merge_engine() {
-        let err = Schema::builder()
-            .column("id", DataType::Int(IntType::new()))
-            .column("op", DataType::VarChar(VarCharType::string_type()))
-            .primary_key(["id"])
-            .option("merge-engine", "partial-update")
-            .option("rowkind.field", "op")
-            .build()
-            .unwrap_err();
-        assert!(
-            matches!(err, crate::Error::ConfigInvalid { ref message }
-                if message.contains("deduplicate")),
-            "got {err:?}"
-        );
+    fn rowkind_field_accepts_supported_merge_engines() {
+        for engine in ["deduplicate", "first-row", "partial-update", 
"aggregation"] {
+            Schema::builder()
+                .column("id", DataType::Int(IntType::new()))
+                .column("op", DataType::VarChar(VarCharType::string_type()))
+                .primary_key(["id"])
+                .option("merge-engine", engine)
+                .option("rowkind.field", "op")
+                .build()
+                .unwrap();
+        }
     }
 
     #[test]
diff --git a/crates/paimon/src/spec/types.rs b/crates/paimon/src/spec/types.rs
index 551dd08e..ab73da7e 100644
--- a/crates/paimon/src/spec/types.rs
+++ b/crates/paimon/src/spec/types.rs
@@ -664,6 +664,11 @@ impl FromStr for BinaryType {
             .fail();
         }
 
+        // Java permits omitted lengths and uses the type's default length.
+        if s == "BINARY" || s == "BINARY NOT NULL" {
+            return Self::with_nullable(s == "BINARY", Self::DEFAULT_LENGTH);
+        }
+
         let (open_bracket, close_bracket) = 
serde_utils::extract_brackets_pos(s, "BinaryType")?;
         let length_str = &s[open_bracket + 1..close_bracket];
         let length = length_str
@@ -1691,6 +1696,11 @@ impl FromStr for VarBinaryType {
             .fail();
         }
 
+        // Java permits omitted lengths and uses the type's default length.
+        if s == "VARBINARY" || s == "VARBINARY NOT NULL" {
+            return Self::try_new(s == "VARBINARY", Self::DEFAULT_LENGTH);
+        }
+
         let (open_bracket, close_bracket) = 
serde_utils::extract_brackets_pos(s, "VarBinaryType")?;
         let length_str = &s[open_bracket + 1..close_bracket];
         let length = length_str
@@ -1791,6 +1801,11 @@ impl FromStr for VarCharType {
             .fail();
         }
 
+        // Java permits omitted lengths and uses the type's default length.
+        if s == "VARCHAR" || s == "VARCHAR NOT NULL" {
+            return Self::with_nullable(s == "VARCHAR", Self::DEFAULT_LENGTH);
+        }
+
         let (open_bracket, close_bracket) = 
serde_utils::extract_brackets_pos(s, "VarCharType")?;
         let length_str = &s[open_bracket + 1..close_bracket];
         let length = length_str
@@ -2878,6 +2893,28 @@ mod tests {
         }
     }
 
+    #[test]
+    fn test_default_length_string_types_match_java() {
+        // Java DataTypeJsonParser uses each type's DEFAULT_LENGTH when 
omitted.
+        let cases = [
+            ("VARCHAR", DataType::VarChar(VarCharType::default())),
+            ("BINARY", DataType::Binary(BinaryType::default())),
+            ("VARBINARY", DataType::VarBinary(VarBinaryType::default())),
+        ];
+        for (name, expected) in cases {
+            for nullable in [true, false] {
+                let suffix = if nullable { "" } else { " NOT NULL" };
+                let expected = expected.copy_with_nullable(nullable).unwrap();
+                for input in [format!("{name}{suffix}"), 
format!("{name}(1){suffix}")] {
+                    let parsed: DataType =
+                        
serde_json::from_value(serde_json::json!(input)).unwrap();
+                    assert_eq!(parsed, expected, "{input}");
+                    assert_eq!(parsed.to_string(), 
format!("{name}(1){suffix}"));
+                }
+            }
+        }
+    }
+
     #[test]
     fn test_string_type_serializes_like_java() {
         let nullable = DataType::VarChar(VarCharType::string_type());
@@ -3299,7 +3336,9 @@ mod tests {
     fn test_datatype_deserialize_rejects_unknown_shapes() {
         assert!(serde_json::from_str::<DataType>("\"TUPLE\"").is_err());
         assert!(serde_json::from_str::<DataType>("\"INT NULLABLE\"").is_err());
-        assert!(serde_json::from_str::<DataType>("\"VARCHAR\"").is_err());
+        for malformed in ["VARCHAR(", "BINARY(", "VARBINARY("] {
+            
assert!(serde_json::from_value::<DataType>(serde_json::json!(malformed)).is_err());
+        }
         
assert!(serde_json::from_str::<DataType>(r#"{"type":"TUPLE","element":"INT"}"#).is_err());
         
assert!(serde_json::from_str::<DataType>(r#"{"element":"INT"}"#).is_err());
         assert!(serde_json::from_str::<DataType>("7").is_err());
diff --git a/crates/paimon/src/table/format_table_write_tests.rs 
b/crates/paimon/src/table/format_table_write_tests.rs
index f12c9e6d..f9bbe50e 100644
--- a/crates/paimon/src/table/format_table_write_tests.rs
+++ b/crates/paimon/src/table/format_table_write_tests.rs
@@ -393,10 +393,31 @@ async fn 
rejected_batch_does_not_stage_or_publish_a_file() {
     .unwrap();
     let error = write.write_arrow_batch(&wrong).await.err().unwrap();
     assert!(error.to_string().contains("expects"));
-    assert!(write.prepare_commit().await.is_err());
+    assert!(write.prepare_commit().await.unwrap().is_empty());
     assert!(visible_files(&table, "dt=a").await.is_empty());
 }
 
+#[tokio::test]
+async fn schema_error_preserves_accepted_rows_and_writer_recovery() {
+    let table = memory_table("format_schema_recovery", true, &[]);
+    let builder = table.new_write_builder();
+    let mut write = builder.new_write().unwrap();
+    write.write_arrow_batch(&batch(&[("a", 1)])).await.unwrap();
+    let error = write
+        .write_arrow_batch(&unpartitioned_batch(&[2]))
+        .await
+        .unwrap_err();
+    assert!(error.to_string().contains("expects 2 columns"));
+    write.write_arrow_batch(&batch(&[("a", 3)])).await.unwrap();
+    builder
+        .new_commit()
+        .commit(write.prepare_commit().await.unwrap())
+        .await
+        .unwrap();
+    assert_eq!(ids(&table).await, [1, 3]);
+    assert_eq!(visible_files(&table, "dt=a").await.len(), 1);
+}
+
 #[tokio::test]
 async fn one_writer_can_prepare_more_than_one_append() {
     let table = memory_table("format_writer_reuse", true, &[]);
diff --git a/crates/paimon/src/table/format_table_writer.rs 
b/crates/paimon/src/table/format_table_writer.rs
index a41f0106..0d3e10c3 100644
--- a/crates/paimon/src/table/format_table_writer.rs
+++ b/crates/paimon/src/table/format_table_writer.rs
@@ -213,6 +213,8 @@ impl FormatTableWriter {
                 source: None,
             });
         }
+        // Reject schema mistakes before entering the mutation/failure path.
+        self.check_schema(batch)?;
         let result = self.write_inner(batch).await;
         if result.is_err() {
             self.failed = true;
@@ -222,7 +224,6 @@ impl FormatTableWriter {
     }
 
     async fn write_inner(&mut self, batch: &RecordBatch) -> Result<()> {
-        self.check_schema(batch)?;
         if batch.num_rows() == 0 {
             return Ok(());
         }
diff --git a/crates/paimon/src/table/kv_file_writer.rs 
b/crates/paimon/src/table/kv_file_writer.rs
index 4ddaa63b..000a6057 100644
--- a/crates/paimon/src/table/kv_file_writer.rs
+++ b/crates/paimon/src/table/kv_file_writer.rs
@@ -44,8 +44,8 @@ use 
crate::table::managed_blob_reference::ManagedBlobReferences;
 use crate::table::managed_blob_writer::ManagedBlobWriteState;
 use crate::table::prepared_files::PreparedFiles;
 use crate::table::sort_merge::{
-    AggregateMergeFunction, BufferedBatch, MergeFunction, MergeResult, 
MergeRow,
-    PartialUpdateMergeFunction,
+    AggregateMergeFunction, BufferedBatch, FirstRowMergeFunction, 
MergeFunction, MergeResult,
+    MergeRow, PartialUpdateMergeFunction,
 };
 use crate::Result;
 use arrow_array::{Array, BooleanArray, Int64Array, Int8Array, RecordBatch, 
UInt32Array};
@@ -1089,40 +1089,80 @@ impl KeyValueFileWriter {
         batch: &RecordBatch,
         sorted_indices: &arrow_array::UInt32Array,
     ) -> Result<Vec<u32>> {
-        let n = sorted_indices.len();
-        if n == 0 {
-            return Ok(vec![]);
+        if sorted_indices.is_empty() {
+            return Ok(Vec::new());
         }
 
-        let rows = self.convert_key_rows(batch)?;
-
-        let mut result: Vec<u32> = Vec::with_capacity(n);
-        // Track the start of the current key group and the candidate winner.
-        let mut group_winner = sorted_indices.value(0);
-
-        for i in 1..n {
-            let cur = sorted_indices.value(i);
-            if rows.row(group_winner as usize) == rows.row(cur as usize) {
-                // Same key group — update winner based on merge engine.
-                match self.config.merge_engine {
-                    // Deduplicate: keep last (highest seq), which is the 
current row
-                    // since we sorted ascending.
-                    MergeEngine::Deduplicate => group_winner = cur,
-                    // FirstRow: keep first (lowest seq), so don't update.
-                    MergeEngine::FirstRow => {}
-                    MergeEngine::PartialUpdate | MergeEngine::Aggregation => 
unreachable!(
-                        "{:?} should use select_flush_indices and skip dedup",
-                        self.config.merge_engine
-                    ),
+        let schema = batch.schema();
+        let value_kinds = if self.config.merge_engine == MergeEngine::FirstRow 
{
+            batch
+                .column_by_name(VALUE_KIND_FIELD_NAME)
+                .map(|column| {
+                    column.as_any().downcast_ref::<Int8Array>().ok_or_else(|| {
+                        crate::Error::DataInvalid {
+                            message: "_VALUE_KIND column must be Int8".into(),
+                            source: None,
+                        }
+                    })
+                })
+                .transpose()?
+        } else {
+            None
+        };
+        let first_row_merge = FirstRowMergeFunction {
+            ignore_delete: 
CoreOptions::new(&self.config.table_options).ignore_delete(),
+        };
+        let key_rows = self.convert_key_rows(batch)?;
+        let mut result = Vec::with_capacity(sorted_indices.len());
+        let mut start = 0;
+        while start < sorted_indices.len() {
+            let first = sorted_indices.value(start);
+            let mut end = start + 1;
+            while end < sorted_indices.len()
+                && key_rows.row(first as usize) == 
key_rows.row(sorted_indices.value(end) as usize)
+            {
+                end += 1;
+            }
+            match self.config.merge_engine {
+                MergeEngine::Deduplicate => 
result.push(sorted_indices.value(end - 1)),
+                // Java's ReducerMergeFunctionWrapper bypasses the merge 
function
+                // for singleton groups. Insert-only batches need no kind 
validation.
+                MergeEngine::FirstRow if end - start == 1 || 
value_kinds.is_none() => {
+                    result.push(first);
                 }
-            } else {
-                // New key group — emit the winner of the previous group.
-                result.push(group_winner);
-                group_winner = cur;
+                MergeEngine::FirstRow => {
+                    let kinds = value_kinds.unwrap();
+                    let rows = sorted_indices.values()[start..end]
+                        .iter()
+                        .map(|&idx| MergeRow {
+                            batch_idx: 0,
+                            row_idx: idx as usize,
+                            // Keep the established user-sequence and arrival 
order.
+                            sequence_number: 0,
+                            user_sequence: None,
+                            value_kind: if kinds.is_null(idx as usize) {
+                                RowKind::Insert as i8
+                            } else {
+                                kinds.value(idx as usize)
+                            },
+                        })
+                        .collect::<Vec<_>>();
+                    // First-row returns a source index without 
reading/materializing values.
+                    match first_row_merge.merge(&rows, &[], &[], &schema)? {
+                        MergeResult::SourceRow { row_idx, .. } => 
result.push(row_idx as u32),
+                        MergeResult::Omit => {}
+                        MergeResult::MaterializedRow(_) | 
MergeResult::MaterializedDeleteRow(_) => {
+                            unreachable!("first-row merge must return an input 
row")
+                        }
+                    }
+                }
+                MergeEngine::PartialUpdate | MergeEngine::Aggregation => 
unreachable!(
+                    "{:?} should use select_flush_indices and skip dedup",
+                    self.config.merge_engine
+                ),
             }
+            start = end;
         }
-        // Emit the last group's winner.
-        result.push(group_winner);
         Ok(result)
     }
 
@@ -1572,6 +1612,115 @@ mod tests {
         assert_eq!(deduped, vec![0, 2]);
     }
 
+    fn first_row_kind_batch(ids: Vec<i32>, kinds: Vec<i8>) -> RecordBatch {
+        let len = ids.len();
+        RecordBatch::try_new(
+            Arc::new(ArrowSchema::new(vec![
+                ArrowField::new("id", ArrowDataType::Int32, false),
+                ArrowField::new("seq", ArrowDataType::Int64, false),
+                ArrowField::new("value", ArrowDataType::Int32, false),
+                ArrowField::new(VALUE_KIND_FIELD_NAME, ArrowDataType::Int8, 
false),
+            ])),
+            vec![
+                Arc::new(Int32Array::from(ids)),
+                Arc::new(Int64Array::from_iter_values(0..len as i64)),
+                Arc::new(Int32Array::from_iter_values(10..10 + len as i32)),
+                Arc::new(Int8Array::from(kinds)),
+            ],
+        )
+        .unwrap()
+    }
+
+    #[tokio::test]
+    async fn test_first_row_flush_rejects_every_retract_in_a_multi_row_group() 
{
+        for retract in [RowKind::Delete, RowKind::UpdateBefore] {
+            for position in 1..=3 {
+                let mut kinds = vec![RowKind::Insert as i8; 5];
+                kinds[position] = retract as i8;
+                let batch = first_row_kind_batch(vec![0, 1, 1, 1, 2], kinds);
+                let mut config = test_write_config(MergeEngine::FirstRow);
+                config.write_buffer_size = i64::MAX;
+                config.input_changelog = true;
+                let mut writer = KeyValueFileWriter::new(
+                    FileIOBuilder::new("memory").build().unwrap(),
+                    config,
+                    0,
+                )
+                .unwrap();
+                // The invalid group spans batches and follows a valid key 
group.
+                writer.write(&batch.slice(0, 3)).await.unwrap();
+                writer.write(&batch.slice(3, 2)).await.unwrap();
+                let error = writer
+                    .prepare_commit()
+                    .await
+                    .err()
+                    .expect("retract must be rejected");
+                assert!(
+                    error
+                        .to_string()
+                        .contains("does not support DELETE or UPDATE_BEFORE"),
+                    "{retract:?} at {position}: {error}"
+                );
+                assert!(writer.written_files.is_empty());
+                assert!(writer.written_changelog_files.is_empty());
+            }
+        }
+    }
+
+    #[tokio::test]
+    async fn 
test_first_row_flush_preserves_singleton_retracts_like_java_reducer() {
+        let batch = first_row_kind_batch(vec![1, 2, 3], vec![1, 3, 2]);
+        let io = FileIOBuilder::new("memory").build().unwrap();
+        let mut writer =
+            KeyValueFileWriter::new(io.clone(), 
test_write_config(MergeEngine::FirstRow), 0)
+                .unwrap();
+        writer.write(&batch).await.unwrap();
+        let prepared = writer.prepare_commit().await.unwrap();
+        let stored = read_kv_file(&io, &prepared.data_files[0]).await;
+        assert_eq!(prepared.data_files[0].delete_row_count, Some(2));
+        assert_eq!(
+            stored
+                .column_by_name(VALUE_KIND_FIELD_NAME)
+                .unwrap()
+                .as_ref(),
+            batch
+                .column_by_name(VALUE_KIND_FIELD_NAME)
+                .unwrap()
+                .as_ref()
+        );
+        assert_eq!(
+            stored.column_by_name("value").unwrap().as_ref(),
+            batch.column_by_name("value").unwrap().as_ref()
+        );
+    }
+
+    #[tokio::test]
+    async fn 
test_first_row_flush_ignore_delete_uses_first_add_and_omits_retract_group() {
+        let batch = first_row_kind_batch(vec![1, 1, 1, 1, 2, 2], vec![3, 2, 0, 
1, 1, 3]);
+        let io = FileIOBuilder::new("memory").build().unwrap();
+        let mut config = test_write_config(MergeEngine::FirstRow);
+        config
+            .table_options
+            .insert("ignore-delete".into(), "true".into());
+        let mut writer = KeyValueFileWriter::new(io.clone(), config, 
0).unwrap();
+        writer.write(&batch).await.unwrap();
+        let prepared = writer.prepare_commit().await.unwrap();
+        assert_eq!(prepared.data_files[0].row_count, 1);
+        assert_eq!(prepared.data_files[0].delete_row_count, Some(0));
+        let stored = read_kv_file(&io, &prepared.data_files[0]).await;
+        assert_eq!(
+            stored.column_by_name("value").unwrap().as_ref(),
+            &Int32Array::from(vec![11]) as &dyn Array
+        );
+        assert_eq!(
+            stored
+                .column_by_name(VALUE_KIND_FIELD_NAME)
+                .unwrap()
+                .as_ref(),
+            &Int8Array::from(vec![RowKind::UpdateAfter as i8]) as &dyn Array
+        );
+    }
+
     fn partial_update_writer() -> KeyValueFileWriter {
         KeyValueFileWriter::new(
             FileIOBuilder::new("memory").build().unwrap(),
diff --git a/crates/paimon/src/table/table_write.rs 
b/crates/paimon/src/table/table_write.rs
index b50a7501..1f28cc6d 100644
--- a/crates/paimon/src/table/table_write.rs
+++ b/crates/paimon/src/table/table_write.rs
@@ -26,7 +26,7 @@ use crate::resource::ResourceContext;
 use crate::spec::PartitionComputer;
 use crate::spec::{
     first_row_supports_changelog_producer, BinaryRow, ChangelogProducer, 
CoreOptions, DataField,
-    DataType, MergeEngine, RowKindFilter, EMPTY_SERIALIZED_ROW, 
POSTPONE_BUCKET,
+    DataType, MergeEngine, RowKind, RowKindFilter, EMPTY_SERIALIZED_ROW, 
POSTPONE_BUCKET,
     VALUE_KIND_FIELD_NAME,
 };
 use crate::table::bucket_assigner::{BucketAssignerEnum, PartitionBucketKey};
@@ -882,7 +882,7 @@ impl TableWrite {
 
     fn enrich_rowkind_batch(&self, batch: &RecordBatch) -> Result<RecordBatch> 
{
         let Some(generator) = &self.row_kind_generator else {
-            return Ok(batch.clone());
+            return self.filter_rowkind_batch(batch);
         };
         if batch
             .schema()
@@ -896,23 +896,37 @@ impl TableWrite {
             });
         }
 
+        let kinds = (0..batch.num_rows())
+            .map(|row| generator.generate(batch, row).map(|kind| 
kind.to_value()))
+            .collect::<Result<Vec<_>>>()?;
+        let batch = Self::add_per_row_value_kind_column(batch, kinds)?;
+        self.filter_rowkind_batch(&batch)
+    }
+
+    /// Both generated row kinds and explicit caller-provided kinds use the
+    /// same Java RowKindFilter before bucket assignment and changelog writing.
+    fn filter_rowkind_batch(&self, batch: &RecordBatch) -> Result<RecordBatch> 
{
+        let Some(filter) = &self.row_kind_filter else {
+            return Ok(batch.clone());
+        };
+        let Some(column) = batch.column_by_name(VALUE_KIND_FIELD_NAME) else {
+            return Ok(batch.clone());
+        };
+        let kinds = column
+            .as_any()
+            .downcast_ref::<arrow_array::Int8Array>()
+            .ok_or_else(|| crate::Error::DataInvalid {
+                message: "_VALUE_KIND column must be Int8".to_string(),
+                source: None,
+            })?;
         let mut keep_rows = Vec::new();
-        let mut kinds = Vec::new();
-        for row in 0..batch.num_rows() {
-            let kind = generator.generate(batch, row)?;
-            if let Some(filter) = &self.row_kind_filter {
-                if !filter.test(kind) {
-                    continue;
-                }
+        for (row, kind) in kinds.iter().enumerate() {
+            let kind = RowKind::from_value(kind.unwrap_or(RowKind::Insert as 
i8))?;
+            if filter.test(kind) {
+                keep_rows.push(row);
             }
-            keep_rows.push(row);
-            kinds.push(kind.to_value());
-        }
-        if keep_rows.is_empty() {
-            return Ok(RecordBatch::new_empty(batch.schema()));
         }
-        let filtered = take_rows(batch, &keep_rows)?;
-        Self::add_per_row_value_kind_column(&filtered, kinds)
+        take_rows(batch, &keep_rows)
     }
 
     fn add_per_row_value_kind_column(
diff --git a/crates/paimon/tests/parquet_write_page_index_test.rs 
b/crates/paimon/tests/parquet_write_page_index_test.rs
new file mode 100644
index 00000000..134d8d8c
--- /dev/null
+++ b/crates/paimon/tests/parquet_write_page_index_test.rs
@@ -0,0 +1,122 @@
+// 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 std::sync::Arc;
+
+use arrow_array::builder::{Int32Builder, ListBuilder};
+use arrow_array::{ArrayRef, Int32Array, RecordBatch};
+use arrow_schema::{Field, Schema as ArrowSchema};
+use common::incremental_helpers::{memory_table, persist_table_schema, 
setup_dirs};
+use paimon::spec::{ArrayType, DataType, IntType, Schema, TableSchema};
+use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
+
+#[tokio::test]
+async fn parquet_page_indexes_follow_write_option_for_data_and_changelog() {
+    for setting in [None, Some("true"), Some("false")] {
+        for changelog in [false, true] {
+            let mut schema = Schema::builder()
+                .column("id", DataType::Int(IntType::new()))
+                .column(
+                    "items",
+                    
DataType::Array(ArrayType::new(DataType::Int(IntType::new()))),
+                );
+            if changelog {
+                schema = schema
+                    .primary_key(["id"])
+                    .option("bucket", "1")
+                    .option("changelog-producer", "input");
+            }
+            if let Some(setting) = setting {
+                schema = schema.option("parquet.write-page-index.enabled", 
setting);
+            }
+            // Read-side pruning is independent of whether the writer emits 
indexes.
+            let schema = schema
+                .option("parquet.filter.columnindex.enabled", "false")
+                .build()
+                .unwrap();
+            let path = "memory:/parquet_page_indexes";
+            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;
+
+            let mut items = ListBuilder::new(Int32Builder::new());
+            items.values().append_value(10);
+            items.values().append_null();
+            items.append(true);
+            items.append(true);
+            items.append(false);
+            let columns: Vec<ArrayRef> = vec![
+                Arc::new(Int32Array::from(vec![1, 2, 3])),
+                Arc::new(items.finish()),
+            ];
+            let batch = RecordBatch::try_new(
+                Arc::new(ArrowSchema::new(vec![
+                    Field::new("id", columns[0].data_type().clone(), false),
+                    Field::new("items", columns[1].data_type().clone(), true),
+                ])),
+                columns,
+            )
+            .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(), 1);
+            assert_eq!(messages[0].new_files.len(), 1);
+            assert_eq!(
+                messages[0].new_changelog_files.len(),
+                usize::from(changelog)
+            );
+
+            for file in messages[0]
+                .new_files
+                .iter()
+                .chain(&messages[0].new_changelog_files)
+            {
+                let file_path = format!("{path}/bucket-0/{}", file.file_name);
+                let bytes = 
file_io.new_input(&file_path).unwrap().read().await.unwrap();
+                let reader = 
ParquetRecordBatchReaderBuilder::try_new(bytes).unwrap();
+                let enabled = setting != Some("false");
+                for row_group in reader.metadata().row_groups() {
+                    for column in row_group.columns() {
+                        assert_eq!(column.column_index_offset().is_some(), 
enabled);
+                        assert_eq!(column.offset_index_offset().is_some(), 
enabled);
+                        assert!(column.statistics().is_some());
+                    }
+                }
+                let ids = reader
+                    .build()
+                    .unwrap()
+                    .flat_map(|batch| {
+                        batch
+                            .unwrap()
+                            .column_by_name("id")
+                            .unwrap()
+                            .as_any()
+                            .downcast_ref::<Int32Array>()
+                            .unwrap()
+                            .values()
+                            .to_vec()
+                    })
+                    .collect::<Vec<_>>();
+                assert_eq!(ids, vec![1, 2, 3]);
+            }
+        }
+    }
+}
diff --git a/crates/paimon/tests/rowkind_field_test.rs 
b/crates/paimon/tests/rowkind_field_test.rs
index 07dd2fb4..d4cd97fa 100644
--- a/crates/paimon/tests/rowkind_field_test.rs
+++ b/crates/paimon/tests/rowkind_field_test.rs
@@ -38,7 +38,7 @@ use arrow_schema::{DataType as ArrowDataType, Field as 
ArrowField, Schema as Arr
 use futures::StreamExt;
 use paimon::arrow::paimon_type_to_arrow;
 use paimon::spec::{
-    ArrayType, DataField, DataType, IntType, RowType, Schema, TableSchema, 
VarBinaryType,
+    ArrayType, DataField, DataType, IntType, RowKind, RowType, Schema, 
TableSchema, VarBinaryType,
     VarCharType, VALUE_KIND_FIELD_NAME,
 };
 
@@ -838,6 +838,94 @@ async fn rowkind_field_ignore_delete() {
     );
 }
 
+#[tokio::test]
+async fn aggregation_ignore_delete_supports_generated_and_explicit_row_kinds() 
{
+    for explicit_kind in [true, false] {
+        for option in [
+            "ignore-delete",
+            "first-row.ignore-delete",
+            "deduplicate.ignore-delete",
+            "partial-update.ignore-delete",
+        ] {
+            let path = 
format!("memory:/rowkind_field/aggregation_{option}_{explicit_kind}");
+            let mut schema = Schema::builder()
+                .column("id", DataType::Int(IntType::new()))
+                .column("value", DataType::Int(IntType::new()))
+                .column("kind", DataType::VarChar(VarCharType::string_type()))
+                .primary_key(["id"])
+                .option("bucket", "1")
+                .option("merge-engine", "aggregation")
+                .option("fields.value.aggregate-function", "sum")
+                .option(option, "true");
+            if !explicit_kind {
+                schema = schema.option("rowkind.field", "kind");
+            }
+            let (file_io, table) =
+                memory_table(&path, TableSchema::new(0, 
&schema.build().unwrap()));
+            setup_dirs(&file_io, &path).await;
+            persist_table_schema(&file_io, &path, table.schema()).await;
+            let make_batch = |ids, values, kinds: Vec<&str>| {
+                let value_kinds = kinds
+                    .iter()
+                    .map(|kind| 
RowKind::from_short_string(kind).unwrap().to_value())
+                    .collect::<Vec<_>>();
+                let batch = make_batch_with_rowkind(ids, values, kinds, 
"kind");
+                if !explicit_kind {
+                    return batch;
+                }
+                let mut fields = batch.schema().fields().to_vec();
+                fields.push(Arc::new(ArrowField::new(
+                    VALUE_KIND_FIELD_NAME,
+                    ArrowDataType::Int8,
+                    false,
+                )));
+                let mut columns = batch.columns().to_vec();
+                columns.push(Arc::new(Int8Array::from(value_kinds)));
+                RecordBatch::try_new(Arc::new(ArrowSchema::new(fields)), 
columns).unwrap()
+            };
+
+            write_batch(
+                &table,
+                &make_batch(
+                    vec![1, 1, 1, 1, 2],
+                    vec![10, 7, 20, 8, 100],
+                    vec!["+I", "-D", "+U", "-U", "-D"],
+                ),
+            )
+            .await;
+            assert_eq!(
+                scan_id_values(&table).await,
+                [(1, 30)].into_iter().collect(),
+                "{option}, explicit={explicit_kind}"
+            );
+
+            // A retract-only commit cannot cancel an earlier sum or add a new 
key.
+            write_batch(
+                &table,
+                &make_batch(vec![1, 2], vec![30, 100], vec!["-D", "-U"]),
+            )
+            .await;
+            assert_eq!(
+                scan_id_values(&table).await,
+                [(1, 30)].into_iter().collect(),
+                "{option}, explicit={explicit_kind}"
+            );
+
+            // The explicit global option takes precedence over a deprecated 
alias.
+            let table = table_with_options(
+                &table,
+                HashMap::from([("ignore-delete".to_string(), 
"false".to_string())]),
+            );
+            write_batch(&table, &make_batch(vec![1], vec![5], 
vec!["-D"])).await;
+            assert_eq!(
+                scan_id_values(&table).await,
+                [(1, 25)].into_iter().collect(),
+                "{option}, explicit={explicit_kind}"
+            );
+        }
+    }
+}
+
 #[tokio::test]
 async fn rowkind_field_ignore_update_before() {
     let table_path = "memory:/rowkind_field/ignore_update_before";

Reply via email to