From 496e4457524f480d53ee3e1f60e6a569759dea35 Mon Sep 17 00:00:00 2001 From: JingsongLi Date: Fri, 2 Oct 2026 18:35:40 +0800 Subject: [PATCH 1/2] [core] Add composite BTree key storage and point lookups --- .../globalindex/CompositeKeySerializer.java | 117 ++++++++ .../globalindex/GlobalIndexKeyExtractor.java | 2 +- .../paimon/globalindex/GlobalIndexReader.java | 6 + .../globalindex/SortedGlobalIndexer.java | 5 +- .../globalindex/btree/BTreeGlobalIndexer.java | 26 +- .../btree/BTreeGlobalIndexerFactory.java | 12 + .../globalindex/btree/BTreeIndexWriter.java | 15 +- .../btree/LazyFilteredBTreeReader.java | 20 ++ .../btree/CompositeBTreeIndexTest.java | 273 ++++++++++++++++++ 9 files changed, 469 insertions(+), 7 deletions(-) create mode 100644 paimon-common/src/main/java/org/apache/paimon/globalindex/CompositeKeySerializer.java create mode 100644 paimon-common/src/test/java/org/apache/paimon/globalindex/btree/CompositeBTreeIndexTest.java diff --git a/paimon-common/src/main/java/org/apache/paimon/globalindex/CompositeKeySerializer.java b/paimon-common/src/main/java/org/apache/paimon/globalindex/CompositeKeySerializer.java new file mode 100644 index 000000000000..0a5f6a294486 --- /dev/null +++ b/paimon-common/src/main/java/org/apache/paimon/globalindex/CompositeKeySerializer.java @@ -0,0 +1,117 @@ +/* + * 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.globalindex; + +import org.apache.paimon.data.GenericRow; +import org.apache.paimon.data.InternalRow; +import org.apache.paimon.memory.MemorySlice; +import org.apache.paimon.memory.MemorySliceInput; +import org.apache.paimon.memory.MemorySliceOutput; +import org.apache.paimon.types.RowType; + +import java.util.Comparator; + +/** Length-delimited tuple keys, compared by their typed components with nulls first. */ +public class CompositeKeySerializer implements KeySerializer { + + private final RowType rowType; + + private final KeySerializer[] serializers; + private final InternalRow.FieldGetter[] getters; + private final Comparator[] comparators; + + @SuppressWarnings("unchecked") + public CompositeKeySerializer(RowType type) { + this.rowType = type; + int count = type.getFieldCount(); + serializers = new KeySerializer[count]; + getters = new InternalRow.FieldGetter[count]; + comparators = new Comparator[count]; + for (int i = 0; i < count; i++) { + serializers[i] = KeySerializer.create(type.getTypeAt(i)); + getters[i] = InternalRow.createFieldGetter(type.getTypeAt(i), i); + comparators[i] = serializers[i].createComparator(); + } + } + + public RowType rowType() { + return rowType; + } + + @Override + public byte[] serialize(Object key) { + InternalRow row = (InternalRow) key; + if (row.getFieldCount() != serializers.length) { + throw new IllegalArgumentException( + "Expected " + + serializers.length + + " composite key fields, but got " + + row.getFieldCount()); + } + MemorySliceOutput output = new MemorySliceOutput(32); + for (int i = 0; i < serializers.length; i++) { + Object value = getters[i].getFieldOrNull(row); + if (value == null) { + output.writeInt(-1); + } else { + byte[] bytes = serializers[i].serialize(value); + output.writeInt(bytes.length); + output.writeBytes(bytes); + } + } + return output.toSlice().copyBytes(); + } + + @Override + public Object deserialize(MemorySlice data) { + MemorySliceInput input = data.toInput(); + GenericRow row = new GenericRow(serializers.length); + for (int i = 0; i < serializers.length; i++) { + int length = input.readInt(); + if (length >= 0) { + row.setField(i, serializers[i].deserialize(input.readSlice(length))); + } else if (length != -1) { + throw new IllegalArgumentException( + "Invalid composite key component length: " + length); + } + } + if (input.available() != 0) { + throw new IllegalArgumentException("Trailing bytes in composite key"); + } + return row; + } + + @Override + public Comparator createComparator() { + return (left, right) -> { + for (int i = 0; i < serializers.length; i++) { + Object a = getters[i].getFieldOrNull((InternalRow) left); + Object b = getters[i].getFieldOrNull((InternalRow) right); + int comparison = + a == null + ? (b == null ? 0 : -1) + : b == null ? 1 : comparators[i].compare(a, b); + if (comparison != 0) { + return comparison; + } + } + return 0; + }; + } +} diff --git a/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexKeyExtractor.java b/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexKeyExtractor.java index 9a0881fa1371..d2d036e20da2 100644 --- a/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexKeyExtractor.java +++ b/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexKeyExtractor.java @@ -25,7 +25,7 @@ import java.io.IOException; import java.io.Serializable; -/** Extracts zero or more normalized index keys from one source-column value. */ +/** Extracts zero or more normalized index keys from one source value or projected tuple. */ public interface GlobalIndexKeyExtractor extends Serializable { /** Type of the normalized keys emitted by {@link #extract(Object, KeyConsumer)}. */ diff --git a/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexReader.java b/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexReader.java index f5d6ab9fad65..4506713b9176 100644 --- a/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexReader.java +++ b/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexReader.java @@ -36,6 +36,12 @@ public interface GlobalIndexReader extends FunctionVisitor>>, Closeable { + /** Point lookup of a full composite key, with literals in index column order. */ + default CompletableFuture> visitCompositeEqual( + List literals) { + return CompletableFuture.completedFuture(Optional.empty()); + } + @Override default CompletableFuture> visitIsNaN(FieldRef fieldRef) { return CompletableFuture.completedFuture(Optional.empty()); diff --git a/paimon-common/src/main/java/org/apache/paimon/globalindex/SortedGlobalIndexer.java b/paimon-common/src/main/java/org/apache/paimon/globalindex/SortedGlobalIndexer.java index 1f3bc4e6034d..f706060a0955 100644 --- a/paimon-common/src/main/java/org/apache/paimon/globalindex/SortedGlobalIndexer.java +++ b/paimon-common/src/main/java/org/apache/paimon/globalindex/SortedGlobalIndexer.java @@ -21,6 +21,9 @@ /** A global indexer whose normalized keys must be sorted before they are written. */ public interface SortedGlobalIndexer extends GlobalIndexer { - /** Defines how source-column values are normalized into the keys consumed by the writer. */ + /** + * Defines how source values or projected tuples are normalized into the keys consumed by the + * writer. + */ GlobalIndexKeyExtractor keyExtractor(); } diff --git a/paimon-common/src/main/java/org/apache/paimon/globalindex/btree/BTreeGlobalIndexer.java b/paimon-common/src/main/java/org/apache/paimon/globalindex/btree/BTreeGlobalIndexer.java index 774b22bf1fce..e4c632d02b89 100644 --- a/paimon-common/src/main/java/org/apache/paimon/globalindex/btree/BTreeGlobalIndexer.java +++ b/paimon-common/src/main/java/org/apache/paimon/globalindex/btree/BTreeGlobalIndexer.java @@ -20,6 +20,7 @@ import org.apache.paimon.compression.BlockCompressionFactory; import org.apache.paimon.compression.CompressOptions; +import org.apache.paimon.globalindex.CompositeKeySerializer; import org.apache.paimon.globalindex.GlobalIndexIOMeta; import org.apache.paimon.globalindex.GlobalIndexKeyExtractor; import org.apache.paimon.globalindex.GlobalIndexReader; @@ -31,6 +32,8 @@ import org.apache.paimon.io.cache.CacheManager; import org.apache.paimon.options.Options; import org.apache.paimon.types.DataField; +import org.apache.paimon.types.DataType; +import org.apache.paimon.types.RowType; import org.apache.paimon.utils.BloomFilter; import org.apache.paimon.utils.LazyField; import org.apache.paimon.utils.Range; @@ -38,6 +41,8 @@ import javax.annotation.Nullable; import java.io.IOException; +import java.util.ArrayList; +import java.util.Collections; import java.util.List; import java.util.concurrent.ExecutorService; @@ -75,8 +80,25 @@ public class BTreeGlobalIndexer implements SortedGlobalIndexer { private final LazyField cacheManager; public BTreeGlobalIndexer(DataField dataField, Options options) { - this.keySerializer = KeySerializer.create(dataField.type()); - this.keyExtractor = GlobalIndexKeyExtractor.identity(dataField.type()); + this(dataField, Collections.emptyList(), options); + } + + public BTreeGlobalIndexer(DataField dataField, List extraFields, Options options) { + List fields = new ArrayList<>(); + fields.add(dataField); + fields.addAll(extraFields); + for (DataField field : fields) { + if (field.type() instanceof RowType) { + throw new UnsupportedOperationException( + "BTree index columns must have scalar types: " + field.name()); + } + } + DataType keyType = extraFields.isEmpty() ? dataField.type() : new RowType(fields); + this.keySerializer = + extraFields.isEmpty() + ? KeySerializer.create(keyType) + : new CompositeKeySerializer((RowType) keyType); + this.keyExtractor = GlobalIndexKeyExtractor.identity(keyType); this.options = options; this.fallbackScanMaxSize = options.get(BTreeIndexOptions.BTREE_INDEX_FALLBACK_SCAN_MAX_SIZE).getBytes(); diff --git a/paimon-common/src/main/java/org/apache/paimon/globalindex/btree/BTreeGlobalIndexerFactory.java b/paimon-common/src/main/java/org/apache/paimon/globalindex/btree/BTreeGlobalIndexerFactory.java index 98d5c36241bd..2cb788dc4581 100644 --- a/paimon-common/src/main/java/org/apache/paimon/globalindex/btree/BTreeGlobalIndexerFactory.java +++ b/paimon-common/src/main/java/org/apache/paimon/globalindex/btree/BTreeGlobalIndexerFactory.java @@ -27,6 +27,7 @@ import org.apache.paimon.predicate.Predicate; import org.apache.paimon.types.DataField; +import java.util.Collections; import java.util.List; /** The {@link GlobalIndexerFactory} for btree index. */ @@ -45,10 +46,21 @@ public List selectFiles( List extraFields, Predicate predicate, List files) { + if (extraFields != null && !extraFields.isEmpty()) { + // Scalar predicates cannot safely prune tuple metadata. + return files; + } return SortedFileMetaSelector.selectFiles( predicate, files, KeySerializer.create(indexField.type())); } + @Override + public GlobalIndexer create( + DataField indexField, List extraFields, Options options) { + return new BTreeGlobalIndexer( + indexField, extraFields == null ? Collections.emptyList() : extraFields, options); + } + @Override public GlobalIndexer create(DataField dataField, Options options) { return new BTreeGlobalIndexer(dataField, options); diff --git a/paimon-common/src/main/java/org/apache/paimon/globalindex/btree/BTreeIndexWriter.java b/paimon-common/src/main/java/org/apache/paimon/globalindex/btree/BTreeIndexWriter.java index fecf600f56ac..238e70df4166 100644 --- a/paimon-common/src/main/java/org/apache/paimon/globalindex/btree/BTreeIndexWriter.java +++ b/paimon-common/src/main/java/org/apache/paimon/globalindex/btree/BTreeIndexWriter.java @@ -20,6 +20,7 @@ import org.apache.paimon.compression.BlockCompressionFactory; import org.apache.paimon.fs.PositionOutputStream; +import org.apache.paimon.globalindex.CompositeKeySerializer; import org.apache.paimon.globalindex.GlobalIndexSingleColumnWriter; import org.apache.paimon.globalindex.KeySerializer; import org.apache.paimon.globalindex.ResultEntry; @@ -148,19 +149,27 @@ public void write(@Nullable Object key, long rowId) { return; } - if (lastKey != null && comparator.compare(key, lastKey) != 0) { + boolean differentKey = lastKey == null || comparator.compare(key, lastKey) != 0; + if (lastKey != null && differentKey) { try { flush(); } catch (IOException e) { throw new RuntimeException("Error in writing btree index files.", e); } } - lastKey = key; + if (differentKey || !(keySerializer instanceof CompositeKeySerializer)) { + // Sorted engine iterators can reuse the row backing a tuple key. + lastKey = + keySerializer instanceof CompositeKeySerializer + ? keySerializer.deserialize( + MemorySlice.wrap(keySerializer.serialize(key))) + : key; + } currentRowIds.add(rowId); // update stats if (firstKey == null) { - firstKey = key; + firstKey = lastKey; } } diff --git a/paimon-common/src/main/java/org/apache/paimon/globalindex/btree/LazyFilteredBTreeReader.java b/paimon-common/src/main/java/org/apache/paimon/globalindex/btree/LazyFilteredBTreeReader.java index 03bbb8a68d6b..baadcc3c7119 100644 --- a/paimon-common/src/main/java/org/apache/paimon/globalindex/btree/LazyFilteredBTreeReader.java +++ b/paimon-common/src/main/java/org/apache/paimon/globalindex/btree/LazyFilteredBTreeReader.java @@ -18,6 +18,8 @@ package org.apache.paimon.globalindex.btree; +import org.apache.paimon.data.GenericRow; +import org.apache.paimon.globalindex.CompositeKeySerializer; import org.apache.paimon.globalindex.GlobalIndexIOMeta; import org.apache.paimon.globalindex.GlobalIndexResult; import org.apache.paimon.globalindex.KeySerializer; @@ -37,6 +39,7 @@ import java.io.IOException; import java.util.Comparator; import java.util.List; +import java.util.Objects; import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutorService; @@ -109,6 +112,23 @@ private Pair fullRangeBounds(List files) { return remaining == 0 ? Pair.of(min, max) : null; } + @Override + public CompletableFuture> visitCompositeEqual( + List literals) { + if (!(keySerializer instanceof CompositeKeySerializer)) { + return CompletableFuture.completedFuture(Optional.empty()); + } + int fieldCount = ((CompositeKeySerializer) keySerializer).rowType().getFieldCount(); + if (literals.size() != fieldCount) { + throw new IllegalArgumentException( + "Expected " + fieldCount + " composite key fields, but got " + literals.size()); + } + if (literals.stream().anyMatch(Objects::isNull)) { + return CompletableFuture.completedFuture(Optional.of(GlobalIndexResult.createEmpty())); + } + return visitEqual((FieldRef) null, GenericRow.of(literals.toArray())); + } + @Override public CompletableFuture> visitEqual( FieldRef fieldRef, Object literal) { diff --git a/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/CompositeBTreeIndexTest.java b/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/CompositeBTreeIndexTest.java new file mode 100644 index 000000000000..d5fccfe14e1e --- /dev/null +++ b/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/CompositeBTreeIndexTest.java @@ -0,0 +1,273 @@ +/* + * 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.globalindex.btree; + +import org.apache.paimon.data.BinaryString; +import org.apache.paimon.data.GenericRow; +import org.apache.paimon.fs.Path; +import org.apache.paimon.fs.PositionOutputStream; +import org.apache.paimon.fs.local.LocalFileIO; +import org.apache.paimon.globalindex.CompositeKeySerializer; +import org.apache.paimon.globalindex.GlobalIndexIOMeta; +import org.apache.paimon.globalindex.GlobalIndexReader; +import org.apache.paimon.globalindex.GlobalIndexSingleColumnWriter; +import org.apache.paimon.globalindex.GlobalIndexer; +import org.apache.paimon.globalindex.KeySerializer; +import org.apache.paimon.globalindex.ResultEntry; +import org.apache.paimon.globalindex.SortedGlobalIndexer; +import org.apache.paimon.globalindex.io.GlobalIndexFileWriter; +import org.apache.paimon.memory.MemorySlice; +import org.apache.paimon.options.Options; +import org.apache.paimon.predicate.PredicateBuilder; +import org.apache.paimon.types.DataField; +import org.apache.paimon.types.DataTypes; +import org.apache.paimon.types.RowType; +import org.apache.paimon.utils.Range; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; + +import java.util.Arrays; +import java.util.Collections; +import java.util.Comparator; +import java.util.List; +import java.util.UUID; +import java.util.concurrent.ExecutorService; + +import static org.apache.paimon.shade.guava30.com.google.common.util.concurrent.MoreExecutors.newDirectExecutorService; +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +/** Tests for typed composite BTree keys. */ +class CompositeBTreeIndexTest { + + @TempDir java.nio.file.Path tempPath; + + @ParameterizedTest + @ValueSource(ints = {1, 2}) + void testMutableTupleKeysPostingListsAndLocalRanges(int version) throws Exception { + RowType type = + new RowType( + Arrays.asList( + new DataField(10, "category", DataTypes.STRING()), + new DataField(20, "item_number", DataTypes.INT()), + new DataField(30, "tag", DataTypes.STRING()))); + Options options = new Options(); + options.set(BTreeIndexOptions.BTREE_INDEX_FILE_VERSION, version); + options.set(BTreeIndexOptions.BTREE_INDEX_BLOOM_FILTER_ENABLED, true); + options.set(BTreeIndexOptions.BTREE_INDEX_COMPRESSION, "lz4"); + GlobalIndexer indexer = + GlobalIndexer.create( + "btree", type.getFields().get(0), type.getFields().subList(1, 3), options); + LocalFileIO io = LocalFileIO.create(); + Path directory = new Path(tempPath.toUri()); + GlobalIndexFileWriter files = + new GlobalIndexFileWriter() { + @Override + public String newFileName(String prefix) { + return prefix + UUID.randomUUID(); + } + + @Override + public PositionOutputStream newOutputStream(String name) + throws java.io.IOException { + return io.newOutputStream(new Path(directory, name), false); + } + }; + GenericRow reused = row("category-a", -1, ""); + GlobalIndexSingleColumnWriter writer = + (GlobalIndexSingleColumnWriter) indexer.createWriter(files); + writer.write(reused, 0); + reused.setField(1, 7); + writer.write(reused, 1); + reused.setField(2, BinaryString.fromString("tag")); + writer.write(reused, 2); + writer.write(reused, 3); + reused.setField(0, BinaryString.fromString("category-b")); + writer.write(reused, 4); + ResultEntry result = writer.finish().get(0); + Path path = new Path(directory, result.fileName()); + GlobalIndexIOMeta meta = + new GlobalIndexIOMeta(path, io.getFileSize(path), result.rowCount(), result.meta()); + assertThat( + new BTreeGlobalIndexerFactory() + .selectFiles( + type.getFields().get(0), + type.getFields().subList(1, 3), + new PredicateBuilder(type) + .equal(0, BinaryString.fromString("category-a")), + Collections.singletonList(meta))) + .containsExactly(meta); + ExecutorService executor = newDirectExecutorService(); + try (GlobalIndexReader reader = + indexer.createReader( + file -> io.newInputStream(file.filePath()), + Collections.singletonList(meta), + 5, + null, + executor)) { + for (List invalid : + Arrays.asList( + Arrays.asList(BinaryString.fromString("category-a"), 7), + Arrays.asList( + BinaryString.fromString("category-a"), + 7, + BinaryString.fromString("tag"), + BinaryString.fromString("extra")), + Arrays.asList(null, 7, BinaryString.fromString("tag"), null))) { + assertThatThrownBy(() -> reader.visitCompositeEqual(invalid)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("Expected 3 composite key fields"); + } + assertThat( + reader.visitCompositeEqual( + Arrays.asList(null, 7, BinaryString.fromString("tag"))) + .get() + .get() + .results() + .toRangeList()) + .isEmpty(); + assertThat( + reader.visitCompositeEqual( + Arrays.asList( + BinaryString.fromString("category-a"), + -1, + BinaryString.fromString(""))) + .get() + .get() + .results() + .toRangeList()) + .containsExactly(new Range(0, 0)); + assertThat( + reader.visitCompositeEqual( + Arrays.asList( + BinaryString.fromString("category-a"), + 7, + BinaryString.fromString("tag"))) + .get() + .get() + .results() + .toRangeList()) + .containsExactly(new Range(2, 3)); + assertThat( + reader.visitCompositeEqual( + Arrays.asList( + BinaryString.fromString("category-a"), + 8, + BinaryString.fromString("tag"))) + .get() + .get() + .results() + .isEmpty()) + .isTrue(); + } + try (GlobalIndexReader reader = + indexer.createReader( + file -> io.newInputStream(file.filePath()), + Collections.singletonList(meta), + 5, + Collections.singletonList(new Range(3, 3)), + executor)) { + assertThat( + reader.visitCompositeEqual( + Arrays.asList( + BinaryString.fromString("category-a"), + 7, + BinaryString.fromString("tag"))) + .get() + .get() + .results() + .toRangeList()) + .containsExactly(new Range(3, 3)); + } finally { + executor.shutdownNow(); + } + } + + @Test + void testCompositeKeysPreserveTypesBoundariesAndNulls() { + RowType type = + new RowType( + Arrays.asList( + new DataField(10, "category", DataTypes.STRING()), + new DataField(20, "item_number", DataTypes.INT()), + new DataField(30, "tag", DataTypes.STRING()))); + SortedGlobalIndexer indexer = + (SortedGlobalIndexer) + GlobalIndexer.create( + "btree", + type.getFields().get(0), + type.getFields().subList(1, 3), + new Options()); + KeySerializer serializer = + new CompositeKeySerializer((RowType) indexer.keyExtractor().keyType()); + Comparator comparator = serializer.createComparator(); + for (GenericRow invalid : + Arrays.asList( + GenericRow.of(BinaryString.fromString("category-a"), 7), + GenericRow.of( + BinaryString.fromString("category-a"), + 7, + BinaryString.fromString("tag"), + BinaryString.fromString("extra")))) { + assertThatThrownBy(() -> serializer.serialize(invalid)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("Expected 3 composite key fields"); + } + GenericRow first = row("category-a", -1, "a\u0000b"); + GenericRow second = row("category-a", 107, ""); + GenericRow nullable = row(null, 107, null); + for (GenericRow key : Arrays.asList(first, second, nullable)) { + Object restored = serializer.deserialize(MemorySlice.wrap(serializer.serialize(key))); + assertThat(restored).isEqualTo(key); + assertThat(comparator.compare(restored, key)).isZero(); + } + assertThat(comparator.compare(first, second)).isNegative(); + assertThat(comparator.compare(nullable, first)).isNegative(); + assertThat(serializer.serialize(row("a", 1, "bc"))) + .isNotEqualTo(serializer.serialize(row("ab", 1, "c"))); + assertThat( + GlobalIndexer.create( + "btree", + type.getFields().get(0), + Collections.emptyList(), + new Options())) + .isInstanceOf(BTreeGlobalIndexer.class); + assertThat(GlobalIndexer.create("btree", type.getFields().get(0), null, new Options())) + .isInstanceOf(BTreeGlobalIndexer.class); + } + + @Test + void testCompositeSupportDoesNotEnablePhysicalRowIndexes() { + DataField field = new DataField(40, "nested", RowType.of(DataTypes.INT())); + for (String indexType : Arrays.asList("btree", "bitmap")) { + assertThatThrownBy(() -> GlobalIndexer.create(indexType, field, new Options())) + .isInstanceOf(UnsupportedOperationException.class); + } + } + + private GenericRow row(String category, int itemNumber, String tag) { + return GenericRow.of( + category == null ? null : BinaryString.fromString(category), + itemNumber, + tag == null ? null : BinaryString.fromString(tag)); + } +} From 365abedce2106777fbdc3445e8731e238ef17c56 Mon Sep 17 00:00:00 2001 From: JingsongLi Date: Fri, 2 Oct 2026 18:54:41 +0800 Subject: [PATCH 2/2] [core] Unify global index factories on ordered field lists --- .../paimon/globalindex/GlobalIndexer.java | 11 +--- .../globalindex/GlobalIndexerFactory.java | 24 ++------- .../GlobalIndexerFactoryUtils.java | 7 +-- .../bitmap/BitmapGlobalIndexerFactory.java | 14 +++--- .../MultiValueGlobalIndexerFactory.java | 9 +++- .../globalindex/btree/BTreeGlobalIndexer.java | 18 +++---- .../btree/BTreeGlobalIndexerFactory.java | 21 ++------ .../fmindex/FMGlobalIndexerFactory.java | 9 +++- .../globalindex/SortedIndexFileMetaTest.java | 8 ++- .../btree/AbstractIndexReaderTest.java | 5 +- .../btree/BTreeBloomFilterTest.java | 3 +- .../btree/BTreeIndexReaderCloseTest.java | 5 +- .../btree/BTreeThreadSafetyTest.java | 3 +- .../btree/CompositeBTreeIndexTest.java | 50 ++++++++++++------- .../LazyFilteredBTreeIndexReaderTest.java | 19 +++++-- .../fmindex/FMGlobalIndexTest.java | 2 +- .../TestFullTextGlobalIndexerFactory.java | 23 +++------ .../TestMultiFieldVectorGlobalIndexer.java | 4 +- .../TestVectorGlobalIndexerFactory.java | 15 ++---- .../DataEvolutionGlobalIndexScanner.java | 7 ++- .../globalindex/GlobalIndexBuilderUtils.java | 16 +----- .../paimon/globalindex/GlobalIndexQuery.java | 12 +++-- .../sorted/SortedGlobalIndexWriter.java | 7 ++- .../index/pkfulltext/PkFullTextIndexFile.java | 8 ++- .../index/pksorted/PkSortedIndexBuilder.java | 4 +- .../index/pksorted/PkSortedIndexFile.java | 5 +- .../pkvector/PkVectorAnnSegmentFile.java | 4 +- .../pkvector/PkVectorAnnSegmentSearcher.java | 6 ++- .../paimon/schema/SchemaValidation.java | 8 +-- .../AbstractDataEvolutionVectorRead.java | 7 +-- .../source/DataEvolutionFullTextRead.java | 8 +-- .../table/source/PrimaryKeyFullTextRead.java | 4 +- .../source/PrimaryKeySortedIndexScan.java | 2 +- .../table/source/RawFullTextReadImpl.java | 3 +- .../globalindex/GlobalIndexQueryTest.java | 17 +++---- .../pksorted/PkSortedIndexBuilderTest.java | 4 +- .../table/BitmapGlobalIndexTableTest.java | 5 +- .../paimon/table/IndexQuerySplitTest.java | 5 +- .../source/FullTextSearchBuilderTest.java | 22 +++++--- .../table/source/VectorSearchBuilderTest.java | 20 ++++---- .../VectorSearchRowFilterExactnessTest.java | 12 +++-- .../index/ESIndexGlobalIndexerFactory.java | 18 +------ .../globalindex/GenericIndexTopoBuilder.java | 2 +- .../globalindex/SortedIndexTopoBuilder.java | 4 +- .../procedure/CreateGlobalIndexProcedure.java | 9 ++-- .../VectorSearchProcedureITCase.java | 5 +- .../NativeFullTextGlobalIndexerFactory.java | 8 ++- .../index/JavaPyNativeFullTextE2ETest.java | 2 +- ...ativeFullTextGlobalIndexerFactoryTest.java | 4 +- .../index/NativeFullTextRowFilterTest.java | 4 +- .../NativePrimaryKeyFullTextIndexTest.java | 6 ++- .../LuminaVectorGlobalIndexerFactory.java | 9 +++- .../lumina/index/JavaPyLuminaE2ETest.java | 9 ++-- .../DefaultGlobalIndexBuilder.java | 2 +- .../sorted/SortedIndexTopoBuilder.java | 3 +- .../procedure/CreateGlobalIndexProcedure.java | 6 +-- .../read/SparkDataEvolutionVectorRead.java | 5 +- .../NativeVectorGlobalIndexerFactory.java | 8 ++- .../java/org/apache/paimon/JavaPyE2ETest.java | 4 +- 59 files changed, 291 insertions(+), 253 deletions(-) diff --git a/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexer.java b/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexer.java index 45a4724bb193..9792602393dd 100644 --- a/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexer.java +++ b/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexer.java @@ -49,14 +49,7 @@ GlobalIndexReader createReader( @Nullable List rowRanges, ExecutorService executor); - static GlobalIndexer create(String type, DataField indexField, Options options) { - GlobalIndexerFactory globalIndexerFactory = GlobalIndexerFactoryUtils.load(type); - return globalIndexerFactory.create(indexField, options); - } - - static GlobalIndexer create( - String type, DataField indexField, List extraFields, Options options) { - GlobalIndexerFactory globalIndexerFactory = GlobalIndexerFactoryUtils.load(type); - return globalIndexerFactory.create(indexField, extraFields, options); + static GlobalIndexer create(String type, List indexFields, Options options) { + return GlobalIndexerFactoryUtils.load(type).create(indexFields, options); } } diff --git a/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexerFactory.java b/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexerFactory.java index 6170515e6df4..2e2d9166e689 100644 --- a/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexerFactory.java +++ b/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexerFactory.java @@ -39,28 +39,10 @@ default boolean supportsFullTextSearch() { * match; the default keeps all files when no safe pruning is available. */ default List selectFiles( - DataField indexField, - List extraFields, - Predicate predicate, - List files) { + List indexFields, Predicate predicate, List files) { return files; } - GlobalIndexer create(DataField indexField, Options options); - - /** - * Creates an indexer over a primary column plus optional extra columns. {@code indexField} is - * the primary column; {@code extraFields} holds the remaining columns and is empty for a - * single-column index. - */ - default GlobalIndexer create( - DataField indexField, List extraFields, Options options) { - if (extraFields != null && !extraFields.isEmpty()) { - throw new UnsupportedOperationException( - String.format( - "Index type '%s' does not support multi-column index, got extra columns: %s", - identifier(), extraFields)); - } - return create(indexField, options); - } + /** Creates an indexer over a non-empty list of columns in index order. */ + GlobalIndexer create(List indexFields, Options options); } diff --git a/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexerFactoryUtils.java b/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexerFactoryUtils.java index a793090d1552..6e045ec79c6e 100644 --- a/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexerFactoryUtils.java +++ b/paimon-common/src/main/java/org/apache/paimon/globalindex/GlobalIndexerFactoryUtils.java @@ -61,13 +61,10 @@ public static GlobalIndexerFactory load(String type) { /** Keep unrecognized index types for the reader to resolve at execution time. */ public static List selectFiles( String type, - DataField indexField, - List extraFields, + List indexFields, Predicate predicate, List files) { GlobalIndexerFactory factory = factories.get(type); - return factory == null - ? files - : factory.selectFiles(indexField, extraFields, predicate, files); + return factory == null ? files : factory.selectFiles(indexFields, predicate, files); } } diff --git a/paimon-common/src/main/java/org/apache/paimon/globalindex/bitmap/BitmapGlobalIndexerFactory.java b/paimon-common/src/main/java/org/apache/paimon/globalindex/bitmap/BitmapGlobalIndexerFactory.java index e573850bb6e9..c7ab53bd3e3d 100644 --- a/paimon-common/src/main/java/org/apache/paimon/globalindex/bitmap/BitmapGlobalIndexerFactory.java +++ b/paimon-common/src/main/java/org/apache/paimon/globalindex/bitmap/BitmapGlobalIndexerFactory.java @@ -41,16 +41,18 @@ public String identifier() { @Override public List selectFiles( - DataField indexField, - List extraFields, - Predicate predicate, - List files) { + List indexFields, Predicate predicate, List files) { return SortedFileMetaSelector.selectFiles( - predicate, files, KeySerializer.create(indexField.type())); + predicate, files, KeySerializer.create(indexFields.get(0).type())); } @Override - public GlobalIndexer create(DataField dataField, Options options) { + public GlobalIndexer create(List indexFields, Options options) { + if (indexFields.size() != 1) { + throw new UnsupportedOperationException( + "Index type '" + identifier() + "' requires exactly one index field."); + } + DataField dataField = indexFields.get(0); return new BitmapGlobalIndexer(dataField, options); } } diff --git a/paimon-common/src/main/java/org/apache/paimon/globalindex/bitmap/MultiValueGlobalIndexerFactory.java b/paimon-common/src/main/java/org/apache/paimon/globalindex/bitmap/MultiValueGlobalIndexerFactory.java index dba47d9d3c2b..f3cde2e9c9b7 100644 --- a/paimon-common/src/main/java/org/apache/paimon/globalindex/bitmap/MultiValueGlobalIndexerFactory.java +++ b/paimon-common/src/main/java/org/apache/paimon/globalindex/bitmap/MultiValueGlobalIndexerFactory.java @@ -23,6 +23,8 @@ import org.apache.paimon.options.Options; import org.apache.paimon.types.DataField; +import java.util.List; + /** Factory for bitmap-backed multivalue indexes on array columns. */ public class MultiValueGlobalIndexerFactory implements GlobalIndexerFactory { @@ -34,7 +36,12 @@ public String identifier() { } @Override - public GlobalIndexer create(DataField indexField, Options options) { + public GlobalIndexer create(List indexFields, Options options) { + if (indexFields.size() != 1) { + throw new UnsupportedOperationException( + "Index type '" + identifier() + "' requires exactly one index field."); + } + DataField indexField = indexFields.get(0); return new MultiValueGlobalIndexer(indexField, options); } } diff --git a/paimon-common/src/main/java/org/apache/paimon/globalindex/btree/BTreeGlobalIndexer.java b/paimon-common/src/main/java/org/apache/paimon/globalindex/btree/BTreeGlobalIndexer.java index e4c632d02b89..451b292f4c2c 100644 --- a/paimon-common/src/main/java/org/apache/paimon/globalindex/btree/BTreeGlobalIndexer.java +++ b/paimon-common/src/main/java/org/apache/paimon/globalindex/btree/BTreeGlobalIndexer.java @@ -41,11 +41,11 @@ import javax.annotation.Nullable; import java.io.IOException; -import java.util.ArrayList; -import java.util.Collections; import java.util.List; import java.util.concurrent.ExecutorService; +import static org.apache.paimon.utils.Preconditions.checkArgument; + /** * The {@link GlobalIndexer} for btree index. We do not build a B-tree directly in memory, instead, * we form a logical B-tree via multi-level metadata over SST files that store the actual data, as @@ -79,23 +79,17 @@ public class BTreeGlobalIndexer implements SortedGlobalIndexer { private final long fallbackScanMaxSize; private final LazyField cacheManager; - public BTreeGlobalIndexer(DataField dataField, Options options) { - this(dataField, Collections.emptyList(), options); - } - - public BTreeGlobalIndexer(DataField dataField, List extraFields, Options options) { - List fields = new ArrayList<>(); - fields.add(dataField); - fields.addAll(extraFields); + public BTreeGlobalIndexer(List fields, Options options) { + checkArgument(!fields.isEmpty(), "BTree index requires at least one field."); for (DataField field : fields) { if (field.type() instanceof RowType) { throw new UnsupportedOperationException( "BTree index columns must have scalar types: " + field.name()); } } - DataType keyType = extraFields.isEmpty() ? dataField.type() : new RowType(fields); + DataType keyType = fields.size() == 1 ? fields.get(0).type() : new RowType(fields); this.keySerializer = - extraFields.isEmpty() + fields.size() == 1 ? KeySerializer.create(keyType) : new CompositeKeySerializer((RowType) keyType); this.keyExtractor = GlobalIndexKeyExtractor.identity(keyType); diff --git a/paimon-common/src/main/java/org/apache/paimon/globalindex/btree/BTreeGlobalIndexerFactory.java b/paimon-common/src/main/java/org/apache/paimon/globalindex/btree/BTreeGlobalIndexerFactory.java index 2cb788dc4581..d0b1140fbad0 100644 --- a/paimon-common/src/main/java/org/apache/paimon/globalindex/btree/BTreeGlobalIndexerFactory.java +++ b/paimon-common/src/main/java/org/apache/paimon/globalindex/btree/BTreeGlobalIndexerFactory.java @@ -27,7 +27,6 @@ import org.apache.paimon.predicate.Predicate; import org.apache.paimon.types.DataField; -import java.util.Collections; import java.util.List; /** The {@link GlobalIndexerFactory} for btree index. */ @@ -42,27 +41,17 @@ public String identifier() { @Override public List selectFiles( - DataField indexField, - List extraFields, - Predicate predicate, - List files) { - if (extraFields != null && !extraFields.isEmpty()) { + List indexFields, Predicate predicate, List files) { + if (indexFields.size() > 1) { // Scalar predicates cannot safely prune tuple metadata. return files; } return SortedFileMetaSelector.selectFiles( - predicate, files, KeySerializer.create(indexField.type())); + predicate, files, KeySerializer.create(indexFields.get(0).type())); } @Override - public GlobalIndexer create( - DataField indexField, List extraFields, Options options) { - return new BTreeGlobalIndexer( - indexField, extraFields == null ? Collections.emptyList() : extraFields, options); - } - - @Override - public GlobalIndexer create(DataField dataField, Options options) { - return new BTreeGlobalIndexer(dataField, options); + public GlobalIndexer create(List indexFields, Options options) { + return new BTreeGlobalIndexer(indexFields, options); } } diff --git a/paimon-common/src/main/java/org/apache/paimon/globalindex/fmindex/FMGlobalIndexerFactory.java b/paimon-common/src/main/java/org/apache/paimon/globalindex/fmindex/FMGlobalIndexerFactory.java index 4b5e37c0567a..3ac045fc3dde 100644 --- a/paimon-common/src/main/java/org/apache/paimon/globalindex/fmindex/FMGlobalIndexerFactory.java +++ b/paimon-common/src/main/java/org/apache/paimon/globalindex/fmindex/FMGlobalIndexerFactory.java @@ -24,6 +24,8 @@ import org.apache.paimon.types.DataField; import org.apache.paimon.types.DataTypeFamily; +import java.util.List; + import static org.apache.paimon.utils.Preconditions.checkArgument; /** Factory for the exact partitioned FM contains index. */ @@ -37,7 +39,12 @@ public String identifier() { } @Override - public GlobalIndexer create(DataField dataField, Options options) { + public GlobalIndexer create(List indexFields, Options options) { + if (indexFields.size() != 1) { + throw new UnsupportedOperationException( + "Index type '" + identifier() + "' requires exactly one index field."); + } + DataField dataField = indexFields.get(0); checkArgument( dataField.type().is(DataTypeFamily.CHARACTER_STRING), "FM index requires a character string column, but field '%s' is %s.", diff --git a/paimon-common/src/test/java/org/apache/paimon/globalindex/SortedIndexFileMetaTest.java b/paimon-common/src/test/java/org/apache/paimon/globalindex/SortedIndexFileMetaTest.java index 537d6724028d..960b98bcf4b0 100644 --- a/paimon-common/src/test/java/org/apache/paimon/globalindex/SortedIndexFileMetaTest.java +++ b/paimon-common/src/test/java/org/apache/paimon/globalindex/SortedIndexFileMetaTest.java @@ -34,6 +34,7 @@ import org.junit.jupiter.api.io.TempDir; import java.io.IOException; +import java.util.Collections; import java.util.List; import static org.assertj.core.api.Assertions.assertThat; @@ -157,8 +158,11 @@ public PositionOutputStream newOutputStream(String fileName) }; GlobalIndexSingleColumnWriter indexWriter = new BTreeGlobalIndexer( - new DataField( - 1, "testField", new VarCharType(VarCharType.MAX_LENGTH)), + Collections.singletonList( + new DataField( + 1, + "testField", + new VarCharType(VarCharType.MAX_LENGTH))), new Options()) .createWriter(fileWriter); for (int index = 0; index < keys.length; index++) { diff --git a/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/AbstractIndexReaderTest.java b/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/AbstractIndexReaderTest.java index c00f26bff5f8..c46472305c04 100644 --- a/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/AbstractIndexReaderTest.java +++ b/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/AbstractIndexReaderTest.java @@ -143,7 +143,10 @@ public PositionOutputStream newOutputStream(String fileName) new Path(new Path(tempPath.toUri()), meta.filePath())); options = new Options(); options.set(BTreeIndexOptions.BTREE_INDEX_CACHE_SIZE, MemorySize.ofMebiBytes(8)); - globalIndexer = new BTreeGlobalIndexer(new DataField(1, "testField", dataType), options); + globalIndexer = + new BTreeGlobalIndexer( + Collections.singletonList(new DataField(1, "testField", dataType)), + options); keySerializer = KeySerializer.create(dataType); comparator = keySerializer.createComparator(); diff --git a/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/BTreeBloomFilterTest.java b/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/BTreeBloomFilterTest.java index 67f2985b58dd..1db4a2c22296 100644 --- a/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/BTreeBloomFilterTest.java +++ b/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/BTreeBloomFilterTest.java @@ -115,7 +115,8 @@ private void assertMissingLookupReadsOnlyBloom( private Fixture writeIndex(Options options) throws IOException { ByteArrayGlobalIndexFileWriter fileWriter = new ByteArrayGlobalIndexFileWriter(); BTreeGlobalIndexer indexer = - new BTreeGlobalIndexer(new DataField(0, "k", new IntType()), options); + new BTreeGlobalIndexer( + Collections.singletonList(new DataField(0, "k", new IntType())), options); GlobalIndexSingleColumnWriter writer = indexer.createWriter(fileWriter); for (int i = 0; i < ENTRY_COUNT; i++) { writer.write(i * 2, i); diff --git a/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/BTreeIndexReaderCloseTest.java b/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/BTreeIndexReaderCloseTest.java index 90db80263a46..11a389bd2f14 100644 --- a/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/BTreeIndexReaderCloseTest.java +++ b/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/BTreeIndexReaderCloseTest.java @@ -43,6 +43,7 @@ import java.io.IOException; import java.io.RandomAccessFile; +import java.util.Collections; import java.util.List; import java.util.UUID; import java.util.concurrent.atomic.AtomicInteger; @@ -82,7 +83,9 @@ public PositionOutputStream newOutputStream(String fileName) }; BTreeGlobalIndexer indexer = - new BTreeGlobalIndexer(new DataField(1, "testField", dataType), new Options()); + new BTreeGlobalIndexer( + Collections.singletonList(new DataField(1, "testField", dataType)), + new Options()); GlobalIndexSingleColumnWriter writer = indexer.createWriter(fileWriter); for (int i = 0; i < RECORD_NUM; i++) { writer.write(i, (long) i); diff --git a/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/BTreeThreadSafetyTest.java b/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/BTreeThreadSafetyTest.java index 83961dc63638..9e6edfbd667b 100644 --- a/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/BTreeThreadSafetyTest.java +++ b/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/BTreeThreadSafetyTest.java @@ -46,6 +46,7 @@ import java.io.IOException; import java.util.ArrayList; +import java.util.Collections; import java.util.Comparator; import java.util.HashSet; import java.util.Iterator; @@ -108,7 +109,7 @@ public PositionOutputStream newOutputStream(String fileName) options.set(BTreeIndexOptions.BTREE_INDEX_CACHE_SIZE, MemorySize.ofMebiBytes(8)); options.set(BTreeIndexOptions.BTREE_INDEX_BLOOM_FILTER_ENABLED, true); DataField dataField = new DataField(1, "id", new IntType()); - globalIndexer = new BTreeGlobalIndexer(dataField, options); + globalIndexer = new BTreeGlobalIndexer(Collections.singletonList(dataField), options); keySerializer = KeySerializer.create(new IntType()); comparator = keySerializer.createComparator(); diff --git a/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/CompositeBTreeIndexTest.java b/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/CompositeBTreeIndexTest.java index d5fccfe14e1e..7a75de7382fa 100644 --- a/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/CompositeBTreeIndexTest.java +++ b/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/CompositeBTreeIndexTest.java @@ -74,9 +74,7 @@ void testMutableTupleKeysPostingListsAndLocalRanges(int version) throws Exceptio options.set(BTreeIndexOptions.BTREE_INDEX_FILE_VERSION, version); options.set(BTreeIndexOptions.BTREE_INDEX_BLOOM_FILTER_ENABLED, true); options.set(BTreeIndexOptions.BTREE_INDEX_COMPRESSION, "lz4"); - GlobalIndexer indexer = - GlobalIndexer.create( - "btree", type.getFields().get(0), type.getFields().subList(1, 3), options); + GlobalIndexer indexer = GlobalIndexer.create("btree", type.getFields(), options); LocalFileIO io = LocalFileIO.create(); Path directory = new Path(tempPath.toUri()); GlobalIndexFileWriter files = @@ -110,8 +108,7 @@ public PositionOutputStream newOutputStream(String name) assertThat( new BTreeGlobalIndexerFactory() .selectFiles( - type.getFields().get(0), - type.getFields().subList(1, 3), + type.getFields(), new PredicateBuilder(type) .equal(0, BinaryString.fromString("category-a")), Collections.singletonList(meta))) @@ -212,11 +209,7 @@ void testCompositeKeysPreserveTypesBoundariesAndNulls() { new DataField(30, "tag", DataTypes.STRING()))); SortedGlobalIndexer indexer = (SortedGlobalIndexer) - GlobalIndexer.create( - "btree", - type.getFields().get(0), - type.getFields().subList(1, 3), - new Options()); + GlobalIndexer.create("btree", type.getFields(), new Options()); KeySerializer serializer = new CompositeKeySerializer((RowType) indexer.keyExtractor().keyType()); Comparator comparator = serializer.createComparator(); @@ -244,22 +237,43 @@ void testCompositeKeysPreserveTypesBoundariesAndNulls() { assertThat(comparator.compare(nullable, first)).isNegative(); assertThat(serializer.serialize(row("a", 1, "bc"))) .isNotEqualTo(serializer.serialize(row("ab", 1, "c"))); - assertThat( + SortedGlobalIndexer scalarIndexer = + (SortedGlobalIndexer) GlobalIndexer.create( "btree", - type.getFields().get(0), - Collections.emptyList(), - new Options())) - .isInstanceOf(BTreeGlobalIndexer.class); - assertThat(GlobalIndexer.create("btree", type.getFields().get(0), null, new Options())) - .isInstanceOf(BTreeGlobalIndexer.class); + Collections.singletonList(type.getFields().get(0)), + new Options()); + assertThat(scalarIndexer.keyExtractor().keyType()).isEqualTo(DataTypes.STRING()); + assertThat(indexer.keyExtractor().keyType()).isEqualTo(type); + } + + @Test + void testFactoryRejectsUnsupportedFieldLists() { + assertThatThrownBy( + () -> GlobalIndexer.create("btree", Collections.emptyList(), new Options())) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("at least one field"); + List fields = + Arrays.asList( + new DataField(10, "category", DataTypes.STRING()), + new DataField(20, "item_number", DataTypes.INT())); + for (String indexType : Arrays.asList("bitmap", "multivalue", "fm")) { + assertThatThrownBy(() -> GlobalIndexer.create(indexType, fields, new Options())) + .isInstanceOf(UnsupportedOperationException.class) + .hasMessageContaining("exactly one index field"); + } } @Test void testCompositeSupportDoesNotEnablePhysicalRowIndexes() { DataField field = new DataField(40, "nested", RowType.of(DataTypes.INT())); for (String indexType : Arrays.asList("btree", "bitmap")) { - assertThatThrownBy(() -> GlobalIndexer.create(indexType, field, new Options())) + assertThatThrownBy( + () -> + GlobalIndexer.create( + indexType, + Collections.singletonList(field), + new Options())) .isInstanceOf(UnsupportedOperationException.class); } } diff --git a/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/LazyFilteredBTreeIndexReaderTest.java b/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/LazyFilteredBTreeIndexReaderTest.java index 83942e37b549..658d74cafbcb 100644 --- a/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/LazyFilteredBTreeIndexReaderTest.java +++ b/paimon-common/src/test/java/org/apache/paimon/globalindex/btree/LazyFilteredBTreeIndexReaderTest.java @@ -232,7 +232,10 @@ fileReader, written, dataNum, null, newDirectExecutorService())) { @TestTemplate public void testFallbackScanDisabledByBudget() throws Exception { options.set(BTreeIndexOptions.BTREE_INDEX_FALLBACK_SCAN_MAX_SIZE, MemorySize.ofBytes(1)); - globalIndexer = new BTreeGlobalIndexer(new DataField(1, "testField", dataType), options); + globalIndexer = + new BTreeGlobalIndexer( + Collections.singletonList(new DataField(1, "testField", dataType)), + options); List written = writeData(); FieldRef ref = new FieldRef(1, "testField", dataType); @@ -264,7 +267,10 @@ public void testFallbackBudgetUsesSelectedFiles() throws Exception { options.set( BTreeIndexOptions.BTREE_INDEX_FALLBACK_SCAN_MAX_SIZE, MemorySize.ofBytes(written.get(1).fileSize())); - globalIndexer = new BTreeGlobalIndexer(new DataField(1, "testField", dataType), options); + globalIndexer = + new BTreeGlobalIndexer( + Collections.singletonList(new DataField(1, "testField", dataType)), + options); FieldRef ref = new FieldRef(1, "testField", dataType); Object min = data.get(0).getKey(); @@ -336,7 +342,10 @@ public void testAllMatchAvoidsOpeningBTreeFiles() throws Exception { @TestTemplate public void testAllMatchRangeDoesNotConsumeScanBudget() throws Exception { options.set(BTreeIndexOptions.BTREE_INDEX_FALLBACK_SCAN_MAX_SIZE, MemorySize.ofBytes(1)); - globalIndexer = new BTreeGlobalIndexer(new DataField(1, "testField", dataType), options); + globalIndexer = + new BTreeGlobalIndexer( + Collections.singletonList(new DataField(1, "testField", dataType)), + options); List written = writeData(); Object min = data.get(0).getKey(); Object max = data.get(dataNum - 1).getKey(); @@ -522,7 +531,9 @@ public void testConcurrentAccess() throws Exception { stressOptions.set(BTreeIndexOptions.BTREE_INDEX_CACHE_SIZE, MemorySize.ofKibiBytes(64)); stressOptions.set(BTreeIndexOptions.BTREE_INDEX_HIGH_PRIORITY_POOL_RATIO, 0.1); BTreeGlobalIndexer stressIndexer = - new BTreeGlobalIndexer(new DataField(1, "testField", dataType), stressOptions); + new BTreeGlobalIndexer( + Collections.singletonList(new DataField(1, "testField", dataType)), + stressOptions); // Inject null values at the tail to test isNull/isNotNull under concurrency for (int i = dataNum - 1; i >= dataNum * 0.9; i--) { diff --git a/paimon-common/src/test/java/org/apache/paimon/globalindex/fmindex/FMGlobalIndexTest.java b/paimon-common/src/test/java/org/apache/paimon/globalindex/fmindex/FMGlobalIndexTest.java index 49eba6f1e0d4..8e7befeb4e88 100644 --- a/paimon-common/src/test/java/org/apache/paimon/globalindex/fmindex/FMGlobalIndexTest.java +++ b/paimon-common/src/test/java/org/apache/paimon/globalindex/fmindex/FMGlobalIndexTest.java @@ -295,7 +295,7 @@ public void testEmptyShardServiceLoadingAndFileCoverageValidation() throws Excep try (GlobalIndexReader reader = createReader(Collections.emptyList(), 0)) { assertRows(reader.visitContains(fieldRef, str("")).join()); } - assertThat(GlobalIndexer.create("fm", dataField, new Options())) + assertThat(GlobalIndexer.create("fm", Collections.singletonList(dataField), new Options())) .isInstanceOf(FMGlobalIndexer.class); List shifted = writeData(Collections.singletonList(str("value")), 1); diff --git a/paimon-common/src/test/java/org/apache/paimon/globalindex/testfulltext/TestFullTextGlobalIndexerFactory.java b/paimon-common/src/test/java/org/apache/paimon/globalindex/testfulltext/TestFullTextGlobalIndexerFactory.java index e8894affbe4d..c70267c78543 100644 --- a/paimon-common/src/test/java/org/apache/paimon/globalindex/testfulltext/TestFullTextGlobalIndexerFactory.java +++ b/paimon-common/src/test/java/org/apache/paimon/globalindex/testfulltext/TestFullTextGlobalIndexerFactory.java @@ -45,22 +45,13 @@ public boolean supportsFullTextSearch() { } @Override - public GlobalIndexer create(DataField field, Options options) { - return new TestFullTextGlobalIndexer(field.type(), options); - } - - @Override - public GlobalIndexer create( - DataField indexField, List extraFields, Options options) { - // Multi-column support: this brute-force backend indexes a single text column, which may be - // the primary field or an extra field. Pick the VARCHAR/STRING column among them. - DataField textField = indexField; - if (!(textField.type() instanceof VarCharType) && extraFields != null) { - for (DataField extra : extraFields) { - if (extra.type() instanceof VarCharType) { - textField = extra; - break; - } + public GlobalIndexer create(List indexFields, Options options) { + // This brute-force backend indexes the first VARCHAR/STRING column. + DataField textField = indexFields.get(0); + for (DataField field : indexFields) { + if (field.type() instanceof VarCharType) { + textField = field; + break; } } return new TestFullTextGlobalIndexer(textField.type(), options); diff --git a/paimon-common/src/test/java/org/apache/paimon/globalindex/testvector/TestMultiFieldVectorGlobalIndexer.java b/paimon-common/src/test/java/org/apache/paimon/globalindex/testvector/TestMultiFieldVectorGlobalIndexer.java index ff18c4af1e5e..3edda5577a80 100644 --- a/paimon-common/src/test/java/org/apache/paimon/globalindex/testvector/TestMultiFieldVectorGlobalIndexer.java +++ b/paimon-common/src/test/java/org/apache/paimon/globalindex/testvector/TestMultiFieldVectorGlobalIndexer.java @@ -42,6 +42,7 @@ import java.io.IOException; import java.util.ArrayList; +import java.util.Collections; import java.util.List; import java.util.Optional; import java.util.concurrent.CompletableFuture; @@ -66,7 +67,8 @@ class TestMultiFieldVectorGlobalIndexer implements VectorGlobalIndexer { this.vectorField = vectorField; this.scalarField = extraFields.get(0); this.vectorIndexer = new TestVectorGlobalIndexer(vectorField.type(), options); - this.scalarIndexer = new BTreeGlobalIndexer(scalarField, options); + this.scalarIndexer = + new BTreeGlobalIndexer(Collections.singletonList(scalarField), options); } @Override diff --git a/paimon-common/src/test/java/org/apache/paimon/globalindex/testvector/TestVectorGlobalIndexerFactory.java b/paimon-common/src/test/java/org/apache/paimon/globalindex/testvector/TestVectorGlobalIndexerFactory.java index 04543f7919e5..e87a8eec3246 100644 --- a/paimon-common/src/test/java/org/apache/paimon/globalindex/testvector/TestVectorGlobalIndexerFactory.java +++ b/paimon-common/src/test/java/org/apache/paimon/globalindex/testvector/TestVectorGlobalIndexerFactory.java @@ -39,16 +39,11 @@ public String identifier() { } @Override - public GlobalIndexer create(DataField field, Options options) { - return new TestVectorGlobalIndexer(field.type(), options); - } - - @Override - public GlobalIndexer create( - DataField indexField, List extraFields, Options options) { - if (extraFields == null || extraFields.isEmpty()) { - return create(indexField, options); + public GlobalIndexer create(List indexFields, Options options) { + if (indexFields.size() == 1) { + return new TestVectorGlobalIndexer(indexFields.get(0).type(), options); } - return new TestMultiFieldVectorGlobalIndexer(indexField, extraFields, options); + return new TestMultiFieldVectorGlobalIndexer( + indexFields.get(0), indexFields.subList(1, indexFields.size()), options); } } diff --git a/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexScanner.java b/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexScanner.java index 5b000ad7546d..4e65b3be0e7f 100644 --- a/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexScanner.java +++ b/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionGlobalIndexScanner.java @@ -419,16 +419,15 @@ public GlobalIndexResult unindexedRows(TopN topN) { private Collection createReaders( GlobalIndexFileReader indexFileReadWrite, IndexMetaFileGroup group, RowType rowType) { - DataField indexField = group.indexField(rowType); - List extraFields = group.extraFields(rowType); + List indexFields = + group.fieldIds.stream().map(rowType::getField).collect(Collectors.toList()); Set readers = new HashSet<>(); for (Map.Entry>> entry : group.metas.entrySet()) { String indexType = entry.getKey(); Map> metas = entry.getValue(); GlobalIndexerFactory globalIndexerFactory = GlobalIndexerFactoryUtils.load(indexType); - GlobalIndexer globalIndexer = - globalIndexerFactory.create(indexField, extraFields, options); + GlobalIndexer globalIndexer = globalIndexerFactory.create(indexFields, options); List> futures = new ArrayList<>(metas.size()); for (Map.Entry> rangeMetas : metas.entrySet()) { 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..f3e6e6fdfbc8 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 @@ -607,21 +607,9 @@ private static List toIndexFileMetas( } public static GlobalIndexWriter createIndexWriter( - FileStoreTable table, String indexType, DataField indexField, Options options) + FileStoreTable table, String indexType, List indexFields, Options options) throws IOException { - GlobalIndexer globalIndexer = GlobalIndexer.create(indexType, indexField, options); - return globalIndexer.createWriter(createGlobalIndexFileReadWrite(table)); - } - - public static GlobalIndexWriter createIndexWriter( - FileStoreTable table, - String indexType, - DataField indexField, - List extraFields, - Options options) - throws IOException { - GlobalIndexer globalIndexer = - GlobalIndexer.create(indexType, indexField, extraFields, options); + GlobalIndexer globalIndexer = GlobalIndexer.create(indexType, indexFields, options); return globalIndexer.createWriter(createGlobalIndexFileReadWrite(table)); } diff --git a/paimon-core/src/main/java/org/apache/paimon/globalindex/GlobalIndexQuery.java b/paimon-core/src/main/java/org/apache/paimon/globalindex/GlobalIndexQuery.java index 7745f8ba2155..7446ce5ed1a1 100644 --- a/paimon-core/src/main/java/org/apache/paimon/globalindex/GlobalIndexQuery.java +++ b/paimon-core/src/main/java/org/apache/paimon/globalindex/GlobalIndexQuery.java @@ -190,7 +190,7 @@ private static GlobalIndexQuery createForIndexedField( for (IndexGroup group : fieldGroups) { List selectedFiles = GlobalIndexerFactoryUtils.selectFiles( - group.type, group.field, group.extraFields, predicate, group.files); + group.type, group.indexFields(), predicate, group.files); if (!selectedFiles.isEmpty()) { selectedGroups.add( selectedFiles == group.files @@ -280,8 +280,7 @@ private GlobalIndexResult evaluateWithExecutor( continue; } GlobalIndexer indexer = - GlobalIndexerFactoryUtils.load(group.type) - .create(group.field, group.extraFields, options); + GlobalIndexerFactoryUtils.load(group.type).create(group.indexFields(), options); try (GlobalIndexReader reader = indexer.createReader( meta -> fileIO.newInputStream(meta.filePath()), @@ -478,6 +477,13 @@ private IndexGroup( this.files = new ArrayList<>(files); } + private List indexFields() { + List fields = new ArrayList<>(1 + extraFields.size()); + fields.add(field); + fields.addAll(extraFields); + return fields; + } + private static List fromMetadata( IndexMetaFileGroup group, RowType rowType, IndexPathFactory pathFactory) { DataField field = group.indexField(rowType); diff --git a/paimon-core/src/main/java/org/apache/paimon/globalindex/sorted/SortedGlobalIndexWriter.java b/paimon-core/src/main/java/org/apache/paimon/globalindex/sorted/SortedGlobalIndexWriter.java index b0bd408d877a..52f33840105a 100644 --- a/paimon-core/src/main/java/org/apache/paimon/globalindex/sorted/SortedGlobalIndexWriter.java +++ b/paimon-core/src/main/java/org/apache/paimon/globalindex/sorted/SortedGlobalIndexWriter.java @@ -86,7 +86,9 @@ public SortedGlobalIndexWriter withIndexField(String indexField) { indexField, table.fullName()); this.indexField = rowType.getField(indexField); - GlobalIndexer indexer = GlobalIndexer.create(indexType, this.indexField, options); + GlobalIndexer indexer = + GlobalIndexer.create( + indexType, Collections.singletonList(this.indexField), options); checkArgument( indexer instanceof SortedGlobalIndexer, "Index algorithm %s does not expose sorted index keys.", @@ -131,7 +133,8 @@ public SortedSingleColumnIndexWriter createTaskWriter(Range rowRange) throws IOE } public GlobalIndexSingleColumnWriter createWriter() throws IOException { - GlobalIndexWriter indexWriter = createIndexWriter(table, indexType, indexField, options); + GlobalIndexWriter indexWriter = + createIndexWriter(table, indexType, Collections.singletonList(indexField), options); if (!(indexWriter instanceof GlobalIndexSingleColumnWriter)) { throw new RuntimeException( "Unexpected implementation, the index writer of " diff --git a/paimon-core/src/main/java/org/apache/paimon/index/pkfulltext/PkFullTextIndexFile.java b/paimon-core/src/main/java/org/apache/paimon/index/pkfulltext/PkFullTextIndexFile.java index 1b2f22fddd23..4a6e02e062bc 100644 --- a/paimon-core/src/main/java/org/apache/paimon/index/pkfulltext/PkFullTextIndexFile.java +++ b/paimon-core/src/main/java/org/apache/paimon/index/pkfulltext/PkFullTextIndexFile.java @@ -63,7 +63,9 @@ public IndexFileMeta build( throws IOException { GlobalIndexer indexer; try { - indexer = GlobalIndexer.create(INDEX_TYPE, textField, indexOptions); + indexer = + GlobalIndexer.create( + INDEX_TYPE, Collections.singletonList(textField), indexOptions); } catch (RuntimeException | Error failure) { IOUtils.closeQuietly(textReader); throw failure; @@ -76,7 +78,9 @@ public IndexFileMeta build(List sources, DataField textField, Options in throws IOException { GlobalIndexer indexer; try { - indexer = GlobalIndexer.create(INDEX_TYPE, textField, indexOptions); + indexer = + GlobalIndexer.create( + INDEX_TYPE, Collections.singletonList(textField), indexOptions); } catch (RuntimeException | Error failure) { closeSources(sources); throw failure; diff --git a/paimon-core/src/main/java/org/apache/paimon/index/pksorted/PkSortedIndexBuilder.java b/paimon-core/src/main/java/org/apache/paimon/index/pksorted/PkSortedIndexBuilder.java index b5412ad49033..2b26739a4805 100644 --- a/paimon-core/src/main/java/org/apache/paimon/index/pksorted/PkSortedIndexBuilder.java +++ b/paimon-core/src/main/java/org/apache/paimon/index/pksorted/PkSortedIndexBuilder.java @@ -42,6 +42,7 @@ import java.io.Closeable; import java.io.IOException; import java.util.ArrayList; +import java.util.Collections; import java.util.Comparator; import java.util.Iterator; import java.util.List; @@ -104,7 +105,8 @@ public IndexFileMeta build(List dataFiles) throws IOException { sourceFiles.add( new PrimaryKeyIndexSourceFile(dataFile.fileName(), dataFile.rowCount())); } - GlobalIndexer indexer = GlobalIndexer.create(indexType, indexField, options); + GlobalIndexer indexer = + GlobalIndexer.create(indexType, Collections.singletonList(indexField), options); checkArgument( indexer instanceof SortedGlobalIndexer, "Index algorithm %s does not expose sorted index keys.", diff --git a/paimon-core/src/main/java/org/apache/paimon/index/pksorted/PkSortedIndexFile.java b/paimon-core/src/main/java/org/apache/paimon/index/pksorted/PkSortedIndexFile.java index fb05e43ba3c6..456f7878f487 100644 --- a/paimon-core/src/main/java/org/apache/paimon/index/pksorted/PkSortedIndexFile.java +++ b/paimon-core/src/main/java/org/apache/paimon/index/pksorted/PkSortedIndexFile.java @@ -39,6 +39,7 @@ import javax.annotation.Nullable; import java.io.IOException; +import java.util.Collections; import java.util.Iterator; import java.util.LinkedHashMap; import java.util.List; @@ -130,7 +131,9 @@ protected GlobalIndexSingleColumnWriter createWriter( Options indexOptions, GlobalIndexFileWriter fileWriter) throws IOException { - GlobalIndexer indexer = GlobalIndexer.create(indexType, indexField, indexOptions); + GlobalIndexer indexer = + GlobalIndexer.create( + indexType, Collections.singletonList(indexField), indexOptions); GlobalIndexWriter writer = indexer.createWriter(fileWriter); checkArgument( writer instanceof GlobalIndexSingleColumnWriter, diff --git a/paimon-core/src/main/java/org/apache/paimon/index/pkvector/PkVectorAnnSegmentFile.java b/paimon-core/src/main/java/org/apache/paimon/index/pkvector/PkVectorAnnSegmentFile.java index e337db2802d2..93ac807aff4c 100644 --- a/paimon-core/src/main/java/org/apache/paimon/index/pkvector/PkVectorAnnSegmentFile.java +++ b/paimon-core/src/main/java/org/apache/paimon/index/pkvector/PkVectorAnnSegmentFile.java @@ -82,7 +82,9 @@ public IndexFileMeta build( } checkArgument(totalRowCount > 0, "An ANN segment must reference at least one source row."); - GlobalIndexer indexer = GlobalIndexer.create(indexType, vectorField, indexOptions); + GlobalIndexer indexer = + GlobalIndexer.create( + indexType, Collections.singletonList(vectorField), indexOptions); checkArgument( indexer instanceof VectorGlobalIndexer, "Index algorithm %s does not implement VectorGlobalIndexer.", diff --git a/paimon-core/src/main/java/org/apache/paimon/index/pkvector/PkVectorAnnSegmentSearcher.java b/paimon-core/src/main/java/org/apache/paimon/index/pkvector/PkVectorAnnSegmentSearcher.java index 30f502455477..14b8ab33346e 100644 --- a/paimon-core/src/main/java/org/apache/paimon/index/pkvector/PkVectorAnnSegmentSearcher.java +++ b/paimon-core/src/main/java/org/apache/paimon/index/pkvector/PkVectorAnnSegmentSearcher.java @@ -181,7 +181,8 @@ CompletableFuture> searchAsync( return CompletableFuture.completedFuture(Collections.emptyList()); } GlobalIndexer indexer = - GlobalIndexer.create(segment.indexType(), vectorField, indexOptions); + GlobalIndexer.create( + segment.indexType(), Collections.singletonList(vectorField), indexOptions); checkArgument( indexer instanceof VectorGlobalIndexer, "Index algorithm %s does not implement VectorGlobalIndexer.", @@ -279,7 +280,8 @@ CompletableFuture>> searchBatchAsync( return CompletableFuture.completedFuture(Collections.unmodifiableList(results)); } GlobalIndexer indexer = - GlobalIndexer.create(segment.indexType(), vectorField, indexOptions); + GlobalIndexer.create( + segment.indexType(), Collections.singletonList(vectorField), indexOptions); checkArgument( indexer instanceof VectorGlobalIndexer, "Index algorithm %s does not implement VectorGlobalIndexer.", 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 ed316a72238d..2da251e5ea0d 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 @@ -1466,25 +1466,25 @@ private static void validatePrimaryKeySortedIndexes(TableSchema schema, CoreOpti for (String column : options.primaryKeyBTreeIndexColumns()) { GlobalIndexer.create( BTreeGlobalIndexerFactory.IDENTIFIER, - fields.get(column), + Collections.singletonList(fields.get(column)), options.primaryKeyBTreeIndexOptions(column)); } for (String column : options.primaryKeyBitmapIndexColumns()) { GlobalIndexer.create( BitmapGlobalIndexerFactory.IDENTIFIER, - fields.get(column), + Collections.singletonList(fields.get(column)), options.primaryKeyBitmapIndexOptions(column)); } for (String column : options.primaryKeyMultiValueIndexColumns()) { GlobalIndexer.create( MultiValueGlobalIndexerFactory.IDENTIFIER, - fields.get(column), + Collections.singletonList(fields.get(column)), options.primaryKeyMultiValueIndexOptions(column)); } for (String column : options.primaryKeyFMIndexColumns()) { GlobalIndexer.create( FMGlobalIndexerFactory.IDENTIFIER, - fields.get(column), + Collections.singletonList(fields.get(column)), options.primaryKeyFMIndexOptions(column)); } } diff --git a/paimon-core/src/main/java/org/apache/paimon/table/source/AbstractDataEvolutionVectorRead.java b/paimon-core/src/main/java/org/apache/paimon/table/source/AbstractDataEvolutionVectorRead.java index 9acef5d9ac45..c638a400a36a 100644 --- a/paimon-core/src/main/java/org/apache/paimon/table/source/AbstractDataEvolutionVectorRead.java +++ b/paimon-core/src/main/java/org/apache/paimon/table/source/AbstractDataEvolutionVectorRead.java @@ -130,7 +130,9 @@ protected GlobalIndexer createGlobalIndexer(List splits) protected GlobalIndexer createGlobalIndexer(String indexType) { return GlobalIndexerFactoryUtils.load(indexType) - .create(vectorColumn, table.coreOptions().toConfiguration()); + .create( + Collections.singletonList(vectorColumn), + table.coreOptions().toConfiguration()); } private GlobalIndexer createGlobalIndexer(String indexType, GlobalIndexMeta meta) { @@ -139,8 +141,7 @@ private GlobalIndexer createGlobalIndexer(String indexType, GlobalIndexMeta meta } return GlobalIndexerFactoryUtils.load(indexType) .create( - meta.getIndexField(table.rowType()), - meta.getExtraFields(table.rowType()), + meta.getIndexedFields(table.rowType()), table.coreOptions().toConfiguration()); } diff --git a/paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionFullTextRead.java b/paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionFullTextRead.java index 33cf8f6c3547..05d16d885268 100644 --- a/paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionFullTextRead.java +++ b/paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionFullTextRead.java @@ -51,6 +51,7 @@ import java.io.IOException; import java.util.ArrayList; +import java.util.Collections; import java.util.Comparator; import java.util.HashMap; import java.util.List; @@ -386,12 +387,13 @@ private GlobalIndexer createIndexer(IndexFileMeta file, DataField textColumn) { if (meta.extraFieldIds() != null) { return GlobalIndexerFactoryUtils.load(indexType) .create( - meta.getIndexField(table.rowType()), - meta.getExtraFields(table.rowType()), + meta.getIndexedFields(table.rowType()), table.coreOptions().toConfiguration()); } return GlobalIndexerFactoryUtils.load(indexType) - .create(textColumn, table.coreOptions().toConfiguration()); + .create( + Collections.singletonList(textColumn), + table.coreOptions().toConfiguration()); } private CompletableFuture> eval( diff --git a/paimon-core/src/main/java/org/apache/paimon/table/source/PrimaryKeyFullTextRead.java b/paimon-core/src/main/java/org/apache/paimon/table/source/PrimaryKeyFullTextRead.java index 18e7f05c245e..80d85b5e8936 100644 --- a/paimon-core/src/main/java/org/apache/paimon/table/source/PrimaryKeyFullTextRead.java +++ b/paimon-core/src/main/java/org/apache/paimon/table/source/PrimaryKeyFullTextRead.java @@ -176,7 +176,9 @@ private ProductionSearch( table.coreOptions().toConfiguration().get(GLOBAL_INDEX_THREAD_NUM)); this.indexer = GlobalIndexer.create( - PkFullTextIndexFile.INDEX_TYPE, textField, definition.options()); + PkFullTextIndexFile.INDEX_TYPE, + Collections.singletonList(textField), + definition.options()); this.archiveReader = meta -> fileIO.newInputStream(meta.filePath()); } diff --git a/paimon-core/src/main/java/org/apache/paimon/table/source/PrimaryKeySortedIndexScan.java b/paimon-core/src/main/java/org/apache/paimon/table/source/PrimaryKeySortedIndexScan.java index dd28bd49bff5..31b35e0a5f13 100644 --- a/paimon-core/src/main/java/org/apache/paimon/table/source/PrimaryKeySortedIndexScan.java +++ b/paimon-core/src/main/java/org/apache/paimon/table/source/PrimaryKeySortedIndexScan.java @@ -114,7 +114,7 @@ static ReaderFactory readerFactory( GlobalIndexer indexer = GlobalIndexer.create( definition.indexType(), - rowType.getField(definition.fieldId()), + Collections.singletonList(rowType.getField(definition.fieldId())), definition.options()); return indexer.createReader(fileReader, ioMetas, totalRowCount, null, executor); }; diff --git a/paimon-core/src/main/java/org/apache/paimon/table/source/RawFullTextReadImpl.java b/paimon-core/src/main/java/org/apache/paimon/table/source/RawFullTextReadImpl.java index 44a7e63a3609..58cb5a59e9c1 100644 --- a/paimon-core/src/main/java/org/apache/paimon/table/source/RawFullTextReadImpl.java +++ b/paimon-core/src/main/java/org/apache/paimon/table/source/RawFullTextReadImpl.java @@ -211,7 +211,8 @@ private Map createRawFullTextIndexes( indexType = checkNotNull(fallbackIndexType); } GlobalIndexer globalIndexer = - GlobalIndexerFactoryUtils.load(indexType).create(textColumn, rawSearchOptions()); + GlobalIndexerFactoryUtils.load(indexType) + .create(Collections.singletonList(textColumn), rawSearchOptions()); try { RawFullTextIndexFileWriter fileWriter = new RawFullTextIndexFileWriter(column); GlobalIndexWriter indexWriter = globalIndexer.createWriter(fileWriter); diff --git a/paimon-core/src/test/java/org/apache/paimon/globalindex/GlobalIndexQueryTest.java b/paimon-core/src/test/java/org/apache/paimon/globalindex/GlobalIndexQueryTest.java index d2821933f7bb..7534931a31c0 100644 --- a/paimon-core/src/test/java/org/apache/paimon/globalindex/GlobalIndexQueryTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/globalindex/GlobalIndexQueryTest.java @@ -39,7 +39,6 @@ import org.apache.paimon.stats.SimpleStats; import org.apache.paimon.table.source.DataSplit; import org.apache.paimon.table.source.SplitSerializer; -import org.apache.paimon.types.DataField; import org.apache.paimon.types.DataType; import org.apache.paimon.types.DataTypes; import org.apache.paimon.types.RowType; @@ -171,8 +170,7 @@ void testBoundedRangeUsesSingleScan(boolean lowerInclusive, boolean upperInclusi GlobalIndexer indexer = mock(GlobalIndexer.class); when(indexer.createReader(any(), anyList(), eq(100L), anyList(), any())).thenReturn(reader); GlobalIndexerFactory factory = mock(GlobalIndexerFactory.class); - when(factory.create(any(DataField.class), anyList(), any(Options.class))) - .thenReturn(indexer); + when(factory.create(anyList(), any(Options.class))).thenReturn(indexer); try (MockedStatic factories = mockStatic(GlobalIndexerFactoryUtils.class)) { factories.when(() -> GlobalIndexerFactoryUtils.load("btree")).thenReturn(factory); @@ -224,7 +222,9 @@ public boolean isExternalPath() { GlobalIndexSingleColumnWriter writer = (GlobalIndexSingleColumnWriter) GlobalIndexer.create( - indexType, rowType.getFields().get(0), new Options()) + indexType, + Collections.singletonList(rowType.getFields().get(0)), + new Options()) .createWriter(io); if (nulls) { writer.write(null, 0); @@ -361,8 +361,7 @@ void testUnsupportedPredicateFailsAtRuntime(String indexType) { when(indexer.createReader(any(), anyList(), anyLong(), anyList(), any())) .thenReturn(reader); GlobalIndexerFactory factory = mock(GlobalIndexerFactory.class); - when(factory.create(any(DataField.class), anyList(), any(Options.class))) - .thenReturn(indexer); + when(factory.create(anyList(), any(Options.class))).thenReturn(indexer); try (MockedStatic factories = mockStatic(GlobalIndexerFactoryUtils.class)) { factories.when(() -> GlobalIndexerFactoryUtils.load(indexType)).thenReturn(factory); @@ -432,8 +431,7 @@ public List toRangeList() { when(indexer.createReader(any(), anyList(), anyLong(), anyList(), any())) .thenReturn(reader); GlobalIndexerFactory factory = mock(GlobalIndexerFactory.class); - when(factory.create(any(DataField.class), anyList(), any(Options.class))) - .thenReturn(indexer); + when(factory.create(anyList(), any(Options.class))).thenReturn(indexer); RowType rowType = RowType.of(DataTypes.INT()); GlobalIndexQuery plan = @@ -508,8 +506,7 @@ void testPushesLocalRangesToIndexer(String indexType) throws Exception { GlobalIndexer indexer = mock(GlobalIndexer.class); when(indexer.createReader(any(), anyList(), eq(100L), anyList(), any())).thenReturn(reader); GlobalIndexerFactory factory = mock(GlobalIndexerFactory.class); - when(factory.create(any(DataField.class), anyList(), any(Options.class))) - .thenReturn(indexer); + when(factory.create(anyList(), any(Options.class))).thenReturn(indexer); FileIO fileIO = mock(FileIO.class); List first = Collections.singletonList(new Range(110, 119)); List second = Collections.singletonList(new Range(130, 139)); diff --git a/paimon-core/src/test/java/org/apache/paimon/index/pksorted/PkSortedIndexBuilderTest.java b/paimon-core/src/test/java/org/apache/paimon/index/pksorted/PkSortedIndexBuilderTest.java index 79d2155149af..5af2f8b00f9b 100644 --- a/paimon-core/src/test/java/org/apache/paimon/index/pksorted/PkSortedIndexBuilderTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/index/pksorted/PkSortedIndexBuilderTest.java @@ -142,7 +142,7 @@ void testBuildsQueryableMultiValueFromUnsortedArrayRows() throws Exception { payload.globalIndexMeta().indexMeta())); ExecutorService executor = newDirectExecutorService(); try (GlobalIndexReader reader = - GlobalIndexer.create("multivalue", tags, options) + GlobalIndexer.create("multivalue", Collections.singletonList(tags), options) .createReader( new GlobalIndexFileReadWrite(fileIO, pathFactory), ioMetas, @@ -391,7 +391,7 @@ private static void assertQuery( payload.globalIndexMeta().indexMeta())); ExecutorService executor = newDirectExecutorService(); try (GlobalIndexReader reader = - GlobalIndexer.create(indexType, field(), options) + GlobalIndexer.create(indexType, Collections.singletonList(field()), options) .createReader( new GlobalIndexFileReadWrite(fileIO, pathFactory), ioMetas, diff --git a/paimon-core/src/test/java/org/apache/paimon/table/BitmapGlobalIndexTableTest.java b/paimon-core/src/test/java/org/apache/paimon/table/BitmapGlobalIndexTableTest.java index fb8f980dccd8..a50eb2d0375e 100644 --- a/paimon-core/src/test/java/org/apache/paimon/table/BitmapGlobalIndexTableTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/table/BitmapGlobalIndexTableTest.java @@ -214,7 +214,10 @@ private CommitMessage buildIndex( Range rowRange = rowRange(split); GlobalIndexWriter indexWriter = GlobalIndexBuilderUtils.createIndexWriter( - table, INDEX_TYPE, indexField, table.coreOptions().toConfiguration()); + table, + INDEX_TYPE, + Collections.singletonList(indexField), + table.coreOptions().toConfiguration()); GlobalIndexSingleColumnWriter writer = (GlobalIndexSingleColumnWriter) indexWriter; InternalRow.FieldGetter fieldGetter = InternalRow.createFieldGetter( diff --git a/paimon-core/src/test/java/org/apache/paimon/table/IndexQuerySplitTest.java b/paimon-core/src/test/java/org/apache/paimon/table/IndexQuerySplitTest.java index 90e48f61aaee..0dff41970940 100644 --- a/paimon-core/src/test/java/org/apache/paimon/table/IndexQuerySplitTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/table/IndexQuerySplitTest.java @@ -834,7 +834,10 @@ public void testFMFallbackAndMixedDistributedIndex() throws Exception { indexOptions.set(FMGlobalIndexOptions.SA_SAMPLE_RATE, 1); GlobalIndexSingleColumnWriter writer = (GlobalIndexSingleColumnWriter) - GlobalIndexer.create("fm", table.rowType().getField("f1"), indexOptions) + GlobalIndexer.create( + "fm", + Collections.singletonList(table.rowType().getField("f1")), + indexOptions) .createWriter(io); for (int i = 0; i < 100; i++) { writer.write(str("a" + i), i); diff --git a/paimon-core/src/test/java/org/apache/paimon/table/source/FullTextSearchBuilderTest.java b/paimon-core/src/test/java/org/apache/paimon/table/source/FullTextSearchBuilderTest.java index 281f3ce32c31..574d4d588834 100644 --- a/paimon-core/src/test/java/org/apache/paimon/table/source/FullTextSearchBuilderTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/table/source/FullTextSearchBuilderTest.java @@ -1616,7 +1616,7 @@ private void buildAndCommitIndexWithFields( GlobalIndexBuilderUtils.createIndexWriter( table, TestFullTextGlobalIndexerFactory.IDENTIFIER, - textField, + Collections.singletonList(textField), options); for (int i = 0; i < documents.length; i++) { writer.write(documents[i], i); @@ -1657,7 +1657,7 @@ private void buildAndCommitSourceBackedIndex(FileStoreTable table, String[] docu GlobalIndexBuilderUtils.createIndexWriter( table, TestFullTextGlobalIndexerFactory.IDENTIFIER, - textField, + Collections.singletonList(textField), options); for (int i = 0; i < documents.length; i++) { writer.write(documents[i], i); @@ -1724,7 +1724,7 @@ private void buildAndCommitIndexRange( GlobalIndexBuilderUtils.createIndexWriter( table, TestFullTextGlobalIndexerFactory.IDENTIFIER, - textField, + Collections.singletonList(textField), options); // Doc ids are file-local (0-based); the global row offset is carried by rowRange.from, // which @@ -1816,7 +1816,7 @@ private void buildAndCommitIndexForColumn( GlobalIndexBuilderUtils.createIndexWriter( table, TestFullTextGlobalIndexerFactory.IDENTIFIER, - textField, + Collections.singletonList(textField), options); for (int i = 0; i < documents.length; i++) { writer.write(documents[i], i); @@ -1856,7 +1856,10 @@ private void buildAndCommitIdBTreeIndexRange(FileStoreTable table, Range rowRang GlobalIndexSingleColumnWriter writer = (GlobalIndexSingleColumnWriter) GlobalIndexBuilderUtils.createIndexWriter( - table, BTreeGlobalIndexerFactory.IDENTIFIER, idField, options); + table, + BTreeGlobalIndexerFactory.IDENTIFIER, + Collections.singletonList(idField), + options); for (long rowId = rowRange.from; rowId <= rowRange.to; rowId++) { writer.write((int) rowId, rowId - rowRange.from); } @@ -1892,7 +1895,10 @@ private void buildAndCommitBTreeIndex(FileStoreTable table, String[] documents) GlobalIndexSingleColumnWriter writer = (GlobalIndexSingleColumnWriter) GlobalIndexBuilderUtils.createIndexWriter( - table, BTreeGlobalIndexerFactory.IDENTIFIER, textField, options); + table, + BTreeGlobalIndexerFactory.IDENTIFIER, + Collections.singletonList(textField), + options); for (int i = 0; i < documents.length; i++) { writer.write(BinaryString.fromString(documents[i]), i); } @@ -1944,7 +1950,7 @@ private void buildAndCommitMultipleIndexFiles(FileStoreTable table, String[] doc GlobalIndexBuilderUtils.createIndexWriter( table, TestFullTextGlobalIndexerFactory.IDENTIFIER, - textField, + Collections.singletonList(textField), options); for (int i = 0; i < mid; i++) { writer1.write(documents[i], i); @@ -1967,7 +1973,7 @@ private void buildAndCommitMultipleIndexFiles(FileStoreTable table, String[] doc GlobalIndexBuilderUtils.createIndexWriter( table, TestFullTextGlobalIndexerFactory.IDENTIFIER, - textField, + Collections.singletonList(textField), options); for (int i = mid; i < documents.length; i++) { writer2.write(documents[i], i - mid); diff --git a/paimon-core/src/test/java/org/apache/paimon/table/source/VectorSearchBuilderTest.java b/paimon-core/src/test/java/org/apache/paimon/table/source/VectorSearchBuilderTest.java index e9c972d97852..13a1fe39a42c 100644 --- a/paimon-core/src/test/java/org/apache/paimon/table/source/VectorSearchBuilderTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/table/source/VectorSearchBuilderTest.java @@ -658,7 +658,7 @@ private void buildAndCommitMultiValueIndex(FileStoreTable table, int rowCount) GlobalIndexBuilderUtils.createIndexWriter( table, MultiValueGlobalIndexerFactory.IDENTIFIER, - tagsField, + Collections.singletonList(tagsField), options); for (int row = 0; row < rowCount; row++) { writer.write(BinaryString.fromString("t" + row), row); @@ -1906,7 +1906,7 @@ private void buildAndCommitIndex(FileStoreTable table, String fieldName, float[] GlobalIndexBuilderUtils.createIndexWriter( table, TestVectorGlobalIndexerFactory.IDENTIFIER, - vectorField, + Collections.singletonList(vectorField), options); for (int i = 0; i < vectors.length; i++) { writer.write(vectors[i], i); @@ -1949,7 +1949,7 @@ private void buildAndCommitMultipleIndexFiles(FileStoreTable table, float[][] ve GlobalIndexBuilderUtils.createIndexWriter( table, TestVectorGlobalIndexerFactory.IDENTIFIER, - vectorField, + Collections.singletonList(vectorField), options); for (int i = 0; i < mid; i++) { writer1.write(vectors[i], i); @@ -1972,7 +1972,7 @@ private void buildAndCommitMultipleIndexFiles(FileStoreTable table, float[][] ve GlobalIndexBuilderUtils.createIndexWriter( table, TestVectorGlobalIndexerFactory.IDENTIFIER, - vectorField, + Collections.singletonList(vectorField), options); for (int i = mid; i < vectors.length; i++) { writer2.write(vectors[i], i - mid); @@ -2128,8 +2128,7 @@ private void buildAndCommitVectorIndexWithFields( GlobalIndexBuilderUtils.createIndexWriter( table, TestVectorGlobalIndexerFactory.IDENTIFIER, - vectorField, - indexFields.subList(1, indexFields.size()), + indexFields, options); for (int i = 0; i < vectors.length; i++) { writer.write( @@ -2144,7 +2143,7 @@ private void buildAndCommitVectorIndexWithFields( GlobalIndexBuilderUtils.createIndexWriter( table, TestVectorGlobalIndexerFactory.IDENTIFIER, - vectorField, + Collections.singletonList(vectorField), options); for (int i = 0; i < vectors.length; i++) { writer.write(vectors[i], i); @@ -2184,7 +2183,10 @@ private void buildAndCommitBTreeIndex(FileStoreTable table, int[] ids, Range row GlobalIndexSingleColumnWriter writer = (GlobalIndexSingleColumnWriter) GlobalIndexBuilderUtils.createIndexWriter( - table, BTreeGlobalIndexerFactory.IDENTIFIER, idField, options); + table, + BTreeGlobalIndexerFactory.IDENTIFIER, + Collections.singletonList(idField), + options); for (int id : ids) { long relativeRowId = id - rowRange.from; writer.write(id, relativeRowId); @@ -2225,7 +2227,7 @@ private void buildAndCommitPartitionedIndex( GlobalIndexBuilderUtils.createIndexWriter( table, TestVectorGlobalIndexerFactory.IDENTIFIER, - vectorField, + Collections.singletonList(vectorField), options); for (int i = 0; i < vectors.length; i++) { writer.write(vectors[i], i); diff --git a/paimon-core/src/test/java/org/apache/paimon/table/source/VectorSearchRowFilterExactnessTest.java b/paimon-core/src/test/java/org/apache/paimon/table/source/VectorSearchRowFilterExactnessTest.java index e203229c30d9..acd97870d832 100644 --- a/paimon-core/src/test/java/org/apache/paimon/table/source/VectorSearchRowFilterExactnessTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/table/source/VectorSearchRowFilterExactnessTest.java @@ -467,7 +467,7 @@ private void buildAndCommitVectorIndex(FileStoreTable table, float[][] vectors, GlobalIndexBuilderUtils.createIndexWriter( table, TestVectorGlobalIndexerFactory.IDENTIFIER, - vectorField, + Collections.singletonList(vectorField), options); for (long rowId = rowRange.from; rowId <= rowRange.to; rowId++) { writer.write(vectors[(int) rowId], rowId - rowRange.from); @@ -491,7 +491,10 @@ private void buildAndCommitNameBTreeIndex(FileStoreTable table, String[] names) GlobalIndexSingleColumnWriter writer = (GlobalIndexSingleColumnWriter) GlobalIndexBuilderUtils.createIndexWriter( - table, BTreeGlobalIndexerFactory.IDENTIFIER, nameField, options); + table, + BTreeGlobalIndexerFactory.IDENTIFIER, + Collections.singletonList(nameField), + options); // The btree writer needs sorted keys. Integer[] order = new Integer[names.length]; for (int i = 0; i < names.length; i++) { @@ -519,7 +522,10 @@ private void buildAndCommitIdBTreeIndex(FileStoreTable table, int rowCount) thro GlobalIndexSingleColumnWriter writer = (GlobalIndexSingleColumnWriter) GlobalIndexBuilderUtils.createIndexWriter( - table, BTreeGlobalIndexerFactory.IDENTIFIER, idField, options); + table, + BTreeGlobalIndexerFactory.IDENTIFIER, + Collections.singletonList(idField), + options); for (int i = 0; i < rowCount; i++) { writer.write(i, i); } diff --git a/paimon-eslib/src/main/java/org/apache/paimon/eslib/index/ESIndexGlobalIndexerFactory.java b/paimon-eslib/src/main/java/org/apache/paimon/eslib/index/ESIndexGlobalIndexerFactory.java index 5e520b7e14d2..98cc1d8deeb5 100644 --- a/paimon-eslib/src/main/java/org/apache/paimon/eslib/index/ESIndexGlobalIndexerFactory.java +++ b/paimon-eslib/src/main/java/org/apache/paimon/eslib/index/ESIndexGlobalIndexerFactory.java @@ -44,21 +44,7 @@ public boolean supportsFullTextSearch() { } @Override - public GlobalIndexer create(DataField field, Options options) { - return new ESIndexGlobalIndexer(java.util.Collections.singletonList(field), options); - } - - @Override - public GlobalIndexer create( - DataField indexField, List extraFields, Options options) { - List fields; - if (extraFields == null || extraFields.isEmpty()) { - fields = java.util.Collections.singletonList(indexField); - } else { - fields = new java.util.ArrayList<>(extraFields.size() + 1); - fields.add(indexField); - fields.addAll(extraFields); - } - return new ESIndexGlobalIndexer(fields, options); + public GlobalIndexer create(List indexFields, Options options) { + return new ESIndexGlobalIndexer(indexFields, options); } } diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/globalindex/GenericIndexTopoBuilder.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/globalindex/GenericIndexTopoBuilder.java index bb14cefa4392..55e005d0a24d 100644 --- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/globalindex/GenericIndexTopoBuilder.java +++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/globalindex/GenericIndexTopoBuilder.java @@ -441,7 +441,7 @@ public void processElement(StreamRecord element) throws Exception long startTime = System.currentTimeMillis(); GlobalIndexWriter indexWriter = - createIndexWriter(table, indexType, indexField, extraFields, mergedOptions); + createIndexWriter(table, indexType, indexedFields, mergedOptions); try { long rowsSeen = 0; 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 67a73da98695..c4d86def5297 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 @@ -183,7 +183,9 @@ public static Optional> buildIndexStream( String buildTaskIdField = buildTaskIdFieldName(dataReadType); GlobalIndexer indexer = GlobalIndexer.create( - indexType, table.rowType().getField(indexColumn), userOptions); + indexType, + Collections.singletonList(table.rowType().getField(indexColumn)), + userOptions); if (!(indexer instanceof SortedGlobalIndexer)) { throw new IllegalArgumentException( "Index algorithm " + indexType + " does not expose sorted index keys."); 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..c0a0ea90d585 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 @@ -119,13 +119,10 @@ public String[] call( 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. - DataField indexField = rowType.getField(indexColumns.get(0)); - List extraFields = - indexColumns.subList(1, indexColumns.size()).stream() - .map(rowType::getField) - .collect(Collectors.toList()); + List indexFields = + indexColumns.stream().map(rowType::getField).collect(Collectors.toList()); try { - GlobalIndexer.create(indexType, indexField, extraFields, userOptions); + GlobalIndexer.create(indexType, indexFields, userOptions); } catch (UnsupportedOperationException e) { throw new IllegalArgumentException( String.format( diff --git a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/procedure/VectorSearchProcedureITCase.java b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/procedure/VectorSearchProcedureITCase.java index 31324f2ec7f3..1cf45a2fa5b9 100644 --- a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/procedure/VectorSearchProcedureITCase.java +++ b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/procedure/VectorSearchProcedureITCase.java @@ -1057,7 +1057,7 @@ private void buildAndCommitVectorIndex(FileStoreTable table, float[][] vectors, GlobalIndexBuilderUtils.createIndexWriter( table, TestVectorGlobalIndexerFactory.IDENTIFIER, - vectorField, + Collections.singletonList(vectorField), options); for (int i = 0; i < vectors.length; i++) { writer.write(vectors[i], i); @@ -1098,8 +1098,7 @@ private void buildAndCommitMultiFieldVectorIndex( GlobalIndexBuilderUtils.createIndexWriter( table, TestVectorGlobalIndexerFactory.IDENTIFIER, - vectorField, - Collections.singletonList(idField), + Arrays.asList(vectorField, idField), options); for (int i = 0; i < vectors.length; i++) { writer.write(i, GenericRow.of(new GenericArray(vectors[i]), (int) (rowRange.from + i))); diff --git a/paimon-full-text/src/main/java/org/apache/paimon/fulltext/index/NativeFullTextGlobalIndexerFactory.java b/paimon-full-text/src/main/java/org/apache/paimon/fulltext/index/NativeFullTextGlobalIndexerFactory.java index 1f1aed8c45ec..7263257346ea 100644 --- a/paimon-full-text/src/main/java/org/apache/paimon/fulltext/index/NativeFullTextGlobalIndexerFactory.java +++ b/paimon-full-text/src/main/java/org/apache/paimon/fulltext/index/NativeFullTextGlobalIndexerFactory.java @@ -23,6 +23,8 @@ import org.apache.paimon.options.Options; import org.apache.paimon.types.DataField; +import java.util.List; + /** Factory for creating native full-text index. */ public class NativeFullTextGlobalIndexerFactory implements GlobalIndexerFactory { @@ -39,7 +41,11 @@ public boolean supportsFullTextSearch() { } @Override - public GlobalIndexer create(DataField field, Options options) { + public GlobalIndexer create(List indexFields, Options options) { + if (indexFields.size() != 1) { + throw new UnsupportedOperationException( + "Index type '" + identifier() + "' requires exactly one index field."); + } return new NativeFullTextGlobalIndexer( new NativeFullTextIndexOptions( options.removePrefix(NativeFullTextIndexOptions.FULL_TEXT_PREFIX).toMap())); diff --git a/paimon-full-text/src/test/java/org/apache/paimon/fulltext/index/JavaPyNativeFullTextE2ETest.java b/paimon-full-text/src/test/java/org/apache/paimon/fulltext/index/JavaPyNativeFullTextE2ETest.java index 7883f0a58454..3ac0c4a3b5e8 100644 --- a/paimon-full-text/src/test/java/org/apache/paimon/fulltext/index/JavaPyNativeFullTextE2ETest.java +++ b/paimon-full-text/src/test/java/org/apache/paimon/fulltext/index/JavaPyNativeFullTextE2ETest.java @@ -190,7 +190,7 @@ private void writeTableWithNativeFullTextIndex( GlobalIndexBuilderUtils.createIndexWriter( table, NativeFullTextGlobalIndexerFactory.IDENTIFIER, - contentField, + Collections.singletonList(contentField), indexOptions); // Write the same text data to the index. diff --git a/paimon-full-text/src/test/java/org/apache/paimon/fulltext/index/NativeFullTextGlobalIndexerFactoryTest.java b/paimon-full-text/src/test/java/org/apache/paimon/fulltext/index/NativeFullTextGlobalIndexerFactoryTest.java index b96c228c2051..8d78bfe28652 100644 --- a/paimon-full-text/src/test/java/org/apache/paimon/fulltext/index/NativeFullTextGlobalIndexerFactoryTest.java +++ b/paimon-full-text/src/test/java/org/apache/paimon/fulltext/index/NativeFullTextGlobalIndexerFactoryTest.java @@ -25,6 +25,8 @@ import org.junit.jupiter.api.Test; +import java.util.Collections; + import static org.assertj.core.api.Assertions.assertThat; /** Tests for {@link NativeFullTextGlobalIndexerFactory}. */ @@ -37,7 +39,7 @@ public void testFactoryAcceptsFullTextSearcherPoolOption() { Options options = new Options(); options.set("full-text.searcher-pool.max-size", "0"); - GlobalIndexer indexer = factory.create(field, options); + GlobalIndexer indexer = factory.create(Collections.singletonList(field), options); assertThat(indexer).isInstanceOf(NativeFullTextGlobalIndexer.class); } diff --git a/paimon-full-text/src/test/java/org/apache/paimon/fulltext/index/NativeFullTextRowFilterTest.java b/paimon-full-text/src/test/java/org/apache/paimon/fulltext/index/NativeFullTextRowFilterTest.java index e2898bdaee3c..a659b1860a7d 100644 --- a/paimon-full-text/src/test/java/org/apache/paimon/fulltext/index/NativeFullTextRowFilterTest.java +++ b/paimon-full-text/src/test/java/org/apache/paimon/fulltext/index/NativeFullTextRowFilterTest.java @@ -450,7 +450,7 @@ private Dataset writeDataset( GlobalIndexBuilderUtils.createIndexWriter( table, NativeFullTextGlobalIndexerFactory.IDENTIFIER, - contentField, + Collections.singletonList(contentField), table.coreOptions().toConfiguration()); for (int i = 0; i < fullTextIndexedRows; i++) { textWriter.write(BinaryString.fromString(contents.get(i)), i); @@ -472,7 +472,7 @@ private Dataset writeDataset( GlobalIndexBuilderUtils.createIndexWriter( table, BTreeGlobalIndexerFactory.IDENTIFIER, - categoryField, + Collections.singletonList(categoryField), table.coreOptions().toConfiguration()); // The btree writer is an SST writer: keys must arrive in sorted order. List rowIdsByCategory = new ArrayList<>(rowCount); diff --git a/paimon-full-text/src/test/java/org/apache/paimon/fulltext/index/NativePrimaryKeyFullTextIndexTest.java b/paimon-full-text/src/test/java/org/apache/paimon/fulltext/index/NativePrimaryKeyFullTextIndexTest.java index a8da5ddebb6a..db0c26b7cc62 100644 --- a/paimon-full-text/src/test/java/org/apache/paimon/fulltext/index/NativePrimaryKeyFullTextIndexTest.java +++ b/paimon-full-text/src/test/java/org/apache/paimon/fulltext/index/NativePrimaryKeyFullTextIndexTest.java @@ -128,7 +128,8 @@ void testPreservesNullOrdinalFiltersDeletedRowsAndOrdersByNativeScore() throws E deletionVector.delete(2); Map deletionVectors = Collections.singletonMap(dataFile.fileName(), deletionVector); - GlobalIndexer indexer = GlobalIndexer.create("full-text", TEXT_FIELD, options); + GlobalIndexer indexer = + GlobalIndexer.create("full-text", Collections.singletonList(TEXT_FIELD), options); PrimaryKeyFullTextBucketSearch search = new PrimaryKeyFullTextBucketSearch( (payload, totalRowCount) -> @@ -169,7 +170,8 @@ void testPersistsTokenizerOptionsInPkArchiveMetadata() throws Exception { } private IndexFileMeta buildArchive(List texts, Options options) throws Exception { - GlobalIndexer indexer = GlobalIndexer.create("full-text", TEXT_FIELD, options); + GlobalIndexer indexer = + GlobalIndexer.create("full-text", Collections.singletonList(TEXT_FIELD), options); GlobalIndexSingleColumnWriter writer = (GlobalIndexSingleColumnWriter) indexer.createWriter(fileWriter()); for (int i = 0; i < texts.size(); i++) { diff --git a/paimon-lumina/src/main/java/org/apache/paimon/lumina/index/LuminaVectorGlobalIndexerFactory.java b/paimon-lumina/src/main/java/org/apache/paimon/lumina/index/LuminaVectorGlobalIndexerFactory.java index 7d9c062feb61..5adc143a1a09 100644 --- a/paimon-lumina/src/main/java/org/apache/paimon/lumina/index/LuminaVectorGlobalIndexerFactory.java +++ b/paimon-lumina/src/main/java/org/apache/paimon/lumina/index/LuminaVectorGlobalIndexerFactory.java @@ -23,6 +23,8 @@ import org.apache.paimon.options.Options; import org.apache.paimon.types.DataField; +import java.util.List; + /** Factory for creating Lumina vector index. */ public class LuminaVectorGlobalIndexerFactory implements GlobalIndexerFactory { @@ -34,7 +36,12 @@ public String identifier() { } @Override - public GlobalIndexer create(DataField field, Options options) { + public GlobalIndexer create(List indexFields, Options options) { + if (indexFields.size() != 1) { + throw new UnsupportedOperationException( + "Index type '" + identifier() + "' requires exactly one index field."); + } + DataField field = indexFields.get(0); Options fieldOptions = LuminaVectorIndexOptions.resolveFieldOptions(field.name(), options); return new LuminaVectorGlobalIndexer(field.type(), fieldOptions); } diff --git a/paimon-lumina/src/test/java/org/apache/paimon/lumina/index/JavaPyLuminaE2ETest.java b/paimon-lumina/src/test/java/org/apache/paimon/lumina/index/JavaPyLuminaE2ETest.java index 3a311b5a58e2..1f14493b282b 100644 --- a/paimon-lumina/src/test/java/org/apache/paimon/lumina/index/JavaPyLuminaE2ETest.java +++ b/paimon-lumina/src/test/java/org/apache/paimon/lumina/index/JavaPyLuminaE2ETest.java @@ -174,7 +174,7 @@ public void testLuminaVectorIndexWrite(String fileFormat) throws Exception { GlobalIndexBuilderUtils.createIndexWriter( table, LuminaVectorGlobalIndexerFactory.IDENTIFIER, - embeddingField, + Collections.singletonList(embeddingField), indexOptions); for (int i = 0; i < vectors.length; i++) { @@ -284,7 +284,7 @@ public void testLuminaVectorWithBTreeIndexWrite() throws Exception { GlobalIndexBuilderUtils.createIndexWriter( table, LuminaVectorGlobalIndexerFactory.IDENTIFIER, - embeddingField, + Collections.singletonList(embeddingField), indexOptions); for (int i = 0; i < vectors.length; i++) { vectorWriter.write(vectors[i], i); @@ -305,7 +305,10 @@ public void testLuminaVectorWithBTreeIndexWrite() throws Exception { GlobalIndexSingleColumnWriter idWriter = (GlobalIndexSingleColumnWriter) GlobalIndexBuilderUtils.createIndexWriter( - table, BTreeGlobalIndexerFactory.IDENTIFIER, idField, indexOptions); + table, + BTreeGlobalIndexerFactory.IDENTIFIER, + Collections.singletonList(idField), + indexOptions); for (int i = 0; i < vectors.length; i++) { idWriter.write(i, i); } diff --git a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/globalindex/DefaultGlobalIndexBuilder.java b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/globalindex/DefaultGlobalIndexBuilder.java index 5734cf84e396..a78ad1d7e19a 100644 --- a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/globalindex/DefaultGlobalIndexBuilder.java +++ b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/globalindex/DefaultGlobalIndexBuilder.java @@ -168,7 +168,7 @@ public CommitMessage build(CloseableIterator data) throws IOExcepti private List writePaimonRows( CloseableIterator rows, LongCounter rowCounter) throws IOException { GlobalIndexWriter indexWriter = - createIndexWriter(table, indexType, indexField, extraFields, options); + createIndexWriter(table, indexType, indexedFields(), options); boolean multiColumn = !extraFields.isEmpty(); try { 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 7431d062f884..3b9a125eb128 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 @@ -131,7 +131,8 @@ public List buildIndex( int maxParallelism = options.get(SortedIndexOptions.SORTED_INDEX_BUILD_MAX_PARALLELISM); List allMessages = new ArrayList<>(); - GlobalIndexer indexer = GlobalIndexer.create(indexType, indexField, options); + GlobalIndexer indexer = + GlobalIndexer.create(indexType, Collections.singletonList(indexField), options); if (!(indexer instanceof SortedGlobalIndexer)) { throw new IllegalArgumentException( "Index algorithm " + indexType + " does not expose sorted index keys."); 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..ed349c65340e 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 @@ -176,11 +176,7 @@ public InternalRow[] call(InternalRow args) { // multi-column throw from GlobalIndexerFactory#create, which happens // before any indexer side effect. try { - GlobalIndexer.create( - indexType, - indexFields.get(0), - indexFields.subList(1, indexFields.size()), - userOptions); + GlobalIndexer.create(indexType, indexFields, userOptions); } catch (UnsupportedOperationException e) { throw new IllegalArgumentException( String.format( diff --git a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/read/SparkDataEvolutionVectorRead.java b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/read/SparkDataEvolutionVectorRead.java index 47e2a1114b23..5b0896a50c5b 100644 --- a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/read/SparkDataEvolutionVectorRead.java +++ b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/read/SparkDataEvolutionVectorRead.java @@ -42,6 +42,7 @@ import java.io.IOException; import java.util.ArrayList; +import java.util.Collections; import java.util.List; import java.util.Map; import java.util.Optional; @@ -126,7 +127,9 @@ protected ScoredGlobalIndexResult readIndexSplitsInSpark( group -> { GlobalIndexer taskGlobalIndexer = GlobalIndexerFactoryUtils.load(indexType) - .create(vectorColumn, table.coreOptions().toConfiguration()); + .create( + Collections.singletonList(vectorColumn), + table.coreOptions().toConfiguration()); IndexPathFactory indexPathFactory = table.store().pathFactory().globalIndexFileFactory(); diff --git a/paimon-vector/src/main/java/org/apache/paimon/vector/index/NativeVectorGlobalIndexerFactory.java b/paimon-vector/src/main/java/org/apache/paimon/vector/index/NativeVectorGlobalIndexerFactory.java index ad11f4bf16eb..01b3a71a1cf8 100644 --- a/paimon-vector/src/main/java/org/apache/paimon/vector/index/NativeVectorGlobalIndexerFactory.java +++ b/paimon-vector/src/main/java/org/apache/paimon/vector/index/NativeVectorGlobalIndexerFactory.java @@ -26,6 +26,7 @@ import org.apache.paimon.types.VectorType; import java.util.LinkedHashMap; +import java.util.List; import java.util.Map; /** Factory for creating vector indexes backed by paimon-vector-index-java. */ @@ -36,7 +37,12 @@ public abstract class NativeVectorGlobalIndexerFactory implements GlobalIndexerF static final double DEFAULT_TRAIN_SAMPLE_RATIO = 1.0; @Override - public GlobalIndexer create(DataField field, Options options) { + public GlobalIndexer create(List indexFields, Options options) { + if (indexFields.size() != 1) { + throw new UnsupportedOperationException( + "Index type '" + identifier() + "' requires exactly one index field."); + } + DataField field = indexFields.get(0); String identifier = identifier(); return new NativeVectorGlobalIndexer( field.type(), diff --git a/paimon-vector/src/test/java/org/apache/paimon/JavaPyE2ETest.java b/paimon-vector/src/test/java/org/apache/paimon/JavaPyE2ETest.java index 8b87bd2bb5cd..7975fc76fb64 100644 --- a/paimon-vector/src/test/java/org/apache/paimon/JavaPyE2ETest.java +++ b/paimon-vector/src/test/java/org/apache/paimon/JavaPyE2ETest.java @@ -157,7 +157,7 @@ public void testVindexVectorIndexWrite() throws Exception { GlobalIndexBuilderUtils.createIndexWriter( table, IvfFlatVectorGlobalIndexerFactory.IDENTIFIER, - embeddingField, + Collections.singletonList(embeddingField), indexOptions); for (int i = 0; i < vectors.length; i++) { @@ -267,7 +267,7 @@ public void testVindexVectorRawFallbackWrite() throws Exception { GlobalIndexBuilderUtils.createIndexWriter( table, IvfFlatVectorGlobalIndexerFactory.IDENTIFIER, - embeddingField, + Collections.singletonList(embeddingField), indexOptions); for (int i = 0; i < indexedVectors.length; i++) {