From 81a3792b60278899185be2e51fc56598a39372cb Mon Sep 17 00:00:00 2001 From: yangjie01 Date: Sat, 26 Sep 2026 18:10:12 +0800 Subject: [PATCH] [core] Reject a conflicting global index over the same primary column The read-time scanner groups global index files by primary field and rejects conflicting column sets, but neither the Flink nor the Spark create_global_index procedure checked primary-column uniqueness: creating two indexes with the same primary column and different column sets succeeded, then every filtered query or TopN on the shared column failed deterministically while grouping the index files. Reject at creation time an existing index over the same primary field with a DIFFERENT column set, scanning all index types table-wide to match the read-time grouping; re-running the creation with the same column set remains the refresh flow and stays allowed. Assisted-by: GLM-5.3 --- .../globalindex/GlobalIndexBuilderUtils.java | 41 +++++++++++++++++ .../GlobalIndexBuilderUtilsTest.java | 46 +++++++++++++++++++ .../procedure/CreateGlobalIndexProcedure.java | 5 ++ .../procedure/CreateGlobalIndexProcedure.java | 3 ++ 4 files changed, 95 insertions(+) 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