diff --git a/paimon-core/src/main/java/org/apache/paimon/globalindex/GlobalIndexBuilderUtils.java b/paimon-core/src/main/java/org/apache/paimon/globalindex/GlobalIndexBuilderUtils.java index aea1efee5dbd..c9304eea4256 100644 --- a/paimon-core/src/main/java/org/apache/paimon/globalindex/GlobalIndexBuilderUtils.java +++ b/paimon-core/src/main/java/org/apache/paimon/globalindex/GlobalIndexBuilderUtils.java @@ -36,6 +36,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.utils.Filter; import org.apache.paimon.utils.Pair; import org.apache.paimon.utils.Range; import org.apache.paimon.utils.RangeHelper; @@ -49,6 +50,7 @@ import java.io.IOException; import java.util.ArrayList; import java.util.Arrays; +import java.util.Collection; import java.util.Collections; import java.util.Comparator; import java.util.HashMap; @@ -132,6 +134,45 @@ public static List unindexedRowRanges( return Range.sortAndMergeOverlap(dataRange.exclude(indexedRanges), true); } + /** + * A primary column can own at most one column set: the read-time scanner groups index files by + * primary field and rejects conflicting column sets, so a second index over the same primary + * column with different columns would make every filtered query fail. Creation procedures must + * reject it here; the same column set is the refresh flow and stays allowed. + */ + public static void checkPrimaryFieldNotIndexed( + FileStoreTable table, DataField indexField, List fields) { + Snapshot snapshot = table.snapshotManager().latestSnapshot(); + if (snapshot == null) { + return; + } + // the read-time grouping is type-agnostic, so the creation check scans all types + checkPrimaryFieldNotIndexed( + table.store().newIndexFileHandler().scan(snapshot, Filter.alwaysTrue()), + indexField, + extraFieldIds(fields)); + } + + static void checkPrimaryFieldNotIndexed( + Collection existingEntries, + DataField indexField, + int[] newExtraFieldIds) { + for (IndexManifestEntry entry : existingEntries) { + GlobalIndexMeta meta = entry.indexFile().globalIndexMeta(); + // re-running the creation with the same column set is the refresh flow; + // only a different column set over the same primary field is rejected + if (meta != null + && meta.indexFieldId() == indexField.id() + && !sameExtraFieldIds(meta.extraFieldIds(), newExtraFieldIds)) { + throw new IllegalArgumentException( + String.format( + "Primary field %s already owns an index with different columns; " + + "a primary column can own at most one column set.", + indexField.name())); + } + } + } + public static List currentIndexEntries( FileStoreTable table, Snapshot snapshot, diff --git a/paimon-core/src/test/java/org/apache/paimon/globalindex/GlobalIndexBuilderUtilsTest.java b/paimon-core/src/test/java/org/apache/paimon/globalindex/GlobalIndexBuilderUtilsTest.java index 063114a99611..4cbdc75d42d1 100644 --- a/paimon-core/src/test/java/org/apache/paimon/globalindex/GlobalIndexBuilderUtilsTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/globalindex/GlobalIndexBuilderUtilsTest.java @@ -24,11 +24,13 @@ import org.apache.paimon.fs.Path; import org.apache.paimon.fs.local.LocalFileIO; import org.apache.paimon.index.DataEvolutionIndexSourceMeta; +import org.apache.paimon.index.GlobalIndexMeta; import org.apache.paimon.index.IndexFileMeta; import org.apache.paimon.index.IndexPathFactory; import org.apache.paimon.io.DataFileMeta; import org.apache.paimon.io.PojoDataFileMeta; import org.apache.paimon.manifest.FileKind; +import org.apache.paimon.manifest.IndexManifestEntry; import org.apache.paimon.manifest.ManifestEntry; import org.apache.paimon.options.Options; import org.apache.paimon.stats.SimpleStats; @@ -36,6 +38,7 @@ import org.apache.paimon.table.source.Split; import org.apache.paimon.types.ArrayType; import org.apache.paimon.types.DataField; +import org.apache.paimon.types.DataTypes; import org.apache.paimon.types.FloatType; import org.apache.paimon.types.IntType; import org.apache.paimon.types.VarCharType; @@ -55,6 +58,7 @@ import java.util.UUID; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; /** Tests for {@link GlobalIndexBuilderUtils}. */ class GlobalIndexBuilderUtilsTest { @@ -366,4 +370,46 @@ private static DataFileMeta createDataFileMeta(long firstRowId, long rowCount) { null, null); } + + @Test + public void testCheckPrimaryFieldNotIndexed() { + DataField vector = new DataField(5, "vec", DataTypes.VECTOR(2, DataTypes.FLOAT())); + DataField text = new DataField(6, "txt", DataTypes.STRING()); + DataField untouched = new DataField(7, "plain", DataTypes.INT()); + GlobalIndexMeta owned = new GlobalIndexMeta(0, 10, vector.id(), new int[] {7}, null, null); + GlobalIndexMeta other = new GlobalIndexMeta(0, 10, text.id(), null, null, null); + IndexManifestEntry ownedEntry = + new IndexManifestEntry( + FileKind.ADD, + BinaryRow.EMPTY_ROW, + 0, + new IndexFileMeta("lumina", "f1", 1L, 1L, owned, null)); + IndexManifestEntry otherEntry = + new IndexManifestEntry( + FileKind.ADD, + BinaryRow.EMPTY_ROW, + 0, + new IndexFileMeta("lumina", "f2", 1L, 1L, other, null)); + + // a different column set over the same primary field is what the read-time scanner + // rejects, so creation must reject it too + assertThatThrownBy( + () -> + GlobalIndexBuilderUtils.checkPrimaryFieldNotIndexed( + Arrays.asList(ownedEntry, otherEntry), + vector, + new int[] {8})) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining( + "Primary field vec already owns an index with different columns; " + + "a primary column can own at most one column set."); + // the same column set is the refresh flow and stays allowed, including the + // single-column form against a null extras entry + GlobalIndexBuilderUtils.checkPrimaryFieldNotIndexed( + Arrays.asList(ownedEntry, otherEntry), vector, new int[] {7}); + GlobalIndexBuilderUtils.checkPrimaryFieldNotIndexed( + Arrays.asList(ownedEntry, otherEntry), text, null); + GlobalIndexBuilderUtils.checkPrimaryFieldNotIndexed( + Arrays.asList(ownedEntry, otherEntry), untouched, new int[] {8}); + } } 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 f93d19a0dcf8..93c9507b3623 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 @@ -20,6 +20,7 @@ import org.apache.paimon.flink.globalindex.GenericIndexTopoBuilder; import org.apache.paimon.flink.globalindex.SortedIndexTopoBuilder; +import org.apache.paimon.globalindex.GlobalIndexBuilderUtils; import org.apache.paimon.globalindex.GlobalIndexer; import org.apache.paimon.options.Options; import org.apache.paimon.partition.PartitionPredicate; @@ -116,6 +117,10 @@ public String[] call( Options userOptions = createUserOptions(table, options); indexType = indexType.toLowerCase(Locale.ROOT).trim(); + GlobalIndexBuilderUtils.checkPrimaryFieldNotIndexed( + table, + rowType.getField(indexColumns.get(0)), + indexColumns.stream().map(rowType::getField).collect(Collectors.toList())); if (indexColumns.size() > 1) { // Fail fast before submitting the job: index types that do not support multi-column // throw from GlobalIndexerFactory#create, which happens before any indexer side effect. diff --git a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CreateGlobalIndexProcedure.java b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CreateGlobalIndexProcedure.java index cff2fece6028..ccb8bb42eeec 100644 --- a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CreateGlobalIndexProcedure.java +++ b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CreateGlobalIndexProcedure.java @@ -18,6 +18,7 @@ package org.apache.paimon.spark.procedure; +import org.apache.paimon.globalindex.GlobalIndexBuilderUtils; import org.apache.paimon.globalindex.GlobalIndexer; import org.apache.paimon.options.Options; import org.apache.paimon.partition.PartitionPredicate; @@ -171,6 +172,8 @@ public InternalRow[] call(InternalRow args) { Options userOptions = createUserOptions(table, optionString); + GlobalIndexBuilderUtils.checkPrimaryFieldNotIndexed( + table, indexFields.get(0), indexFields); if (indexColumns.size() > 1) { // Fail fast before submitting the job: index types that do not support // multi-column throw from GlobalIndexerFactory#create, which happens