diff --git a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java index 4f1cb3111193..21f362a7c69f 100644 --- a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java +++ b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java @@ -100,6 +100,9 @@ public class CoreOptions implements Serializable { public static final String MAP_SHARED_SHREDDING_MAX_COLUMNS = "map.shared-shredding.max-columns"; + public static final String MAP_SHARED_SHREDDING_COLUMN_PLACEMENT_POLICY = + "map.shared-shredding.column-placement-policy"; + public static final String FILE_INDEX = "file-index"; public static final String COLUMNS = "columns"; @@ -5256,6 +5259,18 @@ public int mapSharedShreddingMaxColumns(String fieldName) { return maxColumns; } + public MapSharedShreddingColumnPlacementPolicy mapSharedShreddingColumnPlacementPolicy( + String fieldName) { + return options.get( + key(FIELDS_PREFIX + + "." + + fieldName + + "." + + MAP_SHARED_SHREDDING_COLUMN_PLACEMENT_POLICY) + .enumType(MapSharedShreddingColumnPlacementPolicy.class) + .defaultValue(MapSharedShreddingColumnPlacementPolicy.LRU)); + } + /** MAP storage layout. */ public enum MapStorageLayout implements DescribedEnum { DEFAULT( @@ -5288,6 +5303,40 @@ public InlineElement getDescription() { } } + /** Physical column placement policy for shared-shredding MAP fields. */ + public enum MapSharedShreddingColumnPlacementPolicy implements DescribedEnum { + PLAIN( + "plain", + "Keep each MAP row's input key order and place the first K keys into physical " + + "columns."), + SEQUENTIAL( + "sequential", + "Order keys by their field dictionary IDs and place the first K keys into physical " + + "columns."), + LRU( + "lru", + "Reuse physical columns for recently seen keys and evict the least recently used " + + "column when necessary."); + + private final String value; + private final String description; + + MapSharedShreddingColumnPlacementPolicy(String value, String description) { + this.value = value; + this.description = description; + } + + @Override + public String toString() { + return value; + } + + @Override + public InlineElement getDescription() { + return text(description); + } + } + /** * Action to take when an UPDATE (e.g. via MERGE INTO) modifies columns that are covered by a * global index. diff --git a/paimon-common/src/main/java/org/apache/paimon/data/shredding/LruMapSharedShreddingColumnAllocator.java b/paimon-common/src/main/java/org/apache/paimon/data/shredding/LruMapSharedShreddingColumnAllocator.java new file mode 100644 index 000000000000..ebbf722c8e5d --- /dev/null +++ b/paimon-common/src/main/java/org/apache/paimon/data/shredding/LruMapSharedShreddingColumnAllocator.java @@ -0,0 +1,120 @@ +/* + * 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.data.shredding; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; + +/** + * Allocator that retains resident field-to-column assignments and evicts the least recently used + * column when necessary. + */ +public class LruMapSharedShreddingColumnAllocator extends MapSharedShreddingColumnAllocator { + + private final int[] residentFieldByColumn; + private final long[] lastUsed; + private long lruClock; + + public LruMapSharedShreddingColumnAllocator(int numColumns) { + super(numColumns); + this.residentFieldByColumn = emptyColumnMapping(); + this.lastUsed = new long[numColumns]; + } + + @Override + public RowAllocation allocateRow(List fieldIds) { + List sortedFieldIds = new ArrayList<>(fieldIds); + Collections.sort(sortedFieldIds); + + int[] colToField = emptyColumnMapping(); + int[] nextResidentFieldByColumn = residentFieldByColumn.clone(); + boolean[] usedColumns = new boolean[numColumns]; + List unassignedFields = new ArrayList<>(); + + for (Integer fieldId : sortedFieldIds) { + int column = findResidentColumn(fieldId); + if (column == -1) { + unassignedFields.add(fieldId); + } else { + usedColumns[column] = true; + colToField[column] = fieldId; + } + } + + List overflowFields = new ArrayList<>(); + for (Integer fieldId : unassignedFields) { + int column = selectColumn(usedColumns, nextResidentFieldByColumn); + if (column == -1) { + overflowFields.add(fieldId); + continue; + } + + usedColumns[column] = true; + colToField[column] = fieldId; + nextResidentFieldByColumn[column] = fieldId; + } + + RowAllocation allocation = new RowAllocation(colToField, overflowFields); + updateLastUsed(colToField); + System.arraycopy(nextResidentFieldByColumn, 0, residentFieldByColumn, 0, numColumns); + commitRow(allocation, sortedFieldIds); + return allocation; + } + + private int findResidentColumn(int fieldId) { + for (int column = 0; column < numColumns; column++) { + if (residentFieldByColumn[column] == fieldId) { + return column; + } + } + return -1; + } + + private int selectColumn(boolean[] usedColumns, int[] plannedResidentFieldByColumn) { + int selectedColumn = -1; + long selectedLastUsed = Long.MAX_VALUE; + for (int column = 0; column < numColumns; column++) { + if (usedColumns[column]) { + continue; + } + if (plannedResidentFieldByColumn[column] == -1) { + return column; + } + if (lastUsed[column] < selectedLastUsed) { + selectedColumn = column; + selectedLastUsed = lastUsed[column]; + } + } + return selectedColumn; + } + + private void updateLastUsed(int[] colToField) { + boolean touched = false; + for (int column = 0; column < numColumns; column++) { + if (colToField[column] != -1) { + lastUsed[column] = lruClock; + touched = true; + } + } + if (touched) { + lruClock++; + } + } +} diff --git a/paimon-common/src/main/java/org/apache/paimon/data/shredding/MapSharedShreddingColumnAllocator.java b/paimon-common/src/main/java/org/apache/paimon/data/shredding/MapSharedShreddingColumnAllocator.java index 9a3a9c896556..05a418fb0696 100644 --- a/paimon-common/src/main/java/org/apache/paimon/data/shredding/MapSharedShreddingColumnAllocator.java +++ b/paimon-common/src/main/java/org/apache/paimon/data/shredding/MapSharedShreddingColumnAllocator.java @@ -29,44 +29,43 @@ /** * Per-row physical column allocator for one shared-shredding MAP column. * - *

This is a simple temporary implementation which assigns fields to physical columns by row - * order. A later version will use a more sophisticated LRU-style allocator to improve column reuse - * across rows. + *

Implementations decide the physical column placement for each row. This base class accumulates + * the file-level metadata shared by all placement policies. */ -public class MapSharedShreddingColumnAllocator { +public abstract class MapSharedShreddingColumnAllocator { - private final int numColumns; + protected final int numColumns; private final Map> fieldToColumns = new TreeMap<>(); private final Set overflowFieldSet = new TreeSet<>(); private int maxRowWidth = 0; - public MapSharedShreddingColumnAllocator(int numColumns) { + protected MapSharedShreddingColumnAllocator(int numColumns) { this.numColumns = numColumns; } - public RowAllocation allocateRow(List fieldIds) { + /** Allocates physical columns for one row's field IDs. */ + public abstract RowAllocation allocateRow(List fieldIds); + + /** Commits one row allocation and updates accumulated file-level metadata. */ + protected void commitRow(RowAllocation allocation, List fieldIds) { maxRowWidth = Math.max(maxRowWidth, fieldIds.size()); - int[] colToField = new int[numColumns]; - for (int i = 0; i < numColumns; i++) { - colToField[i] = -1; + for (int column = 0; column < numColumns; column++) { + int fieldId = allocation.colToField[column]; + if (fieldId != -1) { + fieldToColumns.computeIfAbsent(fieldId, ignored -> new TreeSet<>()).add(column); + } } - int assignLimit = Math.min(fieldIds.size(), numColumns); - for (int i = 0; i < assignLimit; i++) { - int fieldId = fieldIds.get(i); - colToField[i] = fieldId; - fieldToColumns.computeIfAbsent(fieldId, ignored -> new TreeSet<>()).add(i); - } + overflowFieldSet.addAll(allocation.overflowFields); + } - List overflowFields = new ArrayList<>(); - for (int i = assignLimit; i < fieldIds.size(); i++) { - int fieldId = fieldIds.get(i); - overflowFields.add(fieldId); - overflowFieldSet.add(fieldId); + protected int[] emptyColumnMapping() { + int[] colToField = new int[numColumns]; + for (int i = 0; i < numColumns; i++) { + colToField[i] = -1; } - - return new RowAllocation(colToField, overflowFields); + return colToField; } public Map> fieldToColumns() { @@ -97,9 +96,9 @@ public static class RowAllocation { private final int[] colToField; private final List overflowFields; - private RowAllocation(int[] colToField, List overflowFields) { - this.colToField = colToField; - this.overflowFields = Collections.unmodifiableList(overflowFields); + RowAllocation(int[] colToField, List overflowFields) { + this.colToField = colToField.clone(); + this.overflowFields = Collections.unmodifiableList(new ArrayList<>(overflowFields)); } public int[] colToField() { diff --git a/paimon-common/src/main/java/org/apache/paimon/data/shredding/MapSharedShreddingRowConverter.java b/paimon-common/src/main/java/org/apache/paimon/data/shredding/MapSharedShreddingRowConverter.java index 2b0bc532c7eb..8a099b12d227 100644 --- a/paimon-common/src/main/java/org/apache/paimon/data/shredding/MapSharedShreddingRowConverter.java +++ b/paimon-common/src/main/java/org/apache/paimon/data/shredding/MapSharedShreddingRowConverter.java @@ -18,6 +18,7 @@ package org.apache.paimon.data.shredding; +import org.apache.paimon.CoreOptions.MapSharedShreddingColumnPlacementPolicy; import org.apache.paimon.data.BinaryString; import org.apache.paimon.data.Blob; import org.apache.paimon.data.Decimal; @@ -42,6 +43,8 @@ import java.util.List; import java.util.Map; +import static org.apache.paimon.utils.Preconditions.checkArgument; + /** Converts logical rows containing shared-shredding MAP fields into physical rows. */ public class MapSharedShreddingRowConverter { @@ -52,7 +55,9 @@ public class MapSharedShreddingRowConverter { private final List shreddingFieldNames; public MapSharedShreddingRowConverter( - RowType logicalType, Map fieldToNumColumns) { + RowType logicalType, + Map fieldToNumColumns, + Map fieldToColumnPlacementPolicy) { this.logicalType = logicalType; this.physicalType = MapSharedShreddingUtils.logicalToPhysicalSchema(logicalType, fieldToNumColumns); @@ -68,7 +73,14 @@ public MapSharedShreddingRowConverter( } MapType mapType = (MapType) field.type(); - ColumnContext context = new ColumnContext(field.name(), numColumns, mapType); + MapSharedShreddingColumnPlacementPolicy placementPolicy = + fieldToColumnPlacementPolicy.get(field.name()); + checkArgument( + placementPolicy != null, + "Missing column placement policy for shared-shredding field '%s'.", + field.name()); + ColumnContext context = + new ColumnContext(field.name(), numColumns, mapType, placementPolicy); contextByFieldName.put(field.name(), context); contextByFieldPos[i] = context; shreddingFieldNames.add(field.name()); @@ -289,14 +301,33 @@ private static class ColumnContext { private final MapSharedShreddingFieldDict dict; private final MapSharedShreddingColumnAllocator allocator; - private ColumnContext(String fieldName, int numColumns, MapType mapType) { + private ColumnContext( + String fieldName, + int numColumns, + MapType mapType, + MapSharedShreddingColumnPlacementPolicy placementPolicy) { this.fieldName = fieldName; this.numColumns = numColumns; this.keyGetter = InternalArray.createElementGetter(mapType.getKeyType()); DataType valueType = mapType.getValueType(); this.valueGetter = InternalArray.createElementGetter(valueType); this.dict = new MapSharedShreddingFieldDict(); - this.allocator = new MapSharedShreddingColumnAllocator(numColumns); + this.allocator = createAllocator(numColumns, placementPolicy); + } + + private static MapSharedShreddingColumnAllocator createAllocator( + int numColumns, MapSharedShreddingColumnPlacementPolicy placementPolicy) { + switch (placementPolicy) { + case PLAIN: + return new PlainMapSharedShreddingColumnAllocator(numColumns); + case SEQUENTIAL: + return new SequentialMapSharedShreddingColumnAllocator(numColumns); + case LRU: + return new LruMapSharedShreddingColumnAllocator(numColumns); + default: + throw new IllegalArgumentException( + "Unknown shared-shredding column placement policy: " + placementPolicy); + } } } } diff --git a/paimon-common/src/main/java/org/apache/paimon/data/shredding/MapSharedShreddingWritePlan.java b/paimon-common/src/main/java/org/apache/paimon/data/shredding/MapSharedShreddingWritePlan.java index af0f614906ed..9555a2f45ed5 100644 --- a/paimon-common/src/main/java/org/apache/paimon/data/shredding/MapSharedShreddingWritePlan.java +++ b/paimon-common/src/main/java/org/apache/paimon/data/shredding/MapSharedShreddingWritePlan.java @@ -18,6 +18,7 @@ package org.apache.paimon.data.shredding; +import org.apache.paimon.CoreOptions.MapSharedShreddingColumnPlacementPolicy; import org.apache.paimon.data.InternalRow; import org.apache.paimon.types.RowType; @@ -35,9 +36,13 @@ public class MapSharedShreddingWritePlan implements ShreddingWritePlan { @Nullable private Map> fieldMetadata; public MapSharedShreddingWritePlan( - RowType logicalRowType, Map fieldToNumColumns) { + RowType logicalRowType, + Map fieldToNumColumns, + Map fieldToColumnPlacementPolicy) { this.logicalRowType = logicalRowType; - this.converter = new MapSharedShreddingRowConverter(logicalRowType, fieldToNumColumns); + this.converter = + new MapSharedShreddingRowConverter( + logicalRowType, fieldToNumColumns, fieldToColumnPlacementPolicy); } @Override diff --git a/paimon-common/src/main/java/org/apache/paimon/data/shredding/MapSharedShreddingWritePlanFactory.java b/paimon-common/src/main/java/org/apache/paimon/data/shredding/MapSharedShreddingWritePlanFactory.java index 24c47207fad5..7818b24baa68 100644 --- a/paimon-common/src/main/java/org/apache/paimon/data/shredding/MapSharedShreddingWritePlanFactory.java +++ b/paimon-common/src/main/java/org/apache/paimon/data/shredding/MapSharedShreddingWritePlanFactory.java @@ -19,6 +19,7 @@ package org.apache.paimon.data.shredding; import org.apache.paimon.CoreOptions; +import org.apache.paimon.CoreOptions.MapSharedShreddingColumnPlacementPolicy; import org.apache.paimon.data.InternalMap; import org.apache.paimon.data.InternalRow; import org.apache.paimon.format.shredding.ShreddingWritePlanFactory; @@ -38,6 +39,7 @@ public class MapSharedShreddingWritePlanFactory implements ShreddingWritePlanFac private final RowType logicalRowType; private final Map fieldToMaxColumns; + private final Map fieldToColumnPlacementPolicy; private final Map fieldToPosition; public MapSharedShreddingWritePlanFactory(RowType logicalRowType, Options options) { @@ -47,8 +49,11 @@ public MapSharedShreddingWritePlanFactory(RowType logicalRowType, Options option MapSharedShreddingUtils.detectShreddingColumns(logicalRowType, coreOptions); this.fieldToMaxColumns = MapSharedShreddingUtils.buildColumnToNumColumns(shreddingFields, coreOptions); + this.fieldToColumnPlacementPolicy = new LinkedHashMap<>(); this.fieldToPosition = new LinkedHashMap<>(); for (String field : shreddingFields) { + fieldToColumnPlacementPolicy.put( + field, coreOptions.mapSharedShreddingColumnPlacementPolicy(field)); fieldToPosition.put(field, logicalRowType.getFieldIndex(field)); } } @@ -97,6 +102,7 @@ public ShreddingWritePlan createWritePlan(List sampleRows) { } // TODO: Infer the column count from recent file metadata instead of current-file samples. - return new MapSharedShreddingWritePlan(logicalRowType, fieldToNumColumns); + return new MapSharedShreddingWritePlan( + logicalRowType, fieldToNumColumns, fieldToColumnPlacementPolicy); } } diff --git a/paimon-common/src/main/java/org/apache/paimon/data/shredding/PlainMapSharedShreddingColumnAllocator.java b/paimon-common/src/main/java/org/apache/paimon/data/shredding/PlainMapSharedShreddingColumnAllocator.java new file mode 100644 index 000000000000..14b778ffc9b4 --- /dev/null +++ b/paimon-common/src/main/java/org/apache/paimon/data/shredding/PlainMapSharedShreddingColumnAllocator.java @@ -0,0 +1,51 @@ +/* + * 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.data.shredding; + +import java.util.ArrayList; +import java.util.List; + +/** Allocator that maps fields to physical columns in the input MAP entry order. */ +public class PlainMapSharedShreddingColumnAllocator extends MapSharedShreddingColumnAllocator { + + public PlainMapSharedShreddingColumnAllocator(int numColumns) { + super(numColumns); + } + + @Override + public RowAllocation allocateRow(List fieldIds) { + RowAllocation allocation = allocateLeadingColumns(fieldIds); + commitRow(allocation, fieldIds); + return allocation; + } + + protected RowAllocation allocateLeadingColumns(List fieldIds) { + int[] colToField = emptyColumnMapping(); + List overflowFields = new ArrayList<>(); + for (int i = 0; i < fieldIds.size(); i++) { + int fieldId = fieldIds.get(i); + if (i < numColumns) { + colToField[i] = fieldId; + } else { + overflowFields.add(fieldId); + } + } + return new RowAllocation(colToField, overflowFields); + } +} diff --git a/paimon-common/src/main/java/org/apache/paimon/data/shredding/SequentialMapSharedShreddingColumnAllocator.java b/paimon-common/src/main/java/org/apache/paimon/data/shredding/SequentialMapSharedShreddingColumnAllocator.java new file mode 100644 index 000000000000..a759402ffe7b --- /dev/null +++ b/paimon-common/src/main/java/org/apache/paimon/data/shredding/SequentialMapSharedShreddingColumnAllocator.java @@ -0,0 +1,41 @@ +/* + * 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.data.shredding; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; + +/** Allocator that orders fields by dictionary ID before filling physical columns. */ +public class SequentialMapSharedShreddingColumnAllocator + extends PlainMapSharedShreddingColumnAllocator { + + public SequentialMapSharedShreddingColumnAllocator(int numColumns) { + super(numColumns); + } + + @Override + public RowAllocation allocateRow(List fieldIds) { + List sortedFieldIds = new ArrayList<>(fieldIds); + Collections.sort(sortedFieldIds); + RowAllocation allocation = allocateLeadingColumns(sortedFieldIds); + commitRow(allocation, sortedFieldIds); + return allocation; + } +} diff --git a/paimon-common/src/test/java/org/apache/paimon/data/shredding/LruMapSharedShreddingColumnAllocatorTest.java b/paimon-common/src/test/java/org/apache/paimon/data/shredding/LruMapSharedShreddingColumnAllocatorTest.java new file mode 100644 index 000000000000..39e4def32dd1 --- /dev/null +++ b/paimon-common/src/test/java/org/apache/paimon/data/shredding/LruMapSharedShreddingColumnAllocatorTest.java @@ -0,0 +1,83 @@ +/* + * 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.data.shredding; + +import org.junit.jupiter.api.Test; + +import java.util.Arrays; +import java.util.Collections; + +import static org.assertj.core.api.Assertions.assertThat; + +/** Tests for {@link LruMapSharedShreddingColumnAllocator}. */ +class LruMapSharedShreddingColumnAllocatorTest { + + @Test + void testAllocatesWithHitRetainEvictAndOverflow() { + LruMapSharedShreddingColumnAllocator allocator = + new LruMapSharedShreddingColumnAllocator(3); + + MapSharedShreddingColumnAllocator.RowAllocation row0 = + allocator.allocateRow(Arrays.asList(0, 1, 2)); + assertThat(row0.colToField()).containsExactly(0, 1, 2); + assertThat(row0.overflowFields()).isEmpty(); + + MapSharedShreddingColumnAllocator.RowAllocation row1 = + allocator.allocateRow(Arrays.asList(0, 1)); + assertThat(row1.colToField()).containsExactly(0, 1, -1); + assertThat(row1.overflowFields()).isEmpty(); + + MapSharedShreddingColumnAllocator.RowAllocation row2 = + allocator.allocateRow(Arrays.asList(3, 4, 5)); + assertThat(row2.colToField()).containsExactly(4, 5, 3); + assertThat(row2.overflowFields()).isEmpty(); + + MapSharedShreddingColumnAllocator.RowAllocation row3 = + allocator.allocateRow(Arrays.asList(0, 3, 4, 5)); + assertThat(row3.colToField()).containsExactly(4, 5, 3); + assertThat(row3.overflowFields()).containsExactly(0); + + assertThat(allocator.maxRowWidth()).isEqualTo(4); + assertThat(allocator.fieldToColumns().get(0)).containsExactly(0); + assertThat(allocator.fieldToColumns().get(1)).containsExactly(1); + assertThat(allocator.fieldToColumns().get(2)).containsExactly(2); + assertThat(allocator.fieldToColumns().get(3)).containsExactly(2); + assertThat(allocator.fieldToColumns().get(4)).containsExactly(0); + assertThat(allocator.fieldToColumns().get(5)).containsExactly(1); + assertThat(allocator.overflowFieldSet()).containsExactly(0); + } + + @Test + void testHandlesEmptyRows() { + LruMapSharedShreddingColumnAllocator allocator = + new LruMapSharedShreddingColumnAllocator(2); + + MapSharedShreddingColumnAllocator.RowAllocation emptyRow = + allocator.allocateRow(Collections.emptyList()); + assertThat(emptyRow.colToField()).containsExactly(-1, -1); + assertThat(emptyRow.overflowFields()).isEmpty(); + assertThat(allocator.maxRowWidth()).isZero(); + + MapSharedShreddingColumnAllocator.RowAllocation row = + allocator.allocateRow(Collections.singletonList(7)); + assertThat(row.colToField()).containsExactly(7, -1); + assertThat(row.overflowFields()).isEmpty(); + assertThat(allocator.maxRowWidth()).isEqualTo(1); + } +} diff --git a/paimon-common/src/test/java/org/apache/paimon/data/shredding/MapSharedShreddingRowConverterTest.java b/paimon-common/src/test/java/org/apache/paimon/data/shredding/MapSharedShreddingRowConverterTest.java index 3fed33f722fb..7960b119a5e2 100644 --- a/paimon-common/src/test/java/org/apache/paimon/data/shredding/MapSharedShreddingRowConverterTest.java +++ b/paimon-common/src/test/java/org/apache/paimon/data/shredding/MapSharedShreddingRowConverterTest.java @@ -18,6 +18,7 @@ package org.apache.paimon.data.shredding; +import org.apache.paimon.CoreOptions.MapSharedShreddingColumnPlacementPolicy; import org.apache.paimon.data.BinaryString; import org.apache.paimon.data.GenericArray; import org.apache.paimon.data.GenericMap; @@ -53,7 +54,7 @@ void testBasicConversion() { DataTypes.FIELD( 1, "tags", DataTypes.MAP(DataTypes.STRING(), DataTypes.BIGINT()))); MapSharedShreddingRowConverter converter = - new MapSharedShreddingRowConverter(logicalType, columns("tags", 3)); + createPlainConverter(logicalType, columns("tags", 3)); assertThat(converter.shreddingFieldNames()).containsExactly("tags"); assertThatThrownBy(() -> converter.shreddingFieldNames().add("metrics")) .isInstanceOf(UnsupportedOperationException.class); @@ -109,7 +110,7 @@ void testBasicMapWithNullValue() { "metrics", DataTypes.MAP(DataTypes.STRING(), DataTypes.BIGINT()))); MapSharedShreddingRowConverter converter = - new MapSharedShreddingRowConverter(logicalType, columns("metrics", 2)); + createPlainConverter(logicalType, columns("metrics", 2)); InternalRow row = converter.convert(GenericRow.of(stringKeyMap("a", null, "b", 20L))); InternalRow metrics = row.getRow(0, 4); @@ -142,7 +143,7 @@ void testOverflowWhenExceedK() { "metrics", DataTypes.MAP(DataTypes.STRING(), DataTypes.BIGINT()))); MapSharedShreddingRowConverter converter = - new MapSharedShreddingRowConverter(logicalType, columns("metrics", 2)); + createPlainConverter(logicalType, columns("metrics", 2)); InternalRow row = converter.convert(GenericRow.of(stringKeyMap("a", 10L, "b", 20L, "c", 30L))); @@ -176,7 +177,7 @@ void testEmptyAndNullMaps() { "metrics", DataTypes.MAP(DataTypes.STRING(), DataTypes.BIGINT()))); MapSharedShreddingRowConverter converter = - new MapSharedShreddingRowConverter(logicalType, columns("metrics", 2)); + createPlainConverter(logicalType, columns("metrics", 2)); InternalRow nullRow = converter.convert(GenericRow.of((InternalMap) null)); assertThat(nullRow.isNullAt(0)).isTrue(); @@ -212,7 +213,7 @@ void testNestedValueStruct() { DataTypes.ROW( DataTypes.FIELD(0, "tags", DataTypes.MAP(DataTypes.STRING(), valueType))); MapSharedShreddingRowConverter converter = - new MapSharedShreddingRowConverter(logicalType, columns("tags", 2)); + createPlainConverter(logicalType, columns("tags", 2)); InternalRow row = converter.convert( @@ -253,7 +254,7 @@ void testNestedValueList() { DataTypes.MAP( DataTypes.STRING(), DataTypes.ARRAY(DataTypes.INT())))); MapSharedShreddingRowConverter converter = - new MapSharedShreddingRowConverter(logicalType, columns("tags", 2)); + createPlainConverter(logicalType, columns("tags", 2)); InternalRow first = converter.convert( @@ -335,7 +336,7 @@ void testNestedValueMap() { DataTypes.FIELD( 1, "nested", DataTypes.MAP(DataTypes.STRING(), innerMapType))); MapSharedShreddingRowConverter converter = - new MapSharedShreddingRowConverter(logicalType, columns("nested", 2)); + createPlainConverter(logicalType, columns("nested", 2)); InternalRow first = converter.convert( @@ -421,7 +422,7 @@ void testNestedComplex() { DataTypes.FIELD(0, "id", DataTypes.INT()), DataTypes.FIELD(1, "data", DataTypes.MAP(DataTypes.STRING(), valueType))); MapSharedShreddingRowConverter converter = - new MapSharedShreddingRowConverter(logicalType, columns("data", 2)); + createPlainConverter(logicalType, columns("data", 2)); InternalRow first = converter.convert( @@ -527,7 +528,7 @@ void testMultipleMapFields() { fieldToNumColumns.put("tags", 2); fieldToNumColumns.put("attrs", 3); MapSharedShreddingRowConverter converter = - new MapSharedShreddingRowConverter(logicalType, fieldToNumColumns); + createPlainConverter(logicalType, fieldToNumColumns); InternalRow first = converter.convert( @@ -620,6 +621,106 @@ void testMultipleMapFields() { 4)); } + @Test + void testSequentialPlacementUsesDictionaryOrder() { + RowType logicalType = + DataTypes.ROW( + DataTypes.FIELD(0, "id", DataTypes.INT()), + DataTypes.FIELD( + 1, "tags", DataTypes.MAP(DataTypes.STRING(), DataTypes.BIGINT()))); + MapSharedShreddingRowConverter converter = + createConverter( + logicalType, "tags", 3, MapSharedShreddingColumnPlacementPolicy.SEQUENTIAL); + + InternalRow first = converter.convert(GenericRow.of(100, stringKeyMap("a", 1L, "b", 2L))); + InternalRow firstTags = first.getRow(1, 5); + assertThat(firstTags.getArray(0).toIntArray()).containsExactly(0, 1, -1); + assertThat(firstTags.getLong(1)).isEqualTo(1L); + assertThat(firstTags.getLong(2)).isEqualTo(2L); + assertThat(firstTags.isNullAt(3)).isTrue(); + assertThat(firstTags.isNullAt(4)).isTrue(); + + InternalRow second = + converter.convert(GenericRow.of(200, stringKeyMap("b", 3L, "c", 4L, "a", 5L))); + InternalRow secondTags = second.getRow(1, 5); + assertThat(secondTags.getArray(0).toIntArray()).containsExactly(0, 1, 2); + assertThat(secondTags.getLong(1)).isEqualTo(5L); + assertThat(secondTags.getLong(2)).isEqualTo(3L); + assertThat(secondTags.getLong(3)).isEqualTo(4L); + assertThat(secondTags.isNullAt(4)).isTrue(); + + assertThat(converter.buildFieldMeta("tags")) + .isEqualTo( + new MapSharedShreddingFieldMeta( + nameToId("a", 0, "b", 1, "c", 2), + fieldToColumns( + 0, Collections.singletonList(0), + 1, Collections.singletonList(1), + 2, Collections.singletonList(2)), + new TreeSet(), + 3, + 3)); + } + + @Test + void testLruPlacementPreservesResidentColumns() { + RowType logicalType = + DataTypes.ROW( + DataTypes.FIELD(0, "id", DataTypes.INT()), + DataTypes.FIELD( + 1, "tags", DataTypes.MAP(DataTypes.STRING(), DataTypes.BIGINT()))); + MapSharedShreddingRowConverter converter = + createConverter( + logicalType, "tags", 3, MapSharedShreddingColumnPlacementPolicy.LRU); + + InternalRow first = + converter.convert(GenericRow.of(1, stringKeyMap("a", 10L, "b", 20L, "c", 30L))); + assertThat(first.getRow(1, 5).getArray(0).toIntArray()).containsExactly(0, 1, 2); + + InternalRow second = converter.convert(GenericRow.of(2, stringKeyMap("a", 40L, "b", 50L))); + assertThat(second.getRow(1, 5).getArray(0).toIntArray()).containsExactly(0, 1, -1); + + InternalRow third = + converter.convert(GenericRow.of(3, stringKeyMap("d", 60L, "e", 70L, "f", 80L))); + InternalRow thirdTags = third.getRow(1, 5); + assertThat(thirdTags.getArray(0).toIntArray()).containsExactly(4, 5, 3); + assertThat(thirdTags.getLong(1)).isEqualTo(70L); + assertThat(thirdTags.getLong(2)).isEqualTo(80L); + assertThat(thirdTags.getLong(3)).isEqualTo(60L); + assertThat(thirdTags.isNullAt(4)).isTrue(); + + InternalRow fourth = + converter.convert( + GenericRow.of(4, stringKeyMap("a", 90L, "d", 100L, "e", 110L, "f", 120L))); + InternalRow fourthTags = fourth.getRow(1, 5); + assertThat(fourthTags.getArray(0).toIntArray()).containsExactly(4, 5, 3); + assertThat(fourthTags.getLong(1)).isEqualTo(110L); + assertThat(fourthTags.getLong(2)).isEqualTo(120L); + assertThat(fourthTags.getLong(3)).isEqualTo(100L); + assertThat(fourthTags.getMap(4)).isEqualTo(intKeyMap(0, 90L)); + + assertThat(converter.buildFieldMeta("tags")) + .isEqualTo( + new MapSharedShreddingFieldMeta( + nameToId( + "a", 0, + "b", 1, + "c", 2, + "d", 3, + "e", 4, + "f", 5), + fieldToColumns( + 0, Collections.singletonList(0), + 1, Collections.singletonList(1), + 2, Collections.singletonList(2), + 3, Collections.singletonList(2), + 4, Collections.singletonList(0), + 5, Collections.singletonList(1)), + new TreeSet<>(Collections.singletonList(0)), + 3, + 4)); + } + @Test void testBuildFieldMetaInvalidFieldName() { RowType logicalType = @@ -628,7 +729,7 @@ void testBuildFieldMetaInvalidFieldName() { DataTypes.FIELD( 1, "tags", DataTypes.MAP(DataTypes.STRING(), DataTypes.BIGINT()))); MapSharedShreddingRowConverter converter = - new MapSharedShreddingRowConverter(logicalType, columns("tags", 3)); + createPlainConverter(logicalType, columns("tags", 3)); assertThat(converter.buildFieldMeta("tags")) .isEqualTo( @@ -652,6 +753,26 @@ private static Map columns(String fieldName, int numColumns) { return columns; } + private static MapSharedShreddingRowConverter createPlainConverter( + RowType logicalType, Map fieldToNumColumns) { + Map policies = new HashMap<>(); + for (String fieldName : fieldToNumColumns.keySet()) { + policies.put(fieldName, MapSharedShreddingColumnPlacementPolicy.PLAIN); + } + return new MapSharedShreddingRowConverter(logicalType, fieldToNumColumns, policies); + } + + private static MapSharedShreddingRowConverter createConverter( + RowType logicalType, + String fieldName, + int numColumns, + MapSharedShreddingColumnPlacementPolicy policy) { + return new MapSharedShreddingRowConverter( + logicalType, + Collections.singletonMap(fieldName, numColumns), + Collections.singletonMap(fieldName, policy)); + } + private static GenericMap stringKeyMap(Object... keyValues) { Map values = new LinkedHashMap<>(); for (int i = 0; i < keyValues.length; i += 2) { diff --git a/paimon-common/src/test/java/org/apache/paimon/data/shredding/MapSharedShreddingWritePlanTest.java b/paimon-common/src/test/java/org/apache/paimon/data/shredding/MapSharedShreddingWritePlanTest.java index ccc323bf5451..1652e8171719 100644 --- a/paimon-common/src/test/java/org/apache/paimon/data/shredding/MapSharedShreddingWritePlanTest.java +++ b/paimon-common/src/test/java/org/apache/paimon/data/shredding/MapSharedShreddingWritePlanTest.java @@ -18,6 +18,7 @@ package org.apache.paimon.data.shredding; +import org.apache.paimon.CoreOptions.MapSharedShreddingColumnPlacementPolicy; import org.apache.paimon.data.BinaryString; import org.apache.paimon.data.GenericMap; import org.apache.paimon.data.GenericRow; @@ -45,7 +46,11 @@ void testConvertAndBuildFieldMetadata() { DataTypes.FIELD( 1, "tags", DataTypes.MAP(DataTypes.STRING(), DataTypes.BIGINT()))); MapSharedShreddingWritePlan writePlan = - new MapSharedShreddingWritePlan(logicalType, Collections.singletonMap("tags", 4)); + new MapSharedShreddingWritePlan( + logicalType, + Collections.singletonMap("tags", 4), + Collections.singletonMap( + "tags", MapSharedShreddingColumnPlacementPolicy.PLAIN)); InternalRow physicalRow = writePlan.toPhysicalRow( @@ -117,12 +122,69 @@ void testFactoryCapsInferredColumnCountAtMaxColumns() { logicalType, Collections.singletonMap("tags", 2))); } + @Test + void testFactoryUsesConfiguredColumnPlacementPolicy() { + RowType logicalType = + DataTypes.ROW( + DataTypes.FIELD( + 0, "tags", DataTypes.MAP(DataTypes.STRING(), DataTypes.INT()))); + MapSharedShreddingWritePlanFactory factory = createFactory(logicalType, 3, "sequential"); + InternalRow first = GenericRow.of(stringKeyMap("a", 1, "b", 2, "c", 6)); + InternalRow second = GenericRow.of(stringKeyMap("b", 3, "d", 4, "a", 5)); + ShreddingWritePlan writePlan = factory.createWritePlan(Collections.singletonList(first)); + + writePlan.toPhysicalRow(first).getRow(0, 5); + InternalRow physicalMap = writePlan.toPhysicalRow(second).getRow(0, 5); + + assertThat(physicalMap.getArray(0).toIntArray()).containsExactly(0, 1, 3); + assertThat(physicalMap.getInt(1)).isEqualTo(5); + assertThat(physicalMap.getInt(2)).isEqualTo(3); + assertThat(physicalMap.getInt(3)).isEqualTo(4); + } + + @Test + void testFactoryUsesLruColumnPlacementByDefault() { + RowType logicalType = + DataTypes.ROW( + DataTypes.FIELD( + 0, "tags", DataTypes.MAP(DataTypes.STRING(), DataTypes.INT()))); + MapSharedShreddingWritePlanFactory factory = createFactory(logicalType, 3); + InternalRow first = GenericRow.of(stringKeyMap("a", 10, "b", 20, "c", 30)); + ShreddingWritePlan writePlan = factory.createWritePlan(Collections.singletonList(first)); + + writePlan.toPhysicalRow(first).getRow(0, 5); + writePlan.toPhysicalRow(GenericRow.of(stringKeyMap("a", 40, "b", 50))).getRow(0, 5); + writePlan + .toPhysicalRow(GenericRow.of(stringKeyMap("d", 60, "e", 70, "f", 80))) + .getRow(0, 5); + InternalRow physicalMap = + writePlan + .toPhysicalRow( + GenericRow.of(stringKeyMap("a", 90, "d", 100, "e", 110, "f", 120))) + .getRow(0, 5); + + assertThat(physicalMap.getArray(0).toIntArray()).containsExactly(4, 5, 3); + assertThat(physicalMap.getInt(1)).isEqualTo(110); + assertThat(physicalMap.getInt(2)).isEqualTo(120); + assertThat(physicalMap.getInt(3)).isEqualTo(100); + assertThat(physicalMap.getMap(4)).isEqualTo(intKeyMap(0, 90)); + } + private static MapSharedShreddingWritePlanFactory createFactory( RowType logicalType, int maxColumns) { + return createFactory(logicalType, maxColumns, null); + } + + private static MapSharedShreddingWritePlanFactory createFactory( + RowType logicalType, int maxColumns, String placementPolicy) { Options options = new Options(); options.setString("fields.tags.map.storage-layout", "shared-shredding"); options.setString( "fields.tags.map.shared-shredding.max-columns", String.valueOf(maxColumns)); + if (placementPolicy != null) { + options.setString( + "fields.tags.map.shared-shredding.column-placement-policy", placementPolicy); + } return new MapSharedShreddingWritePlanFactory(logicalType, options); } @@ -133,4 +195,12 @@ private static GenericMap stringKeyMap(Object... keyValues) { } return new GenericMap(values); } + + private static GenericMap intKeyMap(Object... keyValues) { + Map values = new LinkedHashMap<>(); + for (int i = 0; i < keyValues.length; i += 2) { + values.put(keyValues[i], keyValues[i + 1]); + } + return new GenericMap(values); + } } diff --git a/paimon-common/src/test/java/org/apache/paimon/data/shredding/MapSharedShreddingColumnAllocatorTest.java b/paimon-common/src/test/java/org/apache/paimon/data/shredding/PlainMapSharedShreddingColumnAllocatorTest.java similarity index 65% rename from paimon-common/src/test/java/org/apache/paimon/data/shredding/MapSharedShreddingColumnAllocatorTest.java rename to paimon-common/src/test/java/org/apache/paimon/data/shredding/PlainMapSharedShreddingColumnAllocatorTest.java index 933e1faf6180..c1dc967b8f39 100644 --- a/paimon-common/src/test/java/org/apache/paimon/data/shredding/MapSharedShreddingColumnAllocatorTest.java +++ b/paimon-common/src/test/java/org/apache/paimon/data/shredding/PlainMapSharedShreddingColumnAllocatorTest.java @@ -28,12 +28,13 @@ import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; -/** Tests for {@link MapSharedShreddingColumnAllocator}. */ -class MapSharedShreddingColumnAllocatorTest { +/** Tests for {@link PlainMapSharedShreddingColumnAllocator}. */ +class PlainMapSharedShreddingColumnAllocatorTest { @Test void testBasicAllocation() { - MapSharedShreddingColumnAllocator allocator = new MapSharedShreddingColumnAllocator(3); + PlainMapSharedShreddingColumnAllocator allocator = + new PlainMapSharedShreddingColumnAllocator(3); MapSharedShreddingColumnAllocator.RowAllocation allocation = allocator.allocateRow(Arrays.asList(10, 20)); @@ -44,7 +45,8 @@ void testBasicAllocation() { @Test void testExactlyKFields() { - MapSharedShreddingColumnAllocator allocator = new MapSharedShreddingColumnAllocator(3); + PlainMapSharedShreddingColumnAllocator allocator = + new PlainMapSharedShreddingColumnAllocator(3); MapSharedShreddingColumnAllocator.RowAllocation allocation = allocator.allocateRow(Arrays.asList(0, 1, 2)); @@ -55,7 +57,8 @@ void testExactlyKFields() { @Test void testOverflowWhenExceedK() { - MapSharedShreddingColumnAllocator allocator = new MapSharedShreddingColumnAllocator(2); + PlainMapSharedShreddingColumnAllocator allocator = + new PlainMapSharedShreddingColumnAllocator(2); MapSharedShreddingColumnAllocator.RowAllocation allocation = allocator.allocateRow(Arrays.asList(10, 20, 30, 40)); @@ -68,7 +71,8 @@ void testOverflowWhenExceedK() { @Test void testEmptyRow() { - MapSharedShreddingColumnAllocator allocator = new MapSharedShreddingColumnAllocator(3); + PlainMapSharedShreddingColumnAllocator allocator = + new PlainMapSharedShreddingColumnAllocator(3); MapSharedShreddingColumnAllocator.RowAllocation allocation = allocator.allocateRow(Arrays.asList()); @@ -79,7 +83,8 @@ void testEmptyRow() { @Test void testMaxRowWidthTracked() { - MapSharedShreddingColumnAllocator allocator = new MapSharedShreddingColumnAllocator(3); + PlainMapSharedShreddingColumnAllocator allocator = + new PlainMapSharedShreddingColumnAllocator(3); allocator.allocateRow(Arrays.asList(1, 2)); allocator.allocateRow(Arrays.asList(1, 2, 3, 4, 5)); @@ -90,7 +95,8 @@ void testMaxRowWidthTracked() { @Test void testFieldToColumnsAccumulated() { - MapSharedShreddingColumnAllocator allocator = new MapSharedShreddingColumnAllocator(3); + PlainMapSharedShreddingColumnAllocator allocator = + new PlainMapSharedShreddingColumnAllocator(3); allocator.allocateRow(Arrays.asList(10, 20, 30)); allocator.allocateRow(Arrays.asList(20, 40)); @@ -108,7 +114,8 @@ void testFieldToColumnsAccumulated() { @Test void testOverflowFieldSetAccumulated() { - MapSharedShreddingColumnAllocator allocator = new MapSharedShreddingColumnAllocator(2); + PlainMapSharedShreddingColumnAllocator allocator = + new PlainMapSharedShreddingColumnAllocator(2); allocator.allocateRow(Arrays.asList(1, 2, 3)); allocator.allocateRow(Arrays.asList(4, 5, 6, 7)); @@ -118,21 +125,40 @@ void testOverflowFieldSetAccumulated() { .isInstanceOf(UnsupportedOperationException.class); } - @Test - void testGetNumColumns() { - MapSharedShreddingColumnAllocator allocator = new MapSharedShreddingColumnAllocator(5); - - assertThat(allocator.numColumns()).isEqualTo(5); - } - @Test void testSingleColumnAllocator() { - MapSharedShreddingColumnAllocator allocator = new MapSharedShreddingColumnAllocator(1); + PlainMapSharedShreddingColumnAllocator allocator = + new PlainMapSharedShreddingColumnAllocator(1); MapSharedShreddingColumnAllocator.RowAllocation allocation = allocator.allocateRow(Arrays.asList(10, 20, 30)); assertThat(allocation.colToField()).containsExactly(10); assertThat(allocation.overflowFields()).containsExactly(20, 30); + assertThat(allocator.numColumns()).isEqualTo(1); + } + + @Test + void testUsesInputOrder() { + PlainMapSharedShreddingColumnAllocator allocator = + new PlainMapSharedShreddingColumnAllocator(3); + + MapSharedShreddingColumnAllocator.RowAllocation row0 = + allocator.allocateRow(Arrays.asList(2, 0, 1)); + assertThat(row0.colToField()).containsExactly(2, 0, 1); + assertThat(row0.overflowFields()).isEmpty(); + + MapSharedShreddingColumnAllocator.RowAllocation row1 = + allocator.allocateRow(Arrays.asList(4, 3, 5, 6)); + assertThat(row1.colToField()).containsExactly(4, 3, 5); + assertThat(row1.overflowFields()).containsExactly(6); + + assertThat(allocator.fieldToColumns().get(0)).containsExactly(1); + assertThat(allocator.fieldToColumns().get(1)).containsExactly(2); + assertThat(allocator.fieldToColumns().get(2)).containsExactly(0); + assertThat(allocator.fieldToColumns().get(3)).containsExactly(1); + assertThat(allocator.fieldToColumns().get(4)).containsExactly(0); + assertThat(allocator.fieldToColumns().get(5)).containsExactly(2); + assertThat(allocator.overflowFieldSet()).containsExactly(6); } } diff --git a/paimon-common/src/test/java/org/apache/paimon/data/shredding/SequentialMapSharedShreddingColumnAllocatorTest.java b/paimon-common/src/test/java/org/apache/paimon/data/shredding/SequentialMapSharedShreddingColumnAllocatorTest.java new file mode 100644 index 000000000000..879712b9a4b0 --- /dev/null +++ b/paimon-common/src/test/java/org/apache/paimon/data/shredding/SequentialMapSharedShreddingColumnAllocatorTest.java @@ -0,0 +1,58 @@ +/* + * 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.data.shredding; + +import org.junit.jupiter.api.Test; + +import java.util.Arrays; + +import static org.assertj.core.api.Assertions.assertThat; + +/** Tests for {@link SequentialMapSharedShreddingColumnAllocator}. */ +class SequentialMapSharedShreddingColumnAllocatorTest { + + @Test + void testSortsAndUsesLeadingColumns() { + SequentialMapSharedShreddingColumnAllocator allocator = + new SequentialMapSharedShreddingColumnAllocator(3); + + MapSharedShreddingColumnAllocator.RowAllocation row0 = + allocator.allocateRow(Arrays.asList(1, 2)); + assertThat(row0.colToField()).containsExactly(1, 2, -1); + assertThat(row0.overflowFields()).isEmpty(); + + MapSharedShreddingColumnAllocator.RowAllocation row1 = + allocator.allocateRow(Arrays.asList(2, 3)); + assertThat(row1.colToField()).containsExactly(2, 3, -1); + assertThat(row1.overflowFields()).isEmpty(); + + MapSharedShreddingColumnAllocator.RowAllocation row2 = + allocator.allocateRow(Arrays.asList(7, 4, 6, 5)); + assertThat(row2.colToField()).containsExactly(4, 5, 6); + assertThat(row2.overflowFields()).containsExactly(7); + + assertThat(allocator.fieldToColumns().get(1)).containsExactly(0); + assertThat(allocator.fieldToColumns().get(2)).containsExactly(0, 1); + assertThat(allocator.fieldToColumns().get(3)).containsExactly(1); + assertThat(allocator.fieldToColumns().get(4)).containsExactly(0); + assertThat(allocator.fieldToColumns().get(5)).containsExactly(1); + assertThat(allocator.fieldToColumns().get(6)).containsExactly(2); + assertThat(allocator.overflowFieldSet()).containsExactly(7); + } +} diff --git a/paimon-core/src/main/java/org/apache/paimon/schema/SchemaValidation.java b/paimon-core/src/main/java/org/apache/paimon/schema/SchemaValidation.java index d145704feffb..452ba953f8c0 100644 --- a/paimon-core/src/main/java/org/apache/paimon/schema/SchemaValidation.java +++ b/paimon-core/src/main/java/org/apache/paimon/schema/SchemaValidation.java @@ -677,6 +677,7 @@ private static void validateMapStorageLayout(TableSchema schema, CoreOptions opt fieldName)); } options.mapSharedShreddingMaxColumns(fieldName); + options.mapSharedShreddingColumnPlacementPolicy(fieldName); } if (hasSharedShredding) { diff --git a/paimon-core/src/test/java/org/apache/paimon/CoreOptionsTest.java b/paimon-core/src/test/java/org/apache/paimon/CoreOptionsTest.java index 3399990a8a4b..35b548644cc1 100644 --- a/paimon-core/src/test/java/org/apache/paimon/CoreOptionsTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/CoreOptionsTest.java @@ -126,19 +126,32 @@ public void testMapStorageLayout() { assertThat(options.mapStorageLayout("metrics")) .isEqualTo(CoreOptions.MapStorageLayout.DEFAULT); assertThat(options.mapSharedShreddingMaxColumns("metrics")).isEqualTo(256); + assertThat(options.mapSharedShreddingColumnPlacementPolicy("metrics")) + .isEqualTo(CoreOptions.MapSharedShreddingColumnPlacementPolicy.LRU); conf.setString("fields.metrics.map.storage-layout", "shared-shredding"); conf.setString("fields.metrics.map.shared-shredding.max-columns", "32"); + conf.setString("fields.metrics.map.shared-shredding.column-placement-policy", "sequential"); options = new CoreOptions(conf); assertThat(options.mapStorageLayout("metrics")) .isEqualTo(CoreOptions.MapStorageLayout.SHARED_SHREDDING); assertThat(options.mapSharedShreddingMaxColumns("metrics")).isEqualTo(32); + assertThat(options.mapSharedShreddingColumnPlacementPolicy("metrics")) + .isEqualTo(CoreOptions.MapSharedShreddingColumnPlacementPolicy.SEQUENTIAL); + + conf.setString("fields.metrics.map.shared-shredding.column-placement-policy", "lru"); + options = new CoreOptions(conf); + assertThat(options.mapSharedShreddingColumnPlacementPolicy("metrics")) + .isEqualTo(CoreOptions.MapSharedShreddingColumnPlacementPolicy.LRU); conf = new Options(); conf.setString("fields.metrics.map.storage-layout", "Shared-Shredding"); + conf.setString("fields.metrics.map.shared-shredding.column-placement-policy", "PLAIN"); options = new CoreOptions(conf); assertThat(options.mapStorageLayout("metrics")) .isEqualTo(CoreOptions.MapStorageLayout.SHARED_SHREDDING); + assertThat(options.mapSharedShreddingColumnPlacementPolicy("metrics")) + .isEqualTo(CoreOptions.MapSharedShreddingColumnPlacementPolicy.PLAIN); conf = new Options(); conf.setString("fields.metrics.map.storage-layout", "invalid"); @@ -146,6 +159,15 @@ public void testMapStorageLayout() { assertThatThrownBy(() -> invalidLayoutOptions.mapStorageLayout("metrics")) .hasMessageContaining("invalid"); + conf = new Options(); + conf.setString("fields.metrics.map.shared-shredding.column-placement-policy", "invalid"); + final CoreOptions invalidPlacementPolicyOptions = new CoreOptions(conf); + assertThatThrownBy( + () -> + invalidPlacementPolicyOptions + .mapSharedShreddingColumnPlacementPolicy("metrics")) + .hasMessageContaining("invalid"); + conf = new Options(); conf.setString("fields.metrics.map.shared-shredding.max-columns", "0"); final CoreOptions zeroMaxColumnsOptions = new CoreOptions(conf); diff --git a/paimon-core/src/test/java/org/apache/paimon/append/AppendOnlyWriterTest.java b/paimon-core/src/test/java/org/apache/paimon/append/AppendOnlyWriterTest.java index 47ebf0fbaa46..bd5d1d2cf54f 100644 --- a/paimon-core/src/test/java/org/apache/paimon/append/AppendOnlyWriterTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/append/AppendOnlyWriterTest.java @@ -869,6 +869,9 @@ private Options sharedShreddingOptions(Object... fieldToMaxColumns) { options.setString( "fields." + fieldName + ".map.shared-shredding.max-columns", String.valueOf(fieldToMaxColumns[i + 1])); + options.setString( + "fields." + fieldName + ".map.shared-shredding.column-placement-policy", + "plain"); } options.setString("metadata.stats-mode", "none"); return options; diff --git a/paimon-core/src/test/java/org/apache/paimon/schema/SchemaValidationTest.java b/paimon-core/src/test/java/org/apache/paimon/schema/SchemaValidationTest.java index ec07c17c4fd6..de116cb600da 100644 --- a/paimon-core/src/test/java/org/apache/paimon/schema/SchemaValidationTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/schema/SchemaValidationTest.java @@ -504,6 +504,21 @@ public void testMapStorageLayout() { options, ""))) .hasMessageContaining("options map.shared-shredding.max-columns must > 0"); + + options.put("fields.metrics.map.shared-shredding.max-columns", "1"); + options.put("fields.metrics.map.shared-shredding.column-placement-policy", "invalid"); + assertThatThrownBy( + () -> + validateTableSchema( + new TableSchema( + 1, + fields, + 10, + emptyList(), + emptyList(), + options, + ""))) + .hasMessageContaining("invalid"); } @Test diff --git a/paimon-core/src/test/java/org/apache/paimon/table/MapSharedShreddingTableTest.java b/paimon-core/src/test/java/org/apache/paimon/table/MapSharedShreddingTableTest.java index de5b74572b72..35e86d9a1d23 100644 --- a/paimon-core/src/test/java/org/apache/paimon/table/MapSharedShreddingTableTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/table/MapSharedShreddingTableTest.java @@ -103,6 +103,62 @@ public void testAppendOnlyTableReadWrite(String format) throws Exception { .containsEntry(4, javaMapOf("a", null, "b", 42L, "c", null)); } + @ParameterizedTest + @CsvSource({ + "orc,plain", + "orc,sequential", + "orc,lru", + "parquet,plain", + "parquet,sequential", + "parquet,lru" + }) + public void testColumnPlacementPolicies(String format, String placementPolicy) + throws Exception { + Table table = createTable(format, 3, "metrics"); + catalog.alterTable( + identifier(format), + Collections.singletonList( + SchemaChange.setOption( + "fields.metrics.map.shared-shredding.column-placement-policy", + placementPolicy)), + false); + table = catalog.getTable(identifier(format)); + + write( + table, + GenericRow.of(1, mapOf("a", 10L, "b", 20L, "c", 30L)), + GenericRow.of(2, mapOf("a", 40L, "b", 50L)), + GenericRow.of(3, mapOf("d", 60L)), + GenericRow.of(4, mapOf("a", 70L, "b", 80L, "c", 90L, "d", 100L))); + + FileStoreTable fileStoreTable = (FileStoreTable) table; + List files = currentDataFiles(fileStoreTable); + assertThat(files).hasSize(1); + MapSharedShreddingFieldMeta fieldMeta = + readSharedShreddingFieldMeta(fileStoreTable, files.get(0), "metrics"); + assertThat(fieldMeta.nameToId()).containsOnlyKeys("a", "b", "c", "d"); + assertThat(fieldMeta.numColumns()).isEqualTo(3); + assertThat(fieldMeta.maxRowWidth()).isEqualTo(4); + assertThat(fieldMeta.overflowFieldSet()).hasSize(1); + + if ("lru".equals(placementPolicy)) { + assertThat(fieldMeta.overflowFieldSet()).containsExactly(fieldMeta.nameToId().get("c")); + } else { + assertThat(fieldMeta.overflowFieldSet()).containsExactly(fieldMeta.nameToId().get("d")); + } + + Map> actual = new LinkedHashMap<>(); + for (InternalRow row : read(table)) { + actual.put(row.getInt(0), toJavaMap(row.getMap(1))); + } + assertThat(actual) + .containsOnlyKeys(1, 2, 3, 4) + .containsEntry(1, javaMapOf("a", 10L, "b", 20L, "c", 30L)) + .containsEntry(2, javaMapOf("a", 40L, "b", 50L)) + .containsEntry(3, javaMapOf("d", 60L)) + .containsEntry(4, javaMapOf("a", 70L, "b", 80L, "c", 90L, "d", 100L)); + } + @ParameterizedTest @ValueSource(strings = {"orc", "parquet"}) public void testAppendOnlyTableReadWriteWithTwoMapFields(String format) throws Exception {