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