This is an automated email from the ASF dual-hosted git repository.
sivabalan 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 be216f4a0351 [HUDI-8581] Test schema handler in fg reader and some
refactoring to prevent bugs in the future (#12340)
be216f4a0351 is described below
commit be216f4a0351f202ea79964e8941d68ca33d0aab
Author: Jon Vexler <[email protected]>
AuthorDate: Fri Jul 18 17:08:41 2025 -0400
[HUDI-8581] Test schema handler in fg reader and some refactoring to
prevent bugs in the future (#12340)
---------
Co-authored-by: Jonathan Vexler <=>
Co-authored-by: Lokesh Jain <[email protected]>
Co-authored-by: Lokesh Jain <[email protected]>
---
.../SparkFileFormatInternalRowReaderContext.scala | 2 +-
.../hudi/common/engine/HoodieReaderContext.java | 2 +
.../hudi/common/model/HoodieRecordMerger.java | 1 +
.../table/read/FileGroupReaderSchemaHandler.java | 27 +-
.../common/table/read/HoodieFileGroupReader.java | 9 +-
...java => ParquetRowIndexBasedSchemaHandler.java} | 38 ++-
.../common/table/read/SchemaHandlerTestBase.java | 380 +++++++++++++++++++++
.../read/TestFileGroupReaderSchemaHandler.java | 114 +++++++
.../TestParquetRowIndexBasedSchemaHandler.java | 111 ++++++
.../TestPositionBasedFileGroupRecordBuffer.java | 8 +-
10 files changed, 651 insertions(+), 41 deletions(-)
diff --git
a/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/SparkFileFormatInternalRowReaderContext.scala
b/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/SparkFileFormatInternalRowReaderContext.scala
index 1a08c69cd4bc..521bb3d4876e 100644
---
a/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/SparkFileFormatInternalRowReaderContext.scala
+++
b/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/SparkFileFormatInternalRowReaderContext.scala
@@ -169,7 +169,7 @@ class
SparkFileFormatInternalRowReaderContext(parquetFileReader: SparkParquetRea
HoodieAvroUtils.removeFields(skeletonRequiredSchema, rowIndexColumn))
//If we need to do position based merging with log files we will leave
the row index column at the end
- val dataProjection = if (getHasLogFiles &&
getShouldMergeUseRecordPosition) {
+ val dataProjection = if (getShouldMergeUseRecordPosition) {
getBootstrapProjection(dataRequiredSchema, dataRequiredSchema,
partitionFieldAndValues)
} else {
getBootstrapProjection(dataRequiredSchema,
diff --git
a/hudi-common/src/main/java/org/apache/hudi/common/engine/HoodieReaderContext.java
b/hudi-common/src/main/java/org/apache/hudi/common/engine/HoodieReaderContext.java
index 8a5b162c1356..53a880160e71 100644
---
a/hudi-common/src/main/java/org/apache/hudi/common/engine/HoodieReaderContext.java
+++
b/hudi-common/src/main/java/org/apache/hudi/common/engine/HoodieReaderContext.java
@@ -90,6 +90,8 @@ public abstract class HoodieReaderContext<T> {
private Boolean hasLogFiles = null;
private Boolean hasBootstrapBaseFile = null;
private Boolean needsBootstrapMerge = null;
+
+ // should we do position based merging for mor
private Boolean shouldMergeUseRecordPosition = null;
protected String partitionPath;
protected Option<InstantRange> instantRangeOpt = Option.empty();
diff --git
a/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieRecordMerger.java
b/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieRecordMerger.java
index a893bf49ba6e..a01e7fcfde1a 100644
---
a/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieRecordMerger.java
+++
b/hudi-common/src/main/java/org/apache/hudi/common/model/HoodieRecordMerger.java
@@ -152,6 +152,7 @@ public interface HoodieRecordMerger extends Serializable {
* If false, whenever we have log files, we will need to read all columns
* If true, mor merging can be done without all columns. The columns
required can be configured
* by overriding getMandatoryFieldsForMerging
+ * EventTime based merging and CommitTime based merging are projection
compatible
*/
default boolean isProjectionCompatible() {
return false;
diff --git
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/FileGroupReaderSchemaHandler.java
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/FileGroupReaderSchemaHandler.java
index 816ec21a9fbd..b5f979df8fd5 100644
---
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/FileGroupReaderSchemaHandler.java
+++
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/FileGroupReaderSchemaHandler.java
@@ -26,7 +26,6 @@ import org.apache.hudi.common.engine.HoodieReaderContext;
import org.apache.hudi.common.model.HoodieRecord;
import org.apache.hudi.common.model.HoodieRecordMerger;
import org.apache.hudi.common.table.HoodieTableConfig;
-import org.apache.hudi.common.util.LocalAvroSchemaCache;
import org.apache.hudi.common.util.Option;
import org.apache.hudi.common.util.StringUtils;
import org.apache.hudi.common.util.VisibleForTesting;
@@ -79,15 +78,6 @@ public class FileGroupReaderSchemaHandler<T> {
protected final TypedProperties properties;
- protected final Option<HoodieRecordMerger> recordMerger;
-
- protected final boolean hasBootstrapBaseFile;
- protected boolean needsBootstrapMerge;
-
- protected final boolean needsMORMerge;
-
- private final LocalAvroSchemaCache localAvroSchemaCache;
-
private final Option<Pair<String, String>> customDeleteMarkerKeyValue;
private final boolean hasBuiltInDelete;
private final int hoodieOperationPos;
@@ -100,9 +90,6 @@ public class FileGroupReaderSchemaHandler<T> {
TypedProperties properties) {
this.properties = properties;
this.readerContext = readerContext;
- this.hasBootstrapBaseFile = readerContext.getHasBootstrapBaseFile();
- this.needsMORMerge = readerContext.getHasLogFiles();
- this.recordMerger = readerContext.getRecordMerger();
this.tableSchema = tableSchema;
this.requestedSchema = AvroSchemaCache.intern(requestedSchema);
this.hoodieTableConfig = hoodieTableConfig;
@@ -113,8 +100,6 @@ public class FileGroupReaderSchemaHandler<T> {
this.hoodieOperationPos =
Option.ofNullable(requiredSchema.getField(HoodieRecord.OPERATION_METADATA_FIELD)).map(Schema.Field::pos).orElse(-1);
this.internalSchema = pruneInternalSchema(requiredSchema,
internalSchemaOpt);
this.internalSchemaOpt = getInternalSchemaOpt(internalSchemaOpt);
- readerContext.setNeedsBootstrapMerge(this.needsBootstrapMerge);
- this.localAvroSchemaCache = LocalAvroSchemaCache.getInstance();
}
public Schema getTableSchema() {
@@ -179,7 +164,8 @@ public class FileGroupReaderSchemaHandler<T> {
@VisibleForTesting
Schema generateRequiredSchema() {
boolean hasInstantRange = readerContext.getInstantRange().isPresent();
- if (!needsMORMerge) {
+ //might need to change this if other queries than mor have mandatory fields
+ if (!readerContext.getHasLogFiles()) {
if (hasInstantRange && !findNestedField(requestedSchema,
HoodieRecord.COMMIT_TIME_METADATA_FIELD).isPresent()) {
List<Schema.Field> addedFields = new ArrayList<>();
addedFields.add(getField(tableSchema,
HoodieRecord.COMMIT_TIME_METADATA_FIELD));
@@ -189,14 +175,14 @@ public class FileGroupReaderSchemaHandler<T> {
}
if (hoodieTableConfig.getRecordMergeMode() == RecordMergeMode.CUSTOM) {
- if (!recordMerger.get().isProjectionCompatible()) {
+ if (!readerContext.getRecordMerger().get().isProjectionCompatible()) {
return tableSchema;
}
}
List<Schema.Field> addedFields = new ArrayList<>();
for (String field : getMandatoryFieldsForMerging(
- hoodieTableConfig, properties, tableSchema, recordMerger,
+ hoodieTableConfig, properties, tableSchema,
readerContext.getRecordMerger(),
hasBuiltInDelete, customDeleteMarkerKeyValue, hasInstantRange)) {
if (!findNestedField(requestedSchema, field).isPresent()) {
addedFields.add(getField(tableSchema, field));
@@ -270,8 +256,9 @@ public class FileGroupReaderSchemaHandler<T> {
protected Schema prepareRequiredSchema() {
Schema preReorderRequiredSchema = generateRequiredSchema();
Pair<List<Schema.Field>, List<Schema.Field>> requiredFields =
getDataAndMetaCols(preReorderRequiredSchema);
- this.needsBootstrapMerge = hasBootstrapBaseFile &&
!requiredFields.getLeft().isEmpty() && !requiredFields.getRight().isEmpty();
- return needsBootstrapMerge
+
readerContext.setNeedsBootstrapMerge(readerContext.getHasBootstrapBaseFile()
+ && !requiredFields.getLeft().isEmpty() &&
!requiredFields.getRight().isEmpty());
+ return readerContext.getNeedsBootstrapMerge()
?
createSchemaFromFields(Stream.concat(requiredFields.getLeft().stream(),
requiredFields.getRight().stream()).collect(Collectors.toList()))
: preReorderRequiredSchema;
}
diff --git
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java
index 89e8505af9a8..9f7c8f8e5722 100644
---
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java
+++
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/HoodieFileGroupReader.java
@@ -121,7 +121,12 @@ public final class HoodieFileGroupReader<T> implements
Closeable {
this.metaClient = hoodieTableMetaClient;
this.storage = storage;
this.hoodieBaseFileOption = fileSlice.getBaseFile();
+ readerContext.setHasBootstrapBaseFile(hoodieBaseFileOption.isPresent() &&
hoodieBaseFileOption.get().getBootstrapBaseFile().isPresent());
this.logFiles =
fileSlice.getLogFiles().sorted(HoodieLogFile.getLogFileComparator()).collect(Collectors.toList());
+ readerContext.setHasLogFiles(!this.logFiles.isEmpty());
+ if (readerContext.getHasLogFiles() && start != 0) {
+ throw new IllegalArgumentException("Filegroup reader is doing log file
merge but not reading from the start of the base file");
+ }
this.props = props;
this.start = start;
this.length = length;
@@ -133,7 +138,7 @@ public final class HoodieFileGroupReader<T> implements
Closeable {
readerContext.setTablePath(tablePath);
readerContext.setLatestCommitTime(latestCommitTime);
boolean isSkipMerge = ConfigUtils.getStringWithAltKeys(props,
HoodieReaderConfig.MERGE_TYPE,
true).equalsIgnoreCase(HoodieReaderConfig.REALTIME_SKIP_MERGE);
- readerContext.setShouldMergeUseRecordPosition(shouldUseRecordPosition &&
!isSkipMerge);
+ readerContext.setShouldMergeUseRecordPosition(shouldUseRecordPosition &&
!isSkipMerge && readerContext.getHasLogFiles());
readerContext.setHasLogFiles(!this.logFiles.isEmpty());
readerContext.setPartitionPath(partitionPath);
if (readerContext.getHasLogFiles() && start != 0) {
@@ -141,7 +146,7 @@ public final class HoodieFileGroupReader<T> implements
Closeable {
}
readerContext.setHasBootstrapBaseFile(hoodieBaseFileOption.isPresent() &&
hoodieBaseFileOption.get().getBootstrapBaseFile().isPresent());
readerContext.setSchemaHandler(readerContext.supportsParquetRowIndex()
- ? new PositionBasedSchemaHandler<>(readerContext, dataSchema,
requestedSchema, internalSchemaOpt, tableConfig, props)
+ ? new ParquetRowIndexBasedSchemaHandler<>(readerContext, dataSchema,
requestedSchema, internalSchemaOpt, tableConfig, props)
: new FileGroupReaderSchemaHandler<>(readerContext, dataSchema,
requestedSchema, internalSchemaOpt, tableConfig, props));
this.outputConverter =
readerContext.getSchemaHandler().getOutputConverter();
this.orderingFieldName = readerContext.getMergeMode() ==
RecordMergeMode.COMMIT_TIME_ORDERING
diff --git
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/PositionBasedSchemaHandler.java
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/ParquetRowIndexBasedSchemaHandler.java
similarity index 70%
rename from
hudi-common/src/main/java/org/apache/hudi/common/table/read/PositionBasedSchemaHandler.java
rename to
hudi-common/src/main/java/org/apache/hudi/common/table/read/ParquetRowIndexBasedSchemaHandler.java
index c91073b4f78c..81a4fc016bea 100644
---
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/PositionBasedSchemaHandler.java
+++
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/ParquetRowIndexBasedSchemaHandler.java
@@ -23,6 +23,7 @@ import org.apache.hudi.common.config.TypedProperties;
import org.apache.hudi.common.engine.HoodieReaderContext;
import org.apache.hudi.common.table.HoodieTableConfig;
import org.apache.hudi.common.util.Option;
+import org.apache.hudi.common.util.VisibleForTesting;
import org.apache.hudi.common.util.collection.Pair;
import org.apache.hudi.internal.schema.InternalSchema;
import org.apache.hudi.internal.schema.Types;
@@ -38,34 +39,37 @@ import static
org.apache.hudi.avro.AvroSchemaUtils.appendFieldsToSchemaDedupNest
import static
org.apache.hudi.common.table.read.PositionBasedFileGroupRecordBuffer.ROW_INDEX_TEMPORARY_COLUMN_NAME;
/**
- * This class is responsible for handling the schema for the file group reader
that supports positional merge.
+ * This class is responsible for handling the schema for the file group reader
that supports row index based positional merge.
*/
-public class PositionBasedSchemaHandler<T> extends
FileGroupReaderSchemaHandler<T> {
- public PositionBasedSchemaHandler(HoodieReaderContext<T> readerContext,
- Schema dataSchema,
- Schema requestedSchema,
- Option<InternalSchema> internalSchemaOpt,
- HoodieTableConfig hoodieTableConfig,
- TypedProperties properties) {
+public class ParquetRowIndexBasedSchemaHandler<T> extends
FileGroupReaderSchemaHandler<T> {
+ public ParquetRowIndexBasedSchemaHandler(HoodieReaderContext<T>
readerContext,
+ Schema dataSchema,
+ Schema requestedSchema,
+ Option<InternalSchema>
internalSchemaOpt,
+ HoodieTableConfig hoodieTableConfig,
+ TypedProperties properties) {
super(readerContext, dataSchema, requestedSchema, internalSchemaOpt,
hoodieTableConfig, properties);
+ if (!readerContext.supportsParquetRowIndex()) {
+ throw new IllegalStateException("Using " + this.getClass().getName() + "
but context does not support parquet row index");
+ }
}
@Override
protected Schema prepareRequiredSchema() {
Schema preMergeSchema = super.prepareRequiredSchema();
- return readerContext.getShouldMergeUseRecordPosition() &&
readerContext.getHasLogFiles()
+ return readerContext.getShouldMergeUseRecordPosition()
? addPositionalMergeCol(preMergeSchema)
: preMergeSchema;
}
@Override
protected Option<InternalSchema> getInternalSchemaOpt(Option<InternalSchema>
internalSchemaOpt) {
- return
internalSchemaOpt.map(PositionBasedSchemaHandler::addPositionalMergeCol);
+ return
internalSchemaOpt.map(ParquetRowIndexBasedSchemaHandler::addPositionalMergeCol);
}
@Override
protected InternalSchema doPruneInternalSchema(Schema requiredSchema,
InternalSchema internalSchema) {
- if (!(readerContext.getShouldMergeUseRecordPosition() &&
readerContext.getHasLogFiles())) {
+ if (!readerContext.getShouldMergeUseRecordPosition()) {
return super.doPruneInternalSchema(requiredSchema, internalSchema);
}
@@ -82,20 +86,24 @@ public class PositionBasedSchemaHandler<T> extends
FileGroupReaderSchemaHandler<
@Override
public Pair<List<Schema.Field>,List<Schema.Field>>
getBootstrapRequiredFields() {
Pair<List<Schema.Field>,List<Schema.Field>> dataAndMetaCols =
super.getBootstrapRequiredFields();
- if (readerContext.supportsParquetRowIndex()) {
- if (!dataAndMetaCols.getLeft().isEmpty() &&
!dataAndMetaCols.getRight().isEmpty()) {
+ if (readerContext.getNeedsBootstrapMerge() ||
readerContext.getShouldMergeUseRecordPosition()) {
+ if (!dataAndMetaCols.getLeft().isEmpty()) {
dataAndMetaCols.getLeft().add(getPositionalMergeField());
+ }
+ if (!dataAndMetaCols.getRight().isEmpty()) {
dataAndMetaCols.getRight().add(getPositionalMergeField());
}
}
return dataAndMetaCols;
}
- private static Schema addPositionalMergeCol(Schema input) {
+ @VisibleForTesting
+ static Schema addPositionalMergeCol(Schema input) {
return appendFieldsToSchemaDedupNested(input,
Collections.singletonList(getPositionalMergeField()));
}
- private static Schema.Field getPositionalMergeField() {
+ @VisibleForTesting
+ static Schema.Field getPositionalMergeField() {
return new Schema.Field(ROW_INDEX_TEMPORARY_COLUMN_NAME,
Schema.create(Schema.Type.LONG), "", -1L);
}
diff --git
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/SchemaHandlerTestBase.java
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/SchemaHandlerTestBase.java
new file mode 100644
index 000000000000..e209b91dc303
--- /dev/null
+++
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/SchemaHandlerTestBase.java
@@ -0,0 +1,380 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.hudi.common.table.read;
+
+import org.apache.hudi.avro.HoodieAvroUtils;
+import org.apache.hudi.common.config.RecordMergeMode;
+import org.apache.hudi.common.engine.HoodieReaderContext;
+import org.apache.hudi.common.model.HoodieRecord;
+import org.apache.hudi.common.model.HoodieRecordMerger;
+import org.apache.hudi.common.table.HoodieTableConfig;
+import org.apache.hudi.common.table.HoodieTableVersion;
+import org.apache.hudi.common.testutils.HoodieTestDataGenerator;
+import org.apache.hudi.common.util.Option;
+import org.apache.hudi.common.util.collection.ClosableIterator;
+import org.apache.hudi.common.util.collection.Pair;
+import org.apache.hudi.storage.HoodieStorage;
+import org.apache.hudi.storage.StoragePath;
+
+import org.apache.avro.Schema;
+import org.apache.avro.generic.GenericRecord;
+import org.apache.avro.generic.IndexedRecord;
+import org.junit.jupiter.params.provider.Arguments;
+
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.List;
+import java.util.Map;
+import java.util.function.UnaryOperator;
+import java.util.stream.Stream;
+
+import static
org.apache.hudi.common.config.RecordMergeMode.COMMIT_TIME_ORDERING;
+import static org.apache.hudi.common.config.RecordMergeMode.CUSTOM;
+import static
org.apache.hudi.common.config.RecordMergeMode.EVENT_TIME_ORDERING;
+import static
org.apache.hudi.common.table.read.ParquetRowIndexBasedSchemaHandler.addPositionalMergeCol;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+public abstract class SchemaHandlerTestBase {
+
+ protected static final Schema DATA_SCHEMA =
HoodieAvroUtils.addMetadataFields(HoodieTestDataGenerator.AVRO_SCHEMA);
+ protected static final Schema DATA_SCHEMA_NO_DELETE =
generateProjectionSchema(DATA_SCHEMA.getFields().stream()
+ .map(Schema.Field::name).filter(f ->
!f.equals("_hoodie_is_deleted")).toArray(String[]::new));
+ protected static final Schema DATA_COLS_ONLY_SCHEMA =
generateProjectionSchema("begin_lat", "tip_history", "rider");
+ protected static final Schema META_COLS_ONLY_SCHEMA =
generateProjectionSchema("_hoodie_commit_seqno", "_hoodie_record_key");
+
+ static Stream<Arguments> testMorParams(boolean supportsParquetRowIndex) {
+ Stream.Builder<Arguments> b = Stream.builder();
+ for (boolean mergeUseRecordPosition : new boolean[] {true, false}) {
+ for (boolean hasBuiltInDelete : new boolean[] {true, false}) {
+ b.add(Arguments.of(EVENT_TIME_ORDERING, true, false,
mergeUseRecordPosition, supportsParquetRowIndex, hasBuiltInDelete));
+ b.add(Arguments.of(EVENT_TIME_ORDERING, false, false,
mergeUseRecordPosition, supportsParquetRowIndex, hasBuiltInDelete));
+ b.add(Arguments.of(COMMIT_TIME_ORDERING, false, false,
mergeUseRecordPosition, supportsParquetRowIndex, hasBuiltInDelete));
+ b.add(Arguments.of(CUSTOM, false, true, mergeUseRecordPosition,
supportsParquetRowIndex, hasBuiltInDelete));
+ b.add(Arguments.of(CUSTOM, false, false, mergeUseRecordPosition,
supportsParquetRowIndex, hasBuiltInDelete));
+ }
+ }
+ return b.build();
+ }
+
+ public void testMor(RecordMergeMode mergeMode,
+ boolean hasPrecombine,
+ boolean isProjectionCompatible,
+ boolean mergeUseRecordPosition,
+ boolean supportsParquetRowIndex,
+ boolean hasBuiltInDelete) throws IOException {
+ Schema dataSchema = hasBuiltInDelete ? DATA_SCHEMA : DATA_SCHEMA_NO_DELETE;
+ HoodieTableConfig hoodieTableConfig = mock(HoodieTableConfig.class);
+ setupMORTable(mergeMode, hasPrecombine, hoodieTableConfig);
+ HoodieRecordMerger merger = mockRecordMerger(isProjectionCompatible,
+ isProjectionCompatible ? new String[] {"begin_lat", "begin_lon",
"_hoodie_record_key", "timestamp"} : new String[] {"begin_lat", "begin_lon",
"timestamp"});
+ HoodieReaderContext<String> readerContext =
createReaderContext(hoodieTableConfig, supportsParquetRowIndex, true, false,
mergeUseRecordPosition, merger);
+ readerContext.setRecordMerger(Option.of(merger));
+ Schema requestedSchema = dataSchema;
+ FileGroupReaderSchemaHandler schemaHandler =
createSchemaHandler(readerContext, dataSchema, requestedSchema,
hoodieTableConfig, supportsParquetRowIndex);
+ Schema expectedRequiredFullSchema = supportsParquetRowIndex &&
mergeUseRecordPosition
+ ?
ParquetRowIndexBasedSchemaHandler.addPositionalMergeCol(requestedSchema)
+ : requestedSchema;
+ assertEquals(expectedRequiredFullSchema,
schemaHandler.getRequiredSchema());
+ assertFalse(readerContext.getNeedsBootstrapMerge());
+
+ //read subset of columns
+ requestedSchema = DATA_COLS_ONLY_SCHEMA;
+ schemaHandler = createSchemaHandler(readerContext, dataSchema,
requestedSchema, hoodieTableConfig, supportsParquetRowIndex);
+ Schema expectedRequiredSchema;
+ if (mergeMode == EVENT_TIME_ORDERING && hasPrecombine) {
+ expectedRequiredSchema = generateProjectionSchema(hasBuiltInDelete,
"begin_lat", "tip_history", "rider", "_hoodie_record_key", "timestamp");
+ } else if (mergeMode == EVENT_TIME_ORDERING || mergeMode ==
COMMIT_TIME_ORDERING) {
+ expectedRequiredSchema = generateProjectionSchema(hasBuiltInDelete,
"begin_lat", "tip_history", "rider", "_hoodie_record_key");
+ } else if (mergeMode == CUSTOM && isProjectionCompatible) {
+ expectedRequiredSchema = generateProjectionSchema("begin_lat",
"tip_history", "rider", "begin_lon", "_hoodie_record_key", "timestamp");
+ } else {
+ expectedRequiredSchema = dataSchema;
+ }
+ if (supportsParquetRowIndex && mergeUseRecordPosition) {
+ expectedRequiredSchema = addPositionalMergeCol(expectedRequiredSchema);
+ }
+ assertEquals(expectedRequiredSchema, schemaHandler.getRequiredSchema());
+ assertFalse(readerContext.getNeedsBootstrapMerge());
+ }
+
+ public void testMorBootstrap(RecordMergeMode mergeMode,
+ boolean hasPrecombine,
+ boolean isProjectionCompatible,
+ boolean mergeUseRecordPosition,
+ boolean supportsParquetRowIndex,
+ boolean hasBuiltInDelete) throws IOException {
+ Schema dataSchema = hasBuiltInDelete ? DATA_SCHEMA : DATA_SCHEMA_NO_DELETE;
+ HoodieTableConfig hoodieTableConfig = mock(HoodieTableConfig.class);
+ setupMORTable(mergeMode, hasPrecombine, hoodieTableConfig);
+ HoodieRecordMerger merger = mockRecordMerger(isProjectionCompatible, new
String[] {"begin_lat", "begin_lon", "timestamp"});
+ HoodieReaderContext<String> readerContext =
createReaderContext(hoodieTableConfig, supportsParquetRowIndex, true, true,
mergeUseRecordPosition, merger);
+ readerContext.setRecordMerger(Option.of(merger));
+ Schema requestedSchema = dataSchema;
+ FileGroupReaderSchemaHandler schemaHandler =
createSchemaHandler(readerContext, dataSchema, requestedSchema,
hoodieTableConfig, supportsParquetRowIndex);
+ Schema expectedRequiredFullSchema = supportsParquetRowIndex &&
mergeUseRecordPosition
+ ?
ParquetRowIndexBasedSchemaHandler.addPositionalMergeCol(requestedSchema)
+ : requestedSchema;
+ assertEquals(expectedRequiredFullSchema,
schemaHandler.getRequiredSchema());
+ assertTrue(readerContext.getNeedsBootstrapMerge());
+ Pair<List<Schema.Field>, List<Schema.Field>> bootstrapFields =
schemaHandler.getBootstrapRequiredFields();
+ Pair<List<Schema.Field>, List<Schema.Field>> expectedBootstrapFields =
FileGroupReaderSchemaHandler.getDataAndMetaCols(expectedRequiredFullSchema);
+ if (supportsParquetRowIndex) {
+
expectedBootstrapFields.getLeft().add(ParquetRowIndexBasedSchemaHandler.getPositionalMergeField());
+
expectedBootstrapFields.getRight().add(ParquetRowIndexBasedSchemaHandler.getPositionalMergeField());
+ }
+ assertEquals(expectedBootstrapFields.getLeft(), bootstrapFields.getLeft());
+ assertEquals(expectedBootstrapFields.getRight(),
bootstrapFields.getRight());
+
+ //read subset of columns
+ requestedSchema = generateProjectionSchema("begin_lat", "tip_history",
"_hoodie_record_key", "rider");
+ schemaHandler = createSchemaHandler(readerContext, dataSchema,
requestedSchema, hoodieTableConfig, supportsParquetRowIndex);
+ Schema expectedRequiredSchema;
+ if (mergeMode == EVENT_TIME_ORDERING && hasPrecombine) {
+ expectedRequiredSchema = generateProjectionSchema(hasBuiltInDelete,
"_hoodie_record_key", "begin_lat", "tip_history", "rider", "timestamp");
+ } else if (mergeMode == EVENT_TIME_ORDERING || mergeMode ==
COMMIT_TIME_ORDERING) {
+ expectedRequiredSchema = generateProjectionSchema(hasBuiltInDelete,
"_hoodie_record_key", "begin_lat", "tip_history", "rider");
+ } else if (mergeMode == CUSTOM && isProjectionCompatible) {
+ expectedRequiredSchema = generateProjectionSchema("_hoodie_record_key",
"begin_lat", "tip_history", "rider", "begin_lon", "timestamp");
+ } else {
+ expectedRequiredSchema = dataSchema;
+ }
+ if (supportsParquetRowIndex && mergeUseRecordPosition) {
+ expectedRequiredSchema = addPositionalMergeCol(expectedRequiredSchema);
+ }
+ assertEquals(expectedRequiredSchema, schemaHandler.getRequiredSchema());
+ assertTrue(readerContext.getNeedsBootstrapMerge());
+ bootstrapFields = schemaHandler.getBootstrapRequiredFields();
+ expectedBootstrapFields =
FileGroupReaderSchemaHandler.getDataAndMetaCols(expectedRequiredSchema);
+ if (supportsParquetRowIndex) {
+
expectedBootstrapFields.getLeft().add(ParquetRowIndexBasedSchemaHandler.getPositionalMergeField());
+
expectedBootstrapFields.getRight().add(ParquetRowIndexBasedSchemaHandler.getPositionalMergeField());
+ }
+ assertEquals(expectedBootstrapFields.getLeft(), bootstrapFields.getLeft());
+ assertEquals(expectedBootstrapFields.getRight(),
bootstrapFields.getRight());
+
+ // request only data cols
+ requestedSchema = DATA_COLS_ONLY_SCHEMA;
+ schemaHandler = createSchemaHandler(readerContext, dataSchema,
requestedSchema, hoodieTableConfig, supportsParquetRowIndex);
+ if (mergeMode == EVENT_TIME_ORDERING && hasPrecombine) {
+ expectedRequiredSchema = generateProjectionSchema(hasBuiltInDelete,
"_hoodie_record_key", "begin_lat", "tip_history", "rider", "timestamp");
+ assertTrue(readerContext.getNeedsBootstrapMerge());
+ } else if (mergeMode == EVENT_TIME_ORDERING || mergeMode ==
COMMIT_TIME_ORDERING) {
+ expectedRequiredSchema = generateProjectionSchema(hasBuiltInDelete,
"_hoodie_record_key", "begin_lat", "tip_history", "rider");
+ assertTrue(readerContext.getNeedsBootstrapMerge());
+ } else if (mergeMode == CUSTOM && isProjectionCompatible) {
+ expectedRequiredSchema = generateProjectionSchema("begin_lat",
"tip_history", "rider", "begin_lon", "timestamp");
+ assertFalse(readerContext.getNeedsBootstrapMerge());
+ } else {
+ assertTrue(readerContext.getNeedsBootstrapMerge());
+ expectedRequiredSchema = dataSchema;
+ }
+ if (supportsParquetRowIndex && mergeUseRecordPosition) {
+ expectedRequiredSchema = addPositionalMergeCol(expectedRequiredSchema);
+ }
+ assertEquals(expectedRequiredSchema, schemaHandler.getRequiredSchema());
+ bootstrapFields = schemaHandler.getBootstrapRequiredFields();
+ expectedBootstrapFields =
FileGroupReaderSchemaHandler.getDataAndMetaCols(expectedRequiredSchema);
+ if (supportsParquetRowIndex) {
+ if (mergeMode == CUSTOM && isProjectionCompatible) {
+ if (mergeUseRecordPosition) {
+
expectedBootstrapFields.getRight().add(ParquetRowIndexBasedSchemaHandler.getPositionalMergeField());
+ }
+ } else {
+
expectedBootstrapFields.getLeft().add(ParquetRowIndexBasedSchemaHandler.getPositionalMergeField());
+
expectedBootstrapFields.getRight().add(ParquetRowIndexBasedSchemaHandler.getPositionalMergeField());
+ }
+ }
+ assertEquals(expectedBootstrapFields.getLeft(), bootstrapFields.getLeft());
+ assertEquals(expectedBootstrapFields.getRight(),
bootstrapFields.getRight());
+
+ // request only meta cols
+ requestedSchema = META_COLS_ONLY_SCHEMA;
+ schemaHandler = createSchemaHandler(readerContext, dataSchema,
requestedSchema, hoodieTableConfig, supportsParquetRowIndex);
+ if (mergeMode == EVENT_TIME_ORDERING && hasPrecombine) {
+ expectedRequiredSchema = generateProjectionSchema(hasBuiltInDelete,
"_hoodie_commit_seqno", "_hoodie_record_key", "timestamp");
+ assertTrue(readerContext.getNeedsBootstrapMerge());
+ } else if (mergeMode == EVENT_TIME_ORDERING || mergeMode ==
COMMIT_TIME_ORDERING) {
+ expectedRequiredSchema = generateProjectionSchema(hasBuiltInDelete,
"_hoodie_commit_seqno", "_hoodie_record_key");
+ assertEquals(hasBuiltInDelete, readerContext.getNeedsBootstrapMerge());
+ } else if (mergeMode == CUSTOM && isProjectionCompatible) {
+ expectedRequiredSchema =
generateProjectionSchema("_hoodie_commit_seqno", "_hoodie_record_key",
"begin_lat", "begin_lon", "timestamp");
+ assertTrue(readerContext.getNeedsBootstrapMerge());
+ } else {
+ expectedRequiredSchema = dataSchema;
+ assertTrue(readerContext.getNeedsBootstrapMerge());
+ }
+ if (supportsParquetRowIndex && mergeUseRecordPosition) {
+ expectedRequiredSchema = addPositionalMergeCol(expectedRequiredSchema);
+ }
+ assertEquals(expectedRequiredSchema, schemaHandler.getRequiredSchema());
+ bootstrapFields = schemaHandler.getBootstrapRequiredFields();
+ expectedBootstrapFields =
FileGroupReaderSchemaHandler.getDataAndMetaCols(expectedRequiredSchema);
+ if (supportsParquetRowIndex) {
+ if (!expectedBootstrapFields.getRight().isEmpty()) {
+
expectedBootstrapFields.getLeft().add(ParquetRowIndexBasedSchemaHandler.getPositionalMergeField());
+
expectedBootstrapFields.getRight().add(ParquetRowIndexBasedSchemaHandler.getPositionalMergeField());
+ } else if (mergeUseRecordPosition) {
+
expectedBootstrapFields.getLeft().add(ParquetRowIndexBasedSchemaHandler.getPositionalMergeField());
+ }
+ }
+ assertEquals(expectedBootstrapFields.getLeft(), bootstrapFields.getLeft());
+ assertEquals(expectedBootstrapFields.getRight(),
bootstrapFields.getRight());
+ }
+
+ private static void setupMORTable(RecordMergeMode mergeMode, boolean
hasPrecombine, HoodieTableConfig hoodieTableConfig) {
+ when(hoodieTableConfig.populateMetaFields()).thenReturn(true);
+ when(hoodieTableConfig.getRecordMergeMode()).thenReturn(mergeMode);
+ if (hasPrecombine) {
+ when(hoodieTableConfig.getPreCombineField()).thenReturn("timestamp");
+ } else {
+ when(hoodieTableConfig.getPreCombineField()).thenReturn(null);
+ }
+ if (mergeMode == CUSTOM) {
+ when(hoodieTableConfig.getRecordMergeStrategyId()).thenReturn("asdf");
+ // NOTE: in this test custom doesn't have any meta cols because it is
more interesting of a test case
+
when(hoodieTableConfig.getTableVersion()).thenReturn(HoodieTableVersion.EIGHT);
+ }
+ }
+
+ private static HoodieRecordMerger mockRecordMerger(boolean
isProjectionCompatible, String[] mandatoryFields) throws IOException {
+ HoodieRecordMerger merger = mock(HoodieRecordMerger.class);
+ when(merger.isProjectionCompatible()).thenReturn(isProjectionCompatible);
+ when(merger.merge(any(), any(), any(), any(), any())).thenReturn(null);
+ when(merger.getMandatoryFieldsForMerging(any(), any(),
any())).thenReturn(mandatoryFields);
+ return merger;
+ }
+
+ static HoodieReaderContext<String> createReaderContext(HoodieTableConfig
hoodieTableConfig, boolean supportsParquetRowIndex, boolean hasLogFiles,
+ boolean
hasBootstrapBaseFile, boolean mergeUseRecordPosition, HoodieRecordMerger
merger) {
+ HoodieReaderContext<String> readerContext = new
StubbedReaderContext(hoodieTableConfig, supportsParquetRowIndex);
+ readerContext.setHasLogFiles(hasLogFiles);
+ readerContext.setHasBootstrapBaseFile(hasBootstrapBaseFile);
+ readerContext.setShouldMergeUseRecordPosition(mergeUseRecordPosition);
+ readerContext.setRecordMerger(Option.ofNullable(merger));
+ return readerContext;
+ }
+
+ abstract FileGroupReaderSchemaHandler
createSchemaHandler(HoodieReaderContext<String> readerContext, Schema
dataSchema,
+ Schema
requestedSchema, HoodieTableConfig hoodieTableConfig,
+ boolean
supportsParquetRowIndex);
+
+ static Schema generateProjectionSchema(String... fields) {
+ return HoodieAvroUtils.generateProjectionSchema(DATA_SCHEMA,
Arrays.asList(fields));
+ }
+
+ private static Schema generateProjectionSchema(boolean hasBuiltInDelete,
String... fields) {
+ List<String> fieldList = Arrays.asList(fields);
+ if (hasBuiltInDelete) {
+ fieldList = new ArrayList<>(fieldList);
+ fieldList.add("_hoodie_is_deleted");
+ }
+ return HoodieAvroUtils.generateProjectionSchema(DATA_SCHEMA, fieldList);
+ }
+
+ Schema.Field getField(String fieldName) {
+ return DATA_SCHEMA.getField(fieldName);
+ }
+
+ static class StubbedReaderContext extends HoodieReaderContext<String> {
+ private final boolean supportsParquetRowIndex;
+
+ protected StubbedReaderContext(HoodieTableConfig hoodieTableConfig,
boolean supportsParquetRowIndex) {
+ super(null, hoodieTableConfig, Option.empty(), Option.empty());
+ this.supportsParquetRowIndex = supportsParquetRowIndex;
+ }
+
+ @Override
+ public boolean supportsParquetRowIndex() {
+ return this.supportsParquetRowIndex;
+ }
+
+ @Override
+ public ClosableIterator<String> getFileRecordIterator(StoragePath
filePath, long start, long length, Schema dataSchema, Schema requiredSchema,
HoodieStorage storage) throws IOException {
+ return null;
+ }
+
+ @Override
+ public String convertAvroRecord(IndexedRecord avroRecord) {
+ return "";
+ }
+
+ @Override
+ public GenericRecord convertToAvroRecord(String record, Schema schema) {
+ return null;
+ }
+
+ @Override
+ public String getDeleteRow(String record, String recordKey) {
+ return "";
+ }
+
+ @Override
+ public Option<HoodieRecordMerger> getRecordMerger(RecordMergeMode
mergeMode, String mergeStrategyId, String mergeImplClasses) {
+ return null;
+ }
+
+ @Override
+ public Object getValue(String record, Schema schema, String fieldName) {
+ return null;
+ }
+
+ @Override
+ public String getMetaFieldValue(String record, int pos) {
+ return "";
+ }
+
+ @Override
+ public HoodieRecord<String> constructHoodieRecord(BufferedRecord<String>
bufferedRecord) {
+ return null;
+ }
+
+ @Override
+ public String constructEngineRecord(Schema schema, Map<Integer, Object>
updateValues, BufferedRecord<String> baseRecord) {
+ return "";
+ }
+
+ @Override
+ public String seal(String record) {
+ return "";
+ }
+
+ @Override
+ public String toBinaryRow(Schema avroSchema, String record) {
+ return "";
+ }
+
+ @Override
+ public ClosableIterator<String>
mergeBootstrapReaders(ClosableIterator<String> skeletonFileIterator, Schema
skeletonRequiredSchema, ClosableIterator<String> dataFileIterator,
+ Schema
dataRequiredSchema, List<Pair<String, Object>> requiredPartitionFieldAndValues)
{
+ return null;
+ }
+
+ @Override
+ public UnaryOperator<String> projectRecord(Schema from, Schema to,
Map<String, String> renamedColumns) {
+ return null;
+ }
+ }
+}
diff --git
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestFileGroupReaderSchemaHandler.java
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestFileGroupReaderSchemaHandler.java
new file mode 100644
index 000000000000..4373110cb32f
--- /dev/null
+++
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestFileGroupReaderSchemaHandler.java
@@ -0,0 +1,114 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.hudi.common.table.read;
+
+import org.apache.hudi.common.config.RecordMergeMode;
+import org.apache.hudi.common.config.TypedProperties;
+import org.apache.hudi.common.engine.HoodieReaderContext;
+import org.apache.hudi.common.table.HoodieTableConfig;
+import org.apache.hudi.common.util.Option;
+import org.apache.hudi.common.util.collection.Pair;
+
+import org.apache.avro.Schema;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.Arguments;
+import org.junit.jupiter.params.provider.MethodSource;
+
+import java.io.IOException;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.List;
+import java.util.stream.Stream;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+public class TestFileGroupReaderSchemaHandler extends SchemaHandlerTestBase {
+
+ @Test
+ public void testCow() {
+ HoodieTableConfig hoodieTableConfig = mock(HoodieTableConfig.class);
+ when(hoodieTableConfig.populateMetaFields()).thenReturn(true);
+ HoodieReaderContext<String> readerContext =
createReaderContext(hoodieTableConfig, false, false, false, false, null);
+ Schema requestedSchema = DATA_SCHEMA;
+ FileGroupReaderSchemaHandler schemaHandler =
createSchemaHandler(readerContext, DATA_SCHEMA, requestedSchema,
hoodieTableConfig, false);
+ assertEquals(requestedSchema, schemaHandler.getRequiredSchema());
+
+ //read subset of columns
+ requestedSchema = generateProjectionSchema("begin_lat", "tip_history",
"rider");
+ schemaHandler = createSchemaHandler(readerContext, DATA_SCHEMA,
requestedSchema, hoodieTableConfig, false);
+ assertEquals(requestedSchema, schemaHandler.getRequiredSchema());
+ assertFalse(readerContext.getNeedsBootstrapMerge());
+ }
+
+ @Test
+ public void testCowBootstrap() {
+ HoodieTableConfig hoodieTableConfig = mock(HoodieTableConfig.class);
+ when(hoodieTableConfig.populateMetaFields()).thenReturn(true);
+ HoodieReaderContext<String> readerContext =
createReaderContext(hoodieTableConfig, false, false, true, false, null);
+ Schema requestedSchema = generateProjectionSchema("begin_lat",
"tip_history", "_hoodie_record_key", "rider");
+
+ //meta cols must go first in the required schema
+ FileGroupReaderSchemaHandler schemaHandler =
createSchemaHandler(readerContext, DATA_SCHEMA, requestedSchema,
hoodieTableConfig, false);
+ assertTrue(readerContext.getNeedsBootstrapMerge());
+ Schema expectedRequiredSchema =
generateProjectionSchema("_hoodie_record_key", "begin_lat", "tip_history",
"rider");
+ assertEquals(expectedRequiredSchema, schemaHandler.getRequiredSchema());
+ Pair<List<Schema.Field>, List<Schema.Field>> bootstrapFields =
schemaHandler.getBootstrapRequiredFields();
+ assertEquals(Collections.singletonList(getField("_hoodie_record_key")),
bootstrapFields.getLeft());
+ assertEquals(Arrays.asList(getField("begin_lat"), getField("tip_history"),
getField("rider")), bootstrapFields.getRight());
+ }
+
+ private static Stream<Arguments> testMorParams() {
+ return testMorParams(false);
+ }
+
+ @ParameterizedTest
+ @MethodSource("testMorParams")
+ public void testMor(RecordMergeMode mergeMode,
+ boolean hasPrecombine,
+ boolean isProjectionCompatible,
+ boolean mergeUseRecordPosition,
+ boolean supportsParquetRowIndex,
+ boolean hasBuiltInDelete) throws IOException {
+ super.testMor(mergeMode, hasPrecombine, isProjectionCompatible,
mergeUseRecordPosition, supportsParquetRowIndex, hasBuiltInDelete);
+ }
+
+ @ParameterizedTest
+ @MethodSource("testMorParams")
+ public void testMorBootstrap(RecordMergeMode mergeMode,
+ boolean hasPrecombine,
+ boolean isProjectionCompatible,
+ boolean mergeUseRecordPosition,
+ boolean supportsParquetRowIndex,
+ boolean hasBuiltInDelete) throws IOException {
+ super.testMorBootstrap(mergeMode, hasPrecombine, isProjectionCompatible,
mergeUseRecordPosition, supportsParquetRowIndex, hasBuiltInDelete);
+ }
+
+ @Override
+ FileGroupReaderSchemaHandler createSchemaHandler(HoodieReaderContext<String>
readerContext, Schema dataSchema, Schema requestedSchema, HoodieTableConfig
hoodieTableConfig,
+ boolean
supportsParquetRowIndex) {
+ return new FileGroupReaderSchemaHandler(readerContext, dataSchema,
requestedSchema,
+ Option.empty(), hoodieTableConfig, new TypedProperties());
+ }
+}
diff --git
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestParquetRowIndexBasedSchemaHandler.java
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestParquetRowIndexBasedSchemaHandler.java
new file mode 100644
index 000000000000..22b28c7acf40
--- /dev/null
+++
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestParquetRowIndexBasedSchemaHandler.java
@@ -0,0 +1,111 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.hudi.common.table.read;
+
+import org.apache.hudi.common.config.RecordMergeMode;
+import org.apache.hudi.common.config.TypedProperties;
+import org.apache.hudi.common.engine.HoodieReaderContext;
+import org.apache.hudi.common.table.HoodieTableConfig;
+import org.apache.hudi.common.util.Option;
+import org.apache.hudi.common.util.collection.Pair;
+
+import org.apache.avro.Schema;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.Arguments;
+import org.junit.jupiter.params.provider.MethodSource;
+
+import java.io.IOException;
+import java.util.Arrays;
+import java.util.List;
+import java.util.stream.Stream;
+
+import static
org.apache.hudi.common.table.read.ParquetRowIndexBasedSchemaHandler.getPositionalMergeField;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+public class TestParquetRowIndexBasedSchemaHandler extends
SchemaHandlerTestBase {
+
+ @Test
+ public void testCowBootstrapWithPositionMerge() {
+ HoodieTableConfig hoodieTableConfig = mock(HoodieTableConfig.class);
+ when(hoodieTableConfig.populateMetaFields()).thenReturn(true);
+ HoodieReaderContext<String> readerContext =
createReaderContext(hoodieTableConfig, true, false, true, false, null);
+ Schema requestedSchema = generateProjectionSchema("begin_lat",
"tip_history", "_hoodie_record_key", "rider");
+ FileGroupReaderSchemaHandler schemaHandler =
createSchemaHandler(readerContext, DATA_SCHEMA, requestedSchema,
hoodieTableConfig, true);
+ assertTrue(readerContext.getNeedsBootstrapMerge());
+ //meta cols must go first in the required schema
+ Schema expectedRequiredSchema =
generateProjectionSchema("_hoodie_record_key", "begin_lat", "tip_history",
"rider");
+ assertEquals(expectedRequiredSchema, schemaHandler.getRequiredSchema());
+ Pair<List<Schema.Field>, List<Schema.Field>> bootstrapFields =
schemaHandler.getBootstrapRequiredFields();
+ assertEquals(Arrays.asList(getField("_hoodie_record_key"),
getPositionalMergeField()), bootstrapFields.getLeft());
+ assertEquals(Arrays.asList(getField("begin_lat"), getField("tip_history"),
getField("rider"), getPositionalMergeField()), bootstrapFields.getRight());
+
+ schemaHandler = createSchemaHandler(readerContext, DATA_SCHEMA,
DATA_COLS_ONLY_SCHEMA, hoodieTableConfig, true);
+ assertFalse(readerContext.getNeedsBootstrapMerge());
+ assertEquals(DATA_COLS_ONLY_SCHEMA, schemaHandler.getRequiredSchema());
+ bootstrapFields = schemaHandler.getBootstrapRequiredFields();
+ assertTrue(bootstrapFields.getLeft().isEmpty());
+ assertEquals(Arrays.asList(getField("begin_lat"), getField("tip_history"),
getField("rider")), bootstrapFields.getRight());
+
+ schemaHandler = createSchemaHandler(readerContext, DATA_SCHEMA,
META_COLS_ONLY_SCHEMA, hoodieTableConfig, true);
+ assertFalse(readerContext.getNeedsBootstrapMerge());
+ assertEquals(META_COLS_ONLY_SCHEMA, schemaHandler.getRequiredSchema());
+ bootstrapFields = schemaHandler.getBootstrapRequiredFields();
+ assertEquals(Arrays.asList(getField("_hoodie_commit_seqno"),
getField("_hoodie_record_key")), bootstrapFields.getLeft());
+ assertTrue(bootstrapFields.getRight().isEmpty());
+ }
+
+ private static Stream<Arguments> testMorParams() {
+ return testMorParams(true);
+ }
+
+ @ParameterizedTest
+ @MethodSource("testMorParams")
+ public void testMor(RecordMergeMode mergeMode,
+ boolean hasPrecombine,
+ boolean isProjectionCompatible,
+ boolean mergeUseRecordPosition,
+ boolean supportsParquetRowIndex,
+ boolean hasBuiltInDelete) throws IOException {
+ super.testMor(mergeMode, hasPrecombine, isProjectionCompatible,
mergeUseRecordPosition, supportsParquetRowIndex, hasBuiltInDelete);
+ }
+
+ @ParameterizedTest
+ @MethodSource("testMorParams")
+ public void testMorBootstrap(RecordMergeMode mergeMode,
+ boolean hasPrecombine,
+ boolean isProjectionCompatible,
+ boolean mergeUseRecordPosition,
+ boolean supportsParquetRowIndex,
+ boolean hasBuiltInDelete) throws IOException {
+ super.testMorBootstrap(mergeMode, hasPrecombine, isProjectionCompatible,
mergeUseRecordPosition, supportsParquetRowIndex, hasBuiltInDelete);
+ }
+
+ @Override
+ FileGroupReaderSchemaHandler createSchemaHandler(HoodieReaderContext<String>
readerContext, Schema dataSchema, Schema requestedSchema, HoodieTableConfig
hoodieTableConfig,
+ boolean
supportsParquetRowIndex) {
+ return new ParquetRowIndexBasedSchemaHandler(readerContext, dataSchema,
requestedSchema,
+ Option.empty(), hoodieTableConfig, new TypedProperties());
+ }
+}
diff --git
a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/TestPositionBasedFileGroupRecordBuffer.java
b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/TestPositionBasedFileGroupRecordBuffer.java
index 58e94e9d68fb..a9c42d28fed0 100644
---
a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/TestPositionBasedFileGroupRecordBuffer.java
+++
b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/TestPositionBasedFileGroupRecordBuffer.java
@@ -36,9 +36,10 @@ import org.apache.hudi.common.table.TableSchemaResolver;
import org.apache.hudi.common.table.log.block.HoodieDeleteBlock;
import org.apache.hudi.common.table.log.block.HoodieLogBlock;
import org.apache.hudi.common.table.read.CustomPayloadForTesting;
+import org.apache.hudi.common.table.read.FileGroupReaderSchemaHandler;
import org.apache.hudi.common.table.read.HoodieReadStats;
import org.apache.hudi.common.table.read.PositionBasedFileGroupRecordBuffer;
-import org.apache.hudi.common.table.read.PositionBasedSchemaHandler;
+import org.apache.hudi.common.table.read.ParquetRowIndexBasedSchemaHandler;
import org.apache.hudi.common.testutils.HoodieTestDataGenerator;
import org.apache.hudi.common.testutils.RawTripTestPayload;
import org.apache.hudi.common.testutils.SchemaTestUtil;
@@ -135,8 +136,9 @@ public class TestPositionBasedFileGroupRecordBuffer extends
SparkClientFunctiona
} else {
ctx.setRecordMerger(Option.empty());
}
- ctx.setSchemaHandler(new PositionBasedSchemaHandler<>(ctx, avroSchema,
avroSchema,
- Option.empty(), metaClient.getTableConfig(), new TypedProperties()));
+ ctx.setSchemaHandler(HoodieSparkUtils.gteqSpark3_5()
+ ? new ParquetRowIndexBasedSchemaHandler<>(ctx, avroSchema, avroSchema,
Option.empty(), metaClient.getTableConfig(), new TypedProperties())
+ : new FileGroupReaderSchemaHandler<>(ctx, avroSchema, avroSchema,
Option.empty(), metaClient.getTableConfig(), new TypedProperties()));
TypedProperties props = new TypedProperties();
props.put("hoodie.write.record.merge.mode", mergeMode.name());
props.setProperty(HoodieMemoryConfig.MAX_MEMORY_FOR_MERGE.key(),String.valueOf(HoodieMemoryConfig.MAX_MEMORY_FOR_MERGE.defaultValue()));