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 039ccc006fa [HUDI-8586] Allow partition stats only for primitive types 
(#12350)
039ccc006fa is described below

commit 039ccc006fa4a4373ec725415b617d0b8d51f106
Author: Sagar Sumit <[email protected]>
AuthorDate: Fri Nov 29 18:24:17 2024 +0530

    [HUDI-8586] Allow partition stats only for primitive types (#12350)
    
    - Allow partition stats only for certain data types
    - Make column stats a pre-requisite for partition stats
---
 .../metadata/HoodieBackedTableMetadataWriter.java  |  6 ++
 .../hudi/metadata/HoodieTableMetadataUtil.java     | 50 ++++++++++++++--
 .../java/org/apache/hudi/source/TestFileIndex.java |  1 +
 .../hudi/source/TestIncrementalInputSplits.java    |  1 +
 .../hudi/source/stats/TestColumnStatsIndex.java    |  1 +
 .../apache/hudi/table/ITTestHoodieDataSource.java  |  1 +
 .../apache/hudi/table/TestHoodieTableSource.java   |  1 +
 .../hudi/metadata/TestHoodieTableMetadataUtil.java | 48 ++++++++++++++++
 .../functional/PartitionStatsIndexTestBase.scala   |  1 +
 .../hudi/functional/TestPartitionStatsIndex.scala  | 67 +++++++++++++++++++++-
 .../TestPartitionStatsIndexWithSql.scala           |  5 +-
 .../functional/TestSecondaryIndexPruning.scala     |  6 +-
 .../hudi/dml/TestHoodieTableValuedFunction.scala   |  1 +
 13 files changed, 180 insertions(+), 9 deletions(-)

diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/HoodieBackedTableMetadataWriter.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/HoodieBackedTableMetadataWriter.java
index 34ebdea906c..8c8a0cfed9f 100644
--- 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/HoodieBackedTableMetadataWriter.java
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/HoodieBackedTableMetadataWriter.java
@@ -429,6 +429,12 @@ public abstract class HoodieBackedTableMetadataWriter<I> 
implements HoodieTableM
             fileGroupCountAndRecordsPair = 
initializeExpressionIndexPartition(partitionName, instantTimeForPartition);
             break;
           case PARTITION_STATS:
+            // For PARTITION_STATS, COLUMN_STATS should also be enabled
+            if (!dataWriteConfig.isMetadataColumnStatsIndexEnabled()) {
+              LOG.warn("Skipping partition stats initialization as column 
stats index is not enabled. Please enable {}",
+                  
HoodieMetadataConfig.ENABLE_METADATA_INDEX_COLUMN_STATS.key());
+              continue;
+            }
             fileGroupCountAndRecordsPair = 
initializePartitionStatsIndex(partitionInfoList);
             partitionName = PARTITION_STATS.getPartitionPath();
             break;
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/metadata/HoodieTableMetadataUtil.java
 
b/hudi-common/src/main/java/org/apache/hudi/metadata/HoodieTableMetadataUtil.java
index 12eaa6c5732..b76e30e71ae 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/metadata/HoodieTableMetadataUtil.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/metadata/HoodieTableMetadataUtil.java
@@ -188,6 +188,13 @@ public class HoodieTableMetadataUtil {
   public static final String PARTITION_NAME_SECONDARY_INDEX = 
"secondary_index";
   public static final String PARTITION_NAME_SECONDARY_INDEX_PREFIX = 
"secondary_index_";
 
+  private static final Set<Schema.Type> SUPPORTED_TYPES_PARTITION_STATS = new 
HashSet<>(Arrays.asList(
+      Schema.Type.INT, Schema.Type.LONG, Schema.Type.FLOAT, 
Schema.Type.DOUBLE, Schema.Type.STRING, Schema.Type.BOOLEAN, Schema.Type.NULL, 
Schema.Type.BYTES));
+  private static final Set<String> SUPPORTED_META_FIELDS_PARTITION_STATS = new 
HashSet<>(Arrays.asList(
+      
HoodieRecord.HoodieMetadataField.RECORD_KEY_METADATA_FIELD.getFieldName(),
+      
HoodieRecord.HoodieMetadataField.PARTITION_PATH_METADATA_FIELD.getFieldName(),
+      
HoodieRecord.HoodieMetadataField.COMMIT_TIME_METADATA_FIELD.getFieldName()));
+
   private HoodieTableMetadataUtil() {
   }
 
@@ -388,6 +395,8 @@ public class HoodieTableMetadataUtil {
       
partitionToRecordsMap.put(MetadataPartitionType.COLUMN_STATS.getPartitionPath(),
 metadataColumnStatsRDD);
     }
     if (enabledPartitionTypes.contains(MetadataPartitionType.PARTITION_STATS)) 
{
+      
checkState(MetadataPartitionType.COLUMN_STATS.isMetadataPartitionAvailable(dataMetaClient),
+          "Column stats partition must be enabled to generate partition stats. 
Please enable: " + 
HoodieMetadataConfig.ENABLE_METADATA_INDEX_COLUMN_STATS.key());
       final HoodieData<HoodieRecord> partitionStatsRDD = 
convertMetadataToPartitionStatsRecords(commitMetadata, context, dataMetaClient, 
metadataConfig);
       
partitionToRecordsMap.put(MetadataPartitionType.PARTITION_STATS.getPartitionPath(),
 partitionStatsRDD);
     }
@@ -2427,7 +2436,11 @@ public class HoodieTableMetadataUtil {
       LOG.warn("No columns to index for partition stats index");
       return engineContext.emptyHoodieData();
     }
-    LOG.debug("Indexing following columns for partition stats index: {}", 
columnsToIndex);
+    // filter columns with only supported types
+    final List<String> validColumnsToIndex = columnsToIndex.stream()
+        .filter(col -> SUPPORTED_META_FIELDS_PARTITION_STATS.contains(col) || 
validateDataTypeForPartitionStats(col, lazyWriterSchemaOpt.get().get()))
+        .collect(Collectors.toList());
+    LOG.debug("Indexing following columns for partition stats index: {}", 
validColumnsToIndex);
 
     // Group by partition path and collect file names (BaseFile and LogFiles)
     List<Pair<String, Set<String>>> partitionToFileNames = 
partitionInfoList.stream()
@@ -2444,7 +2457,7 @@ public class HoodieTableMetadataUtil {
       final String partitionPath = partitionInfo.getKey();
       // Step 1: Collect Column Metadata for Each File
       List<List<HoodieColumnRangeMetadata<Comparable>>> fileColumnMetadata = 
partitionInfo.getValue().stream()
-          .map(fileName -> getFileStatsRangeMetadata(partitionPath, fileName, 
dataTableMetaClient, columnsToIndex, false,
+          .map(fileName -> getFileStatsRangeMetadata(partitionPath, fileName, 
dataTableMetaClient, validColumnsToIndex, false,
               metadataConfig.getMaxReaderBufferSize()))
           .collect(Collectors.toList());
 
@@ -2498,6 +2511,11 @@ public class HoodieTableMetadataUtil {
       if (columnsToIndex.isEmpty()) {
         return engineContext.emptyHoodieData();
       }
+      // filter columns with only supported types
+      final List<String> validColumnsToIndex = columnsToIndex.stream()
+          .filter(col -> SUPPORTED_META_FIELDS_PARTITION_STATS.contains(col) 
|| validateDataTypeForPartitionStats(col, writerSchemaOpt.get().get()))
+          .collect(Collectors.toList());
+      LOG.debug("Indexing following columns for partition stats index: {}", 
validColumnsToIndex);
       // Group by partitionPath and then gather write stats lists,
       // where each inner list contains HoodieWriteStat objects that have the 
same partitionPath.
       List<List<HoodieWriteStat>> partitionedWriteStats = 
allWriteStats.stream()
@@ -2508,6 +2526,7 @@ public class HoodieTableMetadataUtil {
 
       int parallelism = Math.max(Math.min(partitionedWriteStats.size(), 
metadataConfig.getPartitionStatsIndexParallelism()), 1);
       boolean shouldScanColStatsForTightBound = 
MetadataPartitionType.COLUMN_STATS.isMetadataPartitionAvailable(dataMetaClient);
+
       HoodieTableMetadata tableMetadata;
       if (shouldScanColStatsForTightBound) {
         tableMetadata = HoodieTableMetadata.create(engineContext, 
dataMetaClient.getStorage(), metadataConfig, 
dataMetaClient.getBasePath().toString());
@@ -2518,7 +2537,7 @@ public class HoodieTableMetadataUtil {
         final String partitionName = 
partitionedWriteStat.get(0).getPartitionPath();
         // Step 1: Collect Column Metadata for Each File part of current 
commit metadata
         List<List<HoodieColumnRangeMetadata<Comparable>>> fileColumnMetadata = 
partitionedWriteStat.stream()
-            .map(writeStat -> translateWriteStatToFileStats(writeStat, 
dataMetaClient, columnsToIndex, tableSchema))
+            .map(writeStat -> translateWriteStatToFileStats(writeStat, 
dataMetaClient, validColumnsToIndex, tableSchema))
             .collect(Collectors.toList());
         if (shouldScanColStatsForTightBound) {
           checkState(tableMetadata != null, "tableMetadata should not be null 
when scanning metadata table");
@@ -2532,7 +2551,7 @@ public class HoodieTableMetadataUtil {
               .collect(Collectors.toSet());
           // Fetch metadata table COLUMN_STATS partition records for above 
files
           List<HoodieColumnRangeMetadata<Comparable>> partitionColumnMetadata =
-              
tableMetadata.getRecordsByKeyPrefixes(generateKeyPrefixes(columnsToIndex, 
partitionName), MetadataPartitionType.COLUMN_STATS.getPartitionPath(), false)
+              
tableMetadata.getRecordsByKeyPrefixes(generateKeyPrefixes(validColumnsToIndex, 
partitionName), MetadataPartitionType.COLUMN_STATS.getPartitionPath(), false)
                   // schema and properties are ignored in getInsertValue, so 
simply pass as null
                   .map(record -> record.getData().getInsertValue(null, null))
                   .filter(Option::isPresent)
@@ -2552,6 +2571,29 @@ public class HoodieTableMetadataUtil {
     }
   }
 
+  /**
+   * Given table schema and field to index, checks if field's data type are 
supported.
+   *
+   * @param columnToIndex column to index
+   * @param tableSchema   table schema
+   * @return true if field's data type is supported, false otherwise
+   */
+  @VisibleForTesting
+  static boolean validateDataTypeForPartitionStats(String columnToIndex, 
Schema tableSchema) {
+    Schema fieldSchema = getNestedFieldSchemaFromWriteSchema(tableSchema, 
columnToIndex);
+    // Exclude fields based on logical type
+    if ((fieldSchema.getType() == Schema.Type.INT || fieldSchema.getType() == 
Schema.Type.LONG)
+        && fieldSchema.getLogicalType() != null) {
+
+      // Skip fields with logical types DATE or TIME_MILLIS for INT, 
TIMESTAMP_MILLIS for LONG
+      String logicalType = fieldSchema.getLogicalType().getName();
+      return !logicalType.equals("date") && 
!logicalType.equals("timestamp-millis") && 
!logicalType.equals("timestamp-micros") && !logicalType.equals("time-millis")
+          && !logicalType.equals("time-micros") && 
!logicalType.equals("local-timestamp-millis") && 
!logicalType.equals("local-timestamp-micros");
+    }
+    // Include other supported primitive types
+    return SUPPORTED_TYPES_PARTITION_STATS.contains(fieldSchema.getType());
+  }
+
   /**
    * Generate key prefixes for each combination of column name in {@param 
columnsToIndex} and {@param partitionName}.
    */
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/TestFileIndex.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/TestFileIndex.java
index a2d5e10ec4a..6d35785e1b3 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/TestFileIndex.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/TestFileIndex.java
@@ -171,6 +171,7 @@ public class TestFileIndex {
     conf.set(METADATA_ENABLED, true);
     conf.set(TABLE_TYPE, tableType.name());
     
conf.setBoolean(HoodieMetadataConfig.ENABLE_METADATA_INDEX_PARTITION_STATS.key(),
 true);
+    
conf.setBoolean(HoodieMetadataConfig.ENABLE_METADATA_INDEX_COLUMN_STATS.key(), 
true);
     if (tableType == HoodieTableType.MERGE_ON_READ) {
       // enable CSI for MOR table to collect col stats for delta write stats,
       // which will be used to construct partition stats then.
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/TestIncrementalInputSplits.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/TestIncrementalInputSplits.java
index 02ea8ad553f..fa0d9b56c57 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/TestIncrementalInputSplits.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/TestIncrementalInputSplits.java
@@ -361,6 +361,7 @@ public class TestIncrementalInputSplits extends 
HoodieCommonTestHarness {
     conf.set(FlinkOptions.READ_DATA_SKIPPING_ENABLED, true);
     conf.set(FlinkOptions.TABLE_TYPE, tableType.name());
     
conf.setBoolean(HoodieMetadataConfig.ENABLE_METADATA_INDEX_PARTITION_STATS.key(),
 true);
+    
conf.setBoolean(HoodieMetadataConfig.ENABLE_METADATA_INDEX_COLUMN_STATS.key(), 
true);
     if (tableType == HoodieTableType.MERGE_ON_READ) {
       // enable CSI for MOR table to collect col stats for delta write stats,
       // which will be used to construct partition stats then.
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/stats/TestColumnStatsIndex.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/stats/TestColumnStatsIndex.java
index 15656d0fa0a..631be02c02c 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/stats/TestColumnStatsIndex.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/stats/TestColumnStatsIndex.java
@@ -55,6 +55,7 @@ public class TestColumnStatsIndex {
     Configuration conf = TestConfigurations.getDefaultConf(path);
     conf.set(FlinkOptions.METADATA_ENABLED, true);
     conf.setString("hoodie.metadata.index.partition.stats.enable", "true");
+    conf.setString("hoodie.metadata.index.column.stats.enable", "true");
     HoodieMetadataConfig metadataConfig = HoodieMetadataConfig.newBuilder()
         .enable(true)
         .withMetadataIndexColumnStats(true)
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/ITTestHoodieDataSource.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/ITTestHoodieDataSource.java
index b2bf266129e..2cc964f1a26 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/ITTestHoodieDataSource.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/ITTestHoodieDataSource.java
@@ -516,6 +516,7 @@ public class ITTestHoodieDataSource {
         .option(FlinkOptions.METADATA_ENABLED, true)
         .option(FlinkOptions.READ_AS_STREAMING, true)
         
.option(HoodieMetadataConfig.ENABLE_METADATA_INDEX_PARTITION_STATS.key(), true)
+        .option(HoodieMetadataConfig.ENABLE_METADATA_INDEX_COLUMN_STATS.key(), 
false)
         .option(FlinkOptions.READ_DATA_SKIPPING_ENABLED, true)
         .option(FlinkOptions.TABLE_TYPE, tableType)
         .option(FlinkOptions.HIVE_STYLE_PARTITIONING, hiveStylePartitioning)
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/TestHoodieTableSource.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/TestHoodieTableSource.java
index 899b3b15667..ff0c849e0c1 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/TestHoodieTableSource.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/TestHoodieTableSource.java
@@ -177,6 +177,7 @@ public class TestHoodieTableSource {
     final String path = tempFile.getAbsolutePath();
     conf = TestConfigurations.getDefaultConf(path);
     
conf.setBoolean(HoodieMetadataConfig.ENABLE_METADATA_INDEX_PARTITION_STATS.key(),
 true);
+    
conf.setBoolean(HoodieMetadataConfig.ENABLE_METADATA_INDEX_COLUMN_STATS.key(), 
true);
     conf.set(FlinkOptions.READ_DATA_SKIPPING_ENABLED, true);
     TestData.writeData(TestData.DATA_SET_INSERT, conf);
     HoodieTableSource hoodieTableSource = createHoodieTableSource(conf);
diff --git 
a/hudi-hadoop-common/src/test/java/org/apache/hudi/metadata/TestHoodieTableMetadataUtil.java
 
b/hudi-hadoop-common/src/test/java/org/apache/hudi/metadata/TestHoodieTableMetadataUtil.java
index daf89e47c72..be516c48a5d 100644
--- 
a/hudi-hadoop-common/src/test/java/org/apache/hudi/metadata/TestHoodieTableMetadataUtil.java
+++ 
b/hudi-hadoop-common/src/test/java/org/apache/hudi/metadata/TestHoodieTableMetadataUtil.java
@@ -44,6 +44,7 @@ import org.apache.hudi.io.storage.HoodieFileWriterFactory;
 import org.apache.hudi.storage.StoragePath;
 import org.apache.hudi.util.Lazy;
 
+import org.apache.avro.LogicalTypes;
 import org.apache.avro.Schema;
 import org.apache.avro.SchemaBuilder;
 import org.junit.jupiter.api.AfterEach;
@@ -64,6 +65,7 @@ import java.util.stream.Collectors;
 
 import static org.apache.hudi.avro.TestHoodieAvroUtils.SCHEMA_WITH_AVRO_TYPES;
 import static 
org.apache.hudi.metadata.HoodieTableMetadataUtil.getFileIDForFileGroup;
+import static 
org.apache.hudi.metadata.HoodieTableMetadataUtil.validateDataTypeForPartitionStats;
 import static 
org.apache.hudi.metadata.HoodieTableMetadataUtil.validateDataTypeForSecondaryIndex;
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertFalse;
@@ -542,4 +544,50 @@ public class TestHoodieTableMetadataUtil extends 
HoodieCommonTestHarness {
       list.add("col_" + i);
     }
   }
+
+  @Test
+  public void testValidateDataTypeForPartitionStats() {
+    // Create a dummy schema with both complex and primitive types
+    Schema schema = SchemaBuilder.record("TestRecord")
+        .fields()
+        .requiredString("stringField")
+        .optionalInt("intField")
+        .optionalBoolean("booleanField")
+        .optionalFloat("floatField")
+        .optionalDouble("doubleField")
+        .optionalLong("longField")
+        .optionalBytes("bytesField")
+        
.name("unionIntField").type().unionOf().nullType().and().intType().endUnion().noDefault()
+        .name("arrayField").type().array().items().stringType().noDefault()
+        .name("mapField").type().map().values().intType().noDefault()
+        .name("structField").type().record("NestedRecord")
+        .fields()
+        .requiredString("nestedString")
+        .endRecord()
+        .noDefault()
+        .endRecord();
+
+    // Test for primitive fields
+    assertTrue(validateDataTypeForPartitionStats("stringField", schema));
+    assertTrue(validateDataTypeForPartitionStats("intField", schema));
+    assertTrue(validateDataTypeForPartitionStats("booleanField", schema));
+    assertTrue(validateDataTypeForPartitionStats("floatField", schema));
+    assertTrue(validateDataTypeForPartitionStats("doubleField", schema));
+    assertTrue(validateDataTypeForPartitionStats("longField", schema));
+    assertTrue(validateDataTypeForPartitionStats("bytesField", schema));
+    assertTrue(validateDataTypeForPartitionStats("unionIntField", schema));
+
+    // Test for complex fields
+    assertFalse(validateDataTypeForPartitionStats("arrayField", schema));
+    assertFalse(validateDataTypeForPartitionStats("mapField", schema));
+    assertFalse(validateDataTypeForPartitionStats("structField", schema));
+
+    // Test for logical types
+    Schema dateFieldSchema = 
LogicalTypes.date().addToSchema(Schema.create(Schema.Type.INT));
+    schema = SchemaBuilder.record("TestRecord")
+        .fields()
+        .name("dateField").type(dateFieldSchema).noDefault()
+        .endRecord();
+    assertFalse(validateDataTypeForPartitionStats("dateField", schema));
+  }
 }
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/PartitionStatsIndexTestBase.scala
 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/PartitionStatsIndexTestBase.scala
index 4e08dbb8cb6..5b706dfec14 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/PartitionStatsIndexTestBase.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/PartitionStatsIndexTestBase.scala
@@ -60,6 +60,7 @@ class PartitionStatsIndexTestBase extends 
HoodieSparkClientTestBase {
   val metadataOpts: Map[String, String] = Map(
     HoodieMetadataConfig.ENABLE.key -> "true",
     HoodieMetadataConfig.ENABLE_METADATA_INDEX_PARTITION_STATS.key -> "true",
+    HoodieMetadataConfig.ENABLE_METADATA_INDEX_COLUMN_STATS.key -> "true",
     HoodieMetadataConfig.COLUMN_STATS_INDEX_FOR_COLUMNS.key -> 
targetColumnsToIndex.mkString(",")
   )
   val commonOpts: Map[String, String] = Map(
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestPartitionStatsIndex.scala
 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestPartitionStatsIndex.scala
index 9e886f0b754..5d032c68ebe 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestPartitionStatsIndex.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestPartitionStatsIndex.scala
@@ -30,7 +30,7 @@ import org.apache.hudi.common.table.HoodieTableMetaClient
 import org.apache.hudi.common.table.timeline.HoodieInstant
 import org.apache.hudi.common.testutils.RawTripTestPayload.recordsToStrings
 import org.apache.hudi.config.{HoodieCleanConfig, HoodieClusteringConfig, 
HoodieCompactionConfig, HoodieLockConfig, HoodieWriteConfig}
-import org.apache.hudi.exception.HoodieWriteConflictException
+import org.apache.hudi.exception.{HoodieException, 
HoodieWriteConflictException}
 import org.apache.hudi.keygen.constant.KeyGeneratorOptions
 import org.apache.hudi.metadata.{HoodieBackedTableMetadata, 
HoodieMetadataFileSystemView, MetadataPartitionType}
 import org.apache.hudi.util.{JFunction, JavaConversions}
@@ -42,6 +42,7 @@ import org.junit.jupiter.api.Assertions.{assertEquals, 
assertFalse, assertTrue}
 import org.junit.jupiter.api.{Tag, Test}
 import org.junit.jupiter.params.ParameterizedTest
 import org.junit.jupiter.params.provider.{Arguments, EnumSource, MethodSource}
+import org.scalatest.Assertions.assertThrows
 
 import java.util.concurrent.Executors
 import java.util.stream.Stream
@@ -58,6 +59,62 @@ class TestPartitionStatsIndex extends 
PartitionStatsIndexTestBase {
 
   val sqlTempTable = "hudi_tbl"
 
+  /**
+   * Test case to validate partition stats cannot be created without column 
stats.
+   */
+  @Test
+  def testPartitionStatsWithoutColumnStats(): Unit = {
+    // remove column stats enable key from commonOpts
+    val hudiOpts = commonOpts - 
HoodieMetadataConfig.ENABLE_METADATA_INDEX_COLUMN_STATS.key
+    // should throw an exception as column stats is required for partition 
stats
+    assertThrows[HoodieException] {
+      doWriteAndValidateDataAndPartitionStats(
+        hudiOpts,
+        operation = DataSourceWriteOptions.INSERT_OPERATION_OPT_VAL,
+        saveMode = SaveMode.Overwrite)
+    }
+  }
+
+  /**
+   * Test case to validate partition stats for a logical type column
+   */
+  @Test
+  def testPartitionStatsWithLogicalType(): Unit = {
+    val hudiOpts = commonOpts ++ Map(
+      DataSourceWriteOptions.TABLE_TYPE.key -> 
HoodieTableType.MERGE_ON_READ.name(),
+      HoodieMetadataConfig.COLUMN_STATS_INDEX_FOR_COLUMNS.key -> "current_date"
+    )
+
+    val records = recordsToStrings(dataGen.generateInserts("001", 
100)).asScala.toList
+    val inputDF = spark.read.json(spark.sparkContext.parallelize(records, 2))
+
+    inputDF.write.partitionBy("partition").format("hudi")
+      .options(hudiOpts)
+      .option(DataSourceWriteOptions.OPERATION.key, 
DataSourceWriteOptions.INSERT_OPERATION_OPT_VAL)
+      .option(KeyGeneratorOptions.URL_ENCODE_PARTITIONING.key, "true")
+      .mode(SaveMode.Overwrite)
+      .save(basePath)
+
+    val snapshot0 = 
spark.read.format("org.apache.hudi").options(hudiOpts).load(basePath)
+    assertEquals(100, snapshot0.count())
+
+    val updateRecords = recordsToStrings(dataGen.generateUniqueUpdates("002", 
50)).asScala.toList
+    val updateDF = 
spark.read.json(spark.sparkContext.parallelize(updateRecords, 2))
+    updateDF.write.format("hudi")
+      .options(hudiOpts)
+      .option(DataSourceWriteOptions.OPERATION.key, 
DataSourceWriteOptions.UPSERT_OPERATION_OPT_VAL)
+      .mode(SaveMode.Append)
+      .save(basePath)
+
+    val readOpts = hudiOpts ++ Map(
+      HoodieMetadataConfig.ENABLE.key -> "true",
+      DataSourceReadOptions.ENABLE_DATA_SKIPPING.key -> "true"
+    )
+    val snapshot1 = 
spark.read.format("org.apache.hudi").options(readOpts).load(basePath)
+    val dataFilter = EqualTo(attribute("current_date"), 
Literal(snapshot1.limit(1).collect().head.getAs("current_date")))
+    verifyFilePruning(readOpts, dataFilter, shouldSkipFiles = false)
+  }
+
   /**
    * Test case to do a write (no updates) and validate the partition stats 
index initialization.
    */
@@ -433,14 +490,18 @@ class TestPartitionStatsIndex extends 
PartitionStatsIndexTestBase {
     readDf.createOrReplaceTempView(sqlTempTable)
   }
 
-  private def verifyFilePruning(opts: Map[String, String], dataFilter: 
Expression): Unit = {
+  private def verifyFilePruning(opts: Map[String, String], dataFilter: 
Expression, shouldSkipFiles: Boolean = true): Unit = {
     // with data skipping
     val commonOpts = opts + ("path" -> basePath)
     metaClient = HoodieTableMetaClient.reload(metaClient)
     var fileIndex = HoodieFileIndex(spark, metaClient, None, commonOpts, 
includeLogFiles = true)
     val filteredPartitionDirectories = fileIndex.listFiles(Seq(), 
Seq(dataFilter))
     val filteredFilesCount = filteredPartitionDirectories.flatMap(s => 
s.files).size
-    assertTrue(filteredFilesCount <= getLatestDataFilesCount(opts))
+    if (shouldSkipFiles) {
+      assertTrue(filteredFilesCount <= getLatestDataFilesCount(opts))
+    } else {
+      assertTrue(filteredFilesCount == getLatestDataFilesCount(opts))
+    }
 
     // with no data skipping
     fileIndex = HoodieFileIndex(spark, metaClient, None, commonOpts + 
(DataSourceReadOptions.ENABLE_DATA_SKIPPING.key -> "false"), includeLogFiles = 
true)
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestPartitionStatsIndexWithSql.scala
 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestPartitionStatsIndexWithSql.scala
index dbd764ce976..1e5b64a4b7d 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestPartitionStatsIndexWithSql.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestPartitionStatsIndexWithSql.scala
@@ -63,6 +63,7 @@ class TestPartitionStatsIndexWithSql extends 
HoodieSparkSqlTestBase {
              |    primaryKey = 'id',
              |    preCombineField = 'ts',
              |    'hoodie.metadata.index.partition.stats.enable' = 'true',
+             |    'hoodie.metadata.index.column.stats.enable' = 'true',
              |    'hoodie.metadata.index.column.stats.column.list' = 'name'
              | )
              | location '$tablePath'
@@ -147,6 +148,7 @@ class TestPartitionStatsIndexWithSql extends 
HoodieSparkSqlTestBase {
              |    primaryKey ='uuid',
              |    preCombineField = 'ts',
              |    hoodie.metadata.index.partition.stats.enable = 'true',
+             |    hoodie.metadata.index.column.stats.enable = 'true',
              |    hoodie.metadata.index.column.stats.column.list = 'rider'
              |)
              |PARTITIONED BY (state)
@@ -255,7 +257,8 @@ class TestPartitionStatsIndexWithSql extends 
HoodieSparkSqlTestBase {
              |  type = '$tableType',
              |  primaryKey = 'id',
              |  preCombineField = 'price',
-             |  hoodie.metadata.index.partition.stats.enable = 'true'
+             |  hoodie.metadata.index.partition.stats.enable = 'true',
+             |  hoodie.metadata.index.column.stats.enable = 'true'
              |)
              |location '$tablePath'
              |""".stripMargin
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestSecondaryIndexPruning.scala
 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestSecondaryIndexPruning.scala
index ff85de96cba..002f60f6fda 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestSecondaryIndexPruning.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestSecondaryIndexPruning.scala
@@ -370,7 +370,11 @@ class TestSecondaryIndexPruning extends 
SparkClientFunctionalTestHarness {
       val sqlTableType = if 
(tableType.equals(HoodieTableType.COPY_ON_WRITE.name())) "cow" else "mor"
       tableName += "test_secondary_index_with_partition_stats_index" + (if 
(isPartitioned) "_partitioned" else "") + sqlTableType
       val partitionedByClause = if (isPartitioned) "partitioned 
by(partition_key_col)" else ""
-      val partitionStatsEnable = if (isPartitioned) 
"'hoodie.metadata.index.partition.stats.enable' = 'true'," else ""
+      val partitionStatsEnable = if (isPartitioned)
+        s"""
+           |'hoodie.metadata.index.partition.stats.enable' = 'true',
+           |'hoodie.metadata.index.column.stats.enable' = 'true',
+        """.stripMargin else ""
       val columnsToIndex = if (isPartitioned) 
"'hoodie.metadata.index.column.stats.column.list' = 'name'," else ""
 
       spark.sql(
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/TestHoodieTableValuedFunction.scala
 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/TestHoodieTableValuedFunction.scala
index 6e214e3fe35..77a62fa7c8a 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/TestHoodieTableValuedFunction.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/dml/TestHoodieTableValuedFunction.scala
@@ -662,6 +662,7 @@ class TestHoodieTableValuedFunction extends 
HoodieSparkSqlTestBase {
                |  preCombineField = 'ts',
                |  hoodie.datasource.write.recordkey.field = 'id',
                |  hoodie.metadata.index.partition.stats.enable = 'true',
+               |  hoodie.metadata.index.column.stats.enable = 'true',
                |  hoodie.metadata.index.column.stats.column.list = 'price',
                |  hoodie.populate.meta.fields = 'false'
                |)

Reply via email to