This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new 869e7253f9 [core] Add diagnostic logs for global-index scans (#8657)
869e7253f9 is described below
commit 869e7253f924a8f2d06805e2309f671383ce4b47
Author: jerry <[email protected]>
AuthorDate: Sat Jul 18 12:07:52 2026 +0800
[core] Add diagnostic logs for global-index scans (#8657)
---
.../paimon/globalindex/UnionGlobalIndexReader.java | 56 +++++++++++++------
.../globalindex/GlobalIndexEvaluatorTest.java | 21 ++++++++
.../paimon/globalindex/DataEvolutionBatchScan.java | 22 +++++++-
.../paimon/globalindex/GlobalIndexScanner.java | 19 ++++++-
.../paimon/table/BtreeGlobalIndexTableTest.java | 63 ++++++++++++++++++++++
5 files changed, 161 insertions(+), 20 deletions(-)
diff --git
a/paimon-common/src/main/java/org/apache/paimon/globalindex/UnionGlobalIndexReader.java
b/paimon-common/src/main/java/org/apache/paimon/globalindex/UnionGlobalIndexReader.java
index 9914fa1115..81c83076c2 100644
---
a/paimon-common/src/main/java/org/apache/paimon/globalindex/UnionGlobalIndexReader.java
+++
b/paimon-common/src/main/java/org/apache/paimon/globalindex/UnionGlobalIndexReader.java
@@ -28,6 +28,7 @@ import java.util.List;
import java.util.Optional;
import java.util.concurrent.CompletableFuture;
import java.util.function.Function;
+import java.util.function.LongConsumer;
/**
* A {@link GlobalIndexReader} that combines results from multiple readers by
performing a union
@@ -36,9 +37,15 @@ import java.util.function.Function;
public class UnionGlobalIndexReader implements GlobalIndexReader {
private final List<GlobalIndexReader> readers;
+ private final LongConsumer durationConsumer;
public UnionGlobalIndexReader(List<GlobalIndexReader> readers) {
+ this(readers, null);
+ }
+
+ UnionGlobalIndexReader(List<GlobalIndexReader> readers, LongConsumer
durationConsumer) {
this.readers = readers;
+ this.durationConsumer = durationConsumer;
}
@Override
@@ -132,6 +139,7 @@ public class UnionGlobalIndexReader implements
GlobalIndexReader {
@Override
public CompletableFuture<Optional<ScoredGlobalIndexResult>>
visitVectorSearch(
VectorSearch vectorSearch) {
+ long start = durationConsumer == null ? 0L : System.nanoTime();
List<CompletableFuture<Optional<ScoredGlobalIndexResult>>> futures =
new ArrayList<>(readers.size());
for (GlobalIndexReader reader : readers) {
@@ -153,33 +161,47 @@ public class UnionGlobalIndexReader implements
GlobalIndexReader {
}
}
return result;
+ })
+ .whenComplete(
+ (ignored, throwable) -> {
+ if (durationConsumer != null) {
+ durationConsumer.accept(System.nanoTime() -
start);
+ }
});
}
private CompletableFuture<Optional<GlobalIndexResult>> unionAsync(
Function<GlobalIndexReader,
CompletableFuture<Optional<GlobalIndexResult>>> visitor) {
+ long start = durationConsumer == null ? 0L : System.nanoTime();
List<CompletableFuture<Optional<GlobalIndexResult>>> futures =
new ArrayList<>(readers.size());
for (GlobalIndexReader reader : readers) {
futures.add(visitor.apply(reader));
}
- return CompletableFuture.allOf(futures.toArray(new
CompletableFuture[0]))
- .thenApply(
- v -> {
- Optional<GlobalIndexResult> result =
Optional.empty();
- for
(CompletableFuture<Optional<GlobalIndexResult>> f : futures) {
- Optional<GlobalIndexResult> current = f.join();
- if (!current.isPresent()) {
- continue;
- }
- if (!result.isPresent()) {
- result = current;
- } else {
- result =
Optional.of(result.get().or(current.get()));
- }
- }
- return result;
- });
+ CompletableFuture<Optional<GlobalIndexResult>> result =
+ CompletableFuture.allOf(futures.toArray(new
CompletableFuture[0]))
+ .thenApply(
+ v -> {
+ Optional<GlobalIndexResult> union =
Optional.empty();
+ for
(CompletableFuture<Optional<GlobalIndexResult>> f :
+ futures) {
+ Optional<GlobalIndexResult> current =
f.join();
+ if (!current.isPresent()) {
+ continue;
+ }
+ if (!union.isPresent()) {
+ union = current;
+ } else {
+ union =
Optional.of(union.get().or(current.get()));
+ }
+ }
+ return union;
+ });
+ if (durationConsumer != null) {
+ return result.whenComplete(
+ (ignored, throwable) ->
durationConsumer.accept(System.nanoTime() - start));
+ }
+ return result;
}
@Override
diff --git
a/paimon-common/src/test/java/org/apache/paimon/globalindex/GlobalIndexEvaluatorTest.java
b/paimon-common/src/test/java/org/apache/paimon/globalindex/GlobalIndexEvaluatorTest.java
index b3ffad5263..b2e75a4b6e 100644
---
a/paimon-common/src/test/java/org/apache/paimon/globalindex/GlobalIndexEvaluatorTest.java
+++
b/paimon-common/src/test/java/org/apache/paimon/globalindex/GlobalIndexEvaluatorTest.java
@@ -23,6 +23,7 @@ import org.apache.paimon.predicate.CompoundPredicate;
import org.apache.paimon.predicate.FieldRef;
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.predicate.PredicateBuilder;
+import org.apache.paimon.predicate.VectorSearch;
import org.apache.paimon.types.DataField;
import org.apache.paimon.types.DataTypes;
import org.apache.paimon.types.RowType;
@@ -681,6 +682,26 @@ class GlobalIndexEvaluatorTest {
assertThat(secondClosed).isTrue();
}
+ @Test
+ void testUnionReaderReportsVectorSearchDuration() {
+ AtomicInteger callbacks = new AtomicInteger();
+ GlobalIndexReader reader =
+ new StubGlobalIndexReader(null) {
+ @Override
+ public
CompletableFuture<Optional<ScoredGlobalIndexResult>> visitVectorSearch(
+ VectorSearch vectorSearch) {
+ return
CompletableFuture.completedFuture(Optional.empty());
+ }
+ };
+ UnionGlobalIndexReader union =
+ new UnionGlobalIndexReader(
+ Collections.singletonList(reader), ignored ->
callbacks.incrementAndGet());
+
+ union.visitVectorSearch(new VectorSearch(new float[] {1}, 1,
"test")).join();
+
+ assertThat(callbacks).hasValue(1);
+ }
+
private static void assertBitmapContainsExactly(
RoaringNavigableMap64 bitmap, long... expected) {
assertThat(bitmap.getLongCardinality()).isEqualTo(expected.length);
diff --git
a/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionBatchScan.java
b/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionBatchScan.java
index 57be42ff73..ebdfb0894f 100644
---
a/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionBatchScan.java
+++
b/paimon-core/src/main/java/org/apache/paimon/globalindex/DataEvolutionBatchScan.java
@@ -264,6 +264,8 @@ public class DataEvolutionBatchScan implements
DataTableScan {
if (result instanceof ScoredGlobalIndexResult) {
scoreGetter = ((ScoredGlobalIndexResult)
result).scoreGetter();
}
+ } else if (filter != null) {
+ LOG.info("Scan table '{}' without global index.",
table.name());
}
}
@@ -292,17 +294,33 @@ public class DataEvolutionBatchScan implements
DataTableScan {
}
PartitionPredicate partitionFilter =
batchScan.snapshotReader().manifestsReader().partitionFilter();
+ long totalStart = System.nanoTime();
Optional<GlobalIndexScanner> optionalScanner =
GlobalIndexScanner.create(table, partitionFilter,
globalIndexFilter);
+ long metadataDuration = System.nanoTime() - totalStart;
if (!optionalScanner.isPresent()) {
return Optional.empty();
}
try (GlobalIndexScanner scanner = optionalScanner.get()) {
+ long lookupStart = System.nanoTime();
Optional<GlobalIndexResult> result =
scanner.scan(globalIndexFilter);
+ long lookupDuration = System.nanoTime() - lookupStart;
if (result.isPresent()) {
- LOG.info("Scan table '{}' with global index.", table.name());
- return
Optional.of(result.get().or(scanner.unindexedRows(globalIndexFilter)));
+ long coverageStart = System.nanoTime();
+ GlobalIndexResult finalResult =
+
result.get().or(scanner.unindexedRows(globalIndexFilter));
+ long coverageDuration = System.nanoTime() - coverageStart;
+ long totalDuration = System.nanoTime() - totalStart;
+ LOG.info(
+ "Scan table '{}' with global index. searchMode='{}',
total={} ms, metadata={} ms, lookup={} ms, coverage={} ms.",
+ table.name(),
+ options.globalIndexSearchMode(),
+ totalDuration / 1_000_000,
+ metadataDuration / 1_000_000,
+ lookupDuration / 1_000_000,
+ coverageDuration / 1_000_000);
+ return Optional.of(finalResult);
}
return Optional.empty();
} catch (IOException e) {
diff --git
a/paimon-core/src/main/java/org/apache/paimon/globalindex/GlobalIndexScanner.java
b/paimon-core/src/main/java/org/apache/paimon/globalindex/GlobalIndexScanner.java
index e9ec45b041..df9a306120 100644
---
a/paimon-core/src/main/java/org/apache/paimon/globalindex/GlobalIndexScanner.java
+++
b/paimon-core/src/main/java/org/apache/paimon/globalindex/GlobalIndexScanner.java
@@ -36,6 +36,9 @@ import org.apache.paimon.utils.Filter;
import org.apache.paimon.utils.Range;
import org.apache.paimon.utils.RoaringNavigableMap64;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
import javax.annotation.Nullable;
import java.io.Closeable;
@@ -63,6 +66,8 @@ import static
org.apache.paimon.utils.Preconditions.checkNotNull;
/** Scanner for shard-based global indexes. */
public class GlobalIndexScanner implements Closeable {
+ private static final Logger LOG =
LoggerFactory.getLogger(GlobalIndexScanner.class);
+
private final Options options;
private final RowType rowType;
private final ExecutorService executor;
@@ -325,7 +330,19 @@ public class GlobalIndexScanner implements Closeable {
for (CompletableFuture<GlobalIndexReader> future : futures) {
unionReader.add(future.join());
}
- readers.add(new UnionGlobalIndexReader(unionReader));
+ readers.add(
+ new UnionGlobalIndexReader(
+ unionReader,
+ duration ->
+ LOG.info(
+ "Global index lookup table='{}',
type='{}', fields='{}', lookup={} ms.",
+ table.name(),
+ indexType,
+ group.fieldIds.stream()
+ .map(rowType::getField)
+ .map(DataField::name)
+
.collect(Collectors.toList()),
+ duration / 1_000_000)));
}
return readers;
diff --git
a/paimon-core/src/test/java/org/apache/paimon/table/BtreeGlobalIndexTableTest.java
b/paimon-core/src/test/java/org/apache/paimon/table/BtreeGlobalIndexTableTest.java
index 5f6a9152fd..cf4ac40ac8 100644
---
a/paimon-core/src/test/java/org/apache/paimon/table/BtreeGlobalIndexTableTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/table/BtreeGlobalIndexTableTest.java
@@ -47,8 +47,15 @@ import org.apache.paimon.types.DataTypes;
import org.apache.paimon.utils.Range;
import org.apache.paimon.utils.RoaringNavigableMap64;
+import org.apache.logging.log4j.Level;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.core.Logger;
+import org.apache.logging.log4j.core.appender.OutputStreamAppender;
+import org.apache.logging.log4j.core.layout.PatternLayout;
import org.junit.jupiter.api.Test;
+import java.io.ByteArrayOutputStream;
+import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
@@ -159,6 +166,62 @@ public class BtreeGlobalIndexTableTest extends
DataEvolutionTestBase {
assertThat(splits).allMatch(split -> split instanceof DataSplit);
}
+ @Test
+ public void testGlobalIndexDiagnosticLogs() throws Exception {
+ write(10L);
+ createIndex("f1");
+
+ FileStoreTable table = (FileStoreTable) catalog.getTable(identifier());
+ ByteArrayOutputStream output = new ByteArrayOutputStream();
+ OutputStreamAppender appender =
+ OutputStreamAppender.newBuilder()
+ .setName("global-index-diagnostic-test")
+ .setTarget(output)
+
.setLayout(PatternLayout.newBuilder().withPattern("%level %msg%n").build())
+ .build();
+ Logger scanLogger = (Logger)
LogManager.getLogger(DataEvolutionBatchScan.class);
+ Logger scannerLogger = (Logger)
LogManager.getLogger(GlobalIndexScanner.class);
+ Level previousScanLevel = scanLogger.getLevel();
+ Level previousScannerLevel = scannerLogger.getLevel();
+
+ appender.start();
+ scanLogger.addAppender(appender);
+ scannerLogger.addAppender(appender);
+ scanLogger.setLevel(Level.INFO);
+ scannerLogger.setLevel(Level.INFO);
+ try {
+ Predicate predicate =
+ new PredicateBuilder(table.rowType()).equal(1,
BinaryString.fromString("a7"));
+ table.newReadBuilder().withFilter(predicate).newScan().plan();
+
+ PredicateBuilder rowIdBuilder =
+ new
PredicateBuilder(SpecialFields.rowTypeWithRowId(table.rowType()));
+ int rowIdIndex = table.rowType().getFieldCount();
+ Predicate mixedRowIdPredicate =
+ PredicateBuilder.or(
+ rowIdBuilder.equal(rowIdIndex, 1L),
+ rowIdBuilder.equal(1,
BinaryString.fromString("a7")));
+
table.newReadBuilder().withFilter(mixedRowIdPredicate).newScan().plan();
+
+ String logs = new String(output.toByteArray(),
StandardCharsets.UTF_8);
+ assertThat(logs)
+ .containsPattern(
+ "INFO Scan table '[^']+' with global index\\. "
+ + "searchMode='fast', total=\\d+ ms,
metadata=\\d+ ms, "
+ + "lookup=\\d+ ms, coverage=\\d+ ms\\.")
+ .containsPattern(
+ "INFO Global index lookup table='[^']+',
type='btree', "
+ + "fields='\\[f1\\]', lookup=\\d+ ms\\.")
+ .contains("INFO Scan table '" + table.name() + "' without
global index.");
+ } finally {
+ scanLogger.setLevel(previousScanLevel);
+ scannerLogger.setLevel(previousScannerLevel);
+ scanLogger.removeAppender(appender);
+ scannerLogger.removeAppender(appender);
+ appender.stop();
+ }
+ }
+
@Test
public void testBTreeGlobalIndexSearchModeControlsUnindexedData() throws
Exception {
write(500L);