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,