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