From 3bb99650d6ad714e82f0cf7ee16e1458f2259a56 Mon Sep 17 00:00:00 2001 From: yangjie01 Date: Sat, 26 Sep 2026 02:46:06 +0800 Subject: [PATCH] [core] Validate ANN segment metric against the built-in metric The metric guard in PkVectorAnnSegmentSearcher compared two values that both derive from the current table config, so it could never fail: changing the metric option after ANN segments were built let searches score old segments with the new metric and silently return wrong distances. Record the effective metric in the segment's vector index metadata when writing it, expose it through VectorGlobalIndexer.segmentMetric, and reject segments whose recorded metric differs from the current one. Legacy segments that record no metric are still accepted. Assisted-by: GLM-5.3 --- .../globalindex/VectorGlobalIndexer.java | 9 ++ .../pkvector/PkVectorAnnSegmentSearcher.java | 23 +++++ .../PkVectorAnnSegmentSearcherMetricTest.java | 94 +++++++++++++++++++ .../index/NativeVectorGlobalIndexWriter.java | 5 +- .../index/NativeVectorGlobalIndexer.java | 11 +++ .../paimon/vector/index/VectorIndexMeta.java | 33 +++++-- .../index/NativeVectorGlobalIndexTest.java | 2 +- .../vector/index/VectorIndexMetaTest.java | 38 ++++++++ 8 files changed, 206 insertions(+), 9 deletions(-) create mode 100644 paimon-core/src/test/java/org/apache/paimon/index/pkvector/PkVectorAnnSegmentSearcherMetricTest.java create mode 100644 paimon-vector/src/test/java/org/apache/paimon/vector/index/VectorIndexMetaTest.java diff --git a/paimon-common/src/main/java/org/apache/paimon/globalindex/VectorGlobalIndexer.java b/paimon-common/src/main/java/org/apache/paimon/globalindex/VectorGlobalIndexer.java index 63166ef7c1f1..35f61d705129 100644 --- a/paimon-common/src/main/java/org/apache/paimon/globalindex/VectorGlobalIndexer.java +++ b/paimon-common/src/main/java/org/apache/paimon/globalindex/VectorGlobalIndexer.java @@ -23,4 +23,13 @@ public interface VectorGlobalIndexer extends GlobalIndexer { /** Returns the metric name used to convert vector distances to comparable scores. */ String metric(); + + /** + * Returns the metric recorded in a segment's index metadata when it was built, or {@code null} + * when the metadata records none (legacy segments or indexers that do not persist it). + * Searchers use it to reject segments built with a different metric than the current one. + */ + default String segmentMetric(byte[] indexMeta) { + return null; + } } 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..0db7f6a43cc5 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 @@ -193,6 +193,7 @@ CompletableFuture> searchAsync( "ANN segment metric %s does not match index reader metric %s.", metric, readerMetric); + checkSegmentMetric(indexer, metric, globalIndexMeta.indexMeta()); GlobalIndexIOMeta ioMeta = new GlobalIndexIOMeta( @@ -291,6 +292,7 @@ CompletableFuture>> searchBatchAsync( "ANN segment metric %s does not match index reader metric %s.", metric, readerMetric); + checkSegmentMetric(indexer, metric, globalIndexMeta.indexMeta()); GlobalIndexIOMeta ioMeta = new GlobalIndexIOMeta( @@ -491,4 +493,25 @@ private FilePosition(String dataFileName, long rowPosition) { this.rowPosition = rowPosition; } } + + /** + * The guard above compares two values from the current config; the segment metadata records + * what the index was actually built with, and a mismatch would mean silently wrong distances. + * Legacy segments record no metric and are not checked. + */ + static void checkSegmentMetric( + GlobalIndexer indexer, String normalizedMetric, byte[] indexMeta) { + if (!(indexer instanceof VectorGlobalIndexer)) { + return; + } + String segmentMetric = ((VectorGlobalIndexer) indexer).segmentMetric(indexMeta); + if (segmentMetric != null) { + String normalized = VectorSearchMetric.normalize(segmentMetric); + checkArgument( + normalizedMetric.equals(normalized), + "ANN segment was built with metric %s but the current metric is %s.", + normalized, + normalizedMetric); + } + } } diff --git a/paimon-core/src/test/java/org/apache/paimon/index/pkvector/PkVectorAnnSegmentSearcherMetricTest.java b/paimon-core/src/test/java/org/apache/paimon/index/pkvector/PkVectorAnnSegmentSearcherMetricTest.java new file mode 100644 index 000000000000..59ca0d646b9e --- /dev/null +++ b/paimon-core/src/test/java/org/apache/paimon/index/pkvector/PkVectorAnnSegmentSearcherMetricTest.java @@ -0,0 +1,94 @@ +/* + * 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.index.pkvector; + +import org.apache.paimon.globalindex.GlobalIndexIOMeta; +import org.apache.paimon.globalindex.GlobalIndexReader; +import org.apache.paimon.globalindex.GlobalIndexWriter; +import org.apache.paimon.globalindex.VectorGlobalIndexer; +import org.apache.paimon.globalindex.io.GlobalIndexFileReader; +import org.apache.paimon.globalindex.io.GlobalIndexFileWriter; + +import org.junit.jupiter.api.Test; + +import java.io.IOException; +import java.util.List; +import java.util.concurrent.ExecutorService; + +import static org.assertj.core.api.Assertions.assertThatCode; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +/** Tests for the segment metric validation in {@link PkVectorAnnSegmentSearcher}. */ +class PkVectorAnnSegmentSearcherMetricTest { + + private static VectorGlobalIndexer indexerWithMetric(String segmentMetric) { + return new VectorGlobalIndexer() { + @Override + public String metric() { + return "inner_product"; + } + + @Override + public String segmentMetric(byte[] indexMeta) { + return segmentMetric; + } + + @Override + public GlobalIndexWriter createWriter(GlobalIndexFileWriter fileWriter) + throws IOException { + throw new UnsupportedOperationException(); + } + + @Override + public GlobalIndexReader createReader( + GlobalIndexFileReader fileReader, + List files, + long totalRowCount, + List rowRanges, + ExecutorService executor) { + throw new UnsupportedOperationException(); + } + }; + } + + @Test + void testSegmentMetricMismatchIsRejected() { + assertThatThrownBy( + () -> + PkVectorAnnSegmentSearcher.checkSegmentMetric( + indexerWithMetric("cosine"), "l2", new byte[] {1})) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("built with metric cosine but the current metric is l2"); + } + + @Test + void testMatchingAndLegacySegmentMetricsPass() { + assertThatCode( + () -> + PkVectorAnnSegmentSearcher.checkSegmentMetric( + indexerWithMetric("cosine"), "cosine", new byte[] {1})) + .doesNotThrowAnyException(); + // legacy segments record no metric and must not be rejected + assertThatCode( + () -> + PkVectorAnnSegmentSearcher.checkSegmentMetric( + indexerWithMetric(null), "l2", new byte[] {1})) + .doesNotThrowAnyException(); + } +} diff --git a/paimon-vector/src/main/java/org/apache/paimon/vector/index/NativeVectorGlobalIndexWriter.java b/paimon-vector/src/main/java/org/apache/paimon/vector/index/NativeVectorGlobalIndexWriter.java index d0f52288de3d..82a20993fb79 100644 --- a/paimon-vector/src/main/java/org/apache/paimon/vector/index/NativeVectorGlobalIndexWriter.java +++ b/paimon-vector/src/main/java/org/apache/paimon/vector/index/NativeVectorGlobalIndexWriter.java @@ -290,7 +290,10 @@ private ResultEntry buildIndex() throws IOException { identifier, System.currentTimeMillis() - buildStart); - VectorIndexMeta meta = new VectorIndexMeta(); + VectorIndexMeta meta = + new VectorIndexMeta( + nativeOptions.getOrDefault( + "metric", NativeVectorGlobalIndexer.DEFAULT_METRIC)); return new ResultEntry(fileName, rowCount, meta.serialize()); } } diff --git a/paimon-vector/src/main/java/org/apache/paimon/vector/index/NativeVectorGlobalIndexer.java b/paimon-vector/src/main/java/org/apache/paimon/vector/index/NativeVectorGlobalIndexer.java index 8d7c3e4ab924..35749834c226 100644 --- a/paimon-vector/src/main/java/org/apache/paimon/vector/index/NativeVectorGlobalIndexer.java +++ b/paimon-vector/src/main/java/org/apache/paimon/vector/index/NativeVectorGlobalIndexer.java @@ -27,6 +27,8 @@ import org.apache.paimon.types.DataType; import org.apache.paimon.utils.Range; +import java.io.IOException; +import java.io.UncheckedIOException; import java.util.List; import java.util.Map; import java.util.Objects; @@ -86,6 +88,15 @@ public GlobalIndexReader createReader( return new NativeVectorGlobalIndexReader(fileReader, files, fieldType, executor); } + @Override + public String segmentMetric(byte[] indexMeta) { + try { + return VectorIndexMeta.deserialize(indexMeta).metric(); + } catch (IOException e) { + throw new UncheckedIOException("Failed to read vector index metadata.", e); + } + } + @Override public String metric() { return options.getOrDefault("metric", DEFAULT_METRIC); diff --git a/paimon-vector/src/main/java/org/apache/paimon/vector/index/VectorIndexMeta.java b/paimon-vector/src/main/java/org/apache/paimon/vector/index/VectorIndexMeta.java index ea18e7efebc9..4ffe3406b23e 100644 --- a/paimon-vector/src/main/java/org/apache/paimon/vector/index/VectorIndexMeta.java +++ b/paimon-vector/src/main/java/org/apache/paimon/vector/index/VectorIndexMeta.java @@ -21,17 +21,20 @@ import org.apache.paimon.shade.jackson2.com.fasterxml.jackson.core.type.TypeReference; import org.apache.paimon.shade.jackson2.com.fasterxml.jackson.databind.ObjectMapper; +import javax.annotation.Nullable; + import java.io.IOException; import java.io.Serializable; -import java.util.Collections; import java.util.LinkedHashMap; import java.util.Map; /** * Metadata for a vector index file. * - *

Serialized as an empty JSON {@code Map}. Search-time parameters are passed - * through {@link org.apache.paimon.predicate.VectorSearch#options()}. + *

Serialized as a JSON {@code Map}; it records the metric the index was built + * with so searches can reject segments built under a different metric. Legacy segments carry an + * empty map. Search-time parameters are passed through {@link + * org.apache.paimon.predicate.VectorSearch#options()}. */ public class VectorIndexMeta implements Serializable { @@ -42,14 +45,30 @@ public class VectorIndexMeta implements Serializable { private static final TypeReference> MAP_TYPE_REF = new TypeReference>() {}; - VectorIndexMeta() {} + private static final String METRIC_KEY = "metric"; + + @Nullable private final String metric; + + VectorIndexMeta(@Nullable String metric) { + this.metric = metric; + } public byte[] serialize() throws IOException { - return OBJECT_MAPPER.writeValueAsBytes(Collections.emptyMap()); + Map data = new LinkedHashMap<>(); + if (metric != null) { + data.put(METRIC_KEY, metric); + } + return OBJECT_MAPPER.writeValueAsBytes(data); } public static VectorIndexMeta deserialize(byte[] data) throws IOException { - Map ignored = OBJECT_MAPPER.readValue(data, MAP_TYPE_REF); - return new VectorIndexMeta(); + Map map = OBJECT_MAPPER.readValue(data, MAP_TYPE_REF); + return new VectorIndexMeta(map.get(METRIC_KEY)); + } + + /** The metric this index was built with, or null for legacy segments that record none. */ + @Nullable + public String metric() { + return metric; } } diff --git a/paimon-vector/src/test/java/org/apache/paimon/vector/index/NativeVectorGlobalIndexTest.java b/paimon-vector/src/test/java/org/apache/paimon/vector/index/NativeVectorGlobalIndexTest.java index 6d0180b1611f..40a81ee8c8bf 100644 --- a/paimon-vector/src/test/java/org/apache/paimon/vector/index/NativeVectorGlobalIndexTest.java +++ b/paimon-vector/src/test/java/org/apache/paimon/vector/index/NativeVectorGlobalIndexTest.java @@ -229,7 +229,7 @@ public void testVectorBatchSizeProtectsSingleJavaArrayAllocation() { @Test public void testMetaSerializationIsEmptyMap() throws IOException { - VectorIndexMeta meta = new VectorIndexMeta(); + VectorIndexMeta meta = new VectorIndexMeta(null); byte[] serialized = meta.serialize(); VectorIndexMeta deserialized = VectorIndexMeta.deserialize(serialized); diff --git a/paimon-vector/src/test/java/org/apache/paimon/vector/index/VectorIndexMetaTest.java b/paimon-vector/src/test/java/org/apache/paimon/vector/index/VectorIndexMetaTest.java new file mode 100644 index 000000000000..428e19b9606f --- /dev/null +++ b/paimon-vector/src/test/java/org/apache/paimon/vector/index/VectorIndexMetaTest.java @@ -0,0 +1,38 @@ +/* + * 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.vector.index; + +import org.junit.jupiter.api.Test; + +import static org.assertj.core.api.Assertions.assertThat; + +/** Tests for {@link VectorIndexMeta} metric persistence. */ +class VectorIndexMetaTest { + + @Test + void testMetricRoundTrip() throws Exception { + byte[] data = new VectorIndexMeta("cosine").serialize(); + assertThat(VectorIndexMeta.deserialize(data).metric()).isEqualTo("cosine"); + + // a null metric writes an empty map, matching legacy segments + byte[] legacy = new VectorIndexMeta(null).serialize(); + assertThat(VectorIndexMeta.deserialize(legacy).metric()).isNull(); + assertThat(VectorIndexMeta.deserialize("{}".getBytes()).metric()).isNull(); + } +}