This is an automated email from the ASF dual-hosted git repository.
yihua 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 eeb65f2a8dc [HUDI-5249] Support column metadata stats index for Spark
(#9623)
eeb65f2a8dc is described below
commit eeb65f2a8dc01f1bd8c687a8ccdfbbfc49c5b7d1
Author: Lin Liu <[email protected]>
AuthorDate: Wed Sep 6 18:58:40 2023 -0700
[HUDI-5249] Support column metadata stats index for Spark (#9623)
The core change is to
1. Use HoodieRecord data type for general process;
2. Support HoodieSparkRecord when extract field value from records;
Co-authored-by: Lin Liu <[email protected]>
---
.../org/apache/hudi/io/HoodieAppendHandle.java | 21 +++++--------
.../java/org/apache/hudi/avro/HoodieAvroUtils.java | 8 ++---
.../hudi/metadata/HoodieTableMetadataUtil.java | 34 +++++++++++-----------
3 files changed, 27 insertions(+), 36 deletions(-)
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieAppendHandle.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieAppendHandle.java
index 65f79c5147e..69262baf913 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieAppendHandle.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieAppendHandle.java
@@ -32,7 +32,6 @@ import org.apache.hudi.common.model.HoodieOperation;
import org.apache.hudi.common.model.HoodiePartitionMetadata;
import org.apache.hudi.common.model.HoodiePayloadProps;
import org.apache.hudi.common.model.HoodieRecord;
-import org.apache.hudi.common.model.HoodieRecord.HoodieRecordType;
import org.apache.hudi.common.model.HoodieWriteStat.RuntimeStats;
import org.apache.hudi.common.model.IOType;
import org.apache.hudi.common.model.MetadataValues;
@@ -57,7 +56,6 @@ import org.apache.hudi.exception.HoodieUpsertException;
import org.apache.hudi.table.HoodieTable;
import org.apache.avro.Schema;
-import org.apache.avro.generic.IndexedRecord;
import org.apache.hadoop.fs.Path;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -68,7 +66,6 @@ import java.util.Collections;
import java.util.HashMap;
import java.util.HashSet;
import java.util.Iterator;
-import java.util.LinkedList;
import java.util.List;
import java.util.Map;
import java.util.Properties;
@@ -401,9 +398,7 @@ public class HoodieAppendHandle<T, I, K, O> extends
HoodieWriteHandle<T, I, K, O
updateWriteStatus(stat, result);
}
- // TODO MetadataColumnStatsIndex for spark record
- // https://issues.apache.org/jira/browse/HUDI-5249
- if (config.isMetadataColumnStatsIndexEnabled() &&
recordMerger.getRecordType() == HoodieRecordType.AVRO) {
+ if (config.isMetadataColumnStatsIndexEnabled()) {
final List<Schema.Field> fieldsToIndex;
// If column stats index is enabled but columns not configured then we
assume that
// all columns should be indexed
@@ -417,15 +412,13 @@ public class HoodieAppendHandle<T, I, K, O> extends
HoodieWriteHandle<T, I, K, O
.collect(Collectors.toList());
}
- List<IndexedRecord> indexedRecords = new LinkedList<>();
- for (HoodieRecord hoodieRecord : recordList) {
- indexedRecords.add(hoodieRecord.toIndexedRecord(writeSchema,
config.getProps()).get().getData());
+ try {
+ Map<String, HoodieColumnRangeMetadata<Comparable>>
columnRangeMetadataMap =
+ collectColumnRangeMetadata(recordList, fieldsToIndex,
stat.getPath(), writeSchemaWithMetaFields);
+ stat.putRecordsStats(columnRangeMetadataMap);
+ } catch (HoodieException e) {
+ throw new HoodieAppendException("Failed to extract append result", e);
}
-
- Map<String, HoodieColumnRangeMetadata<Comparable>>
columnRangesMetadataMap =
- collectColumnRangeMetadata(indexedRecords, fieldsToIndex,
stat.getPath());
-
- stat.putRecordsStats(columnRangesMetadataMap);
}
resetWriteCounts();
diff --git
a/hudi-common/src/main/java/org/apache/hudi/avro/HoodieAvroUtils.java
b/hudi-common/src/main/java/org/apache/hudi/avro/HoodieAvroUtils.java
index 0909ee5555a..24edd5149c4 100644
--- a/hudi-common/src/main/java/org/apache/hudi/avro/HoodieAvroUtils.java
+++ b/hudi-common/src/main/java/org/apache/hudi/avro/HoodieAvroUtils.java
@@ -29,8 +29,6 @@ import org.apache.hudi.avro.model.LongWrapper;
import org.apache.hudi.avro.model.StringWrapper;
import org.apache.hudi.avro.model.TimestampMicrosWrapper;
import org.apache.hudi.common.config.SerializableSchema;
-import org.apache.hudi.common.model.HoodieAvroIndexedRecord;
-import org.apache.hudi.common.model.HoodieAvroRecord;
import org.apache.hudi.common.model.HoodieOperation;
import org.apache.hudi.common.model.HoodieRecord;
import org.apache.hudi.common.util.Option;
@@ -792,12 +790,12 @@ public class HoodieAvroUtils {
* @param schema {@link Schema} instance.
* @return Column value.
*/
- public static Object[] getRecordColumnValues(HoodieAvroRecord record,
+ public static Object[] getRecordColumnValues(HoodieRecord record,
String[] columns,
Schema schema,
boolean
consistentLogicalTimestampEnabled) {
try {
- GenericRecord genericRecord = (GenericRecord) ((HoodieAvroIndexedRecord)
record.toIndexedRecord(schema, new Properties()).get()).getData();
+ GenericRecord genericRecord = (GenericRecord)
(record.toIndexedRecord(schema, new Properties()).get()).getData();
List<Object> list = new ArrayList<>();
for (String col : columns) {
list.add(HoodieAvroUtils.getNestedFieldVal(genericRecord, col, true,
consistentLogicalTimestampEnabled));
@@ -816,7 +814,7 @@ public class HoodieAvroUtils {
* @param schema {@link SerializableSchema} instance.
* @return Column value if a single column, or concatenated String values by
comma.
*/
- public static Object getRecordColumnValues(HoodieAvroRecord record,
+ public static Object getRecordColumnValues(HoodieRecord record,
String[] columns,
SerializableSchema schema,
boolean consistentLogicalTimestampEnabled) {
return getRecordColumnValues(record, columns, schema.get(),
consistentLogicalTimestampEnabled);
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 9367b7b0a07..42499267cd5 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
@@ -19,6 +19,7 @@
package org.apache.hudi.metadata;
import org.apache.hudi.avro.ConvertingGenericData;
+import org.apache.hudi.avro.HoodieAvroUtils;
import org.apache.hudi.avro.model.HoodieCleanMetadata;
import org.apache.hudi.avro.model.HoodieMetadataColumnStats;
import org.apache.hudi.avro.model.HoodieRecordIndexInfo;
@@ -66,8 +67,6 @@ import org.apache.hudi.util.Lazy;
import org.apache.avro.AvroTypeException;
import org.apache.avro.LogicalTypes;
import org.apache.avro.Schema;
-import org.apache.avro.generic.GenericRecord;
-import org.apache.avro.generic.IndexedRecord;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.FileSystem;
import org.apache.hadoop.fs.Path;
@@ -101,7 +100,6 @@ import java.util.stream.Stream;
import static org.apache.hudi.avro.AvroSchemaUtils.resolveNullableSchema;
import static org.apache.hudi.avro.HoodieAvroUtils.addMetadataFields;
-import static
org.apache.hudi.avro.HoodieAvroUtils.convertValueForSpecificDataTypes;
import static
org.apache.hudi.avro.HoodieAvroUtils.getNestedFieldSchemaFromWriteSchema;
import static org.apache.hudi.avro.HoodieAvroUtils.unwrapAvroValueWrapper;
import static
org.apache.hudi.common.table.timeline.HoodieInstantTimeGenerator.MILLIS_INSTANT_ID_LENGTH;
@@ -179,9 +177,8 @@ public class HoodieTableMetadataUtil {
* @return map of {@link HoodieColumnRangeMetadata} for each of the provided
target fields for
* the collection of provided records
*/
- public static Map<String, HoodieColumnRangeMetadata<Comparable>>
collectColumnRangeMetadata(List<IndexedRecord> records,
-
List<Schema.Field> targetFields,
-
String filePath) {
+ public static Map<String, HoodieColumnRangeMetadata<Comparable>>
collectColumnRangeMetadata(
+ List<HoodieRecord> records, List<Schema.Field> targetFields, String
filePath, Schema recordSchema) {
// Helper class to calculate column stats
class ColumnStats {
Object minValue;
@@ -200,23 +197,26 @@ public class HoodieTableMetadataUtil {
targetFields.forEach(field -> {
ColumnStats colStats = allColumnStats.computeIfAbsent(field.name(),
(ignored) -> new ColumnStats());
- GenericRecord genericRecord = (GenericRecord) record;
-
- final Object fieldVal =
convertValueForSpecificDataTypes(field.schema(),
genericRecord.get(field.name()), false);
- final Schema fieldSchema =
getNestedFieldSchemaFromWriteSchema(genericRecord.getSchema(), field.name());
+ Schema fieldSchema = getNestedFieldSchemaFromWriteSchema(recordSchema,
field.name());
+ Object fieldValue;
+ if (record.getRecordType() == HoodieRecordType.AVRO) {
+ fieldValue = HoodieAvroUtils.getRecordColumnValues(record, new
String[]{field.name()}, recordSchema, false)[0];
+ } else if (record.getRecordType() == HoodieRecordType.SPARK) {
+ fieldValue = record.getColumnValues(recordSchema, new
String[]{field.name()}, false)[0];
+ } else {
+ throw new HoodieException(String.format("Unknown record type: %s",
record.getRecordType()));
+ }
colStats.valueCount++;
-
- if (fieldVal != null && canCompare(fieldSchema)) {
+ if (fieldValue != null && canCompare(fieldSchema)) {
// Set the min value of the field
if (colStats.minValue == null
- || ConvertingGenericData.INSTANCE.compare(fieldVal,
colStats.minValue, fieldSchema) < 0) {
- colStats.minValue = fieldVal;
+ || ConvertingGenericData.INSTANCE.compare(fieldValue,
colStats.minValue, fieldSchema) < 0) {
+ colStats.minValue = fieldValue;
}
-
// Set the max value of the field
- if (colStats.maxValue == null ||
ConvertingGenericData.INSTANCE.compare(fieldVal, colStats.maxValue,
fieldSchema) > 0) {
- colStats.maxValue = fieldVal;
+ if (colStats.maxValue == null ||
ConvertingGenericData.INSTANCE.compare(fieldValue, colStats.maxValue,
fieldSchema) > 0) {
+ colStats.maxValue = fieldValue;
}
} else {
colStats.nullCount++;