diff --git a/paimon-core/src/main/java/org/apache/paimon/schema/SchemaManagerUtils.java b/paimon-core/src/main/java/org/apache/paimon/schema/SchemaManagerUtils.java index a71451c865d6..39c99691ba97 100644 --- a/paimon-core/src/main/java/org/apache/paimon/schema/SchemaManagerUtils.java +++ b/paimon-core/src/main/java/org/apache/paimon/schema/SchemaManagerUtils.java @@ -24,6 +24,7 @@ import org.apache.paimon.catalog.Identifier; import org.apache.paimon.fs.FileIO; import org.apache.paimon.fs.Path; +import org.apache.paimon.options.FallbackKey; import org.apache.paimon.schema.ColumnDirectiveUtils.ConvertedColumn; import org.apache.paimon.schema.SchemaChange.AddColumn; import org.apache.paimon.schema.SchemaChange.DropColumn; @@ -610,6 +611,7 @@ static Map applyRenameColumnsToOptions( Map renameMappings = Streams.stream(renameColumns) + .filter(rename -> rename.fieldNames().length == 1) .collect( Collectors.toMap( // currently only non-nested columns are supported @@ -618,10 +620,14 @@ static Map applyRenameColumnsToOptions( // case 1: the option key is fixed and only value may contain field names - // bucket key rename + // bucket key rename; canonical readers trim the csv entries, so the rewrite must + // trim too or a spaced entry never matches the rename String bucketKeysStr = options.get(BUCKET_KEY.key()); if (!StringUtils.isNullOrWhitespaceOnly(bucketKeysStr)) { - List bucketColumns = Arrays.asList(bucketKeysStr.split(",")); + List bucketColumns = + Arrays.stream(bucketKeysStr.split(",")) + .map(String::trim) + .collect(Collectors.toList()); List newBucketColumns = applyNotNestedColumnRename(bucketColumns, renameMappings); newOptions.put(BUCKET_KEY.key(), String.join(",", newBucketColumns)); @@ -630,12 +636,36 @@ static Map applyRenameColumnsToOptions( // sequence field rename String sequenceFieldsStr = options.get(SEQUENCE_FIELD.key()); if (!StringUtils.isNullOrWhitespaceOnly(sequenceFieldsStr)) { - List sequenceFields = Arrays.asList(sequenceFieldsStr.split(",")); + List sequenceFields = + Arrays.stream(sequenceFieldsStr.split(",")) + .map(String::trim) + .collect(Collectors.toList()); List newSequenceFields = applyNotNestedColumnRename(sequenceFields, renameMappings); newOptions.put(SEQUENCE_FIELD.key(), String.join(",", newSequenceFields)); } + // clustering columns rename; also cover the fallback key(s) (e.g. the deprecated + // sink.clustering.by-columns) so a table configured via the old key still follows the + // rename instead of keeping a stale column name the canonical reader would resolve + List clusteringColumnKeys = new ArrayList<>(); + clusteringColumnKeys.add(CLUSTERING_COLUMNS.key()); + for (FallbackKey fallbackKey : CLUSTERING_COLUMNS.fallbackKeys()) { + clusteringColumnKeys.add(fallbackKey.getKey()); + } + for (String clusteringColumnKey : clusteringColumnKeys) { + String clusteringColumnsStr = options.get(clusteringColumnKey); + if (!StringUtils.isNullOrWhitespaceOnly(clusteringColumnsStr)) { + List clusteringColumns = + Arrays.stream(clusteringColumnsStr.split(",")) + .map(String::trim) + .collect(Collectors.toList()); + List newClusteringColumns = + applyNotNestedColumnRename(clusteringColumns, renameMappings); + newOptions.put(clusteringColumnKey, String.join(",", newClusteringColumns)); + } + } + // case 2: the option key is composed of certain fixed prefixes, suffixes, and the field // name, while the option value doesn't contain field names. List> fieldNameToOptionKeys = @@ -661,6 +691,10 @@ static Map applyRenameColumnsToOptions( + MAP_SHARED_SHREDDING_COLUMN_PLACEMENT_POLICY); for (RenameColumn rename : renameColumns) { + if (rename.fieldNames().length > 1) { + // nested renames have no field-scoped option keys + continue; + } String fieldName = rename.fieldNames()[0]; String newFieldName = rename.newName(); @@ -691,12 +725,18 @@ static Map applyRenameColumnsToOptions( key.substring( FIELDS_PREFIX.length() + 1, key.length() - matchedSuffix.length() - 1); - List keyFields = Arrays.asList(keyFieldsStr.split(",")); + List keyFields = + Arrays.stream(keyFieldsStr.split(",")) + .map(String::trim) + .collect(Collectors.toList()); List newKeyFields = applyNotNestedColumnRename(keyFields, renameMappings); String valueFieldsStr = newOptions.remove(key); - List valueFields = Arrays.asList(valueFieldsStr.split(",")); + List valueFields = + Arrays.stream(valueFieldsStr.split(",")) + .map(String::trim) + .collect(Collectors.toList()); List newValueFields = applyNotNestedColumnRename(valueFields, renameMappings); newOptions.put( diff --git a/paimon-core/src/test/java/org/apache/paimon/schema/SchemaManagerUtilsTest.java b/paimon-core/src/test/java/org/apache/paimon/schema/SchemaManagerUtilsTest.java new file mode 100644 index 000000000000..b30aa59d6992 --- /dev/null +++ b/paimon-core/src/test/java/org/apache/paimon/schema/SchemaManagerUtilsTest.java @@ -0,0 +1,125 @@ +/* + * 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. + */ + +package org.apache.paimon.schema; + +import org.junit.jupiter.api.Test; + +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.Map; + +import static org.assertj.core.api.Assertions.assertThat; + +/** Tests for option rewriting in {@link SchemaManagerUtils}. */ +public class SchemaManagerUtilsTest { + + private static Map options(String... pairs) { + Map options = new HashMap<>(); + for (int i = 0; i < pairs.length; i += 2) { + options.put(pairs[i], pairs[i + 1]); + } + return options; + } + + @Test + public void testNestedRenameLeavesRootOptionKeysAlone() { + Map options = + options("fields.v.map.storage-layout", "shared-shredding", "bucket-key", "a"); + + Map rewritten = + SchemaManagerUtils.applyRenameColumnsToOptions( + options, + Collections.singletonList( + SchemaChange.renameColumn( + new String[] {"v", "value", "f1"}, "f100"))); + + // a nested rename must not relocate the root column's options to the new nested name + assertThat(rewritten).containsEntry("fields.v.map.storage-layout", "shared-shredding"); + assertThat(rewritten).doesNotContainKey("fields.f100.map.storage-layout"); + } + + @Test + public void testTwoNestedRenamesUnderOneRootDoNotCrash() { + Map options = options("fields.v.aggregate-function", "last_non_null"); + + Map rewritten = + SchemaManagerUtils.applyRenameColumnsToOptions( + options, + Arrays.asList( + SchemaChange.renameColumn(new String[] {"v", "value", "f1"}, "f1n"), + SchemaChange.renameColumn( + new String[] {"v", "value", "f2"}, "f2n"))); + + assertThat(rewritten).containsEntry("fields.v.aggregate-function", "last_non_null"); + } + + @Test + public void testSpacedCsvOptionEntriesMatchRename() { + Map options = + options( + "bucket-key", "a, b", + "sequence.field", "a, b", + "clustering.columns", "b"); + + Map rewritten = + SchemaManagerUtils.applyRenameColumnsToOptions( + options, Collections.singletonList(SchemaChange.renameColumn("b", "c"))); + + // canonical readers trim the entries, so the rewrite must match them trimmed + assertThat(rewritten).containsEntry("bucket-key", "a,c"); + assertThat(rewritten).containsEntry("sequence.field", "a,c"); + assertThat(rewritten).containsEntry("clustering.columns", "c"); + } + + @Test + public void testClusteringColumnsFollowRename() { + Map options = options("clustering.columns", "c,d"); + + Map rewritten = + SchemaManagerUtils.applyRenameColumnsToOptions( + options, Collections.singletonList(SchemaChange.renameColumn("c", "c2"))); + + assertThat(rewritten).containsEntry("clustering.columns", "c2,d"); + } + + @Test + public void testClusteringColumnsFallbackKeyFollowsRename() { + Map options = options("sink.clustering.by-columns", "c, d"); + + Map rewritten = + SchemaManagerUtils.applyRenameColumnsToOptions( + options, Collections.singletonList(SchemaChange.renameColumn("c", "c2"))); + + // the deprecated fallback key must follow the rename too, otherwise the canonical + // reader (which resolves the fallback) sees a column name that no longer exists + assertThat(rewritten).containsEntry("sink.clustering.by-columns", "c2,d"); + } + + @Test + public void testSequenceGroupValueEntriesMatchRenameTrimmed() { + Map options = options("fields.x.sequence-group", "a, b"); + + Map rewritten = + SchemaManagerUtils.applyRenameColumnsToOptions( + options, Collections.singletonList(SchemaChange.renameColumn("b", "c"))); + + assertThat(rewritten).containsEntry("fields.x.sequence-group", "a,c"); + } +}