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);

Reply via email to