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