yihua commented on code in PR #20107:
URL: https://github.com/apache/hudi/pull/20107#discussion_r4163156384


##########
hudi-common/src/main/java/org/apache/hudi/common/index/vector/VectorIndexRliArbitrator.java:
##########
@@ -0,0 +1,230 @@
+/*
+ * 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.hudi.common.index.vector;
+
+import org.apache.hudi.common.data.HoodieData;
+import org.apache.hudi.common.data.HoodieListData;
+import org.apache.hudi.common.model.HoodieRecordGlobalLocation;
+import org.apache.hudi.common.util.Option;
+import org.apache.hudi.metadata.HoodieTableMetadata;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+import java.util.Set;
+import java.util.stream.Collectors;
+
+/** Applies the RFC-109 record-level-index arbitration contract to vector 
finalists. */
+final class VectorIndexRliArbitrator {
+
+  private VectorIndexRliArbitrator() {
+  }
+
+  /**
+   * The RFC-109 RLI finalist arbiter. Resolves each finalist's current 
location from the
+   * record-level index (one batched {@code readRecordIndexLocationsWithKeys} 
over the distinct
+   * finalist keys) and tags it with a {@link VectorIndexArbiter.Decision} 
plus the resolved
+   * location.
+   *
+   * <p>Unlike {@code VectorIndexMdtSearchUtils.attachRecordLocations}, this 
does <em>not</em> drop candidates: it tags all
+   * of them so callers can tally {@code arbiterExclusions.stale} / {@code 
.deleted} and apply the
+   * mode-specific action (approx: exclude STALE + DELETED; exact: 
key-fallback STALE, exclude
+   * DELETED, positional SERVE). Resolved location semantics:
+   *
+   * <ul>
+   *   <li>{@code SERVE}: the posting's own location when present (positional 
trust), else the RLI
+   *       location.</li>
+   *   <li>{@code STALE}: the RLI current location, so exact mode can 
key-fetch at the live slice.</li>
+   *   <li>{@code DELETED}: {@code null}.</li>
+   * </ul>
+   */
+  static HoodieData<ScoredVectorPostingMatch> 
arbitrateFinalists(HoodieTableMetadata metadataTable,
+                                                                  
HoodieData<ScoredVectorPostingMatch> finalists) {
+    return arbitrateFinalists(metadataTable, finalists, false);
+  }
+
+  static HoodieData<ScoredVectorPostingMatch> arbitrateFinalists(
+      HoodieTableMetadata metadataTable,
+      HoodieData<ScoredVectorPostingMatch> finalists,
+      boolean partitionedRecordIndex) {
+    if (!partitionedRecordIndex) {
+      return arbitrateFinalistsForPartition(metadataTable, finalists, 
Option.empty());
+    }
+    List<String> partitions = 
finalists.map(ScoredVectorPostingMatch::getPartitionPath)
+        .distinct()
+        .collectAsList();
+    HoodieData<ScoredVectorPostingMatch> arbitrated = null;
+    for (String partition : partitions) {
+      HoodieData<ScoredVectorPostingMatch> partitionFinalists = finalists
+          .filter(candidate -> Objects.equals(partition, 
candidate.getPartitionPath()));
+      HoodieData<ScoredVectorPostingMatch> partitionResult = 
arbitrateFinalistsForPartition(
+          metadataTable, partitionFinalists, Option.ofNullable(partition));
+      arbitrated = arbitrated == null ? partitionResult : 
arbitrated.union(partitionResult);
+    }
+    return arbitrated == null ? HoodieListData.eager(Collections.emptyList()) 
: arbitrated;
+  }
+
+  private static HoodieData<ScoredVectorPostingMatch> 
arbitrateFinalistsForPartition(
+      HoodieTableMetadata metadataTable,
+      HoodieData<ScoredVectorPostingMatch> finalists,
+      Option<String> dataTablePartition) {
+    // Resolve current RLI locations for the finalist keys into a bounded 
driver-side map, then
+    // attach per candidate via map(...). The finalist set is a bounded 
candidate pool and
+    // {@code finalists} is already persisted upstream, so the two passes are 
cache hits.
+    //
+    // This deliberately avoids leftOuterJoin: HoodiePairData.leftOuterJoin 
requires both operands to
+    // share the same backing flavor, but readRecordIndexLocationsWithKeys 
returns list-backed pair
+    // data for a single-slice RLI (the common 1-file-group case) and 
RDD-backed for multi-slice.
+    // Joining an RDD-backed finalist set against list-backed locations throws 
ClassCastException.
+    // Attaching via map(...) preserves the finalists' backing (RDD stays RDD, 
list stays list).
+    List<String> distinctKeys = 
finalists.map(ScoredVectorPostingMatch::getRecordKey)
+        .distinct()
+        .collectAsList();
+    Map<String, HoodieRecordGlobalLocation> currentLocations = new HashMap<>();
+    if (!distinctKeys.isEmpty()) {
+      metadataTable.readRecordIndexLocationsWithKeys(
+              HoodieListData.eager(distinctKeys), dataTablePartition)
+          .collectAsList()
+          .forEach(pair -> currentLocations.put(pair.getKey(), 
pair.getValue()));

Review Comment:
   See https://github.com/apache/hudi/pull/19802#issuecomment-5946183886
   
   Wrt partitioned RLI, `readRecordIndexLocationsWithKeys(keys, partition)` has 
an implicit precondition that all keys in a batch hash to the same RLI shard, 
and every current caller (Spark write/read, Flink read) groups the lookup keys 
by partition AND shard before calling. The vector lookup caller here does not 
follow this contract - it invokes the method after grouping the keys by data 
table partition only, without shard ID within the partition, which breaks on 
partitioned RLI with more than one shard per data table partition.
   
   The fix is to group finalist keys by mapRecordKeyToFileGroupIndex within 
each partition, the way the Spark/Flink read path does.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to