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