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 b8073f20 Support Mosaic writes for native Paimon tables (#952)
b8073f20 is described below

commit b8073f20a22a1aff812998dd0d35e67bd433e117
Author: Jingsong Lee <[email protected]>
AuthorDate: Fri Sep 25 12:09:15 2026 +0800

    Support Mosaic writes for native Paimon tables (#952)
---
 crates/paimon/src/arrow/format/mod.rs              |   14 +
 crates/paimon/src/arrow/format/mosaic.rs           |    7 +-
 crates/paimon/src/arrow/format/mosaic_write.rs     | 1215 ++++++++++++++++++++
 crates/paimon/src/table/kv_file_writer.rs          |   48 +-
 crates/paimon/src/table/mod.rs                     |    2 +
 .../paimon/src/table/mosaic_table_write_tests.rs   |  962 ++++++++++++++++
 crates/paimon/src/table/table_write.rs             |   74 ++
 docs/src/getting-started.md                        |   18 +-
 docs/src/sql.md                                    |    7 +-
 9 files changed, 2337 insertions(+), 10 deletions(-)

diff --git a/crates/paimon/src/arrow/format/mod.rs 
b/crates/paimon/src/arrow/format/mod.rs
index 37607d31..24c9ed04 100644
--- a/crates/paimon/src/arrow/format/mod.rs
+++ b/crates/paimon/src/arrow/format/mod.rs
@@ -19,6 +19,7 @@ mod avro;
 mod avro_write;
 pub(crate) mod blob;
 mod mosaic;
+mod mosaic_write;
 mod orc;
 pub(crate) mod parquet;
 mod row;
@@ -448,6 +449,7 @@ fn supported_write_formats() -> Vec<&'static str> {
         ".blob",
         ".avro",
         ".row",
+        ".mosaic",
         #[cfg(feature = "vortex")]
         ".vortex",
     ]
@@ -508,6 +510,18 @@ pub(crate) async fn create_format_writer(
         Ok(Box::new(
             row::RowFormatWriter::new(output, schema, row_type, 
zstd_level).await?,
         ))
+    } else if lower.ends_with(".mosaic") {
+        Ok(Box::new(
+            mosaic_write::MosaicFormatWriter::new(
+                output,
+                schema,
+                compression,
+                zstd_level,
+                write_fields,
+                format_options,
+            )
+            .await?,
+        ))
     } else {
         #[cfg(feature = "vortex")]
         if lower.ends_with(".vortex") {
diff --git a/crates/paimon/src/arrow/format/mosaic.rs 
b/crates/paimon/src/arrow/format/mosaic.rs
index 2c5a2970..f84ee873 100644
--- a/crates/paimon/src/arrow/format/mosaic.rs
+++ b/crates/paimon/src/arrow/format/mosaic.rs
@@ -655,7 +655,10 @@ fn build_file_column_indices(
         .collect()
 }
 
-fn mosaic_value_to_datum(value: &MosaicValue, data_type: &PaimonDataType) -> 
Option<Datum> {
+pub(super) fn mosaic_value_to_datum(
+    value: &MosaicValue,
+    data_type: &PaimonDataType,
+) -> Option<Datum> {
     match (value, data_type) {
         (MosaicValue::Boolean(value), PaimonDataType::Boolean(_)) => 
Some(Datum::Bool(*value)),
         (MosaicValue::TinyInt(value), PaimonDataType::TinyInt(_)) => 
Some(Datum::TinyInt(*value)),
@@ -801,7 +804,7 @@ fn block_on_file_read(
         .map_err(|_| io::Error::other("mosaic async read task was cancelled"))?
 }
 
-fn validate_mosaic_schema(schema: &SchemaRef) -> crate::Result<()> {
+pub(super) fn validate_mosaic_schema(schema: &SchemaRef) -> crate::Result<()> {
     for field in schema.fields() {
         validate_mosaic_arrow_type(field.data_type()).map_err(|message| 
Error::Unsupported {
             message: format!(
diff --git a/crates/paimon/src/arrow/format/mosaic_write.rs 
b/crates/paimon/src/arrow/format/mosaic_write.rs
new file mode 100644
index 00000000..23b6267e
--- /dev/null
+++ b/crates/paimon/src/arrow/format/mosaic_write.rs
@@ -0,0 +1,1215 @@
+// 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.
+
+//! Async table-writer adapter for the synchronous Mosaic encoder. The encoder
+//! buffers one row group; completed blocks are drained after each input batch.
+
+use std::collections::{HashMap, HashSet};
+use std::io;
+
+use arrow_array::RecordBatch;
+use arrow_schema::SchemaRef;
+use async_trait::async_trait;
+use bytes::Bytes;
+use paimon_mosaic_core::writer::{MosaicWriter, OutputFile, WriterOptions};
+
+use super::{FormatFileWriter, FormatWriteResult};
+use crate::io::{FileWrite, OutputFile as PaimonOutputFile};
+use crate::spec::stats::BinaryTableStats;
+use crate::spec::{BinaryRowBuilder, CoreOptions, DataField, DataType, Datum};
+use crate::{Error, Result};
+
+/// A synchronous Mosaic output that only owns bytes the asynchronous adapter
+/// has not yet sent to storage. `pos` includes drained bytes: Mosaic footer
+/// offsets are absolute, not offsets within the pending buffer.
+#[derive(Default)]
+struct PendingOutput {
+    pending: Vec<u8>,
+    position: u64,
+}
+
+impl PendingOutput {
+    fn take(&mut self) -> Bytes {
+        Bytes::from(std::mem::take(&mut self.pending))
+    }
+}
+
+impl OutputFile for PendingOutput {
+    fn write(&mut self, bytes: &[u8]) -> io::Result<()> {
+        self.position = self
+            .position
+            .checked_add(bytes.len() as u64)
+            .ok_or_else(|| {
+                io::Error::new(io::ErrorKind::InvalidInput, "Mosaic file size 
overflow")
+            })?;
+        self.pending.extend_from_slice(bytes);
+        Ok(())
+    }
+
+    fn flush(&mut self) -> io::Result<()> {
+        Ok(())
+    }
+
+    fn pos(&self) -> u64 {
+        self.position
+    }
+}
+
+pub(crate) struct MosaicFormatWriter {
+    output: Box<dyn FileWrite>,
+    encoder: MosaicWriter<PendingOutput>,
+    schema: SchemaRef,
+    physical_schema: SchemaRef,
+    stats_fields: Vec<DataField>,
+    pending_rows: usize,
+}
+
+impl MosaicFormatWriter {
+    pub(crate) async fn new(
+        output: &PaimonOutputFile,
+        schema: SchemaRef,
+        compression: &str,
+        zstd_level: i32,
+        write_fields: Option<&[DataField]>,
+        format_options: Option<&HashMap<String, String>>,
+    ) -> Result<Self> {
+        // Java MosaicWriterFactory accepts only zstd, even though the 
low-level
+        // Mosaic encoder also knows an uncompressed wire format.
+        if !compression.eq_ignore_ascii_case("zstd") {
+            return Err(Error::Unsupported {
+                message: format!(
+                    "Mosaic format only supports zstd compression, but got: 
{compression}"
+                ),
+            });
+        }
+
+        let physical_schema = super::timestamp_millis_schema(&schema);
+        super::mosaic::validate_mosaic_schema(&physical_schema)?;
+        let fields = match write_fields {
+            Some(fields) => fields.to_vec(),
+            None => super::row::row_type_from_arrow_schema(&schema)?,
+        };
+        for field in &fields {
+            validate_paimon_type(field.data_type())?;
+        }
+        let mut options = WriterOptions {
+            zstd_level,
+            ..WriterOptions::default()
+        };
+        if let Some(format_options) = format_options {
+            let core_options = CoreOptions::new(format_options);
+            if let Some(block_size) = core_options.file_block_size()? {
+                if block_size <= 0 {
+                    return Err(Error::ConfigInvalid {
+                        message: "file.block-size for Mosaic must be 
positive".into(),
+                    });
+                }
+                options.row_group_max_size = block_size as u64;
+            }
+            if let Some(raw) = format_options.get("mosaic.num-buckets") {
+                // Java's option has an int type, so values above i32::MAX must
+                // fail here too instead of silently using a wider Rust usize.
+                let buckets = raw.parse::<i32>().map_err(|source| 
Error::ConfigInvalid {
+                    message: format!("Invalid mosaic.num-buckets '{raw}': 
{source}"),
+                })?;
+                if buckets <= 0 {
+                    return Err(Error::ConfigInvalid {
+                        message: "mosaic.num-buckets must be positive".into(),
+                    });
+                }
+                options.num_buckets = buckets as usize;
+            }
+            if let Some(raw) = format_options.get("mosaic.stats-columns") {
+                let mut seen = HashSet::new();
+                options.stats_columns = raw
+                    .split(',')
+                    .map(str::trim)
+                    .filter(|name| !name.is_empty())
+                    // Java's statistics extractor resolves selected fields as
+                    // a set, even if the option lists a name more than once.
+                    .filter(|name| seen.insert(*name))
+                    .map(str::to_owned)
+                    .collect();
+            }
+        }
+        if options.row_group_max_size == 0 {
+            return Err(Error::ConfigInvalid {
+                message: "file.block-size for Mosaic must be positive".into(),
+            });
+        }
+        let stats_fields = options
+            .stats_columns
+            .iter()
+            .map(|name| {
+                fields
+                    .iter()
+                    .find(|field| field.name() == name)
+                    .cloned()
+                    .ok_or_else(|| Error::ConfigInvalid {
+                        message: format!(
+                            "Mosaic statistics column '{name}' is not in the 
file schema"
+                        ),
+                    })
+            })
+            .collect::<Result<Vec<_>>>()?;
+        let encoder = MosaicWriter::new(PendingOutput::default(), 
&physical_schema, options)
+            .map_err(mosaic_write_error)?;
+        Ok(Self {
+            output: output.writer().await?,
+            encoder,
+            schema,
+            physical_schema,
+            stats_fields,
+            pending_rows: 0,
+        })
+    }
+
+    async fn drain(&mut self) -> Result<()> {
+        let bytes = self.encoder.output_mut().take();
+        if !bytes.is_empty() {
+            self.output.write(bytes).await?;
+        }
+        Ok(())
+    }
+
+    fn collect_stats(&self) -> Result<Option<BinaryTableStats>> {
+        if self.stats_fields.is_empty() || self.encoder.num_row_groups() == 0 {
+            return Ok(None);
+        }
+        let positions = self
+            .stats_fields
+            .iter()
+            .map(|field| {
+                self.encoder
+                    .schema()
+                    .columns
+                    .iter()
+                    .position(|column| column.name == field.name())
+                    .expect("statistics field was checked against Mosaic 
schema")
+            })
+            .collect::<Vec<_>>();
+        let mut minima: Vec<Option<Datum>> = vec![None; positions.len()];
+        let mut maxima: Vec<Option<Datum>> = vec![None; positions.len()];
+        let mut null_counts = vec![Some(0_i64); positions.len()];
+        for group in 0..self.encoder.num_row_groups() {
+            for stat in self.encoder.row_group_stats(group) {
+                let Some(output_index) = positions
+                    .iter()
+                    .position(|position| *position == stat.column_index)
+                else {
+                    continue;
+                };
+                let field = &self.stats_fields[output_index];
+                let count =
+                    i64::try_from(stat.null_count).map_err(|source| 
Error::DataInvalid {
+                        message: format!("Mosaic null count for '{}' exceeds 
i64", field.name()),
+                        source: Some(Box::new(source)),
+                    })?;
+                let total = null_counts[output_index].as_mut().unwrap();
+                *total = total.checked_add(count).ok_or_else(|| 
Error::DataInvalid {
+                    message: format!("Mosaic null count for '{}' exceeds i64", 
field.name()),
+                    source: None,
+                })?;
+                if let Some(min) = stat.min.as_ref().and_then(|value| {
+                    super::mosaic::mosaic_value_to_datum(value, 
field.data_type())
+                }) {
+                    if minima[output_index]
+                        .as_ref()
+                        .is_none_or(|current| min < *current)
+                    {
+                        minima[output_index] = Some(min);
+                    }
+                }
+                if let Some(max) = stat.max.as_ref().and_then(|value| {
+                    super::mosaic::mosaic_value_to_datum(value, 
field.data_type())
+                }) {
+                    if maxima[output_index]
+                        .as_ref()
+                        .is_none_or(|current| max > *current)
+                    {
+                        maxima[output_index] = Some(max);
+                    }
+                }
+            }
+        }
+        let mut min_row = BinaryRowBuilder::new(positions.len() as i32);
+        let mut max_row = BinaryRowBuilder::new(positions.len() as i32);
+        for (index, field) in self.stats_fields.iter().enumerate() {
+            match &minima[index] {
+                Some(value) => min_row.write_datum(index, value, 
field.data_type()),
+                None => min_row.set_null_at(index),
+            }
+            match &maxima[index] {
+                Some(value) => max_row.write_datum(index, value, 
field.data_type()),
+                None => max_row.set_null_at(index),
+            }
+        }
+        Ok(Some(BinaryTableStats::new(
+            min_row.build_serialized(),
+            max_row.build_serialized(),
+            null_counts,
+        )))
+    }
+}
+
+/// Match MosaicFileFormat.MosaicRowTypeVisitor before the Arrow conversion
+/// erases the distinction between a Paimon MAP and MULTISET.
+fn validate_paimon_type(data_type: &DataType) -> Result<()> {
+    match data_type {
+        DataType::Array(array) => validate_paimon_type(array.element_type()),
+        DataType::Map(map) => {
+            validate_paimon_type(map.key_type())?;
+            validate_paimon_type(map.value_type())
+        }
+        DataType::Variant(_)
+        | DataType::Blob(_)
+        | DataType::Vector(_)
+        | DataType::Multiset(_)
+        | DataType::Row(_) => Err(Error::Unsupported {
+            message: format!("Mosaic file format does not support type 
{data_type}"),
+        }),
+        _ => Ok(()),
+    }
+}
+
+#[async_trait]
+impl FormatFileWriter for MosaicFormatWriter {
+    async fn write(&mut self, batch: &RecordBatch) -> Result<()> {
+        if batch.schema() != self.schema {
+            return Err(Error::DataInvalid {
+                message: "Mosaic writer input schema differs from its file 
schema".into(),
+                source: None,
+            });
+        }
+        if batch.num_rows() == 0 {
+            return Ok(());
+        }
+        let batch = if self.schema == self.physical_schema {
+            batch.clone()
+        } else {
+            let columns = batch
+                .columns()
+                .iter()
+                .zip(self.physical_schema.fields())
+                .map(|(array, field)| {
+                    if array.data_type() == field.data_type() {
+                        Ok(array.clone())
+                    } else {
+                        arrow_cast::cast(array, 
field.data_type()).map_err(|source| {
+                            Error::DataInvalid {
+                                message: format!(
+                                    "Cannot convert Mosaic column '{}' to its 
storage type",
+                                    field.name()
+                                ),
+                                source: Some(Box::new(source)),
+                            }
+                        })
+                    }
+                })
+                .collect::<Result<Vec<_>>>()?;
+            RecordBatch::try_new(self.physical_schema.clone(), 
columns).map_err(|source| {
+                Error::DataInvalid {
+                    message: "Cannot build Mosaic storage batch".into(),
+                    source: Some(Box::new(source)),
+                }
+            })?
+        };
+        let groups_before = self.encoder.num_row_groups();
+        self.encoder
+            .write_batch(&batch)
+            .map_err(mosaic_write_error)?;
+        self.pending_rows = if self.encoder.num_row_groups() > groups_before {
+            0
+        } else {
+            self.pending_rows.saturating_add(batch.num_rows())
+        };
+        self.drain().await
+    }
+
+    fn num_bytes(&self) -> usize {
+        self.encoder.estimated_file_size() as usize
+    }
+
+    fn in_progress_size(&self) -> usize {
+        if self.pending_rows == 0 {
+            0
+        } else {
+            self.encoder
+                .estimated_file_size()
+                .saturating_sub(self.encoder.output().pos()) as usize
+        }
+    }
+
+    fn pending_rows(&self) -> Option<usize> {
+        Some(self.pending_rows)
+    }
+
+    async fn flush(&mut self) -> Result<()> {
+        // The Mosaic encoder chooses a row-group boundary by file.block-size.
+        // It exposes finalization at close, so there is no intermediate forced
+        // row-group flush. Drain every completed block before returning.
+        self.drain().await
+    }
+
+    async fn close(mut self: Box<Self>) -> Result<FormatWriteResult> {
+        self.encoder.close().map_err(mosaic_write_error)?;
+        let stats = self.collect_stats()?;
+        self.drain().await?;
+        let file_size = self.encoder.output().pos();
+        self.output.close().await?;
+        Ok(match stats {
+            Some(stats) => FormatWriteResult::with_value_stats(
+                file_size,
+                stats,
+                Some(
+                    self.stats_fields
+                        .iter()
+                        .map(|field| field.name().to_owned())
+                        .collect(),
+                ),
+            ),
+            None => FormatWriteResult::new(file_size),
+        })
+    }
+}
+
+fn mosaic_write_error(source: io::Error) -> Error {
+    Error::DataInvalid {
+        message: format!("Failed to write Mosaic file: {source}"),
+        source: Some(Box::new(source)),
+    }
+}
+
+#[cfg(test)]
+mod tests {
+    use super::*;
+    use crate::arrow::build_target_arrow_schema;
+    use crate::arrow::format::FormatFileReader;
+    use crate::io::FileIOBuilder;
+    use crate::spec::{BinaryRow, DataType, IntType, TimestampType, 
VarCharType};
+    use arrow_array::{
+        Array, Int32Array, StringArray, TimestampMillisecondArray, 
TimestampSecondArray,
+    };
+    use futures::TryStreamExt;
+    use std::sync::Arc;
+
+    fn fields() -> Vec<DataField> {
+        vec![
+            DataField::new(0, "id".into(), DataType::Int(IntType::new())),
+            DataField::new(
+                1,
+                "name".into(),
+                DataType::VarChar(VarCharType::string_type()),
+            ),
+        ]
+    }
+
+    #[tokio::test]
+    async fn writes_multiple_batches_and_collects_java_style_stats() {
+        let io = FileIOBuilder::new("memory").build().unwrap();
+        let path = "memory:/mosaic-writer/values.mosaic";
+        let output = io.new_output(path).unwrap();
+        let fields = fields();
+        let schema = build_target_arrow_schema(&fields).unwrap();
+        let options = HashMap::from([
+            ("mosaic.stats-columns".into(), "id, name, id".into()),
+            ("mosaic.num-buckets".into(), "2".into()),
+            ("file.block-size".into(), "32".into()),
+        ]);
+        let mut writer = MosaicFormatWriter::new(
+            &output,
+            schema.clone(),
+            "zstd",
+            1,
+            Some(&fields),
+            Some(&options),
+        )
+        .await
+        .unwrap();
+        let first = RecordBatch::try_new(
+            schema.clone(),
+            vec![
+                Arc::new(Int32Array::from(vec![Some(3), None])),
+                Arc::new(StringArray::from(vec![Some("z"), Some("b")])),
+            ],
+        )
+        .unwrap();
+        let second = RecordBatch::try_new(
+            schema.clone(),
+            vec![
+                Arc::new(Int32Array::from(vec![Some(1), Some(2)])),
+                Arc::new(StringArray::from(vec![Some("a"), None])),
+            ],
+        )
+        .unwrap();
+        writer.write(&first).await.unwrap();
+        writer.write(&second).await.unwrap();
+        let result = Box::new(writer).close().await.unwrap();
+        let value_stats = result.value_stats.unwrap();
+        assert_eq!(value_stats.columns, Some(vec!["id".into(), 
"name".into()]));
+        assert_eq!(value_stats.stats.null_counts(), &vec![Some(1), Some(1)]);
+        let min = 
BinaryRow::from_serialized_bytes(value_stats.stats.min_values()).unwrap();
+        let max = 
BinaryRow::from_serialized_bytes(value_stats.stats.max_values()).unwrap();
+        assert_eq!(min.get_int(0).unwrap(), 1);
+        assert_eq!(max.get_int(0).unwrap(), 3);
+        assert_eq!(min.get_string(1).unwrap(), "a");
+        assert_eq!(max.get_string(1).unwrap(), "z");
+
+        let input = io.new_input(path).unwrap().reader().await.unwrap();
+        let reader = super::super::mosaic::MosaicFormatReader::default();
+        let batches: Vec<_> = reader
+            .read_batch_stream(
+                Box::new(input),
+                result.file_size,
+                &fields,
+                None,
+                Some(2),
+                None,
+            )
+            .await
+            .unwrap()
+            .try_collect()
+            .await
+            .unwrap();
+        let ids = batches
+            .iter()
+            .flat_map(|batch| {
+                let values = batch
+                    .column(0)
+                    .as_any()
+                    .downcast_ref::<Int32Array>()
+                    .unwrap();
+                (0..batch.num_rows())
+                    .map(|row| (!values.is_null(row)).then(|| 
values.value(row)))
+                    .collect::<Vec<_>>()
+            })
+            .collect::<Vec<_>>();
+        assert_eq!(ids, vec![Some(3), None, Some(1), Some(2)]);
+    }
+
+    #[tokio::test]
+    async fn 
stats_preserve_option_order_across_row_groups_and_temporal_types() {
+        use crate::spec::{DateType, DecimalType};
+        use arrow_array::{Date32Array, Decimal128Array, 
TimestampMicrosecondArray};
+
+        let fields = vec![
+            DataField::new(
+                0,
+                "ts".into(),
+                DataType::Timestamp(TimestampType::new(6).unwrap()),
+            ),
+            DataField::new(
+                1,
+                "amount".into(),
+                DataType::Decimal(DecimalType::new(10, 2).unwrap()),
+            ),
+            DataField::new(2, "day".into(), DataType::Date(DateType::new())),
+        ];
+        let schema = build_target_arrow_schema(&fields).unwrap();
+        let options = HashMap::from([
+            ("mosaic.stats-columns".into(), "day, amount, ts".into()),
+            ("file.block-size".into(), "1".into()),
+        ]);
+        let io = FileIOBuilder::new("memory").build().unwrap();
+        let path = "memory:/mosaic-writer/stats-temporal.mosaic";
+        let mut writer = MosaicFormatWriter::new(
+            &io.new_output(path).unwrap(),
+            schema.clone(),
+            "zstd",
+            1,
+            Some(&fields),
+            Some(&options),
+        )
+        .await
+        .unwrap();
+        for (timestamps, amounts, days) in [
+            (
+                vec![Some(3_000_001), None],
+                vec![Some(1234), None],
+                vec![Some(10), None],
+            ),
+            (
+                vec![Some(1_001_234), Some(2_000_000)],
+                vec![Some(-234), Some(500)],
+                vec![Some(0), Some(5)],
+            ),
+        ] {
+            let batch = RecordBatch::try_new(
+                schema.clone(),
+                vec![
+                    Arc::new(TimestampMicrosecondArray::from(timestamps)),
+                    Arc::new(
+                        Decimal128Array::from(amounts)
+                            .with_precision_and_scale(10, 2)
+                            .unwrap(),
+                    ),
+                    Arc::new(Date32Array::from(days)),
+                ],
+            )
+            .unwrap();
+            writer.write(&batch).await.unwrap();
+        }
+        let result = Box::new(writer).close().await.unwrap();
+        let stats = result.value_stats.unwrap();
+        assert_eq!(
+            stats.columns,
+            Some(vec!["day".into(), "amount".into(), "ts".into()])
+        );
+        assert_eq!(stats.stats.null_counts(), &[Some(1), Some(1), Some(1)]);
+        let min = 
BinaryRow::from_serialized_bytes(stats.stats.min_values()).unwrap();
+        let max = 
BinaryRow::from_serialized_bytes(stats.stats.max_values()).unwrap();
+        assert_eq!(min.get_int(0).unwrap(), 0);
+        assert_eq!(max.get_int(0).unwrap(), 10);
+        assert_eq!(min.get_decimal_unscaled(1, 10).unwrap(), -234);
+        assert_eq!(max.get_decimal_unscaled(1, 10).unwrap(), 1234);
+        assert_eq!(min.get_timestamp_raw(2, 6).unwrap(), (1001, 234_000));
+        assert_eq!(max.get_timestamp_raw(2, 6).unwrap(), (3000, 1000));
+    }
+
+    #[tokio::test]
+    async fn timestamp_zero_uses_mosaic_millisecond_storage() {
+        let io = FileIOBuilder::new("memory").build().unwrap();
+        let path = "memory:/mosaic-writer/timestamp.mosaic";
+        let field = DataField::new(
+            0,
+            "ts".into(),
+            DataType::Timestamp(TimestampType::new(0).unwrap()),
+        );
+        let fields = vec![field];
+        let schema = build_target_arrow_schema(&fields).unwrap();
+        let batch = RecordBatch::try_new(
+            schema.clone(),
+            vec![Arc::new(TimestampSecondArray::from(vec![
+                Some(-1),
+                Some(1_700_000_001),
+            ]))],
+        )
+        .unwrap();
+        let mut writer = MosaicFormatWriter::new(
+            &io.new_output(path).unwrap(),
+            schema,
+            "zstd",
+            1,
+            Some(&fields),
+            None,
+        )
+        .await
+        .unwrap();
+        writer.write(&batch).await.unwrap();
+        let result = Box::new(writer).close().await.unwrap();
+        let input = io.new_input(path).unwrap().reader().await.unwrap();
+        let reader = super::super::mosaic::MosaicFormatReader::default();
+        let batches: Vec<_> = reader
+            .read_batch_stream(Box::new(input), result.file_size, &fields, 
None, None, None)
+            .await
+            .unwrap()
+            .try_collect()
+            .await
+            .unwrap();
+        assert_eq!(batches.len(), 1);
+        let values = batches[0]
+            .column(0)
+            .as_any()
+            .downcast_ref::<TimestampMillisecondArray>()
+            .unwrap();
+        assert_eq!(values.values(), &[-1_000, 1_700_000_001_000]);
+    }
+
+    #[tokio::test]
+    async fn nested_timestamp_zero_uses_millisecond_storage() {
+        use crate::spec::ArrayType;
+        use arrow_array::ListArray;
+        use arrow_buffer::{OffsetBuffer, ScalarBuffer};
+        use arrow_schema::DataType as ArrowType;
+
+        let fields = vec![DataField::new(
+            0,
+            "times".into(),
+            DataType::Array(ArrayType::new(DataType::Timestamp(
+                TimestampType::new(0).unwrap(),
+            ))),
+        )];
+        let schema = build_target_arrow_schema(&fields).unwrap();
+        let item = match schema.field(0).data_type() {
+            ArrowType::List(item) => item.clone(),
+            other => panic!("expected list, got {other:?}"),
+        };
+        let array = ListArray::try_new(
+            item,
+            OffsetBuffer::new(ScalarBuffer::from(vec![0, 2, 2, 3])),
+            Arc::new(TimestampSecondArray::from(vec![Some(-1), None, 
Some(42)])),
+            None,
+        )
+        .unwrap();
+        let batch = RecordBatch::try_new(schema.clone(), 
vec![Arc::new(array)]).unwrap();
+        let io = FileIOBuilder::new("memory").build().unwrap();
+        let path = "memory:/mosaic-writer/nested-timestamp.mosaic";
+        let mut writer = MosaicFormatWriter::new(
+            &io.new_output(path).unwrap(),
+            schema,
+            "zstd",
+            1,
+            Some(&fields),
+            None,
+        )
+        .await
+        .unwrap();
+        writer.write(&batch).await.unwrap();
+        let result = Box::new(writer).close().await.unwrap();
+        let input = io.new_input(path).unwrap().reader().await.unwrap();
+        let decoded: Vec<RecordBatch> = 
super::super::mosaic::MosaicFormatReader::default()
+            .read_batch_stream(Box::new(input), result.file_size, &fields, 
None, None, None)
+            .await
+            .unwrap()
+            .try_collect()
+            .await
+            .unwrap();
+        assert_eq!(decoded.len(), 1);
+        let array = decoded[0]
+            .column(0)
+            .as_any()
+            .downcast_ref::<ListArray>()
+            .unwrap();
+        let times = array.values();
+        let times = times
+            .as_any()
+            .downcast_ref::<TimestampMillisecondArray>()
+            .unwrap();
+        assert_eq!(times.values(), &[-1_000, 0, 42_000]);
+        assert!(times.is_null(1));
+        assert!(array.is_valid(0));
+        assert!(array.is_valid(1));
+    }
+
+    #[tokio::test]
+    async fn rejects_unsupported_compression_and_invalid_configuration() {
+        let io = FileIOBuilder::new("memory").build().unwrap();
+        let fields = fields();
+        let schema = build_target_arrow_schema(&fields).unwrap();
+        let output = io
+            .new_output("memory:/mosaic-writer/invalid.mosaic")
+            .unwrap();
+        let error = match MosaicFormatWriter::new(
+            &output,
+            schema.clone(),
+            "snappy",
+            1,
+            Some(&fields),
+            None,
+        )
+        .await
+        {
+            Ok(_) => panic!("snappy must be rejected"),
+            Err(error) => error,
+        };
+        assert!(error.to_string().contains("only supports zstd"));
+
+        for (option, value) in [
+            ("mosaic.num-buckets", "0"),
+            ("mosaic.num-buckets", "-1"),
+            ("mosaic.num-buckets", "2147483648"),
+            ("mosaic.num-buckets", "bad"),
+            ("file.block-size", "0"),
+            ("file.block-size", "-1"),
+            ("mosaic.stats-columns", "missing"),
+        ] {
+            let options = HashMap::from([(option.to_string(), 
value.to_string())]);
+            let result = MosaicFormatWriter::new(
+                &output,
+                schema.clone(),
+                "zstd",
+                1,
+                Some(&fields),
+                Some(&options),
+            )
+            .await;
+            assert!(result.is_err(), "{option}={value} should fail");
+        }
+    }
+
+    #[tokio::test]
+    async fn empty_and_multi_group_files_remain_readable() {
+        let io = FileIOBuilder::new("memory").build().unwrap();
+        let fields = fields();
+        let schema = build_target_arrow_schema(&fields).unwrap();
+        let options = HashMap::from([("file.block-size".into(), "1".into())]);
+        let empty_path = "memory:/mosaic-writer/empty.mosaic";
+        let empty_writer = MosaicFormatWriter::new(
+            &io.new_output(empty_path).unwrap(),
+            schema.clone(),
+            "zstd",
+            1,
+            Some(&fields),
+            Some(&options),
+        )
+        .await
+        .unwrap();
+        let empty = Box::new(empty_writer).close().await.unwrap();
+        let input = io.new_input(empty_path).unwrap().reader().await.unwrap();
+        let empty_batches: Vec<RecordBatch> = 
super::super::mosaic::MosaicFormatReader::default()
+            .read_batch_stream(Box::new(input), empty.file_size, &fields, 
None, None, None)
+            .await
+            .unwrap()
+            .try_collect()
+            .await
+            .unwrap();
+        assert!(empty_batches.is_empty());
+
+        let path = "memory:/mosaic-writer/many-groups.mosaic";
+        let mut writer = MosaicFormatWriter::new(
+            &io.new_output(path).unwrap(),
+            schema.clone(),
+            "zstd",
+            1,
+            Some(&fields),
+            Some(&options),
+        )
+        .await
+        .unwrap();
+        for start in [0, 10, 20] {
+            let batch = RecordBatch::try_new(
+                schema.clone(),
+                vec![
+                    Arc::new(Int32Array::from((start..start + 
10).collect::<Vec<_>>())),
+                    Arc::new(StringArray::from(
+                        (0..10).map(|_| Some("payload")).collect::<Vec<_>>(),
+                    )),
+                ],
+            )
+            .unwrap();
+            writer.write(&batch).await.unwrap();
+            assert_eq!(writer.pending_rows(), Some(0));
+        }
+        let result = Box::new(writer).close().await.unwrap();
+        let input = io.new_input(path).unwrap().reader().await.unwrap();
+        let batches: Vec<RecordBatch> = 
super::super::mosaic::MosaicFormatReader::default()
+            .read_batch_stream(Box::new(input), result.file_size, &fields, 
None, None, None)
+            .await
+            .unwrap()
+            .try_collect()
+            .await
+            .unwrap();
+        assert_eq!(batches.iter().map(RecordBatch::num_rows).sum::<usize>(), 
30);
+    }
+
+    #[tokio::test]
+    async fn writer_reports_buffered_rows_until_a_row_group_is_emitted() {
+        let io = FileIOBuilder::new("memory").build().unwrap();
+        let fields = fields();
+        let schema = build_target_arrow_schema(&fields).unwrap();
+        let path = "memory:/mosaic-writer/buffer-lifetime.mosaic";
+        let options = HashMap::from([("file.block-size".into(), 
"1gb".into())]);
+        let mut writer = MosaicFormatWriter::new(
+            &io.new_output(path).unwrap(),
+            schema.clone(),
+            "zstd",
+            1,
+            Some(&fields),
+            Some(&options),
+        )
+        .await
+        .unwrap();
+        assert_eq!(writer.pending_rows(), Some(0));
+        assert_eq!(writer.in_progress_size(), 0);
+        for value in [1, 2, 3] {
+            let batch = RecordBatch::try_new(
+                schema.clone(),
+                vec![
+                    Arc::new(Int32Array::from(vec![value])),
+                    Arc::new(StringArray::from(vec!["payload"])),
+                ],
+            )
+            .unwrap();
+            writer.write(&batch).await.unwrap();
+            assert_eq!(writer.pending_rows(), Some(value as usize));
+            assert!(writer.in_progress_size() > 0);
+            writer.flush().await.unwrap();
+            assert_eq!(writer.pending_rows(), Some(value as usize));
+        }
+        let result = Box::new(writer).close().await.unwrap();
+        let input = io.new_input(path).unwrap().reader().await.unwrap();
+        let batches: Vec<RecordBatch> = 
super::super::mosaic::MosaicFormatReader::default()
+            .read_batch_stream(Box::new(input), result.file_size, &fields, 
None, None, None)
+            .await
+            .unwrap()
+            .try_collect()
+            .await
+            .unwrap();
+        assert_eq!(batches.iter().map(RecordBatch::num_rows).sum::<usize>(), 
3);
+    }
+
+    #[tokio::test]
+    async fn written_row_groups_support_projection_and_row_selection() {
+        use crate::table::RowRange;
+
+        let io = FileIOBuilder::new("memory").build().unwrap();
+        let fields = fields();
+        let schema = build_target_arrow_schema(&fields).unwrap();
+        let options = HashMap::from([("file.block-size".into(), "1".into())]);
+        let path = "memory:/mosaic-writer/selected-groups.mosaic";
+        let mut writer = MosaicFormatWriter::new(
+            &io.new_output(path).unwrap(),
+            schema.clone(),
+            "zstd",
+            1,
+            Some(&fields),
+            Some(&options),
+        )
+        .await
+        .unwrap();
+        for start in [0, 5, 10] {
+            let batch = RecordBatch::try_new(
+                schema.clone(),
+                vec![
+                    Arc::new(Int32Array::from((start..start + 
5).collect::<Vec<_>>())),
+                    Arc::new(StringArray::from(
+                        (start..start + 5)
+                            .map(|value| format!("name-{value}"))
+                            .collect::<Vec<_>>(),
+                    )),
+                ],
+            )
+            .unwrap();
+            writer.write(&batch).await.unwrap();
+        }
+        let result = Box::new(writer).close().await.unwrap();
+        let input = io.new_input(path).unwrap().reader().await.unwrap();
+        let selected: Vec<RecordBatch> = 
super::super::mosaic::MosaicFormatReader::default()
+            .read_batch_stream(
+                Box::new(input),
+                result.file_size,
+                &fields[1..],
+                None,
+                Some(2),
+                Some(vec![RowRange::new(3, 6), RowRange::new(12, 13)]),
+            )
+            .await
+            .unwrap()
+            .try_collect()
+            .await
+            .unwrap();
+        let names = selected
+            .iter()
+            .flat_map(|batch| {
+                assert_eq!(batch.num_columns(), 1);
+                let values = batch
+                    .column(0)
+                    .as_any()
+                    .downcast_ref::<StringArray>()
+                    .unwrap();
+                values
+                    .iter()
+                    .map(|value| value.unwrap().to_owned())
+                    .collect::<Vec<_>>()
+            })
+            .collect::<Vec<_>>();
+        assert_eq!(
+            names,
+            vec!["name-3", "name-4", "name-5", "name-6", "name-12", "name-13"]
+        );
+    }
+
+    #[tokio::test]
+    async fn rejects_different_input_schema_before_emitting_rows() {
+        let io = FileIOBuilder::new("memory").build().unwrap();
+        let fields = fields();
+        let schema = build_target_arrow_schema(&fields).unwrap();
+        let path = "memory:/mosaic-writer/mismatch.mosaic";
+        let mut writer = MosaicFormatWriter::new(
+            &io.new_output(path).unwrap(),
+            schema,
+            "zstd",
+            1,
+            Some(&fields),
+            None,
+        )
+        .await
+        .unwrap();
+        let wrong_batch = RecordBatch::try_new(
+            Arc::new(arrow_schema::Schema::new(vec![
+                arrow_schema::Field::new("name", arrow_schema::DataType::Utf8, 
true),
+                arrow_schema::Field::new("id", arrow_schema::DataType::Int32, 
true),
+            ])),
+            vec![
+                Arc::new(StringArray::from(vec!["wrong"])),
+                Arc::new(Int32Array::from(vec![1])),
+            ],
+        )
+        .unwrap();
+        let error = writer.write(&wrong_batch).await.unwrap_err();
+        assert!(error.to_string().contains("schema differs"));
+        let result = Box::new(writer).close().await.unwrap();
+        let input = io.new_input(path).unwrap().reader().await.unwrap();
+        let batches: Vec<RecordBatch> = 
super::super::mosaic::MosaicFormatReader::default()
+            .read_batch_stream(Box::new(input), result.file_size, &fields, 
None, None, None)
+            .await
+            .unwrap()
+            .try_collect()
+            .await
+            .unwrap();
+        assert!(batches.is_empty());
+    }
+
+    #[tokio::test]
+    async fn round_trips_java_supported_scalars_and_collection_types() {
+        use crate::spec::{
+            ArrayType, BigIntType, BooleanType, DateType, DecimalType, 
DoubleType, FloatType,
+            LocalZonedTimestampType, MapType, SmallIntType, TimeType, 
TinyIntType,
+        };
+        use arrow_array::{
+            BooleanArray, Date32Array, Decimal128Array, Float32Array, 
Float64Array, Int16Array,
+            Int64Array, Int8Array, ListArray, MapArray, StructArray, 
Time32MillisecondArray,
+            TimestampMicrosecondArray, TimestampMillisecondArray,
+        };
+        use arrow_buffer::{BooleanBuffer, NullBuffer, OffsetBuffer, 
ScalarBuffer};
+        use arrow_schema::{DataType as ArrowType, Field};
+
+        let fields = vec![
+            DataField::new(0, "flag".into(), 
DataType::Boolean(BooleanType::new())),
+            DataField::new(1, "tiny".into(), 
DataType::TinyInt(TinyIntType::new())),
+            DataField::new(2, "small".into(), 
DataType::SmallInt(SmallIntType::new())),
+            DataField::new(3, "big".into(), 
DataType::BigInt(BigIntType::new())),
+            DataField::new(4, "float".into(), 
DataType::Float(FloatType::new())),
+            DataField::new(5, "double".into(), 
DataType::Double(DoubleType::new())),
+            DataField::new(6, "date".into(), DataType::Date(DateType::new())),
+            DataField::new(7, "time".into(), 
DataType::Time(TimeType::new(3).unwrap())),
+            DataField::new(
+                8,
+                "decimal5".into(),
+                DataType::Decimal(DecimalType::new(5, 2).unwrap()),
+            ),
+            DataField::new(
+                9,
+                "decimal20".into(),
+                DataType::Decimal(DecimalType::new(20, 0).unwrap()),
+            ),
+            DataField::new(
+                10,
+                "ts3".into(),
+                DataType::Timestamp(TimestampType::new(3).unwrap()),
+            ),
+            DataField::new(
+                11,
+                "ts6".into(),
+                DataType::Timestamp(TimestampType::new(6).unwrap()),
+            ),
+            DataField::new(
+                12,
+                "ltz6".into(),
+                
DataType::LocalZonedTimestamp(LocalZonedTimestampType::new(6).unwrap()),
+            ),
+            DataField::new(
+                13,
+                "array".into(),
+                DataType::Array(ArrayType::new(DataType::Int(IntType::new()))),
+            ),
+            DataField::new(
+                14,
+                "map".into(),
+                DataType::Map(MapType::new(
+                    DataType::VarChar(VarCharType::string_type()),
+                    DataType::Int(IntType::new()),
+                )),
+            ),
+        ];
+        let schema = build_target_arrow_schema(&fields).unwrap();
+        let list_field = match schema.field(13).data_type() {
+            ArrowType::List(field) => field.clone(),
+            other => panic!("expected list, got {other:?}"),
+        };
+        let list = ListArray::try_new(
+            list_field,
+            OffsetBuffer::new(ScalarBuffer::from(vec![0, 2, 2, 3])),
+            Arc::new(Int32Array::from(vec![Some(10), None, Some(-1)])),
+            Some(NullBuffer::new(BooleanBuffer::from(vec![
+                true, false, true,
+            ]))),
+        )
+        .unwrap();
+        let map_field = match schema.field(14).data_type() {
+            ArrowType::Map(field, _) => field.clone(),
+            other => panic!("expected map, got {other:?}"),
+        };
+        let entry_fields = match map_field.data_type() {
+            ArrowType::Struct(fields) => fields.clone(),
+            other => panic!("expected map entries struct, got {other:?}"),
+        };
+        let entries = StructArray::try_new(
+            entry_fields,
+            vec![
+                Arc::new(StringArray::from(vec!["a", "b"])),
+                Arc::new(Int32Array::from(vec![Some(1), None])),
+            ],
+            None,
+        )
+        .unwrap();
+        let map = MapArray::try_new(
+            Arc::new(Field::clone(&map_field)),
+            OffsetBuffer::new(ScalarBuffer::from(vec![0, 2, 2, 2])),
+            entries,
+            Some(NullBuffer::new(BooleanBuffer::from(vec![
+                true, false, true,
+            ]))),
+            false,
+        )
+        .unwrap();
+        let batch = RecordBatch::try_new(
+            schema.clone(),
+            vec![
+                Arc::new(BooleanArray::from(vec![Some(true), None, 
Some(false)])),
+                Arc::new(Int8Array::from(vec![Some(7), None, Some(-8)])),
+                Arc::new(Int16Array::from(vec![Some(256), None, Some(-512)])),
+                Arc::new(Int64Array::from(vec![Some(9_876_543_210), None, 
Some(-1)])),
+                Arc::new(Float32Array::from(vec![Some(1.5), None, 
Some(-2.5)])),
+                Arc::new(Float64Array::from(vec![Some(3.25), None, 
Some(-4.75)])),
+                Arc::new(Date32Array::from(vec![Some(18_000), None, 
Some(-1)])),
+                Arc::new(Time32MillisecondArray::from(vec![
+                    Some(3_600_000),
+                    None,
+                    Some(86_399_999),
+                ])),
+                Arc::new(
+                    Decimal128Array::from(vec![Some(12_345), None, Some(-678)])
+                        .with_precision_and_scale(5, 2)
+                        .unwrap(),
+                ),
+                Arc::new(
+                    
Decimal128Array::from(vec![Some(12_345_678_901_234_567_890), None, Some(-1)])
+                        .with_precision_and_scale(20, 0)
+                        .unwrap(),
+                ),
+                Arc::new(TimestampMillisecondArray::from(vec![
+                    Some(1_700_000_000_000),
+                    None,
+                    Some(-1),
+                ])),
+                Arc::new(TimestampMicrosecondArray::from(vec![
+                    Some(1_700_000_000_000_001),
+                    None,
+                    Some(-1),
+                ])),
+                Arc::new(
+                    TimestampMicrosecondArray::from(vec![
+                        Some(1_700_000_000_000_001),
+                        None,
+                        Some(-1),
+                    ])
+                    .with_timezone("UTC"),
+                ),
+                Arc::new(list),
+                Arc::new(map),
+            ],
+        )
+        .unwrap();
+        let io = FileIOBuilder::new("memory").build().unwrap();
+        let path = "memory:/mosaic-writer/all-types.mosaic";
+        let mut writer = MosaicFormatWriter::new(
+            &io.new_output(path).unwrap(),
+            schema,
+            "zstd",
+            1,
+            Some(&fields),
+            None,
+        )
+        .await
+        .unwrap();
+        writer.write(&batch).await.unwrap();
+        let result = Box::new(writer).close().await.unwrap();
+        let input = io.new_input(path).unwrap().reader().await.unwrap();
+        let decoded: Vec<RecordBatch> = 
super::super::mosaic::MosaicFormatReader::default()
+            .read_batch_stream(Box::new(input), result.file_size, &fields, 
None, None, None)
+            .await
+            .unwrap()
+            .try_collect()
+            .await
+            .unwrap();
+        assert_eq!(decoded.len(), 1);
+        assert_eq!(decoded[0].num_rows(), 3);
+        assert_eq!(decoded[0].num_columns(), fields.len());
+        for (index, field) in fields.iter().enumerate() {
+            assert_eq!(
+                decoded[0].column(index).to_data(),
+                batch.column(index).to_data(),
+                "round-trip differs for {}",
+                field.name()
+            );
+        }
+    }
+
+    #[test]
+    fn 
rejects_java_unsupported_paimon_types_even_when_arrow_shape_is_supported() {
+        use crate::spec::{ArrayType, BlobType, MapType, MultisetType, 
VariantType, VectorType};
+
+        let int = DataType::Int(IntType::new());
+        let unsupported = [
+            DataType::Blob(BlobType::new()),
+            DataType::Variant(VariantType::new()),
+            DataType::Vector(VectorType::new(3, int.clone()).unwrap()),
+            DataType::Multiset(MultisetType::new(int.clone())),
+        ];
+        for data_type in unsupported {
+            let error = validate_paimon_type(&data_type).unwrap_err();
+            assert!(
+                error.to_string().contains("does not support type"),
+                "{error}"
+            );
+        }
+
+        let nested = DataType::Map(MapType::new(
+            DataType::VarChar(VarCharType::string_type()),
+            
DataType::Array(ArrayType::new(DataType::Multiset(MultisetType::new(
+                int.clone(),
+            )))),
+        ));
+        let error = validate_paimon_type(&nested).unwrap_err();
+        assert!(error.to_string().contains("MULTISET"));
+        let supported = 
DataType::Array(ArrayType::new(DataType::Map(MapType::new(
+            DataType::VarChar(VarCharType::string_type()),
+            int,
+        ))));
+        validate_paimon_type(&supported).unwrap();
+    }
+
+    #[tokio::test]
+    async fn multiset_writer_is_rejected_before_file_creation() {
+        use crate::spec::MultisetType;
+
+        let io = FileIOBuilder::new("memory").build().unwrap();
+        let path = "memory:/mosaic-writer/multiset.mosaic";
+        let fields = vec![DataField::new(
+            0,
+            "bag".into(),
+            
DataType::Multiset(MultisetType::new(DataType::Int(IntType::new()))),
+        )];
+        let schema = build_target_arrow_schema(&fields).unwrap();
+        let error = match MosaicFormatWriter::new(
+            &io.new_output(path).unwrap(),
+            schema,
+            "zstd",
+            1,
+            Some(&fields),
+            None,
+        )
+        .await
+        {
+            Ok(_) => panic!("Mosaic must reject Paimon MULTISET"),
+            Err(error) => error,
+        };
+        assert!(error.to_string().contains("MULTISET"));
+        assert!(!io.new_input(path).unwrap().exists().await.unwrap());
+    }
+}
diff --git a/crates/paimon/src/table/kv_file_writer.rs 
b/crates/paimon/src/table/kv_file_writer.rs
index 9965277b..05ae9aad 100644
--- a/crates/paimon/src/table/kv_file_writer.rs
+++ b/crates/paimon/src/table/kv_file_writer.rs
@@ -35,8 +35,9 @@ use crate::resource::{MemoryReservation, ResourceContext};
 use crate::spec::stats::{compute_column_stats, BinaryTableStats};
 use crate::spec::{
     bucket_path_under, data_file_to_file_index_file_name, 
extract_datum_from_arrow,
-    AggregationConfig, BinaryRowBuilder, CoreOptions, DataField, DataFileMeta, 
DataType,
-    MergeEngine, PartialUpdateConfig, RowKind, SEQUENCE_NUMBER_FIELD_NAME, 
VALUE_KIND_FIELD_NAME,
+    AggregationConfig, BigIntType, BinaryRowBuilder, CoreOptions, DataField, 
DataFileMeta,
+    DataType, MergeEngine, PartialUpdateConfig, RowKind, TinyIntType, 
SEQUENCE_NUMBER_FIELD_ID,
+    SEQUENCE_NUMBER_FIELD_NAME, VALUE_KIND_FIELD_ID, VALUE_KIND_FIELD_NAME,
 };
 use crate::table::data_file_index_writer::FileIndexOptions;
 use crate::table::prepared_files::PreparedFiles;
@@ -94,7 +95,8 @@ pub(crate) struct KeyValueWriteConfig {
     pub primary_key_indices: Vec<usize>,
     /// Paimon DataTypes for each primary key column (same order as 
primary_key_indices).
     pub primary_key_types: Vec<DataType>,
-    /// Logical value fields, in file order, for Parquet footer statistics.
+    /// Logical value fields, used for footer statistics and to retain Paimon
+    /// types when building the physical file schema.
     pub value_fields: Vec<DataField>,
     /// Sequence field column indices in the user schema (empty if not 
configured).
     pub sequence_field_indices: Vec<usize>,
@@ -474,14 +476,16 @@ impl KeyValueFileWriter {
                 .await?,
             )
         } else {
+            let physical_fields =
+                build_physical_fields(&physical_schema, 
&self.config.value_fields)?;
             create_format_writer(
                 &output,
                 physical_schema.clone(),
                 write.file_compression,
                 self.config.file_compression_zstd_level,
                 None,
-                None,
-                None,
+                Some(&physical_fields),
+                Some(&self.config.table_options),
             )
             .await?
         };
@@ -1083,6 +1087,40 @@ pub(crate) fn build_physical_schema(user_schema: 
&ArrowSchema) -> Arc<ArrowSchem
     Arc::new(ArrowSchema::new(physical_fields))
 }
 
+/// Describe the actual file columns while retaining logical Paimon types that
+/// Arrow cannot distinguish (for example MULTISET and MAP).
+fn build_physical_fields(
+    physical_schema: &ArrowSchema,
+    value_fields: &[DataField],
+) -> Result<Vec<DataField>> {
+    physical_schema
+        .fields()
+        .iter()
+        .map(|arrow_field| match arrow_field.name().as_str() {
+            SEQUENCE_NUMBER_FIELD_NAME => Ok(DataField::new(
+                SEQUENCE_NUMBER_FIELD_ID,
+                SEQUENCE_NUMBER_FIELD_NAME.into(),
+                DataType::BigInt(BigIntType::new()),
+            )),
+            VALUE_KIND_FIELD_NAME => Ok(DataField::new(
+                VALUE_KIND_FIELD_ID,
+                VALUE_KIND_FIELD_NAME.into(),
+                DataType::TinyInt(TinyIntType::new()),
+            )),
+            name => value_fields
+                .iter()
+                .find(|field| field.name() == name)
+                .cloned()
+                .ok_or_else(|| crate::Error::DataInvalid {
+                    message: format!(
+                        "Physical file column '{name}' is missing from the 
table schema"
+                    ),
+                    source: None,
+                }),
+        })
+        .collect()
+}
+
 #[cfg(test)]
 mod tests {
     use super::*;
diff --git a/crates/paimon/src/table/mod.rs b/crates/paimon/src/table/mod.rs
index 329bc1df..4afd1520 100644
--- a/crates/paimon/src/table/mod.rs
+++ b/crates/paimon/src/table/mod.rs
@@ -73,6 +73,8 @@ mod kv_file_reader;
 mod kv_file_writer;
 mod lumina_index_build_builder;
 pub(crate) mod merge_tree_split_generator;
+#[cfg(test)]
+mod mosaic_table_write_tests;
 mod object_table;
 mod partition_filter;
 mod partition_row_count;
diff --git a/crates/paimon/src/table/mosaic_table_write_tests.rs 
b/crates/paimon/src/table/mosaic_table_write_tests.rs
new file mode 100644
index 00000000..e26a3471
--- /dev/null
+++ b/crates/paimon/src/table/mosaic_table_write_tests.rs
@@ -0,0 +1,962 @@
+// 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.
+
+//! Verify Mosaic writes through ordinary append and primary-key tables.
+
+use std::sync::Arc;
+
+use arrow_array::{Array, Int32Array, Int8Array, RecordBatch, StringArray, 
TimestampSecondArray};
+use arrow_schema::{DataType as ArrowType, Field, Schema as ArrowSchema};
+use futures::TryStreamExt;
+
+use super::{Table, TableCommit, TableWrite};
+use crate::catalog::Identifier;
+use crate::io::{FileIO, FileIOBuilder};
+use crate::spec::{
+    BinaryRow, DataType, Datum, IntType, PredicateBuilder, Schema, 
TableSchema, TimestampType,
+    VarCharType,
+};
+
+fn memory_io() -> FileIO {
+    FileIOBuilder::new("memory").build().unwrap()
+}
+
+async fn setup_dirs(io: &FileIO, path: &str) {
+    io.mkdirs(&format!("{path}/snapshot/")).await.unwrap();
+    io.mkdirs(&format!("{path}/manifest/")).await.unwrap();
+}
+
+fn table(io: &FileIO, path: &str, primary_key: bool, options: &[(&str, &str)]) 
-> Table {
+    let mut schema = Schema::builder()
+        .column("id", DataType::Int(IntType::new()))
+        .column("value", DataType::Int(IntType::new()))
+        .option("file.format", "mosaic");
+    if primary_key {
+        schema = schema.primary_key(["id"]).option("bucket", "1");
+    }
+    for &(key, value) in options {
+        schema = schema.option(key, value);
+    }
+    let schema = schema.build().unwrap();
+    Table::new(
+        io.clone(),
+        Identifier::new("default", "mosaic_table"),
+        path.to_owned(),
+        TableSchema::new(0, &schema),
+        None,
+    )
+}
+
+fn batch(ids: Vec<i32>, values: Vec<i32>) -> RecordBatch {
+    RecordBatch::try_new(
+        Arc::new(ArrowSchema::new(vec![
+            Field::new("id", ArrowType::Int32, true),
+            Field::new("value", ArrowType::Int32, true),
+        ])),
+        vec![
+            Arc::new(Int32Array::from(ids)),
+            Arc::new(Int32Array::from(values)),
+        ],
+    )
+    .unwrap()
+}
+
+async fn read_pairs(table: &Table) -> Vec<(i32, i32)> {
+    let plan = table.new_read_builder().new_scan().plan().await.unwrap();
+    let read = table.new_read_builder().new_read().unwrap();
+    let batches: Vec<RecordBatch> = read
+        .to_arrow(plan.splits())
+        .unwrap()
+        .try_collect()
+        .await
+        .unwrap();
+    let mut pairs = batches
+        .iter()
+        .flat_map(|batch| {
+            let ids = batch
+                .column(0)
+                .as_any()
+                .downcast_ref::<Int32Array>()
+                .unwrap();
+            let values = batch
+                .column(1)
+                .as_any()
+                .downcast_ref::<Int32Array>()
+                .unwrap();
+            (0..batch.num_rows())
+                .map(|row| (ids.value(row), values.value(row)))
+                .collect::<Vec<_>>()
+        })
+        .collect::<Vec<_>>();
+    pairs.sort_unstable();
+    pairs
+}
+
+#[tokio::test]
+async fn append_write_commit_scan_and_read_mosaic() {
+    let io = memory_io();
+    let path = "memory:/native_mosaic_append";
+    setup_dirs(&io, path).await;
+    let table = table(
+        &io,
+        path,
+        false,
+        &[
+            ("mosaic.num-buckets", "2"),
+            ("file.block-size", "64"),
+            ("mosaic.stats-columns", "id,value"),
+        ],
+    );
+    let mut writer = TableWrite::new(&table, "mosaic-test".into()).unwrap();
+    writer
+        .write_arrow_batch(&batch(vec![3, 1], vec![30, 10]))
+        .await
+        .unwrap();
+    writer
+        .write_arrow_batch(&batch(vec![2, 4], vec![20, 40]))
+        .await
+        .unwrap();
+    let messages = writer.prepare_commit().await.unwrap();
+    assert_eq!(messages.len(), 1);
+    assert_eq!(messages[0].new_files.len(), 1);
+    let file = &messages[0].new_files[0];
+    assert!(file.file_name.ends_with(".mosaic"));
+    assert_eq!(file.row_count, 4);
+    assert_eq!(
+        file.value_stats_cols.as_ref().unwrap(),
+        &vec!["id", "value"]
+    );
+    assert_eq!(file.value_stats.null_counts(), &vec![Some(0), Some(0)]);
+    let min = 
BinaryRow::from_serialized_bytes(file.value_stats.min_values()).unwrap();
+    let max = 
BinaryRow::from_serialized_bytes(file.value_stats.max_values()).unwrap();
+    assert_eq!(min.get_int(0).unwrap(), 1);
+    assert_eq!(max.get_int(0).unwrap(), 4);
+    assert_eq!(min.get_int(1).unwrap(), 10);
+    assert_eq!(max.get_int(1).unwrap(), 40);
+
+    TableCommit::new(table.clone(), "mosaic-test".into())
+        .commit(messages)
+        .await
+        .unwrap();
+    assert_eq!(
+        read_pairs(&table).await,
+        vec![(1, 10), (2, 20), (3, 30), (4, 40)]
+    );
+}
+
+#[tokio::test]
+async fn primary_key_mosaic_deduplicates_across_commits() {
+    let io = memory_io();
+    let path = "memory:/native_mosaic_pk";
+    setup_dirs(&io, path).await;
+    let table = table(
+        &io,
+        path,
+        true,
+        &[
+            ("mosaic.num-buckets", "2"),
+            ("mosaic.stats-columns", "value"),
+        ],
+    );
+    let commit = TableCommit::new(table.clone(), "mosaic-test".into());
+
+    let mut first = TableWrite::new(&table, "mosaic-test".into()).unwrap();
+    first
+        .write_arrow_batch(&batch(vec![3, 1, 2], vec![30, 10, 20]))
+        .await
+        .unwrap();
+    let messages = first.prepare_commit().await.unwrap();
+    assert_eq!(messages.len(), 1);
+    let file = &messages[0].new_files[0];
+    assert!(file.file_name.ends_with(".mosaic"));
+    assert_eq!(file.row_count, 3);
+    assert_eq!(file.value_stats_cols.as_ref().unwrap(), &vec!["value"]);
+    commit.commit(messages).await.unwrap();
+
+    let mut second = TableWrite::new(&table, "mosaic-test".into()).unwrap();
+    second
+        .write_arrow_batch(&batch(vec![2, 4], vec![200, 40]))
+        .await
+        .unwrap();
+    commit
+        .commit(second.prepare_commit().await.unwrap())
+        .await
+        .unwrap();
+    assert_eq!(
+        read_pairs(&table).await,
+        vec![(1, 10), (2, 200), (3, 30), (4, 40)]
+    );
+}
+
+#[tokio::test]
+async fn append_mosaic_without_stats_does_not_claim_pruning_data() {
+    let io = memory_io();
+    let path = "memory:/native_mosaic_no_stats";
+    setup_dirs(&io, path).await;
+    let table = table(&io, path, false, &[]);
+    let mut writer = TableWrite::new(&table, "mosaic-test".into()).unwrap();
+    writer
+        .write_arrow_batch(&batch(vec![1, 2], vec![10, 20]))
+        .await
+        .unwrap();
+    let messages = writer.prepare_commit().await.unwrap();
+    let file = &messages[0].new_files[0];
+    assert_eq!(
+        file.value_stats_cols.as_ref().unwrap(),
+        &Vec::<String>::new()
+    );
+    assert!(file.value_stats.null_counts().is_empty());
+    TableCommit::new(table.clone(), "mosaic-test".into())
+        .commit(messages)
+        .await
+        .unwrap();
+    assert_eq!(read_pairs(&table).await, vec![(1, 10), (2, 20)]);
+}
+
+#[tokio::test]
+async fn partitioned_mosaic_writes_files_to_each_partition() {
+    let io = memory_io();
+    let path = "memory:/native_mosaic_partitioned";
+    setup_dirs(&io, path).await;
+    let schema = Schema::builder()
+        .column("pt", DataType::VarChar(VarCharType::string_type()))
+        .column("id", DataType::Int(IntType::new()))
+        .partition_keys(["pt"])
+        .option("file.format", "mosaic")
+        .build()
+        .unwrap();
+    let table = Table::new(
+        io.clone(),
+        Identifier::new("default", "mosaic_partitioned"),
+        path.to_owned(),
+        TableSchema::new(0, &schema),
+        None,
+    );
+    let batch = RecordBatch::try_new(
+        Arc::new(ArrowSchema::new(vec![
+            Field::new("pt", ArrowType::Utf8, true),
+            Field::new("id", ArrowType::Int32, true),
+        ])),
+        vec![
+            Arc::new(StringArray::from(vec!["east", "west", "east"])),
+            Arc::new(Int32Array::from(vec![1, 2, 3])),
+        ],
+    )
+    .unwrap();
+    let mut writer = TableWrite::new(&table, "mosaic-test".into()).unwrap();
+    writer.write_arrow_batch(&batch).await.unwrap();
+    let messages = writer.prepare_commit().await.unwrap();
+    assert_eq!(messages.len(), 2);
+    assert!(messages.iter().all(|message| message.new_files.len() == 1));
+    assert!(messages
+        .iter()
+        .flat_map(|message| &message.new_files)
+        .all(|file| file.file_name.ends_with(".mosaic")));
+    TableCommit::new(table.clone(), "mosaic-test".into())
+        .commit(messages)
+        .await
+        .unwrap();
+    let plan = table.new_read_builder().new_scan().plan().await.unwrap();
+    let read = table.new_read_builder().new_read().unwrap();
+    let batches: Vec<RecordBatch> = read
+        .to_arrow(plan.splits())
+        .unwrap()
+        .try_collect()
+        .await
+        .unwrap();
+    let mut rows = batches
+        .iter()
+        .flat_map(|batch| {
+            let pt = batch
+                .column(0)
+                .as_any()
+                .downcast_ref::<StringArray>()
+                .unwrap();
+            let id = batch
+                .column(1)
+                .as_any()
+                .downcast_ref::<Int32Array>()
+                .unwrap();
+            (0..batch.num_rows())
+                .map(|row| (pt.value(row).to_owned(), id.value(row)))
+                .collect::<Vec<_>>()
+        })
+        .collect::<Vec<_>>();
+    rows.sort_unstable();
+    assert_eq!(
+        rows,
+        vec![("east".into(), 1), ("east".into(), 3), ("west".into(), 2)]
+    );
+}
+
+#[tokio::test]
+async fn append_mosaic_rolls_files_at_target_size_and_preserves_all_rows() {
+    let io = memory_io();
+    let path = "memory:/native_mosaic_rolling";
+    setup_dirs(&io, path).await;
+    let table = table(
+        &io,
+        path,
+        false,
+        &[("target-file-size", "1b"), ("mosaic.stats-columns", "id")],
+    );
+    let mut writer = TableWrite::new(&table, "mosaic-test".into()).unwrap();
+    writer
+        .write_arrow_batch(&batch(vec![1, 2], vec![10, 20]))
+        .await
+        .unwrap();
+    writer
+        .write_arrow_batch(&batch(vec![3, 4], vec![30, 40]))
+        .await
+        .unwrap();
+    writer
+        .write_arrow_batch(&batch(vec![5, 6], vec![50, 60]))
+        .await
+        .unwrap();
+
+    let messages = writer.prepare_commit().await.unwrap();
+    assert_eq!(messages.len(), 1);
+    let files = &messages[0].new_files;
+    assert_eq!(files.len(), 3);
+    assert_eq!(
+        files.iter().map(|file| file.row_count).collect::<Vec<_>>(),
+        vec![2, 2, 2]
+    );
+    assert!(files.iter().all(|file| file.file_name.ends_with(".mosaic")));
+    for (index, file) in files.iter().enumerate() {
+        let min = 
BinaryRow::from_serialized_bytes(file.value_stats.min_values()).unwrap();
+        let max = 
BinaryRow::from_serialized_bytes(file.value_stats.max_values()).unwrap();
+        assert_eq!(min.get_int(0).unwrap(), (index * 2 + 1) as i32);
+        assert_eq!(max.get_int(0).unwrap(), (index * 2 + 2) as i32);
+    }
+    TableCommit::new(table.clone(), "mosaic-test".into())
+        .commit(messages)
+        .await
+        .unwrap();
+    assert_eq!(
+        read_pairs(&table).await,
+        vec![(1, 10), (2, 20), (3, 30), (4, 40), (5, 50), (6, 60)]
+    );
+}
+
+#[tokio::test]
+async fn append_mosaic_stats_prune_manifest_files() {
+    let io = memory_io();
+    let path = "memory:/native_mosaic_pruning";
+    setup_dirs(&io, path).await;
+    let table = table(&io, path, false, &[("mosaic.stats-columns", "id")]);
+    let commit = TableCommit::new(table.clone(), "mosaic-test".into());
+
+    for (ids, values) in [
+        (vec![1, 2], vec![10, 20]),
+        (vec![100, 101], vec![1000, 1010]),
+    ] {
+        let mut writer = TableWrite::new(&table, 
"mosaic-test".into()).unwrap();
+        writer.write_arrow_batch(&batch(ids, values)).await.unwrap();
+        let messages = writer.prepare_commit().await.unwrap();
+        assert_eq!(
+            messages[0].new_files[0].value_stats_cols.as_ref().unwrap(),
+            &["id"]
+        );
+        commit.commit(messages).await.unwrap();
+    }
+
+    let predicate = PredicateBuilder::new(table.schema().fields())
+        .greater_than("id", Datum::Int(10))
+        .unwrap();
+    let mut read = table.new_read_builder();
+    read.with_filter(predicate);
+    let (plan, trace) = read.new_scan().plan_with_trace().await.unwrap();
+    assert_eq!(trace.final_files, 1, "scan trace: {trace:?}");
+    assert!(trace.manifest_entries_pruned_by_data_stats >= 1);
+    let batches: Vec<RecordBatch> = read
+        .new_read()
+        .unwrap()
+        .to_arrow(plan.splits())
+        .unwrap()
+        .try_collect()
+        .await
+        .unwrap();
+    let ids = batches
+        .iter()
+        .flat_map(|batch| {
+            let array = batch
+                .column(0)
+                .as_any()
+                .downcast_ref::<Int32Array>()
+                .unwrap();
+            (0..array.len())
+                .map(|index| array.value(index))
+                .collect::<Vec<_>>()
+        })
+        .collect::<Vec<_>>();
+    assert_eq!(ids, vec![100, 101]);
+}
+
+#[tokio::test]
+async fn primary_key_mosaic_applies_retract_rows() {
+    let io = memory_io();
+    let path = "memory:/native_mosaic_retract";
+    setup_dirs(&io, path).await;
+    let table = table(&io, path, true, &[]);
+    let commit = TableCommit::new(table.clone(), "mosaic-test".into());
+
+    let mut first = TableWrite::new(&table, "mosaic-test".into()).unwrap();
+    first
+        .write_arrow_batch(&batch(vec![1, 2, 3], vec![10, 20, 30]))
+        .await
+        .unwrap();
+    commit
+        .commit(first.prepare_commit().await.unwrap())
+        .await
+        .unwrap();
+
+    let input = RecordBatch::try_new(
+        Arc::new(ArrowSchema::new(vec![
+            Field::new("id", ArrowType::Int32, true),
+            Field::new("value", ArrowType::Int32, true),
+            Field::new("_VALUE_KIND", ArrowType::Int8, false),
+        ])),
+        vec![
+            Arc::new(Int32Array::from(vec![2, 3, 4])),
+            Arc::new(Int32Array::from(vec![20, 300, 40])),
+            Arc::new(Int8Array::from(vec![3, 0, 0])),
+        ],
+    )
+    .unwrap();
+    let mut second = TableWrite::new(&table, "mosaic-test".into()).unwrap();
+    second.write_arrow_batch(&input).await.unwrap();
+    let messages = second.prepare_commit().await.unwrap();
+    assert_eq!(messages[0].new_files[0].delete_row_count, Some(1));
+    commit.commit(messages).await.unwrap();
+    assert_eq!(read_pairs(&table).await, vec![(1, 10), (3, 300), (4, 40)]);
+}
+
+#[tokio::test]
+async fn primary_key_mosaic_merges_latest_values_across_buffer_flushes() {
+    let io = memory_io();
+    let path = "memory:/native_mosaic_pk_rolling";
+    setup_dirs(&io, path).await;
+    let table = table(
+        &io,
+        path,
+        true,
+        &[
+            ("target-file-size", "1b"),
+            ("write-buffer-size", "1b"),
+            ("mosaic.num-buckets", "2"),
+        ],
+    );
+    let mut writer = TableWrite::new(&table, "mosaic-test".into()).unwrap();
+    writer
+        .write_arrow_batch(&batch(vec![1, 2], vec![10, 20]))
+        .await
+        .unwrap();
+    writer
+        .write_arrow_batch(&batch(vec![2, 3], vec![200, 30]))
+        .await
+        .unwrap();
+    writer
+        .write_arrow_batch(&batch(vec![1, 4], vec![100, 40]))
+        .await
+        .unwrap();
+    let messages = writer.prepare_commit().await.unwrap();
+    assert_eq!(messages.len(), 1);
+    assert!(!messages[0].new_files.is_empty());
+    assert!(messages[0]
+        .new_files
+        .iter()
+        .all(|file| file.file_name.ends_with(".mosaic")));
+    TableCommit::new(table.clone(), "mosaic-test".into())
+        .commit(messages)
+        .await
+        .unwrap();
+    assert_eq!(
+        read_pairs(&table).await,
+        vec![(1, 100), (2, 200), (3, 30), (4, 40)]
+    );
+}
+
+#[tokio::test]
+async fn append_mosaic_timestamp_zero_round_trips_through_table_reader() {
+    let io = memory_io();
+    let path = "memory:/native_mosaic_timestamp_zero";
+    setup_dirs(&io, path).await;
+    let schema = Schema::builder()
+        .column("id", DataType::Int(IntType::new()))
+        .column(
+            "event_time",
+            DataType::Timestamp(TimestampType::new(0).unwrap()),
+        )
+        .option("file.format", "mosaic")
+        .option("mosaic.stats-columns", "id")
+        .build()
+        .unwrap();
+    let table = Table::new(
+        io,
+        Identifier::new("default", "mosaic_time"),
+        path.to_owned(),
+        TableSchema::new(0, &schema),
+        None,
+    );
+    let input = RecordBatch::try_new(
+        Arc::new(ArrowSchema::new(vec![
+            Field::new("id", ArrowType::Int32, true),
+            Field::new(
+                "event_time",
+                ArrowType::Timestamp(arrow_schema::TimeUnit::Second, None),
+                true,
+            ),
+        ])),
+        vec![
+            Arc::new(Int32Array::from(vec![1, 2, 3])),
+            Arc::new(TimestampSecondArray::from(vec![Some(-1), None, 
Some(42)])),
+        ],
+    )
+    .unwrap();
+    let mut writer = TableWrite::new(&table, "mosaic-test".into()).unwrap();
+    writer.write_arrow_batch(&input).await.unwrap();
+    let messages = writer.prepare_commit().await.unwrap();
+    assert_eq!(messages[0].new_files[0].row_count, 3);
+    TableCommit::new(table.clone(), "mosaic-test".into())
+        .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>(), 3);
+    let times = batches[0]
+        .column(1)
+        .as_any()
+        .downcast_ref::<TimestampSecondArray>()
+        .unwrap();
+    assert_eq!(times.value(0), -1);
+    assert!(times.is_null(1));
+    assert_eq!(times.value(2), 42);
+}
+
+#[tokio::test]
+async fn append_mosaic_rejects_invalid_table_options_before_commit() {
+    for (case, options) in [
+        ("compression", vec![("file.compression", "snappy")]),
+        ("buckets", vec![("mosaic.num-buckets", "0")]),
+        ("stats", vec![("mosaic.stats-columns", "missing")]),
+    ] {
+        let io = memory_io();
+        let path = format!("memory:/native_mosaic_invalid_{case}");
+        setup_dirs(&io, &path).await;
+        let table = table(&io, &path, false, &options);
+        let mut writer = TableWrite::new(&table, 
"mosaic-test".into()).unwrap();
+        let error = writer
+            .write_arrow_batch(&batch(vec![1], vec![10]))
+            .await
+            .unwrap_err();
+        assert!(
+            error.to_string().contains(case)
+                || error.to_string().contains("zstd")
+                || error.to_string().contains("positive")
+                || error.to_string().contains("not in the file schema"),
+            "{case}: {error}"
+        );
+    }
+}
+
+#[tokio::test]
+async fn primary_key_input_changelog_writes_mosaic_data_and_changelog_files() {
+    let io = memory_io();
+    let path = "memory:/native_mosaic_input_changelog";
+    setup_dirs(&io, path).await;
+    let table = table(
+        &io,
+        path,
+        true,
+        &[
+            ("changelog-producer", "input"),
+            ("mosaic.stats-columns", "value"),
+        ],
+    );
+    let input = RecordBatch::try_new(
+        Arc::new(ArrowSchema::new(vec![
+            Field::new("id", ArrowType::Int32, true),
+            Field::new("value", ArrowType::Int32, true),
+            Field::new("_VALUE_KIND", ArrowType::Int8, false),
+        ])),
+        vec![
+            Arc::new(Int32Array::from(vec![1, 2, 3])),
+            Arc::new(Int32Array::from(vec![10, 20, 30])),
+            Arc::new(Int8Array::from(vec![0, 3, 0])),
+        ],
+    )
+    .unwrap();
+    let mut writer = TableWrite::new(&table, "mosaic-test".into()).unwrap();
+    writer.write_arrow_batch(&input).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(), 1);
+    assert!(messages[0].new_files[0].file_name.ends_with(".mosaic"));
+    assert!(messages[0].new_changelog_files[0]
+        .file_name
+        .ends_with(".mosaic"));
+    assert_eq!(messages[0].new_files[0].delete_row_count, Some(1));
+    assert_eq!(messages[0].new_changelog_files[0].delete_row_count, Some(1));
+    TableCommit::new(table.clone(), "mosaic-test".into())
+        .commit(messages)
+        .await
+        .unwrap();
+    assert_eq!(read_pairs(&table).await, vec![(1, 10), (3, 30)]);
+}
+
+#[tokio::test]
+async fn append_mosaic_write_can_create_file_index() {
+    let io = memory_io();
+    let path = "memory:/native_mosaic_bitmap_index";
+    setup_dirs(&io, path).await;
+    let table = table(
+        &io,
+        path,
+        false,
+        &[
+            ("file-index.bitmap.columns", "id"),
+            ("file-index.in-manifest-threshold", "0 B"),
+        ],
+    );
+    let mut writer = TableWrite::new(&table, "mosaic-test".into()).unwrap();
+    writer
+        .write_arrow_batch(&batch(vec![1, 3], vec![10, 30]))
+        .await
+        .unwrap();
+    let messages = writer.prepare_commit().await.unwrap();
+    assert_eq!(messages[0].new_files.len(), 1);
+    let file = &messages[0].new_files[0];
+    assert!(file.file_name.ends_with(".mosaic"));
+    assert_eq!(file.extra_files, vec![file.data_file_index_file_name()]);
+    let index_path = format!(
+        "{path}/{}/{}",
+        crate::spec::bucket_dir_name(messages[0].bucket),
+        file.data_file_index_file_name()
+    );
+    assert!(io.exists(&index_path).await.unwrap());
+    TableCommit::new(table.clone(), "mosaic-test".into())
+        .commit(messages)
+        .await
+        .unwrap();
+
+    for (key, expected_rows) in [(2, 0), (3, 1)] {
+        let predicate = PredicateBuilder::new(table.schema().fields())
+            .equal("id", Datum::Int(key))
+            .unwrap();
+        let mut reader = table.new_read_builder();
+        reader.with_filter(predicate);
+        let plan = reader.new_scan().plan().await.unwrap();
+        let batches: Vec<RecordBatch> = reader
+            .new_read()
+            .unwrap()
+            .to_arrow(plan.splits())
+            .unwrap()
+            .try_collect()
+            .await
+            .unwrap();
+        assert_eq!(
+            batches.iter().map(RecordBatch::num_rows).sum::<usize>(),
+            expected_rows
+        );
+    }
+}
+
+#[tokio::test]
+async fn append_mosaic_without_stats_keeps_files_in_filtered_plan() {
+    let io = memory_io();
+    let path = "memory:/native_mosaic_untracked_filter";
+    setup_dirs(&io, path).await;
+    let table = table(&io, path, false, &[]);
+    let commit = TableCommit::new(table.clone(), "mosaic-test".into());
+    for (ids, values) in [
+        (vec![1, 2], vec![10, 20]),
+        (vec![100, 101], vec![1000, 1010]),
+    ] {
+        let mut writer = TableWrite::new(&table, 
"mosaic-test".into()).unwrap();
+        writer.write_arrow_batch(&batch(ids, values)).await.unwrap();
+        commit
+            .commit(writer.prepare_commit().await.unwrap())
+            .await
+            .unwrap();
+    }
+    let predicate = PredicateBuilder::new(table.schema().fields())
+        .greater_than("id", Datum::Int(10))
+        .unwrap();
+    let mut reader = table.new_read_builder();
+    reader.with_filter(predicate);
+    let (_plan, trace) = reader.new_scan().plan_with_trace().await.unwrap();
+    assert_eq!(trace.final_files, 2, "scan trace: {trace:?}");
+    assert_eq!(trace.manifest_entries_pruned_by_data_stats, 0);
+}
+
+#[tokio::test]
+async fn append_mosaic_all_null_statistics_remain_conservative() {
+    let io = memory_io();
+    let path = "memory:/native_mosaic_null_stats";
+    setup_dirs(&io, path).await;
+    let table = table(&io, path, false, &[("mosaic.stats-columns", "value")]);
+    let input = RecordBatch::try_new(
+        Arc::new(ArrowSchema::new(vec![
+            Field::new("id", ArrowType::Int32, true),
+            Field::new("value", ArrowType::Int32, true),
+        ])),
+        vec![
+            Arc::new(Int32Array::from(vec![1, 2])),
+            Arc::new(Int32Array::from(vec![None, None])),
+        ],
+    )
+    .unwrap();
+    let mut writer = TableWrite::new(&table, "mosaic-test".into()).unwrap();
+    writer.write_arrow_batch(&input).await.unwrap();
+    let messages = writer.prepare_commit().await.unwrap();
+    let file = &messages[0].new_files[0];
+    assert_eq!(file.value_stats_cols.as_ref().unwrap(), &["value"]);
+    assert_eq!(file.value_stats.null_counts(), &[Some(2)]);
+    let min = 
BinaryRow::from_serialized_bytes(file.value_stats.min_values()).unwrap();
+    let max = 
BinaryRow::from_serialized_bytes(file.value_stats.max_values()).unwrap();
+    assert!(min.is_null_at(0));
+    assert!(max.is_null_at(0));
+    TableCommit::new(table.clone(), "mosaic-test".into())
+        .commit(messages)
+        .await
+        .unwrap();
+
+    let predicate = PredicateBuilder::new(table.schema().fields())
+        .is_null("value")
+        .unwrap();
+    let mut reader = table.new_read_builder();
+    reader.with_filter(predicate);
+    let (plan, trace) = reader.new_scan().plan_with_trace().await.unwrap();
+    assert_eq!(trace.final_files, 1);
+    let batches: Vec<RecordBatch> = reader
+        .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 values = batch
+            .column(1)
+            .as_any()
+            .downcast_ref::<Int32Array>()
+            .unwrap();
+        assert_eq!(values.null_count(), values.len());
+    }
+}
+
+#[tokio::test]
+async fn partitioned_primary_key_mosaic_keeps_same_id_in_distinct_partitions() 
{
+    let io = memory_io();
+    let path = "memory:/native_mosaic_partitioned_pk";
+    setup_dirs(&io, path).await;
+    let schema = Schema::builder()
+        .column("pt", DataType::VarChar(VarCharType::string_type()))
+        .column("id", DataType::Int(IntType::new()))
+        .column("value", DataType::Int(IntType::new()))
+        .partition_keys(["pt"])
+        .primary_key(["pt", "id"])
+        .option("bucket", "1")
+        .option("file.format", "mosaic")
+        .build()
+        .unwrap();
+    let table = Table::new(
+        io,
+        Identifier::new("default", "mosaic_partitioned_pk"),
+        path.to_owned(),
+        TableSchema::new(0, &schema),
+        None,
+    );
+    let make_input = |partitions: Vec<&str>, values: Vec<i32>| {
+        let count = partitions.len();
+        RecordBatch::try_new(
+            Arc::new(ArrowSchema::new(vec![
+                Field::new("pt", ArrowType::Utf8, true),
+                Field::new("id", ArrowType::Int32, true),
+                Field::new("value", ArrowType::Int32, true),
+            ])),
+            vec![
+                Arc::new(StringArray::from(partitions)),
+                Arc::new(Int32Array::from(vec![1; count])),
+                Arc::new(Int32Array::from(values)),
+            ],
+        )
+        .unwrap()
+    };
+    let commit = TableCommit::new(table.clone(), "mosaic-test".into());
+    let mut first = TableWrite::new(&table, "mosaic-test".into()).unwrap();
+    first
+        .write_arrow_batch(&make_input(vec!["east", "west"], vec![10, 20]))
+        .await
+        .unwrap();
+    let messages = first.prepare_commit().await.unwrap();
+    assert_eq!(messages.len(), 2);
+    commit.commit(messages).await.unwrap();
+
+    let mut second = TableWrite::new(&table, "mosaic-test".into()).unwrap();
+    second
+        .write_arrow_batch(&make_input(vec!["east"], vec![100]))
+        .await
+        .unwrap();
+    commit
+        .commit(second.prepare_commit().await.unwrap())
+        .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();
+    let mut rows = batches
+        .iter()
+        .flat_map(|batch| {
+            let partitions = batch
+                .column(0)
+                .as_any()
+                .downcast_ref::<StringArray>()
+                .unwrap();
+            let values = batch
+                .column(2)
+                .as_any()
+                .downcast_ref::<Int32Array>()
+                .unwrap();
+            (0..batch.num_rows())
+                .map(|row| (partitions.value(row).to_owned(), 
values.value(row)))
+                .collect::<Vec<_>>()
+        })
+        .collect::<Vec<_>>();
+    rows.sort_unstable();
+    assert_eq!(rows, vec![("east".into(), 100), ("west".into(), 20)]);
+}
+
+#[tokio::test]
+async fn append_mosaic_filters_rows_when_predicate_column_has_no_stats() {
+    let io = memory_io();
+    let path = "memory:/native_mosaic_residual_filter";
+    setup_dirs(&io, path).await;
+    let table = table(&io, path, false, &[("mosaic.stats-columns", "id")]);
+    let mut writer = TableWrite::new(&table, "mosaic-test".into()).unwrap();
+    writer
+        .write_arrow_batch(&batch(vec![1, 2, 3, 4], vec![5, 20, 10, 30]))
+        .await
+        .unwrap();
+    let messages = writer.prepare_commit().await.unwrap();
+    assert_eq!(
+        messages[0].new_files[0].value_stats_cols.as_ref().unwrap(),
+        &["id"]
+    );
+    TableCommit::new(table.clone(), "mosaic-test".into())
+        .commit(messages)
+        .await
+        .unwrap();
+
+    let predicate = PredicateBuilder::new(table.schema().fields())
+        .greater_than("value", Datum::Int(10))
+        .unwrap();
+    let mut reader = table.new_read_builder();
+    reader.with_filter(predicate);
+    let (plan, trace) = reader.new_scan().plan_with_trace().await.unwrap();
+    assert_eq!(trace.final_files, 1);
+    let batches: Vec<RecordBatch> = reader
+        .new_read()
+        .unwrap()
+        .to_arrow(plan.splits())
+        .unwrap()
+        .try_collect()
+        .await
+        .unwrap();
+    let mut result = batches
+        .iter()
+        .flat_map(|batch| {
+            let ids = batch
+                .column(0)
+                .as_any()
+                .downcast_ref::<Int32Array>()
+                .unwrap();
+            let values = batch
+                .column(1)
+                .as_any()
+                .downcast_ref::<Int32Array>()
+                .unwrap();
+            (0..batch.num_rows())
+                .map(|row| (ids.value(row), values.value(row)))
+                .collect::<Vec<_>>()
+        })
+        .collect::<Vec<_>>();
+    result.sort_unstable();
+    assert_eq!(result, vec![(2, 20), (4, 30)]);
+}
+
+#[tokio::test]
+async fn primary_key_mosaic_rejects_multiset_from_table_schema() {
+    use crate::arrow::build_target_arrow_schema;
+    use crate::spec::MultisetType;
+    use arrow_array::new_null_array;
+
+    let io = memory_io();
+    let path = "memory:/native_mosaic_pk_multiset";
+    setup_dirs(&io, path).await;
+    let schema = Schema::builder()
+        .column("id", DataType::Int(IntType::new()))
+        .column(
+            "bag",
+            
DataType::Multiset(MultisetType::new(DataType::Int(IntType::new()))),
+        )
+        .primary_key(["id"])
+        .option("bucket", "1")
+        .option("file.format", "mosaic")
+        .build()
+        .unwrap();
+    let table = Table::new(
+        io,
+        Identifier::new("default", "mosaic_pk_multiset"),
+        path.to_owned(),
+        TableSchema::new(0, &schema),
+        None,
+    );
+    let arrow_schema = 
build_target_arrow_schema(table.schema().fields()).unwrap();
+    let input = RecordBatch::try_new(
+        arrow_schema.clone(),
+        vec![
+            Arc::new(Int32Array::from(vec![1])),
+            new_null_array(arrow_schema.field(1).data_type(), 1),
+        ],
+    )
+    .unwrap();
+    let mut writer = TableWrite::new(&table, "mosaic-test".into()).unwrap();
+    writer.write_arrow_batch(&input).await.unwrap();
+    let error = writer.prepare_commit().await.unwrap_err();
+    assert!(error.to_string().contains("MULTISET"), "{error}");
+}
diff --git a/crates/paimon/src/table/table_write.rs 
b/crates/paimon/src/table/table_write.rs
index 9901072c..68be711f 100644
--- a/crates/paimon/src/table/table_write.rs
+++ b/crates/paimon/src/table/table_write.rs
@@ -2389,6 +2389,80 @@ pub(in crate::table) mod tests {
         assert_eq!(collect_i32(&batches, 1), vec![10, 20, 30]);
     }
 
+    #[tokio::test]
+    async fn 
test_primary_key_non_parquet_formats_roundtrip_with_typed_physical_fields() {
+        for format in ["row", "avro"] {
+            let file_io = test_file_io();
+            let table_path = format!("memory:/test_pk_typed_fields_{format}");
+            setup_dirs(&file_io, &table_path).await;
+            let schema = Schema::builder()
+                .column("id", DataType::Int(IntType::new()))
+                .column("value", DataType::Int(IntType::new()))
+                .primary_key(["id"])
+                .option("bucket", "1")
+                .option("file.format", format)
+                .build()
+                .unwrap();
+            let table = Table::new(
+                file_io,
+                Identifier::new("default", "test_pk_typed_fields"),
+                table_path,
+                TableSchema::new(0, &schema),
+                None,
+            );
+            let mut writer = TableWrite::new(&table, 
"test-user".into()).unwrap();
+            writer
+                .write_arrow_batch(&make_batch(vec![2, 1], vec![20, 10]))
+                .await
+                .unwrap();
+            let messages = writer.prepare_commit().await.unwrap();
+            assert_eq!(messages[0].new_files.len(), 1);
+            assert!(messages[0].new_files[0].file_name.ends_with(format));
+            let file = &messages[0].new_files[0];
+            let file_path = format!(
+                "{}/{}/{}",
+                table.location(),
+                bucket_dir_name(messages[0].bucket),
+                file.file_name
+            );
+            let physical_fields = vec![
+                DataField::new(
+                    SEQUENCE_NUMBER_FIELD_ID,
+                    SEQUENCE_NUMBER_FIELD_NAME.into(),
+                    DataType::BigInt(BigIntType::new()),
+                ),
+                DataField::new(
+                    VALUE_KIND_FIELD_ID,
+                    VALUE_KIND_FIELD_NAME.into(),
+                    DataType::TinyInt(TinyIntType::new()),
+                ),
+                table.schema().fields()[0].clone(),
+                table.schema().fields()[1].clone(),
+            ];
+            let format_reader = create_format_reader(&file_path, false, 
&physical_fields).unwrap();
+            let input = table.file_io().new_input(&file_path).unwrap();
+            let stream = format_reader
+                .read_batch_stream(
+                    Box::new(input.reader().await.unwrap()),
+                    file.file_size as u64,
+                    &physical_fields,
+                    None,
+                    None,
+                    None,
+                )
+                .await
+                .unwrap();
+            let batches: Vec<RecordBatch> =
+                futures::TryStreamExt::try_collect(stream).await.unwrap();
+            assert_eq!(collect_i32(&batches, 2), vec![1, 2], "{format}");
+            assert_eq!(collect_i32(&batches, 3), vec![10, 20], "{format}");
+            TableCommit::new(table.clone(), "test-user".into())
+                .commit(messages)
+                .await
+                .unwrap();
+        }
+    }
+
     #[test]
     fn test_allows_append_blob_table() {
         let table = Table::new(
diff --git a/docs/src/getting-started.md b/docs/src/getting-started.md
index 71ed524c..de526dff 100644
--- a/docs/src/getting-started.md
+++ b/docs/src/getting-started.md
@@ -131,7 +131,23 @@ change Page Index generation.
 
 ## Mosaic File Format
 
-Mosaic data file reads are always available. The current Mosaic support is 
read-only: Paimon Rust can read existing `.mosaic` data files, including array 
and map columns, in a Paimon table, but it does not write Mosaic data files yet.
+Mosaic data file reads and writes are available for ordinary Paimon tables. Set
+`file.format=mosaic` when creating an append or primary-key table. Mosaic 
supports
+scalar, array, and map columns; `ROW`, `MULTISET`, `VECTOR`, `VARIANT`, and 
`BLOB`
+columns are unsupported, matching Java Paimon. Mosaic writes use Zstd 
compression.
+
+The writer accepts `mosaic.num-buckets`, `file.block-size`, and a 
comma-separated
+`mosaic.stats-columns` list. Statistics are included in manifest entries for 
the
+selected columns and can prune files during scans. An empty list disables 
writer
+statistics. For example:
+
+```text
+file.format = mosaic
+file.compression = zstd
+mosaic.num-buckets = 4
+file.block-size = 128mb
+mosaic.stats-columns = id,event_time
+```
 
 ## FileIndexes for Append Writes
 
diff --git a/docs/src/sql.md b/docs/src/sql.md
index 031a38bc..259b2b2a 100644
--- a/docs/src/sql.md
+++ b/docs/src/sql.md
@@ -31,7 +31,8 @@ datafusion = "54.0.0"
 tokio = { version = "1", features = ["full"] }
 ```
 
-Mosaic support is always available and currently read-only. SQL queries can 
read existing `.mosaic` files, but Paimon Rust does not write Mosaic data files 
yet.
+Mosaic support is always available. SQL queries can read existing `.mosaic` 
files,
+and writes to ordinary Paimon tables can create them with `file.format=mosaic`.
 
 ## SQL Support Scope
 
@@ -779,7 +780,9 @@ For primary-key tables, records with duplicate keys are 
deduplicated according t
 
 The Mosaic reader supports scalar, temporal, array, and map columns. It uses 
row-group statistics for conservative pruning when they are present. This 
pruning is not row-level filter enforcement; DataFusion still applies SQL 
filters above the reader to produce exact query results.
 
-Unsupported or limited Mosaic areas include writing `.mosaic` files, emitting 
manifest `value_stats` for Mosaic writes, Mosaic bloom filters, and 
Mosaic-specific performance tuning.
+Mosaic writes use Zstd and support `mosaic.num-buckets`, `file.block-size`, and
+`mosaic.stats-columns`. Selected statistics are written to the manifest for 
file
+pruning. Mosaic bloom filters are not supported.
 
 ### INSERT OVERWRITE
 

Reply via email to