Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -610,6 +611,7 @@ static Map<String, String> applyRenameColumnsToOptions(

Map<String, String> renameMappings =
Streams.stream(renameColumns)
.filter(rename -> rename.fieldNames().length == 1)
.collect(
Collectors.toMap(
// currently only non-nested columns are supported
Expand All @@ -618,10 +620,14 @@ static Map<String, String> 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<String> bucketColumns = Arrays.asList(bucketKeysStr.split(","));
List<String> bucketColumns =
Arrays.stream(bucketKeysStr.split(","))
.map(String::trim)
.collect(Collectors.toList());
List<String> newBucketColumns =
applyNotNestedColumnRename(bucketColumns, renameMappings);
newOptions.put(BUCKET_KEY.key(), String.join(",", newBucketColumns));
Expand All @@ -630,12 +636,36 @@ static Map<String, String> applyRenameColumnsToOptions(
// sequence field rename
String sequenceFieldsStr = options.get(SEQUENCE_FIELD.key());
if (!StringUtils.isNullOrWhitespaceOnly(sequenceFieldsStr)) {
List<String> sequenceFields = Arrays.asList(sequenceFieldsStr.split(","));
List<String> sequenceFields =
Arrays.stream(sequenceFieldsStr.split(","))
.map(String::trim)
.collect(Collectors.toList());
List<String> 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<String> 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<String> clusteringColumns =
Arrays.stream(clusteringColumnsStr.split(","))
.map(String::trim)
.collect(Collectors.toList());
List<String> 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<Function<String, String>> fieldNameToOptionKeys =
Expand All @@ -661,6 +691,10 @@ static Map<String, String> 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();

Expand Down Expand Up @@ -691,12 +725,18 @@ static Map<String, String> applyRenameColumnsToOptions(
key.substring(
FIELDS_PREFIX.length() + 1,
key.length() - matchedSuffix.length() - 1);
List<String> keyFields = Arrays.asList(keyFieldsStr.split(","));
List<String> keyFields =
Arrays.stream(keyFieldsStr.split(","))
.map(String::trim)
.collect(Collectors.toList());
List<String> newKeyFields =
applyNotNestedColumnRename(keyFields, renameMappings);

String valueFieldsStr = newOptions.remove(key);
List<String> valueFields = Arrays.asList(valueFieldsStr.split(","));
List<String> valueFields =
Arrays.stream(valueFieldsStr.split(","))
.map(String::trim)
.collect(Collectors.toList());
List<String> newValueFields =
applyNotNestedColumnRename(valueFields, renameMappings);
newOptions.put(
Expand Down
Original file line number Diff line number Diff line change
@@ -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<String, String> options(String... pairs) {
Map<String, String> 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<String, String> options =
options("fields.v.map.storage-layout", "shared-shredding", "bucket-key", "a");

Map<String, String> 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<String, String> options = options("fields.v.aggregate-function", "last_non_null");

Map<String, String> 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<String, String> options =
options(
"bucket-key", "a, b",
"sequence.field", "a, b",
"clustering.columns", "b");

Map<String, String> 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<String, String> options = options("clustering.columns", "c,d");

Map<String, String> rewritten =
SchemaManagerUtils.applyRenameColumnsToOptions(
options, Collections.singletonList(SchemaChange.renameColumn("c", "c2")));

assertThat(rewritten).containsEntry("clustering.columns", "c2,d");
}

@Test
public void testClusteringColumnsFallbackKeyFollowsRename() {
Map<String, String> options = options("sink.clustering.by-columns", "c, d");

Map<String, String> 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<String, String> options = options("fields.x.sequence-group", "a, b");

Map<String, String> rewritten =
SchemaManagerUtils.applyRenameColumnsToOptions(
options, Collections.singletonList(SchemaChange.renameColumn("b", "c")));

assertThat(rewritten).containsEntry("fields.x.sequence-group", "a,c");
}
}
Loading