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

codope 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 a987574c061 [HUDI-7472] prevent MDT partitions from getting dropped 
(#10804)
a987574c061 is described below

commit a987574c0613da9010e6a67d058c1debcefe03c4
Author: bhat-vinay <[email protected]>
AuthorDate: Mon Mar 4 12:36:01 2024 +0530

    [HUDI-7472] prevent MDT partitions from getting dropped (#10804)
    
    The functional index creation code-path creates a HudiTable object and gets 
a metadata writer.
    But, this code path (of creating metadata writer) also deletes the existing 
MDT partitions
    iff the write-config does not contain the relevant MDT/index configs. This 
logic is contained
    within HoodieTable::deleteMetadataIndexIfNecessary.
    
    The existing code in HoodieSparkFunctionalIndexClient::create is the entry 
point for
    functional index creation. This creates a custom write-config in
    HoodieSparkFunctionalIndexClient::buildWriteConfig which is then used to 
create a client for
    the base table (on which the functional index needs to be created). This PR 
fixes the
    issue noted earlier by adding the relevant MDT partitions config in
    HoodieSparkFunctionalIndexClient::buildWriteConfig. A test is also added to 
ensure that creating
    of functional index does not drop the existing MDT partitions.
    
    Co-authored-by: Vinaykumar Bhat <[email protected]>
---
 .../hudi/HoodieSparkFunctionalIndexClient.java     | 26 +++++++++++++++++-----
 .../hudi/command/index/TestFunctionalIndex.scala   | 17 ++++++++++++--
 2 files changed, 36 insertions(+), 7 deletions(-)

diff --git 
a/hudi-spark-datasource/hudi-spark-common/src/main/java/org/apache/hudi/HoodieSparkFunctionalIndexClient.java
 
b/hudi-spark-datasource/hudi-spark-common/src/main/java/org/apache/hudi/HoodieSparkFunctionalIndexClient.java
index 541a0d272a4..e66ad5ac417 100644
--- 
a/hudi-spark-datasource/hudi-spark-common/src/main/java/org/apache/hudi/HoodieSparkFunctionalIndexClient.java
+++ 
b/hudi-spark-datasource/hudi-spark-common/src/main/java/org/apache/hudi/HoodieSparkFunctionalIndexClient.java
@@ -27,7 +27,6 @@ import org.apache.hudi.common.model.WriteConcurrencyMode;
 import org.apache.hudi.common.table.HoodieTableMetaClient;
 import org.apache.hudi.common.util.Option;
 import org.apache.hudi.common.util.ValidationUtils;
-import org.apache.hudi.config.HoodieLockConfig;
 import org.apache.hudi.config.HoodieWriteConfig;
 import org.apache.hudi.exception.HoodieException;
 import org.apache.hudi.exception.HoodieFunctionalIndexException;
@@ -49,6 +48,9 @@ import scala.collection.JavaConverters;
 
 import static org.apache.hudi.HoodieConversionUtils.mapAsScalaImmutableMap;
 import static org.apache.hudi.HoodieConversionUtils.toScalaOption;
+import static 
org.apache.hudi.common.config.HoodieMetadataConfig.ENABLE_METADATA_INDEX_BLOOM_FILTER;
+import static 
org.apache.hudi.common.config.HoodieMetadataConfig.ENABLE_METADATA_INDEX_COLUMN_STATS;
+import static 
org.apache.hudi.common.config.HoodieMetadataConfig.RECORD_INDEX_ENABLE_PROP;
 import static org.apache.hudi.common.util.ValidationUtils.checkArgument;
 
 public class HoodieSparkFunctionalIndexClient extends 
BaseHoodieFunctionalIndexClient {
@@ -122,11 +124,25 @@ public class HoodieSparkFunctionalIndexClient extends 
BaseHoodieFunctionalIndexC
   private static Map<String, String> buildWriteConfig(HoodieTableMetaClient 
metaClient, HoodieFunctionalIndexDefinition indexDefinition) {
     Map<String, String> writeConfig = new HashMap<>();
     if (metaClient.getTableConfig().isMetadataTableAvailable()) {
-      if 
(!writeConfig.containsKey(HoodieLockConfig.LOCK_PROVIDER_CLASS_NAME.key())) {
-        writeConfig.put(HoodieWriteConfig.WRITE_CONCURRENCY_MODE.key(), 
WriteConcurrencyMode.OPTIMISTIC_CONCURRENCY_CONTROL.name());
-        
writeConfig.putAll(JavaConverters.mapAsJavaMapConverter(HoodieCLIUtils.getLockOptions(metaClient.getBasePathV2().toString())).asJava());
-      }
+      writeConfig.put(HoodieWriteConfig.WRITE_CONCURRENCY_MODE.key(), 
WriteConcurrencyMode.OPTIMISTIC_CONCURRENCY_CONTROL.name());
+      
writeConfig.putAll(JavaConverters.mapAsJavaMapConverter(HoodieCLIUtils.getLockOptions(metaClient.getBasePathV2().toString())).asJava());
+
+      // [HUDI-7472] Ensure write-config contains the existing MDT partition 
to prevent those from getting deleted
+      
metaClient.getTableConfig().getMetadataPartitions().forEach(partitionPath -> {
+        if 
(partitionPath.equals(MetadataPartitionType.RECORD_INDEX.getPartitionPath())) {
+          writeConfig.put(RECORD_INDEX_ENABLE_PROP.key(), "true");
+        }
+
+        if 
(partitionPath.equals(MetadataPartitionType.BLOOM_FILTERS.getPartitionPath())) {
+          writeConfig.put(ENABLE_METADATA_INDEX_BLOOM_FILTER.key(), "true");
+        }
+
+        if 
(partitionPath.equals(MetadataPartitionType.COLUMN_STATS.getPartitionPath())) {
+          writeConfig.put(ENABLE_METADATA_INDEX_COLUMN_STATS.key(), "true");
+        }
+      });
     }
+
     
HoodieFunctionalIndexConfig.fromIndexDefinition(indexDefinition).getProps().forEach((key,
 value) -> writeConfig.put(key.toString(), value.toString()));
     return writeConfig;
   }
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/command/index/TestFunctionalIndex.scala
 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/command/index/TestFunctionalIndex.scala
index 34f79fa45b5..8a649bf98a6 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/command/index/TestFunctionalIndex.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/command/index/TestFunctionalIndex.scala
@@ -27,6 +27,7 @@ import org.apache.hudi.common.util.Option
 import org.apache.hudi.hive.HiveSyncConfigHolder._
 import org.apache.hudi.hive.{HiveSyncTool, HoodieHiveSyncClient}
 import org.apache.hudi.hive.testutils.HiveTestUtil
+import org.apache.hudi.metadata.MetadataPartitionType
 import org.apache.hudi.sync.common.HoodieSyncConfig.{META_SYNC_BASE_PATH, 
META_SYNC_DATABASE_NAME, META_SYNC_NO_PARTITION_METADATA, META_SYNC_TABLE_NAME}
 import org.apache.spark.sql.catalyst.analysis.Analyzer
 import org.apache.spark.sql.catalyst.catalog.CatalogTable
@@ -186,7 +187,9 @@ class TestFunctionalIndex extends HoodieSparkSqlTestBase {
                | options (
                |  primaryKey ='id',
                |  type = '$tableType',
-               |  preCombineField = 'ts'
+               |  preCombineField = 'ts',
+               |  hoodie.metadata.record.index.enable = 'true',
+               |  hoodie.datasource.write.recordkey.field = 'id'
                | )
                | partitioned by(ts)
                | location '$basePath'
@@ -195,6 +198,13 @@ class TestFunctionalIndex extends HoodieSparkSqlTestBase {
           spark.sql(s"insert into $tableName values(2, 'a2', 10, 1001)")
           spark.sql(s"insert into $tableName values(3, 'a3', 10, 1002)")
 
+          var metaClient = HoodieTableMetaClient.builder()
+            .setBasePath(basePath)
+            .setConf(spark.sessionState.newHadoopConf())
+            .build()
+
+          
assert(metaClient.getTableConfig.isMetadataPartitionAvailable(MetadataPartitionType.RECORD_INDEX))
+
           val sqlParser: ParserInterface = spark.sessionState.sqlParser
           val analyzer: Analyzer = spark.sessionState.analyzer
 
@@ -212,7 +222,7 @@ class TestFunctionalIndex extends HoodieSparkSqlTestBase {
           
assertResult(false)(resolvedLogicalPlan.asInstanceOf[CreateIndexCommand].ignoreIfExists)
 
           spark.sql(createIndexSql)
-          var metaClient = HoodieTableMetaClient.builder()
+          metaClient = HoodieTableMetaClient.builder()
             .setBasePath(basePath)
             .setConf(spark.sessionState.newHadoopConf())
             .build()
@@ -235,6 +245,9 @@ class TestFunctionalIndex extends HoodieSparkSqlTestBase {
           // Ensure that both the indexes are tracked correctly in metadata 
partition config
           val mdtPartitions = metaClient.getTableConfig.getMetadataPartitions
           assert(mdtPartitions.contains("func_index_name_lower") && 
mdtPartitions.contains("func_index_idx_datestr"))
+
+          // [HUDI-7472] After creating functional index, the existing MDT 
partitions should still be available
+          
assert(metaClient.getTableConfig.isMetadataPartitionAvailable(MetadataPartitionType.RECORD_INDEX))
         }
       }
     }

Reply via email to