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 5a780c4e Support nested ROW data evolution in native table reads and
writes (#960)
5a780c4e is described below
commit 5a780c4eb771bbacd8b63be6f42bffbebab6e341
Author: Jingsong Lee <[email protected]>
AuthorDate: Sat Sep 26 10:14:10 2026 +0800
Support nested ROW data evolution in native table reads and writes (#960)
---
crates/paimon/src/spec/core_options.rs | 8 +
crates/paimon/src/table/data_evolution_fields.rs | 338 +++++++++
crates/paimon/src/table/data_evolution_nested.rs | 676 +++++++++++++++++
crates/paimon/src/table/data_evolution_reader.rs | 175 +++--
crates/paimon/src/table/data_evolution_writer.rs | 911 ++++++++++++++++++++---
crates/paimon/src/table/data_file_reader.rs | 281 +++++--
crates/paimon/src/table/mod.rs | 2 +
crates/paimon/src/table/table_commit.rs | 157 +++-
crates/paimon/src/table/table_scan.rs | 24 +-
9 files changed, 2306 insertions(+), 266 deletions(-)
diff --git a/crates/paimon/src/spec/core_options.rs
b/crates/paimon/src/spec/core_options.rs
index fb517c11..cb69c888 100644
--- a/crates/paimon/src/spec/core_options.rs
+++ b/crates/paimon/src/spec/core_options.rs
@@ -23,6 +23,7 @@ const DELETION_VECTORS_ENABLED_OPTION: &str =
"deletion-vectors.enabled";
const DELETION_VECTORS_MERGE_ON_READ_OPTION: &str =
"deletion-vectors.merge-on-read";
pub(crate) const QUERY_AUTH_ENABLED_OPTION: &str = "query-auth.enabled";
const DATA_EVOLUTION_ENABLED_OPTION: &str = "data-evolution.enabled";
+const DATA_EVOLUTION_NESTED_FIELD_ENABLED_OPTION: &str =
"data-evolution.nested-field.enabled";
const FILE_INDEX_READ_ENABLED_OPTION: &str = "file-index.read.enabled";
const GLOBAL_INDEX_ENABLED_OPTION: &str = "global-index.enabled";
const GLOBAL_INDEX_SEARCH_MODE_OPTION: &str = "global-index.search-mode";
@@ -728,6 +729,13 @@ impl<'a> CoreOptions<'a> {
.unwrap_or(false)
}
+ pub fn data_evolution_nested_field_enabled(&self) -> bool {
+ self.options
+ .get(DATA_EVOLUTION_NESTED_FIELD_ENABLED_OPTION)
+ .map(|value| value.eq_ignore_ascii_case("true"))
+ .unwrap_or(false)
+ }
+
/// Maximum complete FileIndex size stored in the manifest. Default is 500
bytes.
pub(crate) fn file_index_in_manifest_threshold(&self) ->
crate::Result<i64> {
match self.options.get("file-index.in-manifest-threshold") {
diff --git a/crates/paimon/src/table/data_evolution_fields.rs
b/crates/paimon/src/table/data_evolution_fields.rs
new file mode 100644
index 00000000..96476d6b
--- /dev/null
+++ b/crates/paimon/src/table/data_evolution_fields.rs
@@ -0,0 +1,338 @@
+// 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.
+
+//! Physical field selection for data-evolution `write_cols`.
+//!
+//! Java `RowType.projectByPaths` preserves the order of the first occurrence
+//! of each selected field. A top-level name takes precedence over splitting at
+//! a dot, because older schemas can contain literal dotted column names.
+
+use crate::spec::{DataField, DataType, RowType};
+use crate::{Error, Result};
+use indexmap::IndexMap;
+use std::collections::HashSet;
+
+#[derive(Default)]
+struct PathSelection {
+ whole: bool,
+ tails: Vec<String>,
+}
+
+pub(super) fn project_by_paths(fields: &[DataField], paths: &[String]) ->
Result<Vec<DataField>> {
+ let mut selected: IndexMap<String, PathSelection> = IndexMap::new();
+ for path in paths {
+ if fields.iter().any(|field| field.name() == path) {
+ selected.entry(path.clone()).or_default().whole = true;
+ continue;
+ }
+ let Some((head, tail)) = path.split_once('.') else {
+ return Err(unknown_field(path));
+ };
+ if head.is_empty() || tail.is_empty() {
+ return Err(unknown_field(path));
+ }
+ selected
+ .entry(head.to_string())
+ .or_default()
+ .tails
+ .push(tail.to_string());
+ }
+
+ selected
+ .into_iter()
+ .map(|(name, selection)| {
+ let field = fields
+ .iter()
+ .find(|field| field.name() == name)
+ .ok_or_else(|| unknown_field(&name))?;
+ if selection.whole || selection.tails.is_empty() {
+ return Ok(field.clone());
+ }
+ let DataType::Row(row) = field.data_type() else {
+ return Err(Error::DataInvalid {
+ message: format!(
+ "Cannot project nested write path(s) {:?} from non-ROW
field '{name}'",
+ selection.tails
+ ),
+ source: None,
+ });
+ };
+ let projected = project_by_paths(row.fields(), &selection.tails)?;
+ Ok(field_with_type(
+ field,
+ DataType::Row(RowType::with_nullable(
+ field.data_type().is_nullable(),
+ projected,
+ )),
+ ))
+ })
+ .collect()
+}
+
+/// Java `collectLeafPaths` only emits a top-level field or one direct child
+/// of a ROW. A deeper partial ROW could be decoded but cannot be assembled
+/// from independently written files, so reject it before creating data files.
+pub(super) fn validate_write_paths(
+ fields: &[DataField],
+ paths: &[String],
+ nested_field_enabled: bool,
+) -> Result<()> {
+ let mut seen = HashSet::new();
+ let whole = paths
+ .iter()
+ .filter(|path| fields.iter().any(|field| field.name() ==
path.as_str()))
+ .map(String::as_str)
+ .collect::<HashSet<_>>();
+ for path in paths {
+ if !seen.insert(path.as_str()) {
+ return Err(Error::DataInvalid {
+ message: format!("Duplicate data-evolution write path
'{path}'"),
+ source: None,
+ });
+ }
+ if whole.contains(path.as_str()) {
+ // `project_by_paths` must prefer an exact top-level name when
+ // decoding existing files. A string-only write API cannot tell
+ // that name apart from an equally named nested SET target.
+ if nested_field_enabled && ambiguous_nested_write_path(fields,
path) {
+ return Err(Error::Unsupported {
+ message: format!(
+ "Ambiguous data-evolution write path '{path}': it
names both a top-level column and a nested field"
+ ),
+ });
+ }
+ continue;
+ }
+ let Some((head, child)) = path.split_once('.') else {
+ return Err(unknown_field(path));
+ };
+ if child.contains('.') {
+ return Err(Error::Unsupported {
+ message: format!(
+ "Sub-field-level data evolution supports only one level of
partial nesting: '{path}'"
+ ),
+ });
+ }
+ if whole.contains(head) {
+ return Err(Error::DataInvalid {
+ message: format!("Cannot write both whole ROW '{head}' and its
sub-field '{path}'"),
+ source: None,
+ });
+ }
+ let parent = fields
+ .iter()
+ .find(|field| field.name() == head)
+ .ok_or_else(|| unknown_field(path))?;
+ let DataType::Row(row) = parent.data_type() else {
+ return Err(Error::DataInvalid {
+ message: format!("Cannot write nested path '{path}' from
non-ROW '{head}'"),
+ source: None,
+ });
+ };
+ if !row.fields().iter().any(|field| field.name() == child) {
+ return Err(unknown_field(path));
+ }
+ }
+ Ok(())
+}
+
+fn ambiguous_nested_write_path(fields: &[DataField], path: &str) -> bool {
+ let Some((head, child)) = path.split_once('.') else {
+ return false;
+ };
+ if child.contains('.') {
+ return false;
+ }
+ fields
+ .iter()
+ .find(|field| field.name() == head)
+ .and_then(|field| match field.data_type() {
+ DataType::Row(row) => Some(row),
+ _ => None,
+ })
+ .is_some_and(|row| row.fields().iter().any(|field| field.name() ==
child))
+}
+
+pub(super) fn field_with_type(field: &DataField, data_type: DataType) ->
DataField {
+ DataField::new(field.id(), field.name().to_string(), data_type)
+ .with_description(field.description().map(str::to_string))
+ .with_default_value(field.default_value().map(str::to_string))
+}
+
+/// Field identities for row-id conflict detection. A whole ROW write covers
+/// every leaf, while a nested write covers only its selected descendants.
+pub(super) fn write_leaf_ids(
+ fields: &[DataField],
+ paths: Option<&[String]>,
+) -> Result<HashSet<i32>> {
+ let selected = match paths {
+ Some(paths) => project_by_paths(fields, paths)?,
+ None => fields.to_vec(),
+ };
+ let mut ids = HashSet::new();
+ collect_leaf_ids(&selected, &mut ids);
+ Ok(ids)
+}
+
+fn collect_leaf_ids(fields: &[DataField], ids: &mut HashSet<i32>) {
+ for field in fields {
+ if let DataType::Row(row) = field.data_type() {
+ collect_leaf_ids(row.fields(), ids);
+ } else {
+ ids.insert(field.id());
+ }
+ }
+}
+
+fn unknown_field(path: &str) -> Error {
+ Error::DataInvalid {
+ message: format!("Unknown data-evolution write path '{path}'"),
+ source: None,
+ }
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+ use crate::spec::{IntType, VarCharType};
+
+ fn field(id: i32, name: &str, data_type: DataType) -> DataField {
+ DataField::new(id, name.to_string(), data_type)
+ }
+
+ fn fields() -> Vec<DataField> {
+ vec![
+ field(1, "id", DataType::Int(IntType::new())),
+ field(
+ 2,
+ "profile",
+ DataType::Row(RowType::new(vec![
+ field(3, "name",
DataType::VarChar(VarCharType::default())),
+ field(4, "age", DataType::Int(IntType::new())),
+ ])),
+ ),
+ field(5, "profile.name", DataType::Int(IntType::new())),
+ ]
+ }
+
+ #[test]
+ fn projects_partial_row_in_path_order() {
+ let projected = project_by_paths(
+ &fields(),
+ &["profile.age".into(), "id".into(), "profile.name".into()],
+ )
+ .unwrap();
+ assert_eq!(
+ projected.iter().map(DataField::id).collect::<Vec<_>>(),
+ vec![2, 1, 5]
+ );
+ let DataType::Row(row) = projected[0].data_type() else {
+ panic!("expected ROW")
+ };
+ assert_eq!(
+ row.fields().iter().map(DataField::id).collect::<Vec<_>>(),
+ vec![4]
+ );
+ }
+
+ #[test]
+ fn whole_field_overrides_partial_selection() {
+ let projected =
+ project_by_paths(&fields(), &["profile.age".into(),
"profile".into()]).unwrap();
+ assert_eq!(projected.len(), 1);
+ assert_eq!(projected[0], fields()[1]);
+ }
+
+ #[test]
+ fn rejects_unknown_or_non_row_path() {
+ for path in ["missing", "profile.missing", "id.part", "profile."] {
+ assert!(
+ project_by_paths(&fields(), &[path.into()]).is_err(),
+ "{path}"
+ );
+ }
+ }
+
+ #[test]
+ fn whole_row_overlaps_both_nested_children_but_siblings_are_distinct() {
+ let fields = fields().into_iter().take(2).collect::<Vec<_>>();
+ let whole = write_leaf_ids(&fields,
Some(&["profile".into()])).unwrap();
+ let name = write_leaf_ids(&fields,
Some(&["profile.name".into()])).unwrap();
+ let age = write_leaf_ids(&fields,
Some(&["profile.age".into()])).unwrap();
+ assert_eq!(whole, HashSet::from([3, 4]));
+ assert_eq!(name, HashSet::from([3]));
+ assert_eq!(age, HashSet::from([4]));
+ assert!(!whole.is_disjoint(&name));
+ assert!(!whole.is_disjoint(&age));
+ assert!(name.is_disjoint(&age));
+ }
+
+ #[test]
+ fn write_paths_allow_direct_child_row_but_reject_deeper_partial_row() {
+ let nested = vec![field(
+ 1,
+ "profile",
+ DataType::Row(RowType::new(vec![field(
+ 2,
+ "address",
+ DataType::Row(RowType::new(vec![field(
+ 3,
+ "zip",
+ DataType::Int(IntType::new()),
+ )])),
+ )])),
+ )];
+ validate_write_paths(&nested, &["profile.address".into()],
true).unwrap();
+ assert!(matches!(
+ validate_write_paths(&nested, &["profile.address.zip".into()],
true),
+ Err(Error::Unsupported { .. })
+ ));
+ }
+
+ #[test]
+ fn write_paths_reject_duplicate_and_overlapping_selection() {
+ let fields = fields().into_iter().take(2).collect::<Vec<_>>();
+ for paths in [
+ vec!["profile.age".into(), "profile.age".into()],
+ vec!["profile".into(), "profile.age".into()],
+ vec!["profile.age".into(), "profile".into()],
+ ] {
+ assert!(validate_write_paths(&fields, &paths, true).is_err());
+ }
+ validate_write_paths(
+ &fields,
+ &["profile.age".into(), "profile.name".into()],
+ true,
+ )
+ .unwrap();
+ }
+
+ #[test]
+ fn
write_paths_reject_top_level_and_nested_name_collision_only_when_enabled() {
+ let fields = fields();
+ let path = ["profile.name".to_string()];
+ assert!(matches!(
+ validate_write_paths(&fields, &path, true),
+ Err(Error::Unsupported { .. })
+ ));
+ validate_write_paths(&fields, &path, false).unwrap();
+
+ // Physical file decoding keeps Java's exact top-level name precedence.
+ let selected = project_by_paths(&fields, &path).unwrap();
+ assert_eq!(selected[0].id(), 5);
+ }
+}
diff --git a/crates/paimon/src/table/data_evolution_nested.rs
b/crates/paimon/src/table/data_evolution_nested.rs
new file mode 100644
index 00000000..6d65be42
--- /dev/null
+++ b/crates/paimon/src/table/data_evolution_nested.rs
@@ -0,0 +1,676 @@
+// 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.
+
+//! One-level ROW composition for data-evolution column groups.
+//!
+//! Files are ordered newest first. As in Java `DataEvolutionReadPlanner`, the
+//! first file containing each leaf wins. All latest sibling providers also
+//! contribute to the parent ROW's nullness, including siblings omitted by a
+//! read projection. Deeper splits of the same direct subfield are rejected.
+
+use crate::spec::{DataField, DataType, RowType};
+use crate::{Error, Result};
+use arrow_array::{Array, ArrayRef, RecordBatch, StructArray, UInt32Array};
+use arrow_buffer::NullBuffer;
+use arrow_schema::DataType as ArrowDataType;
+use arrow_select::take::take;
+use std::collections::HashSet;
+use std::sync::Arc;
+
+pub(super) struct PhysicalRowSource<'a> {
+ pub source_index: usize,
+ pub fields: &'a [DataField],
+}
+
+#[derive(Debug, Clone, PartialEq, Eq)]
+pub(super) struct NestedSelection {
+ pub anchors: Vec<usize>,
+ pub children: Vec<Option<usize>>,
+ pub whole_source: Option<usize>,
+}
+
+/// Return `None` only when no file supplies any leaf of this ROW.
+pub(super) fn select_nested_row(
+ read_field: &DataField,
+ sources: &[PhysicalRowSource<'_>],
+) -> Result<Option<NestedSelection>> {
+ let DataType::Row(read_row) = read_field.data_type() else {
+ return Err(Error::DataInvalid {
+ message: format!(
+ "Nested data-evolution field '{}' is not a ROW",
+ read_field.name()
+ ),
+ source: None,
+ });
+ };
+
+ let source_rows = sources
+ .iter()
+ .filter_map(|source| {
+ source
+ .fields
+ .iter()
+ .find(|field| field.id() == read_field.id())
+ .and_then(|field| match field.data_type() {
+ DataType::Row(row) => Some((source.source_index,
row.fields())),
+ _ => None,
+ })
+ })
+ .collect::<Vec<_>>();
+ if source_rows.is_empty() {
+ return Ok(None);
+ }
+
+ let mut all_leaf_ids = Vec::new();
+ for (_, fields) in &source_rows {
+ collect_leaf_ids(fields, &mut all_leaf_ids);
+ }
+ let mut anchors = Vec::new();
+ for leaf in all_leaf_ids {
+ if let Some(source) = first_provider(leaf, &source_rows) {
+ if !anchors.contains(&source) {
+ anchors.push(source);
+ }
+ }
+ }
+ if anchors.is_empty() {
+ return Ok(None);
+ }
+
+ let mut children = Vec::with_capacity(read_row.fields().len());
+ let mut all_requested_leaves_covered = true;
+ for child in read_row.fields() {
+ let mut requested_leaves = Vec::new();
+ collect_leaf_ids(std::slice::from_ref(child), &mut requested_leaves);
+ let providers = providers_of(&requested_leaves, &source_rows);
+ all_requested_leaves_covered &= requested_leaves
+ .iter()
+ .all(|leaf| first_provider(*leaf, &source_rows).is_some());
+
+ let providers = if providers.is_empty() {
+ // A projected sub-ROW may contain only recently added descendants.
+ // Read its older siblings from their latest provider to retain the
+ // sub-ROW nullness while missing descendants are NULL-filled.
+ let mut sibling_leaves = Vec::new();
+ for (_, fields) in &source_rows {
+ if let Some(existing) = fields.iter().find(|field| field.id()
== child.id()) {
+ collect_leaf_ids(std::slice::from_ref(existing), &mut
sibling_leaves);
+ }
+ }
+ providers_of(&sibling_leaves, &source_rows)
+ } else {
+ providers
+ };
+ if providers.len() > 1 {
+ return Err(Error::Unsupported {
+ message: format!(
+ "Sub-field-level data evolution does not support splitting
nested sub-field '{}.{}' across files",
+ read_field.name(), child.name()
+ ),
+ });
+ }
+ if providers.is_empty() && !child.data_type().is_nullable() {
+ return Err(Error::DataInvalid {
+ message: format!(
+ "Cannot read non-nullable nested field '{}.{}' without a
provider",
+ read_field.name(),
+ child.name()
+ ),
+ source: None,
+ });
+ }
+ children.push(providers.first().copied());
+ }
+
+ let whole_source = if anchors.len() == 1
+ && all_requested_leaves_covered
+ && children.iter().all(|source| *source == Some(anchors[0]))
+ {
+ Some(anchors[0])
+ } else {
+ None
+ };
+ Ok(Some(NestedSelection {
+ anchors,
+ children,
+ whole_source,
+ }))
+}
+
+fn collect_leaf_ids(fields: &[DataField], output: &mut Vec<i32>) {
+ for field in fields {
+ match field.data_type() {
+ DataType::Row(row) => collect_leaf_ids(row.fields(), output),
+ _ => output.push(field.id()),
+ }
+ }
+}
+
+fn first_provider(leaf_id: i32, sources: &[(usize, &[DataField])]) ->
Option<usize> {
+ sources.iter().find_map(|(index, fields)| {
+ let mut leaves = Vec::new();
+ collect_leaf_ids(fields, &mut leaves);
+ leaves.contains(&leaf_id).then_some(*index)
+ })
+}
+
+fn providers_of(leaves: &[i32], sources: &[(usize, &[DataField])]) ->
Vec<usize> {
+ let mut providers = Vec::new();
+ let mut seen = HashSet::new();
+ for leaf in leaves {
+ if let Some(source) = first_provider(*leaf, sources) {
+ if seen.insert(source) {
+ providers.push(source);
+ }
+ }
+ }
+ providers
+}
+
+/// Source indexes and their parent-field offsets are fixed by the outer read
+/// plan after `select_nested_row` chooses the physical providers.
+#[derive(Debug, Clone)]
+pub(super) struct NestedFieldPlan {
+ pub anchors: Vec<(usize, usize)>,
+ pub children: Vec<Option<(usize, usize)>>,
+}
+
+/// Build the smallest current-schema ROW shape needed from one provider.
+/// Selected children are read whole to retain a nested child's nullness after
+/// schema evolution. A provider used only as a parent-nullness anchor still
+/// needs one physical child decoded; otherwise a projected file with no
+/// requested leaves would produce an all-NULL synthetic parent.
+pub(super) fn source_read_field(
+ requested: &DataField,
+ current: &DataField,
+ source_index: usize,
+ source_fields: &[DataField],
+ selection: &NestedSelection,
+) -> Result<DataField> {
+ let DataType::Row(requested_row) = requested.data_type() else {
+ return Err(Error::DataInvalid {
+ message: "Nested source read requires a ROW
projection".to_string(),
+ source: None,
+ });
+ };
+ let DataType::Row(current_row) = current.data_type() else {
+ return Err(Error::DataInvalid {
+ message: "Nested source read requires a current ROW
field".to_string(),
+ source: None,
+ });
+ };
+ let selected_ids = requested_row
+ .fields()
+ .iter()
+ .zip(&selection.children)
+ .filter_map(|(child, provider)| (*provider ==
Some(source_index)).then_some(child.id()))
+ .collect::<HashSet<_>>();
+ let mut children = current_row
+ .fields()
+ .iter()
+ .filter(|child| selected_ids.contains(&child.id()))
+ .cloned()
+ .collect::<Vec<_>>();
+ if children.is_empty() {
+ let physical = source_fields
+ .iter()
+ .find(|field| field.id() == requested.id())
+ .ok_or_else(|| Error::DataInvalid {
+ message: format!(
+ "Nested source {source_index} does not contain field '{}'",
+ requested.name()
+ ),
+ source: None,
+ })?;
+ let DataType::Row(source_row) = physical.data_type() else {
+ return Err(Error::DataInvalid {
+ message: format!("Nested source {source_index} is not a ROW"),
+ source: None,
+ });
+ };
+ let anchor = source_row
+ .fields()
+ .first()
+ .ok_or_else(|| Error::DataInvalid {
+ message: format!("Nested source {source_index} has no physical
child"),
+ source: None,
+ })?;
+ children.push(
+ current_row
+ .fields()
+ .iter()
+ .find(|field| field.id() == anchor.id())
+ .unwrap_or(anchor)
+ .clone(),
+ );
+ }
+ Ok(super::data_evolution_fields::field_with_type(
+ current,
+ DataType::Row(RowType::with_nullable(
+ current.data_type().is_nullable(),
+ children,
+ )),
+ ))
+}
+
+pub(super) fn assemble_nested_row(
+ plan: &NestedFieldPlan,
+ target_type: &ArrowDataType,
+ cursors: &[Option<(RecordBatch, usize)>],
+ rows: usize,
+) -> Result<ArrayRef> {
+ let ArrowDataType::Struct(fields) = target_type else {
+ return Err(Error::UnexpectedError {
+ message: format!("Nested data-evolution target is not a struct:
{target_type:?}"),
+ source: None,
+ });
+ };
+ if fields.len() != plan.children.len() {
+ return Err(Error::UnexpectedError {
+ message: "Nested data-evolution plan has the wrong number of
children".to_string(),
+ source: None,
+ });
+ }
+
+ let mut anchors = Vec::with_capacity(plan.anchors.len());
+ for &(source, field) in &plan.anchors {
+ anchors.push(source_struct(cursors, source, field, rows)?);
+ }
+ let valid = (0..rows)
+ .map(|row| anchors.iter().any(|anchor| anchor.is_valid(row)))
+ .collect::<Vec<_>>();
+ let parent_nulls = (!valid.iter().all(|value| *value)).then(||
NullBuffer::from(valid));
+
+ let mut children = Vec::with_capacity(fields.len());
+ for (index, child_field) in fields.iter().enumerate() {
+ let Some((source, field)) = plan.children[index] else {
+ children.push(arrow_array::new_null_array(child_field.data_type(),
rows));
+ continue;
+ };
+ let parent = source_struct(cursors, source, field, rows)?;
+ let value =
+ parent
+ .column_by_name(child_field.name())
+ .ok_or_else(|| Error::DataInvalid {
+ message: format!(
+ "Nested data-evolution source {source} is missing
child '{}'",
+ child_field.name()
+ ),
+ source: None,
+ })?;
+ let value = project_child(value, child_field.data_type())?;
+ if parent.null_count() == 0 {
+ children.push(value);
+ } else {
+ // A partial ROW's child bytes are meaningless when that source's
+ // parent is NULL, even if another sibling source keeps the merged
+ // parent non-NULL.
+ let positions = UInt32Array::from_iter(
+ (0..rows).map(|row| parent.is_valid(row).then_some(row as
u32)),
+ );
+ children.push(take(value.as_ref(), &positions,
None).map_err(|error| {
+ Error::UnexpectedError {
+ message: format!("Failed to mask NULL nested
data-evolution source: {error}"),
+ source: Some(Box::new(error)),
+ }
+ })?);
+ }
+ }
+
+ let assembled =
+ StructArray::try_new(fields.clone(), children,
parent_nulls).map_err(|error| {
+ Error::UnexpectedError {
+ message: format!("Failed to compose nested data-evolution ROW:
{error}"),
+ source: Some(Box::new(error)),
+ }
+ })?;
+ Ok(Arc::new(assembled))
+}
+
+/// Providers are read as the current full ROW so older siblings can preserve
+/// nullness. Trim each selected child back to the requested nested projection
+/// before composing the output; Arrow Struct children must match that shape.
+fn project_child(value: &ArrayRef, target: &ArrowDataType) -> Result<ArrayRef>
{
+ if value.data_type() == target {
+ return Ok(value.clone());
+ }
+ if let ArrowDataType::Struct(target_fields) = target {
+ let source = value
+ .as_any()
+ .downcast_ref::<StructArray>()
+ .ok_or_else(|| Error::DataInvalid {
+ message: format!(
+ "Expected a ROW while projecting nested data-evolution
child, got {:?}",
+ value.data_type()
+ ),
+ source: None,
+ })?;
+ let children = target_fields
+ .iter()
+ .map(|field| {
+ let child =
+ source
+ .column_by_name(field.name())
+ .ok_or_else(|| Error::DataInvalid {
+ message: format!(
+ "Nested data-evolution source is missing
projected child '{}'",
+ field.name()
+ ),
+ source: None,
+ })?;
+ project_child(child, field.data_type())
+ })
+ .collect::<Result<Vec<_>>>()?;
+ return StructArray::try_new(target_fields.clone(), children,
source.nulls().cloned())
+ .map(|array| Arc::new(array) as ArrayRef)
+ .map_err(|error| Error::DataInvalid {
+ message: format!("Failed to project nested data-evolution ROW:
{error}"),
+ source: None,
+ });
+ }
+ arrow_cast::cast(value.as_ref(), target).map_err(|error|
Error::DataInvalid {
+ message: format!("Failed to cast nested data-evolution child:
{error}"),
+ source: None,
+ })
+}
+
+fn source_struct(
+ cursors: &[Option<(RecordBatch, usize)>],
+ source: usize,
+ field: usize,
+ rows: usize,
+) -> Result<StructArray> {
+ let (batch, offset) = cursors
+ .get(source)
+ .and_then(Option::as_ref)
+ .ok_or_else(|| Error::UnexpectedError {
+ message: format!("Missing nested data-evolution source {source}"),
+ source: None,
+ })?;
+ let array = batch.column(field).slice(*offset, rows);
+ array
+ .as_any()
+ .downcast_ref::<StructArray>()
+ .cloned()
+ .ok_or_else(|| Error::DataInvalid {
+ message: format!("Nested data-evolution source {source} field
{field} is not a ROW"),
+ source: None,
+ })
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+ use crate::spec::{IntType, RowType, VarCharType};
+ use arrow_array::{Int32Array, StringArray};
+ use arrow_schema::{Field, Schema};
+
+ fn field(id: i32, name: &str, data_type: DataType) -> DataField {
+ DataField::new(id, name.to_string(), data_type)
+ }
+
+ fn name() -> DataField {
+ field(2, "name", DataType::VarChar(VarCharType::default()))
+ }
+
+ fn age() -> DataField {
+ field(3, "age", DataType::Int(IntType::new()))
+ }
+
+ fn row(children: Vec<DataField>) -> DataField {
+ field(1, "profile", DataType::Row(RowType::new(children)))
+ }
+
+ #[test]
+ fn latest_leaf_provider_wins_across_partial_files() {
+ let newer = vec![row(vec![age()])];
+ let base = vec![row(vec![name(), age()])];
+ let sources = [
+ PhysicalRowSource {
+ source_index: 0,
+ fields: &newer,
+ },
+ PhysicalRowSource {
+ source_index: 1,
+ fields: &base,
+ },
+ ];
+ let selected = select_nested_row(&row(vec![name(), age()]), &sources)
+ .unwrap()
+ .unwrap();
+ assert_eq!(selected.anchors, vec![0, 1]);
+ assert_eq!(selected.children, vec![Some(1), Some(0)]);
+ assert_eq!(selected.whole_source, None);
+ }
+
+ #[test]
+ fn omitted_sibling_still_anchors_parent_nullness() {
+ let newer = vec![row(vec![age()])];
+ let base = vec![row(vec![name(), age()])];
+ let sources = [
+ PhysicalRowSource {
+ source_index: 0,
+ fields: &newer,
+ },
+ PhysicalRowSource {
+ source_index: 1,
+ fields: &base,
+ },
+ ];
+ let selected = select_nested_row(&row(vec![name()]), &sources)
+ .unwrap()
+ .unwrap();
+ assert_eq!(selected.anchors, vec![0, 1]);
+ assert_eq!(selected.children, vec![Some(1)]);
+ assert_eq!(selected.whole_source, None);
+ }
+
+ #[test]
+ fn
provider_reads_only_selected_children_and_hidden_anchor_reads_one_sibling() {
+ let current = row(vec![name(), age()]);
+ let requested = row(vec![name()]);
+ let age_file = vec![row(vec![age()])];
+ let base_file = vec![current.clone()];
+ let sources = [
+ PhysicalRowSource {
+ source_index: 0,
+ fields: &age_file,
+ },
+ PhysicalRowSource {
+ source_index: 1,
+ fields: &base_file,
+ },
+ ];
+ let selection = select_nested_row(&requested,
&sources).unwrap().unwrap();
+
+ let anchor = source_read_field(&requested, ¤t, 0, &age_file,
&selection).unwrap();
+ let provider = source_read_field(&requested, ¤t, 1, &base_file,
&selection).unwrap();
+ let DataType::Row(anchor_row) = anchor.data_type() else {
+ panic!("anchor must be ROW")
+ };
+ let DataType::Row(provider_row) = provider.data_type() else {
+ panic!("provider must be ROW")
+ };
+ assert_eq!(anchor_row.fields(), &[age()]);
+ assert_eq!(provider_row.fields(), &[name()]);
+ }
+
+ #[test]
+ fn deep_added_leaf_reads_whole_direct_child_to_keep_its_nullness() {
+ let existing = field(4, "existing", DataType::Int(IntType::new()));
+ let added = field(5, "added", DataType::Int(IntType::new()));
+ let sub = |children| field(2, "sub",
DataType::Row(RowType::new(children)));
+ let current = row(vec![sub(vec![existing.clone(), added.clone()]),
age()]);
+ let requested = row(vec![sub(vec![added])]);
+ let physical = vec![row(vec![sub(vec![existing.clone()]), age()])];
+ let sources = [PhysicalRowSource {
+ source_index: 0,
+ fields: &physical,
+ }];
+ let selection = select_nested_row(&requested,
&sources).unwrap().unwrap();
+ let read = source_read_field(&requested, ¤t, 0, &physical,
&selection).unwrap();
+ let DataType::Row(read_row) = read.data_type() else {
+ panic!("read field must be ROW")
+ };
+ let DataType::Row(current_row) = current.data_type() else {
+ panic!("current field must be ROW")
+ };
+ assert_eq!(read_row.fields(), ¤t_row.fields()[0..1]);
+ }
+
+ #[test]
+ fn a_complete_latest_row_needs_no_composition() {
+ let newest = vec![row(vec![name(), age()])];
+ let older = vec![row(vec![name(), age()])];
+ let sources = [
+ PhysicalRowSource {
+ source_index: 0,
+ fields: &newest,
+ },
+ PhysicalRowSource {
+ source_index: 1,
+ fields: &older,
+ },
+ ];
+ let selected = select_nested_row(&row(vec![name(), age()]), &sources)
+ .unwrap()
+ .unwrap();
+ assert_eq!(selected.anchors, vec![0]);
+ assert_eq!(selected.children, vec![Some(0), Some(0)]);
+ assert_eq!(selected.whole_source, Some(0));
+ }
+
+ #[test]
+ fn rejects_deeper_cross_file_subfield_split() {
+ let street = field(4, "street",
DataType::VarChar(VarCharType::default()));
+ let zip = field(5, "zip", DataType::Int(IntType::new()));
+ let address = |children| field(2, "address",
DataType::Row(RowType::new(children)));
+ let newer = vec![row(vec![address(vec![zip.clone()])])];
+ let older = vec![row(vec![address(vec![street.clone()])])];
+ let read = row(vec![address(vec![street, zip])]);
+ let sources = [
+ PhysicalRowSource {
+ source_index: 0,
+ fields: &newer,
+ },
+ PhysicalRowSource {
+ source_index: 1,
+ fields: &older,
+ },
+ ];
+ assert!(matches!(
+ select_nested_row(&read, &sources),
+ Err(Error::Unsupported { message }) if
message.contains("profile.address")
+ ));
+ }
+
+ #[test]
+ fn rejects_missing_non_nullable_subfield() {
+ let required = field(4, "required",
DataType::Int(IntType::with_nullable(false)));
+ let available = vec![row(vec![name()])];
+ let sources = [PhysicalRowSource {
+ source_index: 0,
+ fields: &available,
+ }];
+ assert!(matches!(
+ select_nested_row(&row(vec![required]), &sources),
+ Err(Error::DataInvalid { message, .. }) if
message.contains("profile.required")
+ ));
+ }
+
+ fn struct_batch(
+ names: Vec<Option<&str>>,
+ ages: Vec<Option<i32>>,
+ valid: Vec<bool>,
+ ) -> RecordBatch {
+ let fields = vec![
+ Arc::new(Field::new("name", ArrowDataType::Utf8, true)),
+ Arc::new(Field::new("age", ArrowDataType::Int32, true)),
+ ];
+ let parent = StructArray::try_new(
+ fields.clone().into(),
+ vec![
+ Arc::new(StringArray::from(names)),
+ Arc::new(Int32Array::from(ages)),
+ ],
+ Some(NullBuffer::from(valid)),
+ )
+ .unwrap();
+ RecordBatch::try_new(
+ Arc::new(Schema::new(vec![Field::new(
+ "profile",
+ ArrowDataType::Struct(fields.into()),
+ true,
+ )])),
+ vec![Arc::new(parent)],
+ )
+ .unwrap()
+ }
+
+ #[test]
+ fn assembled_parent_is_valid_when_any_latest_sibling_is_valid() {
+ let newer = struct_batch(
+ vec![None, None, None],
+ vec![Some(20), Some(21), None],
+ vec![true, true, false],
+ );
+ let older = struct_batch(
+ vec![Some("old"), Some("hidden"), Some("kept")],
+ vec![None, None, None],
+ vec![true, false, true],
+ );
+ let target = newer.schema().field(0).data_type().clone();
+ let plan = NestedFieldPlan {
+ anchors: vec![(0, 0), (1, 0)],
+ children: vec![Some((1, 0)), Some((0, 0))],
+ };
+ let cursors = vec![Some((newer, 0)), Some((older, 0))];
+ let output = assemble_nested_row(&plan, &target, &cursors, 3).unwrap();
+ let output = output.as_any().downcast_ref::<StructArray>().unwrap();
+ assert_eq!(output.null_count(), 0);
+ let names = output
+ .column(0)
+ .as_any()
+ .downcast_ref::<StringArray>()
+ .unwrap();
+ let ages = output
+ .column(1)
+ .as_any()
+ .downcast_ref::<Int32Array>()
+ .unwrap();
+ assert_eq!(names.value(0), "old");
+ assert!(names.is_null(1)); // Older parent was NULL despite retained
child bytes.
+ assert_eq!(names.value(2), "kept");
+ assert_eq!(ages.value(0), 20);
+ assert_eq!(ages.value(1), 21);
+ assert!(ages.is_null(2));
+ }
+
+ #[test]
+ fn assembled_parent_null_when_every_source_parent_is_null() {
+ let newer = struct_batch(vec![None], vec![Some(20)], vec![false]);
+ let older = struct_batch(vec![Some("stale")], vec![None], vec![false]);
+ let target = newer.schema().field(0).data_type().clone();
+ let plan = NestedFieldPlan {
+ anchors: vec![(0, 0), (1, 0)],
+ children: vec![Some((1, 0)), Some((0, 0))],
+ };
+ let output =
+ assemble_nested_row(&plan, &target, &[Some((newer, 0)),
Some((older, 0))], 1).unwrap();
+ assert!(output.is_null(0));
+ }
+}
diff --git a/crates/paimon/src/table/data_evolution_reader.rs
b/crates/paimon/src/table/data_evolution_reader.rs
index a8da83c5..f0f435af 100644
--- a/crates/paimon/src/table/data_evolution_reader.rs
+++ b/crates/paimon/src/table/data_evolution_reader.rs
@@ -524,7 +524,7 @@ impl DataEvolutionReader {
&file_meta,
)
.await?;
- let data_fields = raw_file_physical_fields(
+ let (data_fields, file_schema_nested_enabled) =
raw_file_physical_fields(
&self.schema_manager,
self.table_schema_id,
&self.table_fields,
@@ -609,7 +609,13 @@ impl DataEvolutionReader {
let mut row_id_cursor = file_base_row_id;
let mut row_id_offset: usize = 0;
- let mut stream =
raw_file_reader.read_single_file_stream(
+ let mut stream = raw_file_reader.clone()
+ .with_nested_field_enabled(
+ CoreOptions::new(&self.table_options)
+ .data_evolution_nested_field_enabled()
+ || file_schema_nested_enabled,
+ )
+ .read_single_file_stream(
&split,
file_meta,
data_fields,
@@ -973,19 +979,24 @@ impl DataEvolutionReader {
let target_schema = build_target_arrow_schema(&read_type)?;
Ok(try_stream! {
- let file_infos = load_file_infos(
+ let (file_infos, file_schema_nested_enabled) = load_file_infos(
&schema_manager,
table_schema_id,
&table_fields,
&prepared_group.files,
)
.await?;
+ let nested_enabled = CoreOptions::new(&table_options)
+ .data_evolution_nested_field_enabled()
+ || file_schema_nested_enabled;
let source_plan = build_source_plan_with_row_id_pushdown(
&prepared_group,
&file_infos,
&read_type,
+ &table_fields,
&blob_descriptor_fields,
row_ranges.is_some(),
+ nested_enabled,
)?;
let active_source_indices: Vec<usize> = source_plan
@@ -1057,6 +1068,7 @@ impl DataEvolutionReader {
blob_parallelism,
source_parquet_read_budget.clone(),
Arc::clone(&table_options),
+ nested_enabled,
mosaic_prefetch,
read_timing.clone(),
anchor_deletion_vector.as_ref(),
@@ -1137,6 +1149,15 @@ impl DataEvolutionReader {
for (idx, provider) in
source_plan.column_plan.iter().enumerate() {
let target_field = &target_schema.fields()[idx];
+ if let Some(nested) = &source_plan.nested_plan[idx] {
+
columns.push(super::data_evolution_nested::assemble_nested_row(
+ nested,
+ target_field.data_type(),
+ &source_cursors,
+ rows_to_emit,
+ )?);
+ continue;
+ }
let array = provider
.and_then(|(source_idx, field_offset)| {
source_cursors[source_idx].as_ref().map(|(batch,
offset)| {
@@ -1180,40 +1201,23 @@ async fn raw_file_physical_fields(
table_schema_id: i64,
table_fields: &[DataField],
file: &DataFileMeta,
-) -> crate::Result<Option<Vec<DataField>>> {
- let schema_fields = if file.schema_id == table_schema_id {
- None
+) -> crate::Result<(Option<Vec<DataField>>, bool)> {
+ let (schema_fields, file_schema_nested_enabled) = if file.schema_id ==
table_schema_id {
+ (None, false)
} else {
- Some(
- schema_manager
- .schema(file.schema_id)
- .await?
- .fields()
- .to_vec(),
+ let schema = schema_manager.schema(file.schema_id).await?;
+ (
+ Some(schema.fields().to_vec()),
+ schema.core_options().data_evolution_nested_field_enabled(),
)
};
let Some(write_cols) = file.write_cols.as_ref() else {
- return Ok(schema_fields);
+ return Ok((schema_fields, file_schema_nested_enabled));
};
let fields = schema_fields.as_deref().unwrap_or(table_fields);
- let written_fields = write_cols
- .iter()
- .map(|name| {
- fields
- .iter()
- .find(|field| field.name() == name)
- .cloned()
- .ok_or_else(|| Error::DataInvalid {
- message: format!(
- "Failed to resolve write column '{}' in
raw-convertible file '{}'",
- name, file.file_name
- ),
- source: None,
- })
- })
- .collect::<crate::Result<Vec<_>>>()?;
- Ok(Some(written_fields))
+ let written_fields =
super::data_evolution_fields::project_by_paths(fields, write_cols)?;
+ Ok((Some(written_fields), file_schema_nested_enabled))
}
async fn resolve_descriptor_columns(
@@ -1615,6 +1619,7 @@ fn open_source_stream(
blob_parallelism: usize,
parquet_read_budget: Option<Arc<ReadBudget>>,
table_options: Arc<HashMap<String, String>>,
+ nested_field_enabled: bool,
mosaic_prefetch: MosaicPrefetchOptions,
read_timing: Option<Arc<DataFileReadTiming>>,
anchor_deletion_vector: Option<&DeletionVectorContext>,
@@ -1686,6 +1691,7 @@ fn open_source_stream(
.with_blob_parallelism(blob_parallelism)
.with_parquet_read_budget(parquet_read_budget)
.with_table_options(table_options)
+ .with_nested_field_enabled(nested_field_enabled)
.with_mosaic_prefetch(mosaic_prefetch)
.with_read_timing(read_timing);
@@ -2173,8 +2179,9 @@ async fn load_file_infos(
table_schema_id: i64,
table_fields: &[DataField],
files: &[DataFileMeta],
-) -> crate::Result<Vec<ResolvedFileInfo>> {
+) -> crate::Result<(Vec<ResolvedFileInfo>, bool)> {
let mut infos = Vec::with_capacity(files.len());
+ let mut nested_enabled = false;
for file in files {
let (field_ids, data_fields, effective_fields_owned);
@@ -2184,6 +2191,8 @@ async fn load_file_infos(
effective_fields_owned = None;
} else {
let data_schema = schema_manager.schema(file.schema_id).await?;
+ nested_enabled |=
+
CoreOptions::new(data_schema.options()).data_evolution_nested_field_enabled();
let fields = data_schema.fields().to_vec();
field_ids = resolve_field_ids(file, &fields)?;
data_fields = Some(fields.clone());
@@ -2207,27 +2216,17 @@ async fn load_file_infos(
});
}
- Ok(infos)
+ Ok((infos, nested_enabled))
}
fn resolve_field_ids(file: &DataFileMeta, fields: &[DataField]) ->
crate::Result<Vec<i32>> {
match &file.write_cols {
- Some(write_cols) => write_cols
- .iter()
- .map(|name| {
- fields
- .iter()
- .find(|field| field.name() == name)
- .map(|field| field.id())
- .ok_or_else(|| Error::DataInvalid {
- message: format!(
- "Failed to resolve write column '{}' in file '{}'",
- name, file.file_name
- ),
- source: None,
- })
- })
- .collect(),
+ Some(write_cols) => Ok(
+ super::data_evolution_fields::project_by_paths(fields, write_cols)?
+ .iter()
+ .map(DataField::id)
+ .collect(),
+ ),
None => Ok(fields.iter().map(|field| field.id()).collect()),
}
}
@@ -2280,6 +2279,7 @@ fn normalize_vector_write_cols(
struct SourcePlan {
sources: Vec<FieldSource>,
column_plan: Vec<Option<(usize, usize)>>,
+ nested_plan: Vec<Option<super::data_evolution_nested::NestedFieldPlan>>,
}
#[cfg(test)]
@@ -2293,8 +2293,10 @@ fn build_source_plan(
prepared_group,
file_infos,
read_type,
+ read_type,
blob_descriptor_fields,
false,
+ false,
)
}
@@ -2302,8 +2304,10 @@ fn build_source_plan_with_row_id_pushdown(
prepared_group: &PreparedMergeGroup,
file_infos: &[ResolvedFileInfo],
read_type: &[DataField],
+ table_fields: &[DataField],
blob_descriptor_fields: &HashSet<String>,
row_id_pushdown: bool,
+ nested_enabled: bool,
) -> crate::Result<SourcePlan> {
let mut sources = Vec::new();
let mut normal_providers: HashMap<i32, usize> = HashMap::new(); //
field_id -> source_idx
@@ -2311,6 +2315,7 @@ fn build_source_plan_with_row_id_pushdown(
let mut vector_bunch_indices: HashMap<(i64, String, Vec<String>), usize> =
HashMap::new();
let mut blob_source_indices: HashMap<i32, usize> = HashMap::new();
let mut expected_blob_row_count: Option<i64> = None;
+ let mut normal_sources: Vec<(usize, Vec<DataField>)> = Vec::new();
for (file_idx, file) in prepared_group.files.iter().enumerate() {
let info = &file_infos[file_idx];
@@ -2403,6 +2408,17 @@ fn build_source_plan_with_row_id_pushdown(
} else {
expected_blob_row_count = Some(file.row_count);
let source_idx = sources.len();
+ if nested_enabled {
+ let effective_fields =
info.data_fields.as_deref().unwrap_or(table_fields);
+ let physical_fields = match file.write_cols.as_deref() {
+ Some(write_cols) =>
super::data_evolution_fields::project_by_paths(
+ effective_fields,
+ write_cols,
+ )?,
+ None => effective_fields.to_vec(),
+ };
+ normal_sources.push((source_idx, physical_fields));
+ }
sources.push(FieldSource::DataFile {
file: Box::new(file.clone()),
data_fields: info.data_fields.clone(),
@@ -2416,7 +2432,65 @@ fn build_source_plan_with_row_id_pushdown(
}
let mut column_plan = Vec::with_capacity(read_type.len());
+ let mut nested_plan = Vec::with_capacity(read_type.len());
for field in read_type {
+ if nested_enabled && matches!(field.data_type(), DataType::Row(_)) {
+ let physical_sources = normal_sources
+ .iter()
+ .map(
+ |(source_index, fields)|
super::data_evolution_nested::PhysicalRowSource {
+ source_index: *source_index,
+ fields,
+ },
+ )
+ .collect::<Vec<_>>();
+ if let Some(selection) =
+ super::data_evolution_nested::select_nested_row(field,
&physical_sources)?
+ {
+ if selection.whole_source.is_none() {
+ let current_field = table_fields
+ .iter()
+ .find(|candidate| candidate.id() == field.id())
+ .unwrap_or(field);
+ let mut offsets = HashMap::new();
+ for source in &selection.anchors {
+ let physical_fields = normal_sources
+ .iter()
+ .find(|(index, _)| index == source)
+ .map(|(_, fields)| fields.as_slice())
+ .ok_or_else(|| Error::DataInvalid {
+ message: format!(
+ "Missing physical nested source {source}
for '{}'",
+ field.name()
+ ),
+ source: None,
+ })?;
+ let source_field =
super::data_evolution_nested::source_read_field(
+ field,
+ current_field,
+ *source,
+ physical_fields,
+ &selection,
+ )?;
+ offsets.insert(*source,
sources[*source].add_read_field(source_field));
+ }
+
nested_plan.push(Some(super::data_evolution_nested::NestedFieldPlan {
+ anchors: selection
+ .anchors
+ .iter()
+ .map(|source| (*source, offsets[source]))
+ .collect(),
+ children: selection
+ .children
+ .iter()
+ .map(|source| source.map(|source| (source,
offsets[&source])))
+ .collect(),
+ }));
+ column_plan.push(None);
+ continue;
+ }
+ }
+ }
let source_idx = if field.data_type().is_blob_file_field()
&& !blob_descriptor_fields.contains(field.name())
{
@@ -2436,6 +2510,7 @@ fn build_source_plan_with_row_id_pushdown(
if let Some(source_idx) = source_idx {
let field_offset =
sources[source_idx].add_read_field(field.clone());
column_plan.push(Some((source_idx, field_offset)));
+ nested_plan.push(None);
} else if !field.data_type().is_nullable() {
return Err(Error::DataInvalid {
message: format!(
@@ -2446,6 +2521,7 @@ fn build_source_plan_with_row_id_pushdown(
});
} else {
column_plan.push(None);
+ nested_plan.push(None);
}
}
@@ -2485,6 +2561,7 @@ fn build_source_plan_with_row_id_pushdown(
Ok(SourcePlan {
sources,
column_plan,
+ nested_plan,
})
}
@@ -3986,8 +4063,10 @@ mod tests {
&prepared_group,
&file_infos,
&read_type,
+ &read_type,
&HashSet::new(),
true,
+ false,
)
.unwrap_err();
diff --git a/crates/paimon/src/table/data_evolution_writer.rs
b/crates/paimon/src/table/data_evolution_writer.rs
index 98f60f41..842b48cd 100644
--- a/crates/paimon/src/table/data_evolution_writer.rs
+++ b/crates/paimon/src/table/data_evolution_writer.rs
@@ -41,7 +41,9 @@ use crate::table::stats_filter::group_by_overlapping_row_id;
use crate::table::DataSplitBuilder;
use crate::table::Table;
use crate::Result;
-use arrow_array::{Array, ArrayRef, Int64Array, RecordBatch};
+use arrow_array::{Array, ArrayRef, Int64Array, RecordBatch, StructArray};
+use arrow_buffer::NullBuffer;
+use arrow_schema::DataType as ArrowDataType;
use arrow_select::concat::concat_batches;
use arrow_select::interleave::interleave;
use bytes::Bytes;
@@ -74,6 +76,7 @@ const MANIFEST_DIR: &str = "manifest";
pub struct DataEvolutionWriter {
table: Table,
update_columns: Vec<String>,
+ write_fields: Vec<DataField>,
matched_batches: Vec<RecordBatch>,
}
@@ -108,16 +111,35 @@ impl DataEvolutionWriter {
});
}
+ super::data_evolution_fields::validate_write_paths(
+ schema.fields(),
+ &update_columns,
+ core_options.data_evolution_nested_field_enabled(),
+ )?;
+ let write_fields =
+ super::data_evolution_fields::project_by_paths(schema.fields(),
&update_columns)?;
+ if !core_options.data_evolution_nested_field_enabled()
+ && update_columns
+ .iter()
+ .any(|column| !schema.fields().iter().any(|field| field.name()
== column))
+ {
+ return Err(crate::Error::DataInvalid {
+ message: "Nested data-evolution write paths require
data-evolution.nested-field.enabled=true"
+ .to_string(),
+ source: None,
+ });
+ }
let partition_keys = schema.partition_keys();
let blob_descriptor_fields = core_options.blob_descriptor_fields();
for col in &update_columns {
- if partition_keys.contains(col) {
+ let top_level =
DataEvolutionPartialWriter::top_level_write_name(col, schema.fields());
+ if partition_keys.iter().any(|key| key == top_level) {
return Err(crate::Error::Unsupported {
message: format!("Cannot update partition column '{col}'
in MERGE INTO"),
});
}
- if let Some(field) = schema.fields().iter().find(|f| f.name() ==
col) {
- if field.data_type().is_blob_type() &&
!blob_descriptor_fields.contains(col) {
+ if let Some(field) = schema.fields().iter().find(|f| f.name() ==
top_level) {
+ if field.data_type().is_blob_type() &&
!blob_descriptor_fields.contains(top_level) {
return Err(crate::Error::Unsupported {
message: format!(
"Cannot update raw-data BLOB column '{col}' in
MERGE INTO. \
@@ -131,6 +153,7 @@ impl DataEvolutionWriter {
Ok(Self {
table: table.clone(),
update_columns,
+ write_fields,
matched_batches: Vec::new(),
})
}
@@ -236,7 +259,7 @@ impl DataEvolutionWriter {
let row_count = file_range.row_count as usize;
// Read original columns from the entire file group (base +
partial-column files).
- let col_refs: Vec<&str> = self.update_columns.iter().map(|s|
s.as_str()).collect();
+ let col_refs: Vec<&str> =
self.write_fields.iter().map(DataField::name).collect();
let mut rb = self.table.new_read_builder();
rb.with_projection(&col_refs)?;
let read = rb.new_read()?;
@@ -280,97 +303,41 @@ impl DataEvolutionWriter {
});
}
- // Apply updates using 2-array interleave: [original_col,
updates_col].
- // Matched rows are gathered into a single contiguous update array
first,
- // avoiding O(N) array clones for every row in the file.
- let mut new_columns: Vec<ArrayRef> =
Vec::with_capacity(self.update_columns.len());
-
- // Sort matched rows by offset for contiguous iteration
+ // Keep one physical ROW column per top-level field, even when the
+ // update names identify several nested leaves of that ROW.
let mut sorted_matches: Vec<(usize, usize, usize)> = matched_rows
.iter()
.map(|m| (m.offset, m.batch_idx, m.row_idx))
.collect();
sorted_matches.sort_by_key(|(offset, _, _)| *offset);
-
- for (col_idx, col_name) in self.update_columns.iter().enumerate() {
- let original_col = original_batch.column(col_idx);
- let original_dtype = original_col.data_type();
-
- // Gather update values into a single array (one entry per
matched row, in offset order)
- let update_indices: Vec<(usize, usize)> = sorted_matches
- .iter()
- .map(|&(_, batch_idx, row_idx)| (batch_idx, row_idx))
- .collect();
-
- // Collect unique batch arrays, cast if needed
- let mut batch_arrays: Vec<ArrayRef> = Vec::new();
- let mut batch_id_map: HashMap<usize, usize> = HashMap::new();
- let mut interleave_src: Vec<(usize, usize)> =
- Vec::with_capacity(update_indices.len());
-
- for &(batch_idx, row_idx) in &update_indices {
- let arr_idx = match batch_id_map.get(&batch_idx) {
- Some(&idx) => idx,
- None => {
- let src_col =
-
matched_column(&self.matched_batches[batch_idx], col_name)?;
- let casted = if src_col.data_type() !=
original_dtype {
- arrow_cast::cast(src_col.as_ref(),
original_dtype).map_err(|e| {
- crate::Error::DataInvalid {
- message: format!("Failed to cast
column {col_name}: {e}"),
- source: None,
- }
- })?
- } else {
- src_col
- };
- let idx = batch_arrays.len();
- batch_arrays.push(casted);
- batch_id_map.insert(batch_idx, idx);
- idx
- }
- };
- interleave_src.push((arr_idx, row_idx));
- }
-
- let update_col = if batch_arrays.len() == 1 &&
interleave_src.len() == 1 {
- // Single update value — just slice
- let (_, row_idx) = interleave_src[0];
- batch_arrays[0].slice(row_idx, 1)
- } else {
- let refs: Vec<&dyn Array> = batch_arrays.iter().map(|a|
a.as_ref()).collect();
- interleave(&refs, &interleave_src).map_err(|e|
crate::Error::DataInvalid {
- message: format!("Failed to gather update values for
{col_name}: {e}"),
- source: None,
- })?
- };
-
- // Build final indices: 2 sources — [0] = original, [1] =
update_col
- let mut indices: Vec<(usize, usize)> =
Vec::with_capacity(row_count);
- let mut match_pos = 0;
- for row in 0..row_count {
- if match_pos < sorted_matches.len() &&
sorted_matches[match_pos].0 == row {
- indices.push((1, match_pos));
- match_pos += 1;
- } else {
- indices.push((0, row));
- }
- }
-
- let sources: [&dyn Array; 2] = [original_col.as_ref(),
update_col.as_ref()];
- let new_col =
- interleave(&sources, &indices).map_err(|e|
crate::Error::DataInvalid {
- message: format!("Failed to interleave column
{col_name}: {e}"),
- source: None,
- })?;
- new_columns.push(new_col);
+ if sorted_matches.windows(2).any(|rows| rows[0].0 == rows[1].0) {
+ return Err(crate::Error::DataInvalid {
+ message: "Multiple updates target the same row
ID".to_string(),
+ source: None,
+ });
}
-
- let updated_batch = RecordBatch::try_new(original_batch.schema(),
new_columns)
- .map_err(|e| crate::Error::DataInvalid {
+ let new_columns = self
+ .write_fields
+ .iter()
+ .enumerate()
+ .map(|(index, field)| {
+ apply_field_updates(
+ field,
+ field.name(),
+ original_batch.column(index),
+ &self.update_columns,
+ &self.matched_batches,
+ &sorted_matches,
+ )
+ })
+ .collect::<Result<Vec<_>>>()?;
+ let write_schema =
crate::arrow::build_target_arrow_schema(&self.write_fields)?;
+ let updated_batch = RecordBatch::try_new(write_schema,
new_columns).map_err(|e| {
+ crate::Error::DataInvalid {
message: format!("Failed to create updated batch: {e}"),
source: None,
- })?;
+ }
+ })?;
writer
.write_partial_batch(
@@ -930,6 +897,138 @@ fn matched_column(batch: &RecordBatch, col: &str) ->
Result<ArrayRef> {
Ok(batch.column(idx).clone())
}
+/// Apply matched values to one projected field. The projection may be a ROW
+/// containing only the leaves named in `update_columns`; its untouched sibling
+/// leaves remain available from the read batch but are absent from this file.
+fn apply_field_updates(
+ write_field: &DataField,
+ path: &str,
+ original: &ArrayRef,
+ update_columns: &[String],
+ matched_batches: &[RecordBatch],
+ matches: &[(usize, usize, usize)],
+) -> Result<ArrayRef> {
+ if update_columns.iter().any(|column| column == path) {
+ return interleave_updated_leaf(path, original, matched_batches,
matches);
+ }
+
+ let DataType::Row(row) = write_field.data_type() else {
+ return Err(crate::Error::DataInvalid {
+ message: format!("No update value was supplied for '{path}'"),
+ source: None,
+ });
+ };
+ let source = original
+ .as_any()
+ .downcast_ref::<StructArray>()
+ .ok_or_else(|| crate::Error::DataInvalid {
+ message: format!("Expected ROW values for nested update '{path}'"),
+ source: None,
+ })?;
+ let children =
+ row.fields()
+ .iter()
+ .map(|child| {
+ let child_path = format!("{path}.{}", child.name());
+ let source_child =
source.column_by_name(child.name()).ok_or_else(|| {
+ crate::Error::DataInvalid {
+ message: format!("Nested update source is missing
'{child_path}'"),
+ source: None,
+ }
+ })?;
+ apply_field_updates(
+ child,
+ &child_path,
+ source_child,
+ update_columns,
+ matched_batches,
+ matches,
+ )
+ })
+ .collect::<Result<Vec<_>>>()?;
+ let ArrowDataType::Struct(fields) =
+ crate::arrow::paimon_type_to_arrow(write_field.data_type())?
+ else {
+ unreachable!("ROW type must map to Arrow Struct")
+ };
+ let mut valid = (0..original.len())
+ .map(|row| source.is_valid(row))
+ .collect::<Vec<_>>();
+ for &(offset, _, _) in matches {
+ valid[offset] = true;
+ }
+ let struct_array = StructArray::try_new(fields, children,
Some(NullBuffer::from(valid)))
+ .map_err(|error| crate::Error::DataInvalid {
+ message: format!("Failed to build nested update '{path}':
{error}"),
+ source: None,
+ })?;
+ Ok(Arc::new(struct_array))
+}
+
+/// Gather matched rows into one array, then use a two-array interleave over
+/// the file's physical row positions. Casts follow the existing top-level
+/// MERGE update behavior.
+fn interleave_updated_leaf(
+ path: &str,
+ original: &ArrayRef,
+ matched_batches: &[RecordBatch],
+ matches: &[(usize, usize, usize)],
+) -> Result<ArrayRef> {
+ let mut batch_arrays = Vec::<ArrayRef>::new();
+ let mut batch_id_map = HashMap::<usize, usize>::new();
+ let mut gather = Vec::with_capacity(matches.len());
+ for &(_, batch_idx, row_idx) in matches {
+ let array_idx = if let Some(index) = batch_id_map.get(&batch_idx) {
+ *index
+ } else {
+ let value = matched_column(&matched_batches[batch_idx], path)?;
+ let value = if value.data_type() == original.data_type() {
+ value
+ } else {
+ arrow_cast::cast(value.as_ref(),
original.data_type()).map_err(|error| {
+ crate::Error::DataInvalid {
+ message: format!("Failed to cast column {path}:
{error}"),
+ source: None,
+ }
+ })?
+ };
+ let index = batch_arrays.len();
+ batch_arrays.push(value);
+ batch_id_map.insert(batch_idx, index);
+ index
+ };
+ gather.push((array_idx, row_idx));
+ }
+ let updates = if gather.len() == 1 {
+ batch_arrays[gather[0].0].slice(gather[0].1, 1)
+ } else {
+ let arrays = batch_arrays
+ .iter()
+ .map(|array| array.as_ref())
+ .collect::<Vec<_>>();
+ interleave(&arrays, &gather).map_err(|error| crate::Error::DataInvalid
{
+ message: format!("Failed to gather update values for {path}:
{error}"),
+ source: None,
+ })?
+ };
+ let mut indices = Vec::with_capacity(original.len());
+ let mut matched = 0;
+ for row in 0..original.len() {
+ if matched < matches.len() && matches[matched].0 == row {
+ indices.push((1, matched));
+ matched += 1;
+ } else {
+ indices.push((0, row));
+ }
+ }
+ interleave(&[original.as_ref(), updates.as_ref()],
&indices).map_err(|error| {
+ crate::Error::DataInvalid {
+ message: format!("Failed to interleave column {path}: {error}"),
+ source: None,
+ }
+ })
+}
+
struct FileRowRange {
first_row_id: i64,
last_row_id: i64,
@@ -1064,34 +1163,57 @@ impl DataEvolutionPartialWriter {
core_options: &CoreOptions<'_>,
) -> Result<Vec<PartialWriteSet>> {
let vector_file_format = core_options.vector_file_format();
+ if !core_options.data_evolution_nested_field_enabled()
+ && write_columns
+ .iter()
+ .any(|column| !fields.iter().any(|field| field.name() ==
column))
+ {
+ return Err(crate::Error::DataInvalid {
+ message: "Nested data-evolution write paths require
data-evolution.nested-field.enabled=true"
+ .to_string(),
+ source: None,
+ });
+ }
+ super::data_evolution_fields::validate_write_paths(
+ fields,
+ write_columns,
+ core_options.data_evolution_nested_field_enabled(),
+ )?;
+ let projected = super::data_evolution_fields::project_by_paths(fields,
write_columns)?;
let mut normal_fields = Vec::new();
- let mut normal_columns = Vec::new();
let mut normal_indices = Vec::new();
let mut vector_fields = Vec::new();
- let mut vector_columns = Vec::new();
let mut vector_indices = Vec::new();
- for (idx, column) in write_columns.iter().enumerate() {
- let field = fields
- .iter()
- .find(|field| field.name() == column)
- .cloned()
- .ok_or_else(|| crate::Error::DataInvalid {
- message: format!("Unknown data-evolution write column
'{column}'"),
- source: None,
- })?;
-
+ for (idx, field) in projected.into_iter().enumerate() {
if vector_file_format.is_some() && matches!(field.data_type(),
DataType::Vector(_)) {
vector_fields.push(field);
- vector_columns.push(column.clone());
vector_indices.push(idx);
} else {
normal_fields.push(field);
- normal_columns.push(column.clone());
normal_indices.push(idx);
}
}
+ let normal_columns = write_columns
+ .iter()
+ .filter(|path| {
+ normal_fields
+ .iter()
+ .any(|field| field.name() ==
Self::top_level_write_name(path, fields))
+ })
+ .cloned()
+ .collect::<Vec<_>>();
+ let vector_columns = write_columns
+ .iter()
+ .filter(|path| {
+ vector_fields
+ .iter()
+ .any(|field| field.name() ==
Self::top_level_write_name(path, fields))
+ })
+ .cloned()
+ .collect::<Vec<_>>();
+
let mut write_sets = Vec::new();
if !normal_fields.is_empty() {
let schema =
crate::arrow::build_target_arrow_schema(&normal_fields)?;
@@ -1125,6 +1247,14 @@ impl DataEvolutionPartialWriter {
Ok(write_sets)
}
+ fn top_level_write_name<'a>(path: &'a str, fields: &[DataField]) -> &'a
str {
+ if fields.iter().any(|field| field.name() == path) {
+ path
+ } else {
+ path.split_once('.').map_or(path, |(head, _)| head)
+ }
+ }
+
/// Write a partial-column batch for a specific partition, bucket, and row
ID range.
///
/// The `batch` must contain only the columns specified in `write_columns`.
@@ -1259,8 +1389,12 @@ mod tests {
use super::*;
use crate::catalog::Identifier;
use crate::io::FileIOBuilder;
- use crate::spec::{DataType, FloatType, IntType, Schema, TableSchema,
VarCharType, VectorType};
- use arrow_array::StringArray;
+ use crate::spec::{
+ DataField, DataType, FloatType, IntType, RowType, Schema, TableSchema,
VarCharType,
+ VectorType,
+ };
+ use arrow_array::{Int32Array, StringArray, StructArray};
+ use arrow_buffer::NullBuffer;
use arrow_schema::{DataType as ArrowDataType, Field as ArrowField, Schema
as ArrowSchema};
use std::sync::Arc;
@@ -1330,6 +1464,575 @@ mod tests {
TableSchema::new(0, &schema)
}
+ fn test_nested_data_evolution_schema() -> TableSchema {
+ let profile = DataType::Row(RowType::new(vec![
+ DataField::new(
+ 1,
+ "name".to_string(),
+ DataType::VarChar(VarCharType::string_type()),
+ ),
+ DataField::new(2, "age".to_string(),
DataType::Int(IntType::new())),
+ ]));
+ let schema = Schema::builder()
+ .column("id", DataType::Int(IntType::new()))
+ .column("profile", profile)
+ .option("bucket", "-1")
+ .option("data-evolution.enabled", "true")
+ .option("data-evolution.nested-field.enabled", "true")
+ .option("row-tracking.enabled", "true")
+ .build()
+ .unwrap();
+ TableSchema::new(0, &schema)
+ }
+
+ #[test]
+ fn nested_update_rejects_ambiguous_top_level_name() {
+ let profile = DataType::Row(RowType::new(vec![DataField::new(
+ 0,
+ "age".to_string(),
+ DataType::Int(IntType::new()),
+ )]));
+ let schema = Schema::builder()
+ .column("profile", profile)
+ .column("profile.age", DataType::Int(IntType::new()))
+ .option("bucket", "-1")
+ .option("data-evolution.enabled", "true")
+ .option("data-evolution.nested-field.enabled", "true")
+ .option("row-tracking.enabled", "true")
+ .build()
+ .unwrap();
+ let table = Table::new(
+ test_file_io(),
+ Identifier::new("default", "ambiguous_nested_update"),
+ "memory:/ambiguous_nested_update".to_string(),
+ TableSchema::new(0, &schema),
+ None,
+ );
+
+ let error = DataEvolutionWriter::new(&table,
vec!["profile.age".into()])
+ .err()
+ .expect("ambiguous write path must fail before writing");
+ assert!(error
+ .to_string()
+ .contains("Ambiguous data-evolution write path"));
+ }
+
+ fn nested_profile_batch(
+ fields: &[DataField],
+ columns: Vec<ArrayRef>,
+ valid: Vec<bool>,
+ ) -> RecordBatch {
+ let schema = crate::arrow::build_target_arrow_schema(fields).unwrap();
+ let ArrowDataType::Struct(children) = schema.field(0).data_type() else
{
+ panic!("expected profile ROW")
+ };
+ let profile =
+ StructArray::try_new(children.clone(), columns,
Some(NullBuffer::from(valid))).unwrap();
+ RecordBatch::try_new(schema, vec![Arc::new(profile)]).unwrap()
+ }
+
+ fn full_nested_batch(table: &Table) -> RecordBatch {
+ let profile = nested_profile_batch(
+ &table.schema().fields()[1..2],
+ vec![
+ Arc::new(StringArray::from(vec![Some("alice"), Some("bob"),
None])),
+ Arc::new(Int32Array::from(vec![Some(10), Some(20), None])),
+ ],
+ vec![true, true, false],
+ );
+ let schema =
crate::arrow::build_target_arrow_schema(table.schema().fields()).unwrap();
+ RecordBatch::try_new(
+ schema,
+ vec![
+ Arc::new(Int32Array::from(vec![1, 2, 3])),
+ profile.column(0).clone(),
+ ],
+ )
+ .unwrap()
+ }
+
+ async fn commit_full_nested_batch(table: &Table) {
+ let builder = table.new_write_builder();
+ let mut writer = builder.new_write().unwrap();
+ writer
+ .write_arrow_batch(&full_nested_batch(table))
+ .await
+ .unwrap();
+ let messages = writer.prepare_commit().await.unwrap();
+ builder.new_commit().commit(messages).await.unwrap();
+ }
+
+ async fn commit_partial_nested_batch(
+ table: &Table,
+ path: &str,
+ values: ArrayRef,
+ valid: Vec<bool>,
+ check_from_snapshot: i64,
+ ) -> DataFileMeta {
+ let fields = super::super::data_evolution_fields::project_by_paths(
+ table.schema().fields(),
+ &[path.to_string()],
+ )
+ .unwrap();
+ let batch = nested_profile_batch(&fields, vec![values], valid);
+ let mut writer = DataEvolutionPartialWriter::new(table,
vec![path.to_string()]).unwrap();
+ writer
+ .write_partial_batch(
+ crate::spec::EMPTY_SERIALIZED_ROW.clone(),
+ 0,
+ 0,
+ check_from_snapshot,
+ batch,
+ )
+ .await
+ .unwrap();
+ let messages = writer.prepare_commit().await.unwrap();
+ let file = messages[0].new_files[0].clone();
+ assert_eq!(file.write_cols, Some(vec![path.to_string()]));
+ table
+ .new_write_builder()
+ .new_commit()
+ .commit(messages)
+ .await
+ .unwrap();
+ file
+ }
+
+ #[tokio::test]
+ async fn nested_partial_parquet_files_merge_latest_children() {
+ let file_io = test_file_io();
+ let path = "memory:/test_de_nested_partial";
+ setup_dirs(&file_io, path).await;
+ let table = Table::new(
+ file_io,
+ Identifier::new("default", "test_de_nested_partial"),
+ path.to_string(),
+ test_nested_data_evolution_schema(),
+ None,
+ );
+ commit_full_nested_batch(&table).await;
+ let age_file = commit_partial_nested_batch(
+ &table,
+ "profile.age",
+ Arc::new(Int32Array::from(vec![Some(11), Some(21), None])),
+ vec![true, true, false],
+ 1,
+ )
+ .await;
+ assert_eq!(age_file.first_row_id, Some(0));
+ commit_partial_nested_batch(
+ &table,
+ "profile.name",
+ Arc::new(StringArray::from(vec![Some("ALICE"), Some("BOB"),
None])),
+ vec![true, true, false],
+ 2,
+ )
+ .await;
+
+ let plan = table.new_read_builder().new_scan().plan().await.unwrap();
+ let batches = table
+ .new_read_builder()
+ .new_read()
+ .unwrap()
+ .to_arrow(plan.splits())
+ .unwrap()
+ .try_collect::<Vec<RecordBatch>>()
+ .await
+ .unwrap();
+ let mut actual = Vec::new();
+ for batch in batches {
+ let ids = batch
+ .column(0)
+ .as_any()
+ .downcast_ref::<Int32Array>()
+ .unwrap();
+ let profiles = batch
+ .column(1)
+ .as_any()
+ .downcast_ref::<StructArray>()
+ .unwrap();
+ let names = profiles
+ .column(0)
+ .as_any()
+ .downcast_ref::<StringArray>()
+ .unwrap();
+ let ages = profiles
+ .column(1)
+ .as_any()
+ .downcast_ref::<Int32Array>()
+ .unwrap();
+ for row in 0..batch.num_rows() {
+ actual.push((
+ ids.value(row),
+ profiles.is_valid(row).then(||
names.value(row).to_string()),
+ profiles.is_valid(row).then(|| ages.value(row)),
+ ));
+ }
+ }
+ actual.sort_by_key(|row| row.0);
+ assert_eq!(
+ actual,
+ vec![
+ (1, Some("ALICE".to_string()), Some(11)),
+ (2, Some("BOB".to_string()), Some(21)),
+ (3, None, None),
+ ]
+ );
+ }
+
+ #[tokio::test]
+ async fn row_id_update_writes_two_nested_leaves_and_revives_null_parent() {
+ let file_io = test_file_io();
+ let path = "memory:/test_de_nested_row_id_update";
+ setup_dirs(&file_io, path).await;
+ let table = Table::new(
+ file_io,
+ Identifier::new("default", "test_de_nested_row_id_update"),
+ path.to_string(),
+ test_nested_data_evolution_schema(),
+ None,
+ );
+ commit_full_nested_batch(&table).await;
+
+ let matched = RecordBatch::try_new(
+ Arc::new(ArrowSchema::new(vec![
+ ArrowField::new("_ROW_ID", ArrowDataType::Int64, false),
+ ArrowField::new("profile.age", ArrowDataType::Int32, true),
+ ArrowField::new("profile.name", ArrowDataType::Utf8, true),
+ ])),
+ vec![
+ Arc::new(Int64Array::from(vec![0, 2])),
+ Arc::new(Int32Array::from(vec![Some(99), Some(33)])),
+ Arc::new(StringArray::from(vec![Some("ALICE"), Some("new")])),
+ ],
+ )
+ .unwrap();
+ let mut writer =
+ DataEvolutionWriter::new(&table, vec!["profile.age".into(),
"profile.name".into()])
+ .unwrap();
+ writer.add_matched_batch(matched).unwrap();
+ let messages = writer.prepare_commit().await.unwrap();
+ assert_eq!(
+ messages[0].new_files[0].write_cols,
+ Some(vec!["profile.age".into(), "profile.name".into()])
+ );
+ table
+ .new_write_builder()
+ .new_commit()
+ .commit(messages)
+ .await
+ .unwrap();
+
+ let plan = table.new_read_builder().new_scan().plan().await.unwrap();
+ let batches = table
+ .new_read_builder()
+ .new_read()
+ .unwrap()
+ .to_arrow(plan.splits())
+ .unwrap()
+ .try_collect::<Vec<RecordBatch>>()
+ .await
+ .unwrap();
+ let mut actual = Vec::new();
+ for batch in batches {
+ let ids = batch
+ .column(0)
+ .as_any()
+ .downcast_ref::<Int32Array>()
+ .unwrap();
+ let profiles = batch
+ .column(1)
+ .as_any()
+ .downcast_ref::<StructArray>()
+ .unwrap();
+ let names = profiles
+ .column(0)
+ .as_any()
+ .downcast_ref::<StringArray>()
+ .unwrap();
+ let ages = profiles
+ .column(1)
+ .as_any()
+ .downcast_ref::<Int32Array>()
+ .unwrap();
+ for row in 0..batch.num_rows() {
+ actual.push((
+ ids.value(row),
+ names.value(row).to_string(),
+ ages.value(row),
+ ));
+ assert!(profiles.is_valid(row));
+ }
+ }
+ actual.sort_by_key(|row| row.0);
+ assert_eq!(
+ actual,
+ vec![
+ (1, "ALICE".into(), 99),
+ (2, "bob".into(), 20),
+ (3, "new".into(), 33),
+ ]
+ );
+ }
+
+ #[tokio::test]
+ async fn
nested_projection_uses_sibling_nullness_and_masks_null_source_parent() {
+ let file_io = test_file_io();
+ let path = "memory:/test_de_nested_projected_nullness";
+ setup_dirs(&file_io, path).await;
+ let table = Table::new(
+ file_io,
+ Identifier::new("default", "test_de_nested_projected_nullness"),
+ path.to_string(),
+ test_nested_data_evolution_schema(),
+ None,
+ );
+ commit_full_nested_batch(&table).await;
+ // Row 0's newest age file has a NULL parent. Its older age must not
+ // reappear, although the older name keeps the composed parent valid.
+ // Row 2 has the inverse: only the new age parent is valid.
+ commit_partial_nested_batch(
+ &table,
+ "profile.age",
+ Arc::new(Int32Array::from(vec![Some(99), Some(21), Some(30)])),
+ vec![false, true, true],
+ 1,
+ )
+ .await;
+
+ for (path, expected_parent, expected_values) in [
+ (
+ "profile.name",
+ vec![true, true, true],
+ vec![Some("alice"), Some("bob"), None],
+ ),
+ (
+ "profile.age",
+ vec![true, true, true],
+ vec![None, Some("21"), Some("30")],
+ ),
+ ] {
+ let projection =
super::super::data_evolution_fields::project_by_paths(
+ table.schema().fields(),
+ &[path.to_string()],
+ )
+ .unwrap();
+ let mut builder = table.new_read_builder();
+ builder.with_read_type(projection);
+ let plan = builder.new_scan().plan().await.unwrap();
+ let batches = builder
+ .new_read()
+ .unwrap()
+ .to_arrow(plan.splits())
+ .unwrap()
+ .try_collect::<Vec<RecordBatch>>()
+ .await
+ .unwrap();
+ let mut actual = Vec::new();
+ for batch in batches {
+ let parent = batch
+ .column(0)
+ .as_any()
+ .downcast_ref::<StructArray>()
+ .unwrap();
+ for row in 0..batch.num_rows() {
+ let value = if parent.column(0).is_null(row) {
+ None
+ } else if path.ends_with("name") {
+ Some(
+ parent
+ .column(0)
+ .as_any()
+ .downcast_ref::<StringArray>()
+ .unwrap()
+ .value(row)
+ .to_string(),
+ )
+ } else {
+ Some(
+ parent
+ .column(0)
+ .as_any()
+ .downcast_ref::<Int32Array>()
+ .unwrap()
+ .value(row)
+ .to_string(),
+ )
+ };
+ actual.push((parent.is_valid(row), value));
+ }
+ }
+ assert_eq!(
+ actual,
+ expected_parent
+ .into_iter()
+ .zip(
+ expected_values
+ .into_iter()
+ .map(|value| value.map(str::to_string))
+ )
+ .collect::<Vec<_>>()
+ );
+ }
+ }
+
+ #[tokio::test]
+ async fn projected_new_deep_leaf_preserves_both_parent_null_buffers() {
+ let file_io = test_file_io();
+ let path = "memory:/test_de_new_deep_leaf";
+ setup_dirs(&file_io, path).await;
+ let old_schema = Schema::builder()
+ .column("id", DataType::Int(IntType::new()))
+ .column(
+ "profile",
+ DataType::Row(RowType::new(vec![
+ DataField::new(
+ 0,
+ "sub".into(),
+ DataType::Row(RowType::new(vec![DataField::new(
+ 0,
+ "existing".into(),
+ DataType::Int(IntType::new()),
+ )])),
+ ),
+ DataField::new(
+ 0,
+ "other".into(),
+ DataType::VarChar(VarCharType::string_type()),
+ ),
+ ])),
+ )
+ .option("bucket", "-1")
+ .option("data-evolution.enabled", "true")
+ .option("data-evolution.nested-field.enabled", "true")
+ .option("row-tracking.enabled", "true")
+ .build()
+ .unwrap();
+ let old_table_schema = TableSchema::new(0, &old_schema);
+ let old_table = Table::new(
+ file_io.clone(),
+ Identifier::new("default", "test_de_new_deep_leaf"),
+ path.to_string(),
+ old_table_schema.clone(),
+ None,
+ );
+ let arrow_schema =
+
crate::arrow::build_target_arrow_schema(old_table.schema().fields()).unwrap();
+ let ArrowDataType::Struct(profile_fields) =
arrow_schema.field(1).data_type() else {
+ panic!("profile must be ROW")
+ };
+ let ArrowDataType::Struct(sub_fields) = profile_fields[0].data_type()
else {
+ panic!("sub must be ROW")
+ };
+ let sub = StructArray::try_new(
+ sub_fields.clone(),
+ vec![Arc::new(Int32Array::from(vec![Some(10), None]))],
+ Some(NullBuffer::from(vec![true, false])),
+ )
+ .unwrap();
+ let profile = StructArray::try_new(
+ profile_fields.clone(),
+ vec![
+ Arc::new(sub),
+ Arc::new(StringArray::from(vec![Some("one"), Some("two")])),
+ ],
+ Some(NullBuffer::from(vec![true, true])),
+ )
+ .unwrap();
+ let old_batch = RecordBatch::try_new(
+ arrow_schema,
+ vec![Arc::new(Int32Array::from(vec![1, 2])), Arc::new(profile)],
+ )
+ .unwrap();
+ let builder = old_table.new_write_builder();
+ let mut writer = builder.new_write().unwrap();
+ writer.write_arrow_batch(&old_batch).await.unwrap();
+ builder
+ .new_commit()
+ .commit(writer.prepare_commit().await.unwrap())
+ .await
+ .unwrap();
+
+ let schema_path = old_table.schema_manager().schema_path(0);
+ file_io
+ .mkdirs(schema_path.rsplit_once('/').unwrap().0)
+ .await
+ .unwrap();
+ file_io
+ .new_output(&schema_path)
+ .unwrap()
+ .write(Bytes::from(serde_json::to_vec(&old_table_schema).unwrap()))
+ .await
+ .unwrap();
+
+ let mut new_fields = old_schema.fields().to_vec();
+ let DataType::Row(profile_type) = new_fields[1].data_type() else {
+ panic!("profile must be ROW")
+ };
+ let mut profile_children = profile_type.fields().to_vec();
+ let DataType::Row(sub_type) = profile_children[0].data_type() else {
+ panic!("sub must be ROW")
+ };
+ let mut sub_children = sub_type.fields().to_vec();
+ sub_children.push(DataField::new(
+ old_table_schema.highest_field_id() + 1,
+ "added".into(),
+ DataType::Int(IntType::new()),
+ ));
+ profile_children[0] =
super::super::data_evolution_fields::field_with_type(
+ &profile_children[0],
+ DataType::Row(RowType::with_nullable(true, sub_children)),
+ );
+ new_fields[1] = super::super::data_evolution_fields::field_with_type(
+ &new_fields[1],
+ DataType::Row(RowType::with_nullable(true, profile_children)),
+ );
+ let new_schema = old_schema
+ .copy(RowType::with_nullable(false, new_fields))
+ .unwrap();
+ let current = Table::new(
+ file_io,
+ Identifier::new("default", "test_de_new_deep_leaf"),
+ path.to_string(),
+ TableSchema::new(1, &new_schema),
+ None,
+ );
+ let projection = super::super::data_evolution_fields::project_by_paths(
+ current.schema().fields(),
+ &["profile.sub.added".to_string()],
+ )
+ .unwrap();
+ let mut read_builder = current.new_read_builder();
+ read_builder.with_read_type(projection);
+ let plan = read_builder.new_scan().plan().await.unwrap();
+ let batches = read_builder
+ .new_read()
+ .unwrap()
+ .to_arrow(plan.splits())
+ .unwrap()
+ .try_collect::<Vec<RecordBatch>>()
+ .await
+ .unwrap();
+ assert_eq!(batches.iter().map(RecordBatch::num_rows).sum::<usize>(),
2);
+ let profile = batches[0]
+ .column(0)
+ .as_any()
+ .downcast_ref::<StructArray>()
+ .unwrap();
+ let sub = profile
+ .column(0)
+ .as_any()
+ .downcast_ref::<StructArray>()
+ .unwrap();
+ let added =
sub.column(0).as_any().downcast_ref::<Int32Array>().unwrap();
+ assert!(profile.is_valid(0));
+ assert!(profile.is_valid(1));
+ assert!(sub.is_valid(0));
+ assert!(sub.is_null(1));
+ assert!(added.is_null(0));
+ assert!(added.is_null(1));
+ }
+
fn test_table(file_io: &FileIO, table_path: &str) -> Table {
Table::new(
file_io.clone(),
diff --git a/crates/paimon/src/table/data_file_reader.rs
b/crates/paimon/src/table/data_file_reader.rs
index 1b175246..16d6bfa1 100644
--- a/crates/paimon/src/table/data_file_reader.rs
+++ b/crates/paimon/src/table/data_file_reader.rs
@@ -125,6 +125,7 @@ pub(crate) struct DataFileReader {
batch_size: Option<usize>,
parquet_read_budget: Option<Arc<ReadBudget>>,
table_options: Arc<HashMap<String, String>>,
+ nested_field_enabled: bool,
mosaic_prefetch: MosaicPrefetchOptions,
read_timing: Option<Arc<DataFileReadTiming>>,
}
@@ -152,6 +153,7 @@ impl DataFileReader {
batch_size: None,
parquet_read_budget: None,
table_options: Arc::new(HashMap::new()),
+ nested_field_enabled: false,
mosaic_prefetch: MosaicPrefetchOptions::default(),
read_timing: None,
}
@@ -191,6 +193,13 @@ impl DataFileReader {
options: impl Into<Arc<HashMap<String, String>>>,
) -> Self {
self.table_options = options.into();
+ self.nested_field_enabled =
crate::spec::CoreOptions::new(&self.table_options)
+ .data_evolution_nested_field_enabled();
+ self
+ }
+
+ pub(crate) fn with_nested_field_enabled(mut self, enabled: bool) -> Self {
+ self.nested_field_enabled = enabled;
self
}
@@ -528,21 +537,22 @@ impl DataFileReader {
let target_schema = build_target_arrow_schema(&read_type)?;
let file_fields = data_fields.clone().unwrap_or_else(||
table_fields.clone());
- // What the reader is asked for.
- let projected_read_fields: Vec<DataField> = if let Some(ref df) =
data_fields {
- read_data_fields(df, &read_type)?
- } else {
- read_type
- .iter()
- .filter(|field| field.name() != ROW_ID_FIELD_NAME)
- .cloned()
- .collect()
- };
let data_schema_fields = data_schema_fields_for_file(
&file_fields,
file_meta.write_cols.as_deref(),
data_schema_fields.as_deref(),
)?;
+ // What the reader is asked for.
+ let projected_read_fields: Vec<DataField> =
+ if data_fields.is_some() || file_meta.write_cols.is_some() {
+ read_data_fields(&data_schema_fields, &read_type,
self.nested_field_enabled)?
+ } else {
+ read_type
+ .iter()
+ .filter(|field| field.name() != ROW_ID_FIELD_NAME)
+ .cloned()
+ .collect()
+ };
let path_to_read = split.data_file_path(&file_meta);
let configured_reader = create_format_reader_with_budget(
&path_to_read,
@@ -561,14 +571,15 @@ impl DataFileReader {
// The decoded batch is described by `format_read_fields`, so map
// `read_type` onto *that* list: its entries carry the types the
columns
// actually come back as, which is what reconciling them needs.
- let (index_mapping, source_fields) = if data_fields.is_some() {
- (
- create_index_mapping(&read_type, &format_read_fields),
- Some(format_read_fields.clone()),
- )
- } else {
- (None, None)
- };
+ let (index_mapping, source_fields) =
+ if data_fields.is_some() || file_meta.write_cols.is_some() {
+ (
+ create_index_mapping(&read_type, &format_read_fields),
+ Some(format_read_fields.clone()),
+ )
+ } else {
+ (None, None)
+ };
// Remap predicates from table-level to file-level indices.
let file_predicates = if row_id_residual {
@@ -824,18 +835,19 @@ impl DataFileReader {
let target_schema = build_target_arrow_schema(&read_type)?;
let file_fields = data_fields.clone().unwrap_or_else(||
table_fields.clone());
- // What the reader is asked for.
- let projected_read_fields: Vec<DataField> = if let Some(ref df) =
data_fields {
- read_data_fields(df, &read_type)?
- } else {
- read_type
- .iter()
- .filter(|field| field.name() != ROW_ID_FIELD_NAME)
- .cloned()
- .collect()
- };
let data_schema_fields =
data_schema_fields_for_file(&file_fields,
file_meta.write_cols.as_deref(), None)?;
+ // What the reader is asked for.
+ let projected_read_fields: Vec<DataField> =
+ if data_fields.is_some() || file_meta.write_cols.is_some() {
+ read_data_fields(&data_schema_fields, &read_type,
self.nested_field_enabled)?
+ } else {
+ read_type
+ .iter()
+ .filter(|field| field.name() != ROW_ID_FIELD_NAME)
+ .cloned()
+ .collect()
+ };
let path_to_read = split.data_file_path(&file_meta);
let configured_reader = create_format_reader_with_budget(
&path_to_read,
@@ -854,14 +866,15 @@ impl DataFileReader {
// The decoded batch is described by `format_read_fields`, so map
// `read_type` onto *that* list: its entries carry the types the
columns
// actually come back as, which is what reconciling them needs.
- let (index_mapping, source_fields) = if data_fields.is_some() {
- (
- create_index_mapping(&read_type, &format_read_fields),
- Some(format_read_fields.clone()),
- )
- } else {
- (None, None)
- };
+ let (index_mapping, source_fields) =
+ if data_fields.is_some() || file_meta.write_cols.is_some() {
+ (
+ create_index_mapping(&read_type, &format_read_fields),
+ Some(format_read_fields.clone()),
+ )
+ } else {
+ (None, None)
+ };
// Remap predicates from table-level to file-level indices.
let file_predicates = {
@@ -1019,6 +1032,7 @@ fn project_file_batch(
fn read_data_fields(
all_data_fields: &[DataField],
expected_fields: &[DataField],
+ nested_field_enabled: bool,
) -> crate::Result<Vec<DataField>> {
let mut read_fields = Vec::new();
for data_field in all_data_fields {
@@ -1026,9 +1040,11 @@ fn read_data_fields(
.iter()
.find(|field| field.id() == data_field.id())
{
- if let Some(pruned_type) =
- prune_data_type(expected.data_type(), data_field.data_type())?
- {
+ if let Some(pruned_type) = prune_data_type(
+ expected.data_type(),
+ data_field.data_type(),
+ nested_field_enabled,
+ )? {
read_fields.push(data_field_with_type(data_field,
pruned_type));
}
}
@@ -1036,7 +1052,11 @@ fn read_data_fields(
Ok(read_fields)
}
-fn prune_data_type(read_type: &DataType, data_type: &DataType) ->
crate::Result<Option<DataType>> {
+fn prune_data_type(
+ read_type: &DataType,
+ data_type: &DataType,
+ nested_field_enabled: bool,
+) -> crate::Result<Option<DataType>> {
match read_type {
DataType::Row(read_row) if is_variant_extraction_row_type(read_type)
=> {
Ok(Some(DataType::Row(read_row.clone())))
@@ -1052,21 +1072,29 @@ fn prune_data_type(read_type: &DataType, data_type:
&DataType) -> crate::Result<
.iter()
.find(|field| field.id() == read_field.id())
{
- if let Some(pruned_type) =
- prune_data_type(read_field.data_type(),
data_field.data_type())?
- {
+ if let Some(pruned_type) = prune_data_type(
+ read_field.data_type(),
+ data_field.data_type(),
+ nested_field_enabled,
+ )? {
fields.push(data_field_with_type(data_field,
pruned_type));
}
}
}
if fields.is_empty() {
- Ok(None)
- } else {
- Ok(Some(DataType::Row(crate::spec::RowType::with_nullable(
- read_type.is_nullable(),
- fields,
- ))))
+ if !nested_field_enabled {
+ return Ok(None);
+ }
+ if let Some(anchor) = data_row.fields().first().cloned() {
+ fields.push(anchor);
+ } else {
+ return Ok(None);
+ }
}
+ Ok(Some(DataType::Row(crate::spec::RowType::with_nullable(
+ read_type.is_nullable(),
+ fields,
+ ))))
}
// ARRAY and MAP are deliberately NOT descended, even though Java's
// `pruneDataType` does: the pruned type is also what the Vortex
reader is
@@ -1090,19 +1118,7 @@ fn data_schema_fields_for_file(
data_schema_fields: Option<&[DataField]>,
) -> crate::Result<Vec<DataField>> {
if let Some(write_cols) = write_cols {
- return write_cols
- .iter()
- .map(|name| {
- file_fields
- .iter()
- .find(|field| field.name() == name)
- .cloned()
- .ok_or_else(|| Error::DataInvalid {
- message: format!("write column '{name}' is absent from
the file schema"),
- source: None,
- })
- })
- .collect();
+ return super::data_evolution_fields::project_by_paths(file_fields,
write_cols);
}
Ok(data_schema_fields.unwrap_or(file_fields).to_vec())
}
@@ -1507,7 +1523,7 @@ mod row_tests {
));
let expected_field = field(1, "v", extraction_type.clone());
- let read_fields = read_data_fields(&[data_field],
&[expected_field]).unwrap();
+ let read_fields = read_data_fields(&[data_field], &[expected_field],
false).unwrap();
assert_eq!(read_fields.len(), 1);
assert!(is_variant_extraction_row_type(read_fields[0].data_type()));
@@ -1519,12 +1535,145 @@ mod row_tests {
let data_field = field(1, "n", DataType::Int(IntType::new()));
let expected_field = field(1, "n",
DataType::BigInt(BigIntType::new()));
- let read_fields = read_data_fields(&[data_field],
&[expected_field]).unwrap();
+ let read_fields = read_data_fields(&[data_field], &[expected_field],
false).unwrap();
assert_eq!(read_fields.len(), 1);
assert_eq!(read_fields[0].data_type(), &DataType::Int(IntType::new()));
}
+ #[test]
+ fn projected_added_leaf_reads_older_siblings_for_each_row_null_buffer() {
+ let existing = field(3, "existing", DataType::Int(IntType::new()));
+ let old_sub = field(
+ 2,
+ "sub",
+ DataType::Row(RowType::new(vec![
+ existing.clone(),
+ field(6, "unneeded", DataType::Int(IntType::new())),
+ ])),
+ );
+ let file_profile = field(
+ 1,
+ "profile",
+ DataType::Row(RowType::new(vec![
+ old_sub.clone(),
+ field(4, "other", DataType::Int(IntType::new())),
+ ])),
+ );
+ let read_profile = field(
+ 1,
+ "profile",
+ DataType::Row(RowType::new(vec![field(
+ 2,
+ "sub",
+ DataType::Row(RowType::new(vec![field(
+ 5,
+ "added",
+ DataType::Int(IntType::new()),
+ )])),
+ )])),
+ );
+
+ let pruned = read_data_fields(&[file_profile], &[read_profile],
true).unwrap();
+ let DataType::Row(profile) = pruned[0].data_type() else {
+ panic!("profile should remain a ROW")
+ };
+ assert_eq!(
+ profile.fields(),
+ &[field(2, "sub", DataType::Row(RowType::new(vec![existing])))]
+ );
+ }
+
+ #[test]
+ fn disabled_nested_mode_does_not_preserve_deep_hidden_anchor() {
+ let file_profile = field(
+ 1,
+ "profile",
+ DataType::Row(RowType::new(vec![field(
+ 2,
+ "sub",
+ DataType::Row(RowType::new(vec![field(
+ 3,
+ "old",
+ DataType::Int(IntType::new()),
+ )])),
+ )])),
+ );
+ let read_profile = field(
+ 1,
+ "profile",
+ DataType::Row(RowType::new(vec![field(
+ 2,
+ "sub",
+ DataType::Row(RowType::new(vec![field(
+ 4,
+ "added",
+ DataType::Int(IntType::new()),
+ )])),
+ )])),
+ );
+
+ assert!(read_data_fields(&[file_profile], &[read_profile], false)
+ .unwrap()
+ .is_empty());
+ }
+
+ #[test]
+ fn projected_missing_direct_child_reads_a_physical_sibling_as_anchor() {
+ let age = field(3, "age", DataType::Int(IntType::new()));
+ let file_profile = field(1, "profile",
DataType::Row(RowType::new(vec![age.clone()])));
+ let read_profile = field(
+ 1,
+ "profile",
+ DataType::Row(RowType::new(vec![field(
+ 2,
+ "name",
+ DataType::VarChar(VarCharType::string_type()),
+ )])),
+ );
+
+ let pruned = read_data_fields(&[file_profile], &[read_profile],
true).unwrap();
+ let DataType::Row(profile) = pruned[0].data_type() else {
+ panic!("profile should remain a ROW")
+ };
+ assert_eq!(profile.fields(), &[age]);
+ }
+
+ #[test]
+ fn disabled_nested_mode_does_not_read_hidden_sibling_anchor() {
+ let file_profile = field(
+ 1,
+ "profile",
+ DataType::Row(RowType::new(vec![field(
+ 2,
+ "old",
+ DataType::Int(IntType::new()),
+ )])),
+ );
+ let projected = field(
+ 1,
+ "profile",
+ DataType::Row(RowType::new(vec![field(
+ 3,
+ "added",
+ DataType::Int(IntType::new()),
+ )])),
+ );
+ assert!(read_data_fields(
+ std::slice::from_ref(&file_profile),
+ std::slice::from_ref(&projected),
+ false,
+ )
+ .unwrap()
+ .is_empty());
+ assert_eq!(
+ read_data_fields(&[file_profile], &[projected], true)
+ .unwrap()
+ .len(),
+ 1
+ );
+ }
+
#[tokio::test]
async fn row_projection_reads_full_file_schema_before_projecting() {
let fields = vec![
@@ -3673,7 +3822,7 @@ mod prune_container_tests {
])));
let read = DataType::Array(ArrayType::new(row(vec![f(3, "lang",
str_t())])));
- assert_eq!(prune_data_type(&read, &data).unwrap().unwrap(), data);
+ assert_eq!(prune_data_type(&read, &data, false).unwrap().unwrap(),
data);
}
#[test]
@@ -3684,6 +3833,6 @@ mod prune_container_tests {
));
let read = DataType::Map(MapType::new(str_t(), row(vec![f(5, "v",
str_t())])));
- assert_eq!(prune_data_type(&read, &data).unwrap().unwrap(), data);
+ assert_eq!(prune_data_type(&read, &data, false).unwrap().unwrap(),
data);
}
}
diff --git a/crates/paimon/src/table/mod.rs b/crates/paimon/src/table/mod.rs
index 1b918cca..1f29237b 100644
--- a/crates/paimon/src/table/mod.rs
+++ b/crates/paimon/src/table/mod.rs
@@ -37,6 +37,8 @@ mod chunk_shuffle;
mod commit_message;
mod consumer_manager;
pub(crate) mod cow_writer;
+mod data_evolution_fields;
+mod data_evolution_nested;
mod data_evolution_reader;
pub mod data_evolution_writer;
mod data_file_index_writer;
diff --git a/crates/paimon/src/table/table_commit.rs
b/crates/paimon/src/table/table_commit.rs
index 3e4b9c08..ec3049dc 100644
--- a/crates/paimon/src/table/table_commit.rs
+++ b/crates/paimon/src/table/table_commit.rs
@@ -2888,40 +2888,46 @@ impl TableCommit {
.fields()
.to_vec()
};
- let field_id_by_name = fields
- .iter()
- .map(|field| (field.name().to_string(), field.id()))
- .collect::<HashMap<_, _>>();
-
- let mut field_ids = HashSet::new();
- match file.write_cols.as_ref() {
- None => {
- field_ids.extend(
- fields
- .iter()
- .filter(|field| !is_system_field(field.name()))
- .map(|field| field.id()),
- );
- }
- Some(write_cols) => {
- for col in write_cols {
- if is_system_field(col) {
- continue;
- }
- let Some(field_id) = field_id_by_name.get(col) else {
- return Err(crate::Error::DataInvalid {
- message: format!(
- "Cannot find write column '{}' in schema {}.",
- col, file.schema_id
- ),
- source: None,
- });
- };
- field_ids.insert(*field_id);
- }
+ let data_fields = fields
+ .into_iter()
+ .filter(|field| !is_system_field(field.name()))
+ .collect::<Vec<_>>();
+ let write_cols = file.write_cols.as_ref().map(|cols| {
+ cols.iter()
+ .filter(|col| !is_system_field(col))
+ .cloned()
+ .collect::<Vec<_>>()
+ });
+ if self
+ .table
+ .schema()
+ .core_options()
+ .data_evolution_nested_field_enabled()
+ {
+ return super::data_evolution_fields::write_leaf_ids(
+ &data_fields,
+ write_cols.as_deref(),
+ );
+ }
+ let mut ids = HashSet::new();
+ if let Some(write_cols) = write_cols {
+ for col in write_cols {
+ let field = data_fields
+ .iter()
+ .find(|field| field.name() == col)
+ .ok_or_else(|| crate::Error::DataInvalid {
+ message: format!(
+ "Cannot find write column '{}' in schema {}.",
+ col, file.schema_id
+ ),
+ source: None,
+ })?;
+ ids.insert(field.id());
}
+ } else {
+ ids.extend(data_fields.iter().map(|field| field.id()));
}
- Ok(field_ids)
+ Ok(ids)
}
/// Assign row tracking metadata: snapshot ID as sequence number, and
@@ -3790,6 +3796,93 @@ mod tests {
TableSchema::new(0, &schema)
}
+ #[tokio::test]
+ async fn nested_row_id_conflicts_follow_java_leaf_identity() {
+ use crate::spec::{IntType, RowType, Schema};
+
+ let profile = DataType::Row(RowType::new(vec![
+ crate::spec::DataField::new(0, "name".into(),
DataType::Int(IntType::new())),
+ crate::spec::DataField::new(0, "age".into(),
DataType::Int(IntType::new())),
+ ]));
+ let schema = Schema::builder()
+ .column("id", DataType::Int(IntType::new()))
+ .column("profile", profile)
+ .option("bucket", "-1")
+ .option("data-evolution.enabled", "true")
+ .option("data-evolution.nested-field.enabled", "true")
+ .option("row-tracking.enabled", "true")
+ .build()
+ .unwrap();
+ let table = Table::new(
+ test_file_io(),
+ Identifier::new("default", "nested_conflict"),
+ "memory:/nested_conflict".into(),
+ TableSchema::new(0, &schema),
+ None,
+ );
+ let commit = table.new_write_builder().new_commit();
+ let ids_for = |path: Option<&str>| {
+ let mut file = test_data_file("f.parquet", 10);
+ file.write_cols = path.map(|path| vec![path.to_string()]);
+ file
+ };
+ let name = commit
+ .write_field_ids(&ids_for(Some("profile.name")))
+ .await
+ .unwrap();
+ let age = commit
+ .write_field_ids(&ids_for(Some("profile.age")))
+ .await
+ .unwrap();
+ let whole = commit
+ .write_field_ids(&ids_for(Some("profile")))
+ .await
+ .unwrap();
+ let full = commit.write_field_ids(&ids_for(None)).await.unwrap();
+ assert!(name.is_disjoint(&age));
+ assert!(!name.is_disjoint(&whole));
+ assert!(!age.is_disjoint(&whole));
+ assert!(whole.is_subset(&full));
+ }
+
+ #[tokio::test]
+ async fn disabled_nested_mode_rejects_dotted_row_id_write_path() {
+ use crate::spec::{IntType, RowType, Schema};
+
+ let schema = Schema::builder()
+ .column(
+ "profile",
+ DataType::Row(RowType::new(vec![crate::spec::DataField::new(
+ 0,
+ "age".into(),
+ DataType::Int(IntType::new()),
+ )])),
+ )
+ .option("bucket", "-1")
+ .option("data-evolution.enabled", "true")
+ .option("row-tracking.enabled", "true")
+ .build()
+ .unwrap();
+ let table = Table::new(
+ test_file_io(),
+ Identifier::new("default", "nested_disabled"),
+ "memory:/nested_disabled".into(),
+ TableSchema::new(0, &schema),
+ None,
+ );
+ let mut file = test_data_file("f.parquet", 10);
+ file.write_cols = Some(vec!["profile.age".into()]);
+ let err = table
+ .new_write_builder()
+ .new_commit()
+ .write_field_ids(&file)
+ .await
+ .unwrap_err();
+ assert!(err
+ .to_string()
+ .contains("Cannot find write column 'profile.age'"));
+ }
+
fn test_partitioned_schema() -> TableSchema {
use crate::spec::{DataType, IntType, Schema, VarCharType};
let schema = Schema::builder()
diff --git a/crates/paimon/src/table/table_scan.rs
b/crates/paimon/src/table/table_scan.rs
index 3358203c..e975dda7 100644
--- a/crates/paimon/src/table/table_scan.rs
+++ b/crates/paimon/src/table/table_scan.rs
@@ -940,11 +940,6 @@ async fn resolve_data_file_field_ids(
schema.fields()
};
- let field_id_by_name = fields
- .iter()
- .map(|field| (field.name(), field.id()))
- .collect::<HashMap<_, _>>();
-
let mut field_ids = HashSet::new();
match file.write_cols.as_ref() {
None => {
@@ -960,17 +955,14 @@ async fn resolve_data_file_field_ids(
if is_system_field_name(col) {
continue;
}
- let Some(field_id) = field_id_by_name.get(col.as_str()) else {
- return Err(crate::Error::DataInvalid {
- message: format!(
- "Cannot find write column '{}' in schema {}.",
- col, file.schema_id
- ),
- source: None,
- });
- };
- if !is_system_field_id(*field_id) {
- field_ids.insert(*field_id);
+ let projected = super::data_evolution_fields::project_by_paths(
+ fields,
+ std::slice::from_ref(col),
+ )?;
+ for field in projected {
+ if !is_system_field_id(field.id()) {
+ field_ids.insert(field.id());
+ }
}
}
}