diff --git a/docs/docs/multimodal-table/global-index/btree.mdx b/docs/docs/multimodal-table/global-index/btree.mdx index d73fdbfc9b0f..5a31c6b005e1 100644 --- a/docs/docs/multimodal-table/global-index/btree.mdx +++ b/docs/docs/multimodal-table/global-index/btree.mdx @@ -145,6 +145,47 @@ print(pa_table) +## Composite BTree Indexes + +Spark and Flink can create a composite index by listing columns in key order: + +```sql +CALL sys.create_global_index( + `table` => 'db.my_table', + index_column => 'category,item_number', + index_type => 'btree' +); + +SELECT * FROM my_table +WHERE category = 'category-a' + AND item_number = 107; +``` + +The index stores typed tuples and one row-ID posting list for each distinct tuple. +Equality conditions on every key column use a composite point lookup, regardless +of their order in the SQL predicate. This avoids expanding each column's posting +list before intersecting them. Additional predicates remain row filters, and OR +branches can each select a composite lookup. When several definitions match, +selection prefers the one with more key columns. + +Eligible composite point queries honor `global-index.query-in-reader.enabled`. +With reader evaluation enabled, supported index conditions narrow the candidate +rows and remaining conditions are evaluated against the data. + +Composite indexes can coexist with single-column indexes on the same columns. +Composite lookup currently requires equality on all key columns; conditions on +only some key columns, ranges and IN predicates use available single-column +indexes or ordinary scans. Vector and full-text search pre-filters use +single-column indexes or their data fallback. + +In `full` and `detail` modes, rows outside the selected composite index's coverage +remain eligible for data filtering. Calling the procedure again with the same +ordered column list fills missing index ranges. Under the `IGNORE` column-update +policy, updates to any key column invalidate the affected composite index ranges; +calling the procedure again rebuilds those ranges. +To drop the composite definition, pass the same ordered column list to +`sys.drop_global_index`; independent single-column definitions remain in place. + ## BTree Options | Option | Default | Description | diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/globalindex/SortedIndexTopoBuilder.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/globalindex/SortedIndexTopoBuilder.java index 54c010c22879..f40a22426598 100644 --- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/globalindex/SortedIndexTopoBuilder.java +++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/globalindex/SortedIndexTopoBuilder.java @@ -64,6 +64,7 @@ import org.apache.paimon.types.DataType; import org.apache.paimon.types.DataTypes; import org.apache.paimon.types.RowType; +import org.apache.paimon.utils.InternalRowUtils; import org.apache.paimon.utils.Range; import org.apache.flink.api.common.functions.Partitioner; @@ -150,12 +151,34 @@ public static Optional> buildIndexStream( PartitionPredicate partitionPredicate, Options userOptions) throws Exception { + List> definitions = new ArrayList<>(); + for (String name : indexColumns) { + definitions.add(Collections.singletonList(name)); + } + return buildIndexDefinitions( + env, + indexScannerSupplier, + table, + definitions, + indexType, + partitionPredicate, + userOptions); + } + + private static Optional> buildIndexDefinitions( + StreamExecutionEnvironment env, + Supplier indexScannerSupplier, + FileStoreTable table, + List> definitions, + String indexType, + PartitionPredicate partitionPredicate, + Options userOptions) + throws Exception { List> allStreams = new ArrayList<>(); - for (String indexColumn : indexColumns) { + for (List indexColumns : definitions) { + String indexColumn = indexColumns.get(0); SortedGlobalIndexScanner indexScanner = - indexScannerSupplier - .get() - .withIndexFields(Collections.singletonList(indexColumn)); + indexScannerSupplier.get().withIndexFields(indexColumns); if (partitionPredicate != null) { indexScanner = indexScanner.withPartitionPredicate(partitionPredicate); } @@ -176,18 +199,12 @@ public static Optional> buildIndexStream( continue; } - // 2. Select necessary columns (index field + ROW_ID) - List selectedColumns = new ArrayList<>(); - selectedColumns.add(indexColumn); - - RowType dataReadType = - SpecialFields.rowTypeWithRowId(table.rowType().project(selectedColumns)); + // 2. Select the ordered index fields and ROW_ID. + RowType sourceReadType = table.rowType().project(indexColumns); + RowType dataReadType = SpecialFields.rowTypeWithRowId(sourceReadType); String buildTaskIdField = buildTaskIdFieldName(dataReadType); GlobalIndexer indexer = - GlobalIndexer.create( - indexType, - Collections.singletonList(table.rowType().getField(indexColumn)), - userOptions); + GlobalIndexer.create(indexType, sourceReadType.getFields(), userOptions); if (!(indexer instanceof SortedGlobalIndexer)) { throw new IllegalArgumentException( "Index algorithm " + indexType + " does not expose sorted index keys."); @@ -200,7 +217,6 @@ public static Optional> buildIndexStream( int indexFieldPos = sortReadType.getFieldIndex(indexColumn); int rowIdPos = sortReadType.getFieldIndex(SpecialFields.ROW_ID.name()); DataType indexFieldType = sortReadType.getTypeAt(indexFieldPos); - DataType sourceFieldType = table.rowType().getField(indexColumn).type(); // 3. Calculate maximum parallelism bound long recordsPerRange = @@ -220,7 +236,7 @@ public static Optional> buildIndexStream( BinaryRow partition = partitionEntry.getKey(); byte[] partitionBytes = binaryRowSerializer.serializeToBytes(partition); Map> ranges = - keyExtractor.isIdentity() + keyExtractor.isIdentity() && sourceReadType.getFieldCount() == 1 ? partitionEntry.getValue() : shardSplitsByRowRange(partitionEntry.getValue(), recordsPerRange); for (Map.Entry> entry : ranges.entrySet()) { @@ -249,14 +265,14 @@ public static Optional> buildIndexStream( splitTasks, readBuilder, new SortedGlobalIndexWriter(table, indexType, userOptions) - .withIndexFields(Collections.singletonList(indexColumn)), + .withIndexFields(indexColumns), scanResult.scanSnapshotId(), partitionFieldSize, taskIdPos, indexFieldPos, rowIdPos, indexFieldType, - sourceFieldType, + sourceReadType, keyExtractor, coreOptions, sortReadType, @@ -306,19 +322,22 @@ private static List createDeleteCommittables( public static void buildIndexAndExecute( StreamExecutionEnvironment env, FileStoreTable table, - String indexColumn, + List indexColumns, String indexType, PartitionPredicate partitionPredicate, Options userOptions) throws Exception { - if (buildIndex( - env, - () -> new SortedGlobalIndexScanner(table, indexType, userOptions), - table, - Collections.singletonList(indexColumn), - indexType, - partitionPredicate, - userOptions)) { + Optional> written = + buildIndexDefinitions( + env, + () -> new SortedGlobalIndexScanner(table, indexType, userOptions), + table, + Collections.singletonList(indexColumns), + indexType, + partitionPredicate, + userOptions); + if (written.isPresent()) { + commit(table, written.get(), CoreOptions.createCommitUser(userOptions)); env.execute("Create " + indexType + " global index for table: " + table.name()); } } @@ -335,7 +354,7 @@ protected static DataStream executeForBuildTasks( int indexFieldPos, int rowIdPos, DataType indexFieldType, - DataType sourceFieldType, + RowType sourceReadType, GlobalIndexKeyExtractor keyExtractor, CoreOptions coreOptions, RowType readType, @@ -354,14 +373,14 @@ protected static DataStream executeForBuildTasks( .transform( "Read Data", InternalTypeInfo.fromRowType(readType), - new ReadDataOperator(readBuilder, keyExtractor, sourceFieldType)) + new ReadDataOperator(readBuilder, keyExtractor, sourceReadType)) .setParallelism(parallelism); DataStream sortedStream = sortRows( env, rowDataStream, - keyExtractor.isIdentity(), + keyExtractor.isIdentity() && sourceReadType.getFieldCount() == 1, taskIdPos, indexFieldPos, coreOptions, @@ -516,25 +535,34 @@ private static class ReadDataOperator private final ReadBuilder readBuilder; private final GlobalIndexKeyExtractor keyExtractor; - private final DataType sourceFieldType; + private final RowType sourceReadType; private transient TableRead tableRead; private transient InternalRow.FieldGetter sourceFieldGetter; + private transient InternalRow.FieldGetter[] compositeGetters; public ReadDataOperator( ReadBuilder readBuilder, GlobalIndexKeyExtractor keyExtractor, - DataType sourceFieldType) { + RowType sourceReadType) { this.readBuilder = readBuilder; this.keyExtractor = keyExtractor; - this.sourceFieldType = sourceFieldType; + this.sourceReadType = sourceReadType; } @Override public void open() throws Exception { super.open(); this.tableRead = readBuilder.newRead(); - this.sourceFieldGetter = InternalRow.createFieldGetter(sourceFieldType, 0); + if (sourceReadType.getFieldCount() > 1) { + compositeGetters = new InternalRow.FieldGetter[sourceReadType.getFieldCount()]; + for (int i = 0; i < compositeGetters.length; i++) { + compositeGetters[i] = + InternalRow.createFieldGetter(sourceReadType.getTypeAt(i), i); + } + } else { + sourceFieldGetter = InternalRow.createFieldGetter(sourceReadType.getTypeAt(0), 0); + } } @Override @@ -546,10 +574,20 @@ public void processElement(StreamRecord element) throws Excepti try { InternalRow row; while ((row = batch.next()) != null) { - long rowId = row.getLong(1); + long rowId = row.getLong(sourceReadType.getFieldCount()); + Object sourceValue; + if (compositeGetters == null) { + sourceValue = sourceFieldGetter.getFieldOrNull(row); + } else { + GenericRow tuple = new GenericRow(compositeGetters.length); + for (int i = 0; i < compositeGetters.length; i++) { + tuple.setField(i, compositeGetters[i].getFieldOrNull(row)); + } + sourceValue = InternalRowUtils.copy(tuple, sourceReadType); + } boolean[] emitted = new boolean[1]; keyExtractor.extract( - sourceFieldGetter.getFieldOrNull(row), + sourceValue, key -> { emitted[0] = true; output.collect( diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/procedure/CreateGlobalIndexProcedure.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/procedure/CreateGlobalIndexProcedure.java index c0a0ea90d585..bb38de474f17 100644 --- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/procedure/CreateGlobalIndexProcedure.java +++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/procedure/CreateGlobalIndexProcedure.java @@ -135,7 +135,7 @@ public String[] call( SortedIndexTopoBuilder.buildIndexAndExecute( procedureContext.getExecutionEnvironment(), table, - indexColumns.get(0), + indexColumns, indexType, partitionPredicate, userOptions); diff --git a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/SortedGlobalIndexITCase.java b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/SortedGlobalIndexITCase.java index 70be19ff6dc9..0ae2afbe7a94 100644 --- a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/SortedGlobalIndexITCase.java +++ b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/SortedGlobalIndexITCase.java @@ -20,15 +20,28 @@ import org.apache.paimon.catalog.Catalog; import org.apache.paimon.data.BinaryString; +import org.apache.paimon.data.GenericRow; +import org.apache.paimon.data.InternalRow; import org.apache.paimon.globalindex.IndexedSplit; import org.apache.paimon.index.DataEvolutionIndexSourceMeta; import org.apache.paimon.index.IndexFileMeta; +import org.apache.paimon.io.DataFileMeta; import org.apache.paimon.manifest.IndexManifestEntry; import org.apache.paimon.predicate.Predicate; import org.apache.paimon.predicate.PredicateBuilder; +import org.apache.paimon.reader.RecordReader; import org.apache.paimon.table.FileStoreTable; +import org.apache.paimon.table.SpecialFields; +import org.apache.paimon.table.sink.BatchTableCommit; +import org.apache.paimon.table.sink.BatchTableWrite; +import org.apache.paimon.table.sink.BatchWriteBuilder; +import org.apache.paimon.table.sink.CommitMessage; +import org.apache.paimon.table.sink.CommitMessageImpl; import org.apache.paimon.table.source.ReadBuilder; import org.apache.paimon.table.source.TableScan; +import org.apache.paimon.types.RowType; +import org.apache.paimon.utils.InternalRowUtils; +import org.apache.paimon.utils.Range; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; @@ -36,6 +49,10 @@ import org.apache.flink.types.Row; import org.junit.jupiter.api.Test; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.Comparator; import java.util.List; import java.util.Map; import java.util.Set; @@ -80,6 +97,189 @@ public void testBTreeIndex() throws Catalog.TableNotExistException { assertThat(sql("SELECT * FROM T WHERE id = 100")).containsOnly(Row.of(100, "name_100")); } + @Test + public void testCompositeBTreeIndex() throws Exception { + tEnv.getConfig().set(TableConfigOptions.TABLE_DML_SYNC, true); + sql( + "CREATE TABLE T_COMPOSITE (id INT, category STRING, item_number INT) WITH (" + + "'row-tracking.enabled' = 'true', 'data-evolution.enabled' = 'true', " + + "'global-index.column-update-action' = 'IGNORE', " + + "'sorted-index.records-per-file' = '13', 'sorted-index.build.max-parallelism' = '4', " + + "'btree-index.bloom-filter.enabled' = 'true')"); + insertCompositeRows(0, 40); + sql( + "INSERT INTO T_COMPOSITE VALUES " + + "(100, CAST(NULL AS STRING), 7), (101, 'category-a', CAST(NULL AS INT))"); + buildBTreeIndexForTable("T_COMPOSITE", "category"); + buildBTreeIndexForTable("T_COMPOSITE", "item_number"); + // Reverse the schema order to verify the index preserves the requested key order. + buildBTreeIndexForTable("T_COMPOSITE", "item_number, category"); + List compositeEntries = + paimonTable("T_COMPOSITE").store().newIndexFileHandler().scanEntries().stream() + .filter( + entry -> + entry.indexFile() + .globalIndexMeta() + .getIndexedFieldIds() + .size() + == 2) + .collect(Collectors.toList()); + assertThat(compositeEntries).isNotEmpty(); + assertThat(compositeEntries) + .allSatisfy( + entry -> { + assertThat(entry.indexFile().globalIndexMeta().getIndexedFieldIds()) + .containsExactly(2, 1); + assertThat(entry.indexFile().globalIndexMeta().rowRange().count()) + .isLessThanOrEqualTo(13); + }); + assertThat(compositeEntries.stream().mapToLong(entry -> entry.indexFile().rowCount()).sum()) + .isEqualTo(42); + for (boolean inReader : Arrays.asList(false, true)) { + sql( + "ALTER TABLE T_COMPOSITE SET ('global-index.query-in-reader.enabled' = '" + + inReader + + "', 'scalar-index.search-mode' = 'fast')"); + assertThat( + sql( + "SELECT id FROM T_COMPOSITE WHERE item_number = 7 AND category = 'category-a'")) + .containsExactlyInAnyOrder(Row.of(7), Row.of(27)); + assertThat( + sql( + "SELECT id FROM T_COMPOSITE WHERE category = 'category-a' AND item_number = 7")) + .containsExactlyInAnyOrder(Row.of(7), Row.of(27)); + assertThat( + sql( + "SELECT id FROM T_COMPOSITE WHERE category = 'absent' AND item_number = 7")) + .isEmpty(); + } + insertCompositeRows(40, 60); + buildBTreeIndexForTable("T_COMPOSITE", "item_number,category"); + assertThat( + sql( + "SELECT id FROM T_COMPOSITE WHERE category = 'category-a' AND item_number = 7")) + .containsExactlyInAnyOrder(Row.of(7), Row.of(27), Row.of(47)); + updateCompositeColumn(7, "item_number", 107); + buildBTreeIndexForTable("T_COMPOSITE", "item_number,category"); + assertThat( + sql( + "SELECT id FROM T_COMPOSITE WHERE category = 'category-a' AND item_number = 107")) + .containsExactly(Row.of(7)); + updateCompositeColumn(27, "category", BinaryString.fromString("category-c")); + buildBTreeIndexForTable("T_COMPOSITE", "item_number,category"); + assertThat( + sql( + "SELECT id FROM T_COMPOSITE WHERE category = 'category-a' AND item_number = 7")) + .containsExactly(Row.of(47)); + assertThat( + sql( + "SELECT id FROM T_COMPOSITE WHERE category = 'category-c' AND item_number = 7")) + .containsExactly(Row.of(27)); + long snapshot = paimonTable("T_COMPOSITE").snapshotManager().latestSnapshot().id(); + buildBTreeIndexForTable("T_COMPOSITE", "item_number,category"); + assertThat(paimonTable("T_COMPOSITE").snapshotManager().latestSnapshot().id()) + .isEqualTo(snapshot); + sql( + "CALL sys.drop_global_index(`table` => 'default.T_COMPOSITE', index_column => 'item_number,category', index_type => 'btree')"); + List remaining = + paimonTable("T_COMPOSITE").store().newIndexFileHandler().scanEntries(); + assertThat(remaining) + .isNotEmpty() + .allSatisfy( + entry -> + assertThat(entry.indexFile().globalIndexMeta().getIndexedFieldIds()) + .hasSize(1)); + assertThat( + remaining.stream() + .map(entry -> entry.indexFile().globalIndexMeta().indexFieldId()) + .distinct() + .collect(Collectors.toList())) + .containsExactlyInAnyOrder(1, 2); + } + + private void insertCompositeRows(int from, int to) { + String values = + IntStream.range(from, to) + .mapToObj( + i -> + String.format( + "(%d, '%s', %d)", + i, + (i / 10) % 2 == 0 ? "category-a" : "category-b", + i % 10)) + .collect(Collectors.joining(",")); + sql("INSERT INTO T_COMPOSITE VALUES " + values); + } + + private void updateCompositeColumn(int id, String column, Object value) throws Exception { + FileStoreTable table = paimonTable("T_COMPOSITE"); + RowType readType = + RowType.of( + table.rowType().getField("id"), + table.rowType().getField(column), + SpecialFields.ROW_ID); + ReadBuilder readBuilder = table.newReadBuilder().withReadType(readType); + List rows = new ArrayList<>(); + try (RecordReader reader = + readBuilder.newRead().createReader(readBuilder.newScan().plan())) { + reader.forEachRemaining( + row -> rows.add((InternalRow) InternalRowUtils.copy(row, readType))); + } + List rowIds = + rows.stream() + .filter(row -> row.getInt(0) == id) + .map(row -> row.getLong(2)) + .collect(Collectors.toList()); + assertThat(rowIds).hasSize(1); + long rowId = rowIds.get(0); + List ranges = + table.store().newScan().plan().files().stream() + .map(entry -> entry.file().nonNullRowIdRange()) + .filter(range -> range.from <= rowId && rowId <= range.to) + .distinct() + .collect(Collectors.toList()); + assertThat(ranges).hasSize(1); + Range range = ranges.get(0); + List updatedRows = + rows.stream() + .filter(row -> range.from <= row.getLong(2) && row.getLong(2) <= range.to) + .sorted(Comparator.comparingLong(row -> row.getLong(2))) + .collect(Collectors.toList()); + assertThat(updatedRows).hasSize((int) range.count()); + InternalRow.FieldGetter getter = InternalRow.createFieldGetter(readType.getTypeAt(1), 1); + BatchWriteBuilder builder = table.newBatchWriteBuilder(); + try (BatchTableWrite write = + builder.newWrite() + .withWriteType( + table.rowType() + .project(Collections.singletonList(column))); + BatchTableCommit commit = builder.newCommit()) { + // Data Evolution partial-column writes preserve the original data file's row range. + for (InternalRow row : updatedRows) { + write.write( + GenericRow.of(row.getInt(0) == id ? value : getter.getFieldOrNull(row))); + } + List messages = write.prepareCommit(); + List files = + messages.stream() + .flatMap( + message -> + ((CommitMessageImpl) message) + .newFilesIncrement().newFiles().stream()) + .collect(Collectors.toList()); + assertThat(files).hasSize(1); + assertThat(files.get(0).rowCount()).isEqualTo(range.count()); + for (CommitMessage message : messages) { + List newFiles = + ((CommitMessageImpl) message).newFilesIncrement().newFiles(); + if (!newFiles.isEmpty()) { + newFiles.set(0, newFiles.get(0).assignFirstRowId(range.from)); + } + } + commit.commit(messages); + } + } + @Test public void testBitmapIndex() throws Catalog.TableNotExistException { sql( diff --git a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/globalindex/sorted/SortedIndexTopoBuilder.java b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/globalindex/sorted/SortedIndexTopoBuilder.java index e1f23c9fa14f..54a361f0023c 100644 --- a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/globalindex/sorted/SortedIndexTopoBuilder.java +++ b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/globalindex/sorted/SortedIndexTopoBuilder.java @@ -45,6 +45,7 @@ import org.apache.paimon.table.source.DataSplit; import org.apache.paimon.table.source.Split; import org.apache.paimon.types.DataField; +import org.apache.paimon.types.DataType; import org.apache.paimon.types.DataTypes; import org.apache.paimon.types.RowType; import org.apache.paimon.utils.InstantiationUtil; @@ -72,6 +73,7 @@ import java.util.Map; import java.util.NoSuchElementException; import java.util.Optional; +import java.util.stream.Collectors; import static org.apache.paimon.globalindex.GlobalIndexBuilderUtils.groupSplitsByRange; import static org.apache.paimon.globalindex.GlobalIndexBuilderUtils.shardSplitsByRowRange; @@ -101,9 +103,38 @@ public List buildIndex( DataField indexField, Options options) throws IOException { + return buildIndex( + spark, + relation, + partitionPredicate, + table, + indexType, + readType, + indexField, + Collections.emptyList(), + options); + } + + @Override + public List buildIndex( + SparkSession spark, + DataSourceV2Relation relation, + PartitionPredicate partitionPredicate, + FileStoreTable table, + String indexType, + RowType readType, + DataField indexField, + List extraFields, + Options options) + throws IOException { + extraFields = extraFields == null ? Collections.emptyList() : extraFields; + List indexFields = new ArrayList<>(); + indexFields.add(indexField); + indexFields.addAll(extraFields); + List indexNames = + indexFields.stream().map(DataField::name).collect(Collectors.toList()); SortedGlobalIndexScanner indexScanner = - new SortedGlobalIndexScanner(table, indexType, options) - .withIndexFields(Collections.singletonList(indexField.name())); + new SortedGlobalIndexScanner(table, indexType, options).withIndexFields(indexNames); if (partitionPredicate != null) { indexScanner = indexScanner.withPartitionPredicate(partitionPredicate); } @@ -131,8 +162,7 @@ public List buildIndex( int maxParallelism = options.get(SortedIndexOptions.SORTED_INDEX_BUILD_MAX_PARALLELISM); List allMessages = new ArrayList<>(); - GlobalIndexer indexer = - GlobalIndexer.create(indexType, Collections.singletonList(indexField), options); + GlobalIndexer indexer = GlobalIndexer.create(indexType, indexFields, options); if (!(indexer instanceof SortedGlobalIndexer)) { throw new IllegalArgumentException( "Index algorithm " + indexType + " does not expose sorted index keys."); @@ -144,8 +174,7 @@ public List buildIndex( final int partitionKeyNum = table.partitionKeys().size(); BinaryRowSerializer binaryRowSerializer = new BinaryRowSerializer(partitionKeyNum); SortedGlobalIndexWriter indexWriter = - new SortedGlobalIndexWriter(table, indexType, options) - .withIndexFields(Collections.singletonList(indexField.name())); + new SortedGlobalIndexWriter(table, indexType, options).withIndexFields(indexNames); final byte[] serializedWriter = InstantiationUtil.serializeObject(indexWriter); if (keyExtractor.isIdentity()) { List buildTasks = new ArrayList<>(); @@ -177,7 +206,14 @@ public List buildIndex( taskInputs.add( selected.select( functions.col(taskIdField), - functions.col(indexField.name()), + extraFields.isEmpty() + ? functions.col(indexField.name()) + : functions + .struct( + indexNames.stream() + .map(functions::col) + .toArray(Column[]::new)) + .alias(indexField.name()), functions.col(SpecialFields.ROW_ID.name()))); } } @@ -420,10 +456,7 @@ private static String buildTaskIdFieldName(RowType readType) { } private static RowType normalizedReadType( - RowType readType, - String taskIdField, - DataField sourceField, - org.apache.paimon.types.DataType keyType) { + RowType readType, String taskIdField, DataField sourceField, DataType keyType) { return RowType.of( new DataField(BUILD_TASK_ID_FIELD_ID, taskIdField, DataTypes.BIGINT().notNull()), new DataField(sourceField.id(), sourceField.name(), keyType), @@ -454,8 +487,7 @@ private static class SortedTaskInput { private InternalRow next; - private SortedTaskInput( - Iterator input, org.apache.paimon.types.DataType keyType) { + private SortedTaskInput(Iterator input, DataType keyType) { this.input = input; this.keyGetter = InternalRow.createFieldGetter(keyType, 1); advance(); diff --git a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CompositeBTreeIndexProcedureTest.scala b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CompositeBTreeIndexProcedureTest.scala new file mode 100644 index 000000000000..334b68471c39 --- /dev/null +++ b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CompositeBTreeIndexProcedureTest.scala @@ -0,0 +1,152 @@ +/* + * 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.spark.procedure + +import org.apache.paimon.options.Options +import org.apache.paimon.spark.PaimonSparkTestBase +import org.apache.paimon.spark.globalindex.sorted.SortedIndexTopoBuilder +import org.apache.paimon.types.DataField + +import org.apache.spark.sql.Row + +import java.util.Collections + +import scala.collection.JavaConverters._ + +class CompositeBTreeIndexProcedureTest extends PaimonSparkTestBase { + + test("nullable extra fields preserve single column empty builds") { + withTable("T") { + sql("""CREATE TABLE T (id INT) + |TBLPROPERTIES ('bucket' = '-1', 'row-tracking.enabled' = 'true', + |'data-evolution.enabled' = 'true') + |""".stripMargin) + val table = loadTable("T") + for (extras <- Seq(null, Collections.emptyList[DataField]())) { + assert( + new SortedIndexTopoBuilder() + .buildIndex( + null, + null, + null, + table, + "btree", + table.rowType(), + table.rowType().getField("id"), + extras, + new Options()) + .isEmpty) + } + } + } + + test("composite btree creation, incremental build and component refresh") { + withTable("T", "S", "C") { + sql("""CREATE TABLE T (id INT, category STRING, item_number INT) + |TBLPROPERTIES ('bucket' = '-1', 'row-tracking.enabled' = 'true', + |'data-evolution.enabled' = 'true', 'global-index.column-update-action' = 'IGNORE', + |'sorted-index.records-per-file' = '13', 'btree-index.bloom-filter.enabled' = 'true') + |""".stripMargin) + def insert(from: Int, to: Int): Unit = { + val values = (from until to) + .map { + i => + val categoryValue = if ((i / 10) % 2 == 0) "category-a" else "category-b" + s"($i, '$categoryValue', ${i % 10})" + } + .mkString(",") + sql(s"INSERT INTO T VALUES $values") + } + def build(columns: String): Unit = { + sql( + s"CALL sys.create_global_index(table => 'test.T', index_column => '$columns', index_type => 'btree')") + } + insert(0, 40) + sql("INSERT INTO T VALUES (100, NULL, 7), (101, 'category-a', NULL)") + build("category") + build("item_number") + build("category,item_number") + val indexes = loadTable("T").store().newIndexFileHandler().scanEntries().asScala + val composite = + indexes.filter(_.indexFile().globalIndexMeta().getIndexedFieldIds().size() == 2) + assert(composite.nonEmpty) + assert( + composite.forall( + _.indexFile().globalIndexMeta().getIndexedFieldIds().asScala.toSeq == Seq(1, 2))) + assert(composite.map(_.indexFile().rowCount()).sum == 42L) + for (inReader <- Seq(false, true)) { + sql( + s"ALTER TABLE T SET TBLPROPERTIES ('global-index.query-in-reader.enabled' = '$inReader', 'scalar-index.search-mode' = 'fast')") + checkAnswer( + sql("SELECT id FROM T WHERE item_number = 7 AND category = 'category-a'"), + Seq(Row(7), Row(27))) + checkAnswer( + sql("SELECT id FROM T WHERE category = 'absent' AND item_number = 7"), + Seq.empty) + } + insert(40, 60) + build("category,item_number") + checkAnswer( + sql("SELECT id FROM T WHERE item_number = 7 AND category = 'category-a'"), + Seq(Row(7), Row(27), Row(47))) + sql("CREATE TABLE S (id INT, item_number INT)") + sql("INSERT INTO S VALUES (7, 107)") + sql( + "MERGE INTO T USING S ON T.id = S.id WHEN MATCHED THEN UPDATE SET T.item_number = S.item_number") + build("category,item_number") + checkAnswer( + sql("SELECT id FROM T WHERE category = 'category-a' AND item_number = 107"), + Seq(Row(7))) + checkAnswer( + sql("SELECT id FROM T WHERE category = 'category-a' AND item_number = 7"), + Seq(Row(27), Row(47))) + sql("CREATE TABLE C (id INT, category STRING)") + sql("INSERT INTO C VALUES (27, 'category-c')") + sql( + "MERGE INTO T USING C ON T.id = C.id WHEN MATCHED THEN UPDATE SET T.category = C.category") + build("category,item_number") + checkAnswer( + sql("SELECT id FROM T WHERE category = 'category-a' AND item_number = 7"), + Seq(Row(47))) + checkAnswer( + sql("SELECT id FROM T WHERE category = 'category-c' AND item_number = 7"), + Seq(Row(27))) + val snapshot = loadTable("T").snapshotManager().latestSnapshot().id() + build("category,item_number") + assert(loadTable("T").snapshotManager().latestSnapshot().id() == snapshot) + sql( + "CALL sys.drop_global_index(table => 'test.T', index_column => 'category,item_number', index_type => 'btree')") + assert( + loadTable("T") + .store() + .newIndexFileHandler() + .scanEntries() + .asScala + .forall(_.indexFile().globalIndexMeta().getIndexedFieldIds().size() == 1)) + assert( + loadTable("T") + .store() + .newIndexFileHandler() + .scanEntries() + .asScala + .map(_.indexFile().globalIndexMeta().indexFieldId()) + .toSet == Set(1, 2)) + } + } +}