This is an automated email from the ASF dual-hosted git repository.

vinoth pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/hudi.git


The following commit(s) were added to refs/heads/master by this push:
     new 93d9c25  [MINOR] Improve code readability by passing in the 
fileComparisonsRDD in bloom index (#2319)
93d9c25 is described below

commit 93d9c25aee37d36c03e4c0edfc18db7da819ce2d
Author: Danny Chan <[email protected]>
AuthorDate: Tue Dec 15 14:35:24 2020 +0800

    [MINOR] Improve code readability by passing in the fileComparisonsRDD in 
bloom index (#2319)
---
 .../hudi/index/bloom/SparkHoodieBloomIndex.java       | 19 +++++++++----------
 1 file changed, 9 insertions(+), 10 deletions(-)

diff --git 
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/index/bloom/SparkHoodieBloomIndex.java
 
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/index/bloom/SparkHoodieBloomIndex.java
index 894b41b..7316043 100644
--- 
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/index/bloom/SparkHoodieBloomIndex.java
+++ 
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/index/bloom/SparkHoodieBloomIndex.java
@@ -122,13 +122,15 @@ public class SparkHoodieBloomIndex<T extends 
HoodieRecordPayload> extends SparkH
 
     // Step 3: Obtain a RDD, for each incoming record, that already exists, 
with the file id,
     // that contains it.
+    JavaRDD<Tuple2<String, HoodieKey>> fileComparisonsRDD =
+        explodeRecordRDDWithFileComparisons(partitionToFileInfo, 
partitionRecordKeyPairRDD);
     Map<String, Long> comparisonsPerFileGroup =
-        computeComparisonsPerFileGroup(recordsPerPartition, 
partitionToFileInfo, partitionRecordKeyPairRDD);
+        computeComparisonsPerFileGroup(recordsPerPartition, 
partitionToFileInfo, fileComparisonsRDD);
     int inputParallelism = partitionRecordKeyPairRDD.partitions().size();
     int joinParallelism = Math.max(inputParallelism, 
config.getBloomIndexParallelism());
     LOG.info("InputParallelism: ${" + inputParallelism + "}, IndexParallelism: 
${"
         + config.getBloomIndexParallelism() + "}");
-    return findMatchingFilesForRecordKeys(partitionToFileInfo, 
partitionRecordKeyPairRDD, joinParallelism, hoodieTable,
+    return findMatchingFilesForRecordKeys(fileComparisonsRDD, joinParallelism, 
hoodieTable,
         comparisonsPerFileGroup);
   }
 
@@ -137,14 +139,12 @@ public class SparkHoodieBloomIndex<T extends 
HoodieRecordPayload> extends SparkH
    */
   private Map<String, Long> computeComparisonsPerFileGroup(final Map<String, 
Long> recordsPerPartition,
                                                            final Map<String, 
List<BloomIndexFileInfo>> partitionToFileInfo,
-                                                           JavaPairRDD<String, 
String> partitionRecordKeyPairRDD) {
-
+                                                           final 
JavaRDD<Tuple2<String, HoodieKey>> fileComparisonsRDD) {
     Map<String, Long> fileToComparisons;
     if (config.getBloomIndexPruneByRanges()) {
       // we will just try exploding the input and then count to determine 
comparisons
       // FIX(vc): Only do sampling here and extrapolate?
-      fileToComparisons = 
explodeRecordRDDWithFileComparisons(partitionToFileInfo, 
partitionRecordKeyPairRDD)
-          .mapToPair(t -> t).countByKey();
+      fileToComparisons = fileComparisonsRDD.mapToPair(t -> t).countByKey();
     } else {
       fileToComparisons = new HashMap<>();
       partitionToFileInfo.forEach((key, value) -> {
@@ -252,11 +252,10 @@ public class SparkHoodieBloomIndex<T extends 
HoodieRecordPayload> extends SparkH
    * Make sure the parallelism is atleast the groupby parallelism for tagging 
location
    */
   JavaPairRDD<HoodieKey, HoodieRecordLocation> findMatchingFilesForRecordKeys(
-      final Map<String, List<BloomIndexFileInfo>> partitionToFileIndexInfo,
-      JavaPairRDD<String, String> partitionRecordKeyPairRDD, int 
shuffleParallelism, HoodieTable hoodieTable,
+      JavaRDD<Tuple2<String, HoodieKey>> fileComparisonsRDD,
+      int shuffleParallelism,
+      HoodieTable hoodieTable,
       Map<String, Long> fileGroupToComparisons) {
-    JavaRDD<Tuple2<String, HoodieKey>> fileComparisonsRDD =
-        explodeRecordRDDWithFileComparisons(partitionToFileIndexInfo, 
partitionRecordKeyPairRDD);
 
     if (config.useBloomIndexBucketizedChecking()) {
       Partitioner partitioner = new 
BucketizedBloomCheckPartitioner(shuffleParallelism, fileGroupToComparisons,

Reply via email to