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 276133b814e [HUDI-8003] Add OverwriteWithLatestHiveRecordMerger and
restructure existing merger classes (#11649)
276133b814e is described below
commit 276133b814ef97aa427260fd2d0feb70723ae1b2
Author: Jon Vexler <[email protected]>
AuthorDate: Tue Sep 10 22:05:21 2024 -0400
[HUDI-8003] Add OverwriteWithLatestHiveRecordMerger and restructure
existing merger classes (#11649)
Co-authored-by: Jonathan Vexler <=>
---
hudi-client/hudi-java-client/pom.xml | 8 +
.../hadoop/TestHoodieFileGroupReaderOnHive.java | 318 +++++++++++++++++++++
.../hudi/testutils/ArrayWritableTestUtil.java | 194 +++++++++++++
.../testutils/HoodieJavaClientTestHarness.java | 2 +-
...rdMerger.java => DefaultSparkRecordMerger.java} | 10 +-
.../org/apache/hudi/HoodieSparkRecordMerger.java | 107 +------
...a => OverwriteWithLatestSparkRecordMerger.java} | 7 +-
.../hudi/BaseSparkInternalRowReaderContext.java | 15 +-
.../table/read/TestHoodieFileGroupReaderBase.java | 16 +-
...ordMerger.java => DefaultHiveRecordMerger.java} | 10 +-
.../hudi/hadoop/HiveHoodieReaderContext.java | 66 ++---
.../HoodieFileGroupReaderBasedRecordReader.java | 42 ++-
.../apache/hudi/hadoop/HoodieHiveRecordMerger.java | 51 +---
.../hudi/hadoop/HoodieParquetInputFormat.java | 8 +-
.../OverwriteWithLatestHiveRecordMerger.java | 8 +-
.../hadoop/utils/HoodieArrayWritableAvroUtils.java | 12 -
.../utils/HoodieRealtimeRecordReaderUtils.java | 2 +-
.../hudi/hadoop/utils/ObjectInspectorCache.java | 16 +-
...odieSparkValidateDuplicateKeyRecordMerger.scala | 11 +-
.../hudi/TestHoodieMergeHandleWithSparkMerger.java | 10 +-
...stHoodiePositionBasedFileGroupRecordBuffer.java | 7 +-
.../read/TestHoodieFileGroupReaderOnSpark.scala | 11 +-
.../TestSpark35RecordPositionMetadataColumn.scala | 2 +-
.../apache/hudi/functional/CommonOptionUtils.scala | 4 +-
.../apache/hudi/functional/TestCOWDataSource.scala | 2 +-
.../functional/TestDataSourceForBootstrap.scala | 4 +-
.../TestHoodieMultipleBaseFileFormat.scala | 4 +-
.../apache/hudi/functional/TestMORDataSource.scala | 8 +-
.../TestSparkDataSourceDAGExecution.scala | 4 +-
.../ReadAndWriteWithoutAvroBenchmark.scala | 12 +-
.../sql/hudi/common/HoodieSparkSqlTestBase.scala | 6 +-
.../apache/spark/sql/hudi/ddl/TestSpark3DDL.scala | 4 +-
.../deltastreamer/TestHoodieDeltaStreamer.java | 4 +-
33 files changed, 688 insertions(+), 297 deletions(-)
diff --git a/hudi-client/hudi-java-client/pom.xml
b/hudi-client/hudi-java-client/pom.xml
index c6b36bd6bed..95bb71e54e4 100644
--- a/hudi-client/hudi-java-client/pom.xml
+++ b/hudi-client/hudi-java-client/pom.xml
@@ -69,6 +69,14 @@
<type>test-jar</type>
<scope>test</scope>
</dependency>
+ <dependency>
+ <groupId>org.apache.hudi</groupId>
+ <artifactId>hudi-io</artifactId>
+ <version>${project.version}</version>
+ <classifier>tests</classifier>
+ <type>test-jar</type>
+ <scope>test</scope>
+ </dependency>
<dependency>
<groupId>org.apache.hudi</groupId>
<artifactId>hudi-hadoop-common</artifactId>
diff --git
a/hudi-client/hudi-java-client/src/test/java/org/apache/hudi/hadoop/TestHoodieFileGroupReaderOnHive.java
b/hudi-client/hudi-java-client/src/test/java/org/apache/hudi/hadoop/TestHoodieFileGroupReaderOnHive.java
new file mode 100644
index 00000000000..b29668d0d53
--- /dev/null
+++
b/hudi-client/hudi-java-client/src/test/java/org/apache/hudi/hadoop/TestHoodieFileGroupReaderOnHive.java
@@ -0,0 +1,318 @@
+/*
+ * 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.hadoop;
+
+import org.apache.hudi.avro.HoodieAvroUtils;
+import org.apache.hudi.client.HoodieJavaWriteClient;
+import org.apache.hudi.client.common.HoodieJavaEngineContext;
+import org.apache.hudi.common.config.HoodieMemoryConfig;
+import org.apache.hudi.common.config.HoodieReaderConfig;
+import org.apache.hudi.common.config.RecordMergeMode;
+import org.apache.hudi.common.engine.EngineType;
+import org.apache.hudi.common.engine.HoodieReaderContext;
+import org.apache.hudi.common.model.DefaultHoodieRecordPayload;
+import org.apache.hudi.common.model.HoodieRecord;
+import org.apache.hudi.common.model.HoodieRecordMerger;
+import org.apache.hudi.common.model.OverwriteWithLatestAvroPayload;
+import org.apache.hudi.common.table.HoodieTableMetaClient;
+import org.apache.hudi.common.table.read.CustomPayloadForTesting;
+import org.apache.hudi.common.table.read.TestHoodieFileGroupReaderBase;
+import org.apache.hudi.common.testutils.HoodieTestDataGenerator;
+import org.apache.hudi.common.testutils.HoodieTestUtils;
+import org.apache.hudi.common.testutils.minicluster.HdfsTestService;
+import org.apache.hudi.config.HoodieWriteConfig;
+import org.apache.hudi.hadoop.hive.HoodieCombineHiveInputFormat;
+import org.apache.hudi.hadoop.realtime.HoodieParquetRealtimeInputFormat;
+import org.apache.hudi.testutils.ArrayWritableTestUtil;
+import org.apache.hudi.hadoop.utils.ObjectInspectorCache;
+import org.apache.hudi.storage.HoodieStorage;
+import org.apache.hudi.storage.StorageConfiguration;
+import org.apache.hudi.storage.hadoop.HoodieHadoopStorage;
+import org.apache.hudi.testutils.HoodieJavaClientTestHarness;
+
+import org.apache.avro.Schema;
+import org.apache.hadoop.conf.Configuration;
+import org.apache.hadoop.fs.FileSystem;
+import org.apache.hadoop.fs.Path;
+import org.apache.hadoop.hive.metastore.api.hive_metastoreConstants;
+import org.apache.hadoop.hive.ql.exec.Utilities;
+import org.apache.hadoop.hive.ql.exec.mr.ExecMapper;
+import org.apache.hadoop.hive.ql.io.parquet.MapredParquetInputFormat;
+import org.apache.hadoop.hive.ql.plan.MapredWork;
+import org.apache.hadoop.hive.ql.plan.PartitionDesc;
+import org.apache.hadoop.hive.ql.plan.TableDesc;
+import org.apache.hadoop.hive.serde2.ColumnProjectionUtils;
+import org.apache.hadoop.io.ArrayWritable;
+import org.apache.hadoop.io.NullWritable;
+import org.apache.hadoop.mapred.FileInputFormat;
+import org.apache.hadoop.mapred.InputSplit;
+import org.apache.hadoop.mapred.JobConf;
+import org.apache.hadoop.mapred.RecordReader;
+import org.apache.hadoop.mapred.Reporter;
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Disabled;
+
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.HashMap;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.stream.Collectors;
+
+import static org.apache.hadoop.hive.ql.exec.Utilities.HAS_MAP_WORK;
+import static org.apache.hadoop.hive.ql.exec.Utilities.MAPRED_MAPPER_CLASS;
+import static
org.apache.hudi.hadoop.HoodieFileGroupReaderBasedRecordReader.getRecordKeyField;
+import static
org.apache.hudi.hadoop.HoodieFileGroupReaderBasedRecordReader.getStoredPartitionFieldNames;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+
+public class TestHoodieFileGroupReaderOnHive extends
TestHoodieFileGroupReaderBase<ArrayWritable> {
+
+ @Override
+ @Disabled("[HUDI-8072]")
+ public void testReadLogFilesOnlyInMergeOnReadTable(RecordMergeMode
recordMergeMode, String logDataBlockFormat) throws Exception {
+ }
+
+ private static final String PARTITION_COLUMN = "datestr";
+ private static JobConf baseJobConf;
+ private static HdfsTestService hdfsTestService;
+ private static HoodieStorage storage;
+ private static FileSystem fs;
+ private static StorageConfiguration<Configuration> storageConf;
+
+ //currently always true. If we ever have a test with a nonpartitioned table,
the usages of this should be tied together
+ private static final boolean USE_FAKE_PARTITION = true;
+
+ @BeforeAll
+ public static void setUpClass() throws IOException, InterruptedException {
+ // Append is not supported in LocalFileSystem. HDFS needs to be setup.
+ hdfsTestService = new HdfsTestService();
+ fs = hdfsTestService.start(true).getFileSystem();
+ storageConf = HoodieTestUtils.getDefaultStorageConf();
+ baseJobConf = new JobConf(storageConf.unwrap());
+ baseJobConf.set(HoodieMemoryConfig.MAX_DFS_STREAM_BUFFER_SIZE.key(),
String.valueOf(1024 * 1024));
+ fs.setConf(baseJobConf);
+ storage = new HoodieHadoopStorage(fs);
+ }
+
+ @AfterAll
+ public static void tearDownClass() throws IOException {
+ hdfsTestService.stop();
+ if (fs != null) {
+ fs.close();
+ storage.close();
+ }
+ }
+
+ @Override
+ public StorageConfiguration<?> getStorageConf() {
+ return storageConf;
+ }
+
+ @Override
+ public String getBasePath() {
+ return tempDir.toAbsolutePath() + "/myTable";
+ }
+
+ @Override
+ public HoodieReaderContext<ArrayWritable> getHoodieReaderContext(String
tablePath, Schema avroSchema, StorageConfiguration<?> storageConf) {
+ HoodieFileGroupReaderBasedRecordReader.HiveReaderCreator readerCreator =
(inputSplit, jobConf) -> new
MapredParquetInputFormat().getRecordReader(inputSplit, jobConf, null);
+ HoodieTableMetaClient metaClient =
HoodieTableMetaClient.builder().setConf(storageConf).setBasePath(tablePath).build();
+ JobConf jobConf = new JobConf(storageConf.unwrapAs(Configuration.class));
+ setupJobconf(jobConf);
+ return new HiveHoodieReaderContext(readerCreator,
getRecordKeyField(metaClient),
+ getStoredPartitionFieldNames(new
JobConf(storageConf.unwrapAs(Configuration.class)), avroSchema),
+ new ObjectInspectorCache(avroSchema, jobConf));
+ }
+
+ @Override
+ public String getRecordPayloadForMergeMode(RecordMergeMode mergeMode) {
+ switch (mergeMode) {
+ case EVENT_TIME_ORDERING:
+ return DefaultHoodieRecordPayload.class.getName();
+ case OVERWRITE_WITH_LATEST:
+ return OverwriteWithLatestAvroPayload.class.getName();
+ case CUSTOM:
+ default:
+ return CustomPayloadForTesting.class.getName();
+ }
+ }
+
+ @Override
+ public void commitToTable(List<HoodieRecord> recordList, String operation,
Map<String, String> writeConfigs) {
+ HoodieWriteConfig writeConfig = HoodieWriteConfig.newBuilder()
+ .withEngineType(EngineType.JAVA)
+ .withEmbeddedTimelineServerEnabled(false)
+ .withProps(writeConfigs)
+ .withPath(getBasePath())
+ .withSchema(HoodieTestDataGenerator.TRIP_EXAMPLE_SCHEMA)
+ .build();
+
+ HoodieJavaClientTestHarness.TestJavaTaskContextSupplier
taskContextSupplier = new
HoodieJavaClientTestHarness.TestJavaTaskContextSupplier();
+ HoodieJavaEngineContext context = new
HoodieJavaEngineContext(getStorageConf(), taskContextSupplier);
+ //init table if not exists
+ Path basePath = new Path(getBasePath());
+ try {
+ try (FileSystem lfs = basePath.getFileSystem(baseJobConf)) {
+ boolean basepathExists = lfs.exists(basePath);
+ boolean operationIsInsert = operation.equalsIgnoreCase("insert");
+ if (!basepathExists || operationIsInsert) {
+ if (basepathExists) {
+ lfs.delete(new Path(getBasePath()), true);
+ }
+ String recordMergerStrategy = "";
+ if
(RecordMergeMode.valueOf(writeConfigs.get("hoodie.record.merge.mode")).equals(RecordMergeMode.OVERWRITE_WITH_LATEST))
{
+ recordMergerStrategy =
HoodieRecordMerger.OVERWRITE_MERGER_STRATEGY_UUID;
+ } else if
(RecordMergeMode.valueOf(writeConfigs.get("hoodie.record.merge.mode")).equals(RecordMergeMode.EVENT_TIME_ORDERING))
{
+ recordMergerStrategy =
HoodieRecordMerger.DEFAULT_MERGER_STRATEGY_UUID;
+ } else if
(RecordMergeMode.valueOf(writeConfigs.get("hoodie.record.merge.mode")).equals(RecordMergeMode.CUSTOM))
{
+ //match the behavior of spark for now, but this should be a config
+ recordMergerStrategy =
HoodieRecordMerger.DEFAULT_MERGER_STRATEGY_UUID;
+ }
+ Map<String, Object> initConfigs = new HashMap<>(writeConfigs);
+ HoodieTableMetaClient.withPropertyBuilder()
+
.setTableType(writeConfigs.getOrDefault("hoodie.datasource.write.table.type",
"MERGE_ON_READ"))
+ .setTableName(writeConfigs.get("hoodie.table.name"))
+
.setPartitionFields(writeConfigs.getOrDefault("hoodie.datasource.write.partitionpath.field",
""))
+
.setRecordMergeMode(RecordMergeMode.valueOf(writeConfigs.get("hoodie.record.merge.mode")))
+ .setRecordMergerStrategy(recordMergerStrategy)
+ .set(initConfigs).initTable(storageConf, getBasePath());
+ }
+ }
+ } catch (IOException e) {
+ throw new RuntimeException(e);
+ }
+
+ HoodieJavaWriteClient writeClient = new HoodieJavaWriteClient(context,
writeConfig);
+ String instantTime = writeClient.createNewInstantTime();
+ writeClient.startCommitWithTime(instantTime);
+ if (operation.toLowerCase().equals("insert")) {
+ writeClient.insert(recordList, instantTime);
+ } else {
+ writeClient.upsert(recordList, instantTime);
+ }
+ }
+
+ @Override
+ public void validateRecordsInFileGroup(String tablePath, List<ArrayWritable>
actualRecordList, Schema schema, String fileGroupId) {
+
assertEquals(HoodieAvroUtils.addMetadataFields(HoodieTestDataGenerator.AVRO_SCHEMA),
schema);
+ try {
+ //prepare fg reader records to be compared to the baseline reader
+ HoodieReaderContext<ArrayWritable> readerContext =
getHoodieReaderContext(tablePath, schema, storageConf);
+ Map<String, ArrayWritable> recordMap = new HashMap<>();
+ for (ArrayWritable record : actualRecordList) {
+ recordMap.put(readerContext.getRecordKey(record, schema), record);
+ }
+
+ RecordReader<NullWritable, ArrayWritable> reader =
createRecordReader(tablePath);
+ // use reader to read log file.
+ NullWritable key = reader.createKey();
+ ArrayWritable value = reader.createValue();
+ while (reader.next(key, value)) {
+ if (readerContext.getValue(value, schema,
HoodieRecord.FILENAME_METADATA_FIELD).toString().contains(fileGroupId)) {
+ //only evaluate records from the specified filegroup. Maybe there is
a way to get
+ //hive to do this?
+ ArrayWritable compVal =
recordMap.remove(readerContext.getRecordKey(value, schema));
+ assertNotNull(compVal);
+ ArrayWritableTestUtil.assertArrayWritableEqual(schema, value,
compVal, USE_FAKE_PARTITION);
+ }
+ key = reader.createKey();
+ value = reader.createValue();
+ }
+ reader.close();
+ assertEquals(0, recordMap.size());
+ } catch (IOException e) {
+ throw new RuntimeException(e);
+ }
+ }
+
+ private RecordReader<NullWritable, ArrayWritable> createRecordReader(String
tablePath) throws IOException {
+ JobConf jobConf = new JobConf(baseJobConf);
+ jobConf.set(HoodieReaderConfig.FILE_GROUP_READER_ENABLED.key(), "false");
+
+ TableDesc tblDesc = Utilities.defaultTd;
+ // Set the input format
+ tblDesc.setInputFileFormatClass(HoodieParquetRealtimeInputFormat.class);
+ LinkedHashMap<Path, PartitionDesc> pt = new LinkedHashMap<>();
+ LinkedHashMap<Path, ArrayList<String>> talias = new LinkedHashMap<>();
+
+ PartitionDesc partDesc = new PartitionDesc(tblDesc, null);
+
+ pt.put(new Path(tablePath), partDesc);
+
+ ArrayList<String> arrayList = new ArrayList<>();
+ arrayList.add(tablePath);
+ talias.put(new Path(tablePath), arrayList);
+
+ MapredWork mrwork = new MapredWork();
+ mrwork.getMapWork().setPathToPartitionInfo(pt);
+ mrwork.getMapWork().setPathToAliases(talias);
+
+ Path mapWorkPath = new Path(tablePath);
+ Utilities.setMapRedWork(jobConf, mrwork, mapWorkPath);
+
+ // Add three partition path to InputPaths
+ Path[] partitionDirArray = new
Path[HoodieTestDataGenerator.DEFAULT_PARTITION_PATHS.length];
+ Arrays.stream(HoodieTestDataGenerator.DEFAULT_PARTITION_PATHS).map(s ->
new Path(tablePath, s)).collect(Collectors.toList()).toArray(partitionDirArray);
+ FileInputFormat.setInputPaths(jobConf, partitionDirArray);
+ jobConf.set(HAS_MAP_WORK, "true");
+ // The following config tells Hive to choose ExecMapper to read the
MAP_WORK
+ jobConf.set(MAPRED_MAPPER_CLASS, ExecMapper.class.getName());
+ // setting the split size to be 3 to create one split for 3 file groups
+
jobConf.set(org.apache.hadoop.mapreduce.lib.input.FileInputFormat.SPLIT_MAXSIZE,
"128000000");
+ setupJobconf(jobConf);
+
+ HoodieCombineHiveInputFormat combineHiveInputFormat = new
HoodieCombineHiveInputFormat();
+ InputSplit[] splits = combineHiveInputFormat.getSplits(jobConf, 1);
+
+ assertEquals(1, splits.length);
+ return combineHiveInputFormat.getRecordReader(splits[0], jobConf,
Reporter.NULL);
+ }
+
+ private void setupJobconf(JobConf jobConf) {
+ Schema schema =
HoodieAvroUtils.addMetadataFields(HoodieTestDataGenerator.AVRO_SCHEMA);
+ List<Schema.Field> fields = schema.getFields();
+ setHiveColumnNameProps(fields, jobConf, USE_FAKE_PARTITION);
+ jobConf.set("columns.types","string,string,string,string,string," +
HoodieTestDataGenerator.TRIP_HIVE_COLUMN_TYPES + ",string");
+ }
+
+ private void setHiveColumnNameProps(List<Schema.Field> fields, JobConf
jobConf, boolean isPartitioned) {
+ String names =
fields.stream().map(Schema.Field::name).collect(Collectors.joining(","));
+ String positions = fields.stream().map(f ->
String.valueOf(f.pos())).collect(Collectors.joining(","));
+ jobConf.set(ColumnProjectionUtils.READ_COLUMN_NAMES_CONF_STR, names);
+ jobConf.set(ColumnProjectionUtils.READ_COLUMN_IDS_CONF_STR, positions);
+
+ String hiveOrderedColumnNames = fields.stream().filter(field ->
!field.name().equalsIgnoreCase(PARTITION_COLUMN))
+ .map(Schema.Field::name).collect(Collectors.joining(","));
+ if (isPartitioned) {
+ hiveOrderedColumnNames += "," + PARTITION_COLUMN;
+ jobConf.set(hive_metastoreConstants.META_TABLE_PARTITION_COLUMNS,
PARTITION_COLUMN);
+ }
+ jobConf.set(hive_metastoreConstants.META_TABLE_COLUMNS,
hiveOrderedColumnNames);
+ }
+
+ @Override
+ public Comparable getComparableUTF8String(String value) {
+ return value;
+ }
+}
diff --git
a/hudi-client/hudi-java-client/src/test/java/org/apache/hudi/testutils/ArrayWritableTestUtil.java
b/hudi-client/hudi-java-client/src/test/java/org/apache/hudi/testutils/ArrayWritableTestUtil.java
new file mode 100644
index 00000000000..bbca2671127
--- /dev/null
+++
b/hudi-client/hudi-java-client/src/test/java/org/apache/hudi/testutils/ArrayWritableTestUtil.java
@@ -0,0 +1,194 @@
+/*
+ * 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.testutils;
+
+import org.apache.avro.LogicalTypes;
+import org.apache.avro.Schema;
+import org.apache.hadoop.hive.serde2.io.DateWritable;
+import org.apache.hadoop.hive.serde2.io.HiveDecimalWritable;
+import org.apache.hadoop.io.ArrayWritable;
+import org.apache.hadoop.io.BooleanWritable;
+import org.apache.hadoop.io.BytesWritable;
+import org.apache.hadoop.io.DoubleWritable;
+import org.apache.hadoop.io.FloatWritable;
+import org.apache.hadoop.io.IntWritable;
+import org.apache.hadoop.io.LongWritable;
+import org.apache.hadoop.io.NullWritable;
+import org.apache.hadoop.io.Text;
+import org.apache.hadoop.io.Writable;
+
+import java.util.HashMap;
+import java.util.Map;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+
+public class ArrayWritableTestUtil {
+ public static void assertArrayWritableEqual(Schema schema, ArrayWritable
expected, ArrayWritable actual, boolean isPartitioned) {
+ assertArrayWritableEqualInternal(schema, expected, actual, isPartitioned);
+ }
+
+ private static void assertArrayWritableEqualInternal(Schema schema, Writable
expected, Writable actual, boolean ignoreOneExtraCol) {
+ switch (schema.getType()) {
+ case RECORD: {
+ assertInstanceOf(ArrayWritable.class, expected);
+ assertInstanceOf(ArrayWritable.class, actual);
+ //adjust for fake partition
+ int expectedLen = ((ArrayWritable) expected).get().length -
(ignoreOneExtraCol ? 1 : 0);
+ int actualLen = ((ArrayWritable) actual).get().length;
+ assertEquals(expectedLen, actualLen);
+ assertEquals(schema.getFields().size(), expectedLen);
+ for (Schema.Field field : schema.getFields()) {
+ assertArrayWritableEqualInternal(field.schema(), ((ArrayWritable)
expected).get()[field.pos()], ((ArrayWritable) actual).get()[field.pos()],
false);
+ }
+ break;
+ }
+ case ARRAY: {
+ assertInstanceOf(ArrayWritable.class, expected);
+ assertInstanceOf(ArrayWritable.class, actual);
+ int expectedLen = ((ArrayWritable) expected).get().length;
+ int actualLen = ((ArrayWritable) actual).get().length;
+ assertEquals(expectedLen, actualLen);
+ for (int i = 0; i < expectedLen; i++) {
+ assertArrayWritableEqualInternal(schema.getElementType(),
((ArrayWritable) expected).get()[i], ((ArrayWritable) expected).get()[i],
false);
+ }
+ break;
+ }
+ case MAP: {
+ assertInstanceOf(ArrayWritable.class, expected);
+ assertInstanceOf(ArrayWritable.class, actual);
+ int expectedLen = ((ArrayWritable) expected).get().length;
+ int actualLen = ((ArrayWritable) actual).get().length;
+ assertEquals(expectedLen, actualLen);
+ Map<Writable, Writable> expectedMap = new HashMap<>(expectedLen);
+ Map<Writable, Writable> actualMap = new HashMap<>(actualLen);
+ for (int i = 0; i < expectedLen; i++) {
+ Writable expectedKV = ((ArrayWritable) expected).get()[i];
+ assertInstanceOf(ArrayWritable.class, expectedKV);
+ assertEquals(2, ((ArrayWritable) expectedKV).get().length);
+ expectedMap.put(((ArrayWritable) expectedKV).get()[0],
((ArrayWritable) expectedKV).get()[1]);
+ Writable actualKV = ((ArrayWritable) actual).get()[i];
+ assertInstanceOf(ArrayWritable.class, actualKV);
+ assertEquals(2, ((ArrayWritable) actualKV).get().length);
+ actualMap.put(((ArrayWritable) actualKV).get()[0], ((ArrayWritable)
actualKV).get()[1]);
+ }
+
+ for (Writable key : expectedMap.keySet()) {
+ Writable expectedValue = expectedMap.get(key);
+ assertNotNull(expectedValue);
+ Writable actualValue = actualMap.remove(key);
+ assertNotNull(actualValue);
+ assertArrayWritableEqualInternal(schema.getValueType(),
expectedValue, actualValue, false);
+ }
+ assertEquals(0, actualMap.size());
+ break;
+ }
+ case UNION:
+ if (schema.getTypes().size() == 2
+ && schema.getTypes().get(0).getType() == Schema.Type.NULL) {
+ assertArrayWritableEqualInternal(schema.getTypes().get(1), expected,
actual, false);
+ } else if (schema.getTypes().size() == 2
+ && schema.getTypes().get(1).getType() == Schema.Type.NULL) {
+ assertArrayWritableEqualInternal(schema.getTypes().get(0), expected,
actual, false);
+ } else if (schema.getTypes().size() == 1) {
+ assertArrayWritableEqualInternal(schema.getTypes().get(0), expected,
actual, false);
+ } else {
+ throw new IllegalStateException("Union has more than 2 types or one
type is not null: " + schema);
+ }
+ break;
+
+ default:
+ assertWritablePrimaryType(schema, expected, actual);
+ }
+ }
+
+ private static void assertWritablePrimaryType(Schema schema, Writable
expected, Writable actual) {
+ switch (schema.getType()) {
+ case NULL:
+ assertInstanceOf(NullWritable.class, expected);
+ assertInstanceOf(NullWritable.class, actual);
+ assertEquals(expected, actual);
+ break;
+
+ case BOOLEAN:
+ assertInstanceOf(BooleanWritable.class, expected);
+ assertInstanceOf(BooleanWritable.class, actual);
+ assertEquals(expected, actual);
+ break;
+
+ case INT:
+ if (schema.getLogicalType() instanceof LogicalTypes.Date) {
+ assertInstanceOf(DateWritable.class, expected);
+ assertInstanceOf(DateWritable.class, actual);
+ } else {
+ assertInstanceOf(IntWritable.class, expected);
+ assertInstanceOf(IntWritable.class, actual);
+ }
+ assertEquals(expected, actual);
+ break;
+
+ case LONG:
+ assertInstanceOf(LongWritable.class, expected);
+ assertInstanceOf(LongWritable.class, actual);
+ assertEquals(expected, actual);
+ break;
+
+ case FLOAT:
+ assertInstanceOf(FloatWritable.class, expected);
+ assertInstanceOf(FloatWritable.class, actual);
+ assertEquals(expected, actual);
+ break;
+
+ case DOUBLE:
+ assertInstanceOf(DoubleWritable.class, expected);
+ assertInstanceOf(DoubleWritable.class, actual);
+ assertEquals(expected, actual);
+ break;
+
+ case BYTES:
+ case ENUM:
+ assertInstanceOf(BytesWritable.class, expected);
+ assertInstanceOf(BytesWritable.class, actual);
+ assertEquals(expected, actual);
+ break;
+
+ case STRING:
+ assertInstanceOf(Text.class, expected);
+ assertInstanceOf(Text.class, actual);
+ assertEquals(expected, actual);
+ break;
+
+ case FIXED:
+ if (schema.getLogicalType() instanceof LogicalTypes.Decimal) {
+ assertInstanceOf(HiveDecimalWritable.class, expected);
+ assertInstanceOf(HiveDecimalWritable.class, actual);
+ } else {
+ assertEquals(expected.getClass(), actual.getClass());
+ }
+ assertEquals(expected, actual);
+ break;
+
+ default:
+ assertEquals(expected.getClass(), actual.getClass());
+ assertEquals(expected, actual);
+ }
+ }
+}
diff --git
a/hudi-client/hudi-java-client/src/test/java/org/apache/hudi/testutils/HoodieJavaClientTestHarness.java
b/hudi-client/hudi-java-client/src/test/java/org/apache/hudi/testutils/HoodieJavaClientTestHarness.java
index 9195ba9f165..eeebd7960a3 100644
---
a/hudi-client/hudi-java-client/src/test/java/org/apache/hudi/testutils/HoodieJavaClientTestHarness.java
+++
b/hudi-client/hudi-java-client/src/test/java/org/apache/hudi/testutils/HoodieJavaClientTestHarness.java
@@ -154,7 +154,7 @@ public abstract class HoodieJavaClientTestHarness extends
HoodieWriterClientTest
cleanupExecutorService();
}
- public class TestJavaTaskContextSupplier extends TaskContextSupplier {
+ public static class TestJavaTaskContextSupplier extends TaskContextSupplier {
int partitionId = 0;
int stageId = 0;
long attemptId = 0;
diff --git
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/HoodieSparkRecordMerger.java
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/DefaultSparkRecordMerger.java
similarity index 96%
copy from
hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/HoodieSparkRecordMerger.java
copy to
hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/DefaultSparkRecordMerger.java
index a488e7174e5..1b0197a3330 100644
---
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/HoodieSparkRecordMerger.java
+++
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/DefaultSparkRecordMerger.java
@@ -33,7 +33,10 @@ import org.apache.avro.Schema;
import java.io.IOException;
-public class HoodieSparkRecordMerger implements HoodieRecordMerger {
+/**
+ * Record merger for spark that implements the default merger strategy
+ */
+public class DefaultSparkRecordMerger extends HoodieSparkRecordMerger {
@Override
public String getMergingStrategy() {
@@ -115,9 +118,4 @@ public class HoodieSparkRecordMerger implements
HoodieRecordMerger {
(HoodieSparkRecord) older, oldSchema, (HoodieSparkRecord) newer,
newSchema, readerSchema, props));
}
}
-
- @Override
- public HoodieRecordType getRecordType() {
- return HoodieRecordType.SPARK;
- }
}
\ No newline at end of file
diff --git
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/HoodieSparkRecordMerger.java
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/HoodieSparkRecordMerger.java
index a488e7174e5..18fa76004af 100644
---
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/HoodieSparkRecordMerger.java
+++
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/HoodieSparkRecordMerger.java
@@ -19,105 +19,24 @@
package org.apache.hudi;
-import org.apache.hudi.common.config.TypedProperties;
import org.apache.hudi.common.model.HoodieRecord;
-import org.apache.hudi.common.model.HoodieRecord.HoodieRecordType;
import org.apache.hudi.common.model.HoodieRecordMerger;
-import org.apache.hudi.common.model.HoodieSparkRecord;
-import org.apache.hudi.common.util.Option;
-import org.apache.hudi.common.util.ValidationUtils;
-import org.apache.hudi.common.util.collection.Pair;
-import org.apache.hudi.merge.SparkRecordMergingUtils;
-
-import org.apache.avro.Schema;
-
-import java.io.IOException;
-
-public class HoodieSparkRecordMerger implements HoodieRecordMerger {
+import org.apache.hudi.exception.HoodieException;
+public abstract class HoodieSparkRecordMerger implements HoodieRecordMerger {
@Override
- public String getMergingStrategy() {
- return HoodieRecordMerger.DEFAULT_MERGER_STRATEGY_UUID;
+ public HoodieRecord.HoodieRecordType getRecordType() {
+ return HoodieRecord.HoodieRecordType.SPARK;
}
- @Override
- public Option<Pair<HoodieRecord, Schema>> merge(HoodieRecord older, Schema
oldSchema, HoodieRecord newer, Schema newSchema, TypedProperties props) throws
IOException {
- ValidationUtils.checkArgument(older.getRecordType() ==
HoodieRecordType.SPARK);
- ValidationUtils.checkArgument(newer.getRecordType() ==
HoodieRecordType.SPARK);
-
- if (newer instanceof HoodieSparkRecord) {
- HoodieSparkRecord newSparkRecord = (HoodieSparkRecord) newer;
- if (newSparkRecord.isDeleted()) {
- // Delete record
- return Option.empty();
- }
- } else {
- if (newer.getData() == null) {
- // Delete record
- return Option.empty();
- }
+ static HoodieRecordMerger getRecordMerger(String mergerStrategy) {
+ switch (mergerStrategy) {
+ case DEFAULT_MERGER_STRATEGY_UUID:
+ return new DefaultSparkRecordMerger();
+ case OVERWRITE_MERGER_STRATEGY_UUID:
+ return new OverwriteWithLatestSparkRecordMerger();
+ default:
+ throw new HoodieException("This merger strategy UUID is not supported:
" + mergerStrategy);
}
-
- if (older instanceof HoodieSparkRecord) {
- HoodieSparkRecord oldSparkRecord = (HoodieSparkRecord) older;
- if (oldSparkRecord.isDeleted()) {
- // use natural order for delete record
- return Option.of(Pair.of(newer, newSchema));
- }
- } else {
- if (older.getData() == null) {
- // use natural order for delete record
- return Option.of(Pair.of(newer, newSchema));
- }
- }
- if (older.getOrderingValue(oldSchema,
props).compareTo(newer.getOrderingValue(newSchema, props)) > 0) {
- return Option.of(Pair.of(older, oldSchema));
- } else {
- return Option.of(Pair.of(newer, newSchema));
- }
- }
-
- @Override
- public Option<Pair<HoodieRecord, Schema>> partialMerge(HoodieRecord older,
Schema oldSchema, HoodieRecord newer, Schema newSchema, Schema readerSchema,
TypedProperties props) throws IOException {
- ValidationUtils.checkArgument(older.getRecordType() ==
HoodieRecordType.SPARK);
- ValidationUtils.checkArgument(newer.getRecordType() ==
HoodieRecordType.SPARK);
-
- if (newer instanceof HoodieSparkRecord) {
- HoodieSparkRecord newSparkRecord = (HoodieSparkRecord) newer;
- if (newSparkRecord.isDeleted()) {
- // Delete record
- return Option.empty();
- }
- } else {
- if (newer.getData() == null) {
- // Delete record
- return Option.empty();
- }
- }
-
- if (older instanceof HoodieSparkRecord) {
- HoodieSparkRecord oldSparkRecord = (HoodieSparkRecord) older;
- if (oldSparkRecord.isDeleted()) {
- // use natural order for delete record
- return Option.of(Pair.of(newer, newSchema));
- }
- } else {
- if (older.getData() == null) {
- // use natural order for delete record
- return Option.of(Pair.of(newer, newSchema));
- }
- }
- if (older.getOrderingValue(oldSchema,
props).compareTo(newer.getOrderingValue(newSchema, props)) > 0) {
- return Option.of(SparkRecordMergingUtils.mergePartialRecords(
- (HoodieSparkRecord) newer, newSchema, (HoodieSparkRecord) older,
oldSchema, readerSchema, props));
- } else {
- return Option.of(SparkRecordMergingUtils.mergePartialRecords(
- (HoodieSparkRecord) older, oldSchema, (HoodieSparkRecord) newer,
newSchema, readerSchema, props));
- }
- }
-
- @Override
- public HoodieRecordType getRecordType() {
- return HoodieRecordType.SPARK;
}
-}
\ No newline at end of file
+}
diff --git
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/OverwriteWithLatestSparkMerger.java
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/OverwriteWithLatestSparkRecordMerger.java
similarity index 89%
copy from
hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/OverwriteWithLatestSparkMerger.java
copy to
hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/OverwriteWithLatestSparkRecordMerger.java
index 611f045f645..4cd052e8356 100644
---
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/OverwriteWithLatestSparkMerger.java
+++
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/OverwriteWithLatestSparkRecordMerger.java
@@ -31,7 +31,7 @@ import java.io.IOException;
/**
* Spark merger that always chooses the newer record
*/
-public class OverwriteWithLatestSparkMerger extends HoodieSparkRecordMerger {
+public class OverwriteWithLatestSparkRecordMerger extends
HoodieSparkRecordMerger {
@Override
public String getMergingStrategy() {
@@ -43,4 +43,9 @@ public class OverwriteWithLatestSparkMerger extends
HoodieSparkRecordMerger {
return Option.of(Pair.of(newer, newSchema));
}
+ @Override
+ public HoodieRecord.HoodieRecordType getRecordType() {
+ return null;
+ }
+
}
diff --git
a/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRowReaderContext.java
b/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRowReaderContext.java
index 72b6276b458..817e83fce44 100644
---
a/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRowReaderContext.java
+++
b/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/BaseSparkInternalRowReaderContext.java
@@ -28,7 +28,9 @@ import org.apache.hudi.common.model.HoodieRecordMerger;
import org.apache.hudi.common.model.HoodieSparkRecord;
import org.apache.hudi.common.util.ConfigUtils;
import org.apache.hudi.common.util.Option;
-import org.apache.hudi.exception.HoodieException;
+import org.apache.hudi.storage.HoodieStorage;
+import org.apache.hudi.storage.HoodieStorageUtils;
+import org.apache.hudi.storage.StorageConfiguration;
import org.apache.avro.Schema;
import org.apache.spark.sql.HoodieInternalRowUtils;
@@ -44,8 +46,6 @@ import java.util.function.UnaryOperator;
import scala.Function1;
import static
org.apache.hudi.common.model.HoodieRecord.RECORD_KEY_METADATA_FIELD;
-import static
org.apache.hudi.common.model.HoodieRecordMerger.DEFAULT_MERGER_STRATEGY_UUID;
-import static
org.apache.hudi.common.model.HoodieRecordMerger.OVERWRITE_MERGER_STRATEGY_UUID;
import static org.apache.spark.sql.HoodieInternalRowUtils.getCachedSchema;
/**
@@ -56,14 +56,7 @@ public abstract class BaseSparkInternalRowReaderContext
extends HoodieReaderCont
@Override
public HoodieRecordMerger getRecordMerger(String mergerStrategy) {
- switch (mergerStrategy) {
- case DEFAULT_MERGER_STRATEGY_UUID:
- return new HoodieSparkRecordMerger();
- case OVERWRITE_MERGER_STRATEGY_UUID:
- return new OverwriteWithLatestSparkMerger();
- default:
- throw new HoodieException("The merger strategy UUID is not supported:
" + mergerStrategy);
- }
+ return HoodieSparkRecordMerger.getRecordMerger(mergerStrategy);
}
@Override
diff --git
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestHoodieFileGroupReaderBase.java
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestHoodieFileGroupReaderBase.java
index a53583e63da..0a6b8f3937d 100644
---
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestHoodieFileGroupReaderBase.java
+++
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/TestHoodieFileGroupReaderBase.java
@@ -29,6 +29,7 @@ import org.apache.hudi.common.engine.HoodieEngineContext;
import org.apache.hudi.common.engine.HoodieLocalEngineContext;
import org.apache.hudi.common.engine.HoodieReaderContext;
import org.apache.hudi.common.model.FileSlice;
+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.HoodieTableMetaClient;
@@ -68,7 +69,6 @@ import static
org.apache.hudi.common.table.HoodieTableConfig.RECORD_MERGER_STRAT
import static org.apache.hudi.common.table.HoodieTableConfig.RECORD_MERGE_MODE;
import static
org.apache.hudi.common.table.read.HoodieBaseFileGroupRecordBuffer.compareTo;
import static
org.apache.hudi.common.testutils.HoodieTestUtils.getLogFileListFromFileSlice;
-import static
org.apache.hudi.common.testutils.RawTripTestPayload.recordsToStrings;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.params.provider.Arguments.arguments;
@@ -90,7 +90,7 @@ public abstract class TestHoodieFileGroupReaderBase<T> {
public abstract String getRecordPayloadForMergeMode(RecordMergeMode
mergeMode);
- public abstract void commitToTable(List<String> recordList, String operation,
+ public abstract void commitToTable(List<HoodieRecord> recordList, String
operation,
Map<String, String> writeConfigs);
public abstract void validateRecordsInFileGroup(String tablePath,
@@ -105,7 +105,7 @@ public abstract class TestHoodieFileGroupReaderBase<T> {
Map<String, String> writeConfigs = new
HashMap<>(getCommonConfigs(RecordMergeMode.EVENT_TIME_ORDERING));
// Prepare a table for initializing reader context
try (HoodieTestDataGenerator dataGen = new
HoodieTestDataGenerator(0xDEEF)) {
- commitToTable(recordsToStrings(dataGen.generateInserts("001", 1)),
BULK_INSERT.value(), writeConfigs);
+ commitToTable(dataGen.generateInserts("001", 1), BULK_INSERT.value(),
writeConfigs);
}
StorageConfiguration<?> storageConf = getStorageConf();
String tablePath = getBasePath();
@@ -184,17 +184,17 @@ public abstract class TestHoodieFileGroupReaderBase<T> {
try (HoodieTestDataGenerator dataGen = new
HoodieTestDataGenerator(0xDEEF)) {
// One commit; reading one file group containing a base file only
- commitToTable(recordsToStrings(dataGen.generateInserts("001", 100)),
INSERT.value(), writeConfigs);
+ commitToTable(dataGen.generateInserts("001", 100), INSERT.value(),
writeConfigs);
validateOutputFromFileGroupReader(
getStorageConf(), getBasePath(), dataGen.getPartitionPaths(), true,
0, recordMergeMode);
// Two commits; reading one file group containing a base file and a log
file
- commitToTable(recordsToStrings(dataGen.generateUpdates("002", 100)),
UPSERT.value(), writeConfigs);
+ commitToTable(dataGen.generateUpdates("002", 100), UPSERT.value(),
writeConfigs);
validateOutputFromFileGroupReader(
getStorageConf(), getBasePath(), dataGen.getPartitionPaths(), true,
1, recordMergeMode);
// Three commits; reading one file group containing a base file and two
log files
- commitToTable(recordsToStrings(dataGen.generateUpdates("003", 100)),
UPSERT.value(), writeConfigs);
+ commitToTable(dataGen.generateUpdates("003", 100), UPSERT.value(),
writeConfigs);
validateOutputFromFileGroupReader(
getStorageConf(), getBasePath(), dataGen.getPartitionPaths(), true,
2, recordMergeMode);
}
@@ -210,12 +210,12 @@ public abstract class TestHoodieFileGroupReaderBase<T> {
try (HoodieTestDataGenerator dataGen = new
HoodieTestDataGenerator(0xDEEF)) {
// One commit; reading one file group containing a base file only
- commitToTable(recordsToStrings(dataGen.generateInserts("001", 100)),
INSERT.value(), writeConfigs);
+ commitToTable(dataGen.generateInserts("001", 100), INSERT.value(),
writeConfigs);
validateOutputFromFileGroupReader(
getStorageConf(), getBasePath(), dataGen.getPartitionPaths(), false,
1, recordMergeMode);
// Two commits; reading one file group containing a base file and a log
file
- commitToTable(recordsToStrings(dataGen.generateUpdates("002", 100)),
UPSERT.value(), writeConfigs);
+ commitToTable(dataGen.generateUpdates("002", 100), UPSERT.value(),
writeConfigs);
validateOutputFromFileGroupReader(
getStorageConf(), getBasePath(), dataGen.getPartitionPaths(), false,
2, recordMergeMode);
}
diff --git
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieHiveRecordMerger.java
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/DefaultHiveRecordMerger.java
similarity index 92%
copy from
hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieHiveRecordMerger.java
copy to
hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/DefaultHiveRecordMerger.java
index 17a4738569e..7d353f7ee10 100644
---
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieHiveRecordMerger.java
+++
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/DefaultHiveRecordMerger.java
@@ -30,7 +30,10 @@ import org.apache.avro.Schema;
import java.io.IOException;
-public class HoodieHiveRecordMerger implements HoodieRecordMerger {
+/**
+ * Record merger for hive that implements the default merger strategy
+ */
+public class DefaultHiveRecordMerger extends HoodieHiveRecordMerger {
@Override
public Option<Pair<HoodieRecord, Schema>> merge(HoodieRecord older, Schema
oldSchema, HoodieRecord newer, Schema newSchema, TypedProperties props) throws
IOException {
ValidationUtils.checkArgument(older.getRecordType() ==
HoodieRecord.HoodieRecordType.HIVE);
@@ -59,11 +62,6 @@ public class HoodieHiveRecordMerger implements
HoodieRecordMerger {
}
}
- @Override
- public HoodieRecord.HoodieRecordType getRecordType() {
- return HoodieRecord.HoodieRecordType.HIVE;
- }
-
@Override
public String getMergingStrategy() {
return HoodieRecordMerger.DEFAULT_MERGER_STRATEGY_UUID;
diff --git
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HiveHoodieReaderContext.java
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HiveHoodieReaderContext.java
index d2b6ccd8ee8..e31b1756ec9 100644
---
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HiveHoodieReaderContext.java
+++
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HiveHoodieReaderContext.java
@@ -25,13 +25,10 @@ import org.apache.hudi.common.model.HoodieEmptyRecord;
import org.apache.hudi.common.model.HoodieKey;
import org.apache.hudi.common.model.HoodieRecord;
import org.apache.hudi.common.model.HoodieRecordMerger;
-import org.apache.hudi.common.table.HoodieTableMetaClient;
import org.apache.hudi.common.util.Option;
import org.apache.hudi.common.util.StringUtils;
-import org.apache.hudi.common.util.ValidationUtils;
import org.apache.hudi.common.util.collection.ClosableIterator;
import org.apache.hudi.common.util.collection.CloseableMappingIterator;
-import org.apache.hudi.exception.HoodieException;
import org.apache.hudi.hadoop.utils.HoodieArrayWritableAvroUtils;
import org.apache.hudi.hadoop.utils.HoodieRealtimeRecordReaderUtils;
import org.apache.hudi.hadoop.utils.ObjectInspectorCache;
@@ -41,7 +38,10 @@ import org.apache.hudi.storage.StoragePathInfo;
import org.apache.avro.Schema;
import org.apache.avro.generic.IndexedRecord;
+import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.Path;
+import org.apache.hadoop.hive.ql.io.sarg.ConvertAstToSearchArg;
+import org.apache.hadoop.hive.ql.plan.TableScanDesc;
import org.apache.hadoop.hive.serde.serdeConstants;
import org.apache.hadoop.hive.serde2.ColumnProjectionUtils;
import org.apache.hadoop.hive.serde2.typeinfo.TypeInfo;
@@ -52,7 +52,6 @@ import org.apache.hadoop.mapred.FileSplit;
import org.apache.hadoop.mapred.InputSplit;
import org.apache.hadoop.mapred.JobConf;
import org.apache.hadoop.mapred.RecordReader;
-import org.apache.hadoop.mapred.Reporter;
import java.io.IOException;
import java.util.Arrays;
@@ -65,18 +64,11 @@ import java.util.function.UnaryOperator;
import java.util.stream.Collectors;
import java.util.stream.Stream;
-import static
org.apache.hudi.common.model.HoodieRecordMerger.DEFAULT_MERGER_STRATEGY_UUID;
-import static
org.apache.hudi.hadoop.utils.HoodieInputFormatUtils.getPartitionFieldNames;
-
/**
* {@link HoodieReaderContext} for Hive-specific {@link
HoodieFileGroupReaderBasedRecordReader}.
*/
public class HiveHoodieReaderContext extends
HoodieReaderContext<ArrayWritable> {
protected final HoodieFileGroupReaderBasedRecordReader.HiveReaderCreator
readerCreator;
- protected final InputSplit split;
- protected final JobConf jobConf;
- protected final Reporter reporter;
- protected final Schema writerSchema;
protected final Map<String, TypeInfo> columnTypeMap;
private final ObjectInspectorCache objectInspectorCache;
private RecordReader<NullWritable, ArrayWritable> firstRecordReader = null;
@@ -87,39 +79,17 @@ public class HiveHoodieReaderContext extends
HoodieReaderContext<ArrayWritable>
private final String recordKeyField;
protected
HiveHoodieReaderContext(HoodieFileGroupReaderBasedRecordReader.HiveReaderCreator
readerCreator,
- InputSplit split,
- JobConf jobConf,
- Reporter reporter,
- Schema writerSchema,
- HoodieTableMetaClient metaClient) {
+ String recordKeyField,
+ List<String> partitionCols,
+ ObjectInspectorCache objectInspectorCache)
{
this.readerCreator = readerCreator;
- this.split = split;
- this.jobConf = jobConf;
- this.reporter = reporter;
- this.writerSchema = writerSchema;
- this.partitionCols = getPartitionFieldNames(jobConf).stream().filter(n ->
writerSchema.getField(n) != null).collect(Collectors.toList());
+ this.partitionCols = partitionCols;
this.partitionColSet = new HashSet<>(this.partitionCols);
- String tableName = metaClient.getTableConfig().getTableName();
- recordKeyField = getRecordKeyField(metaClient);
- this.objectInspectorCache =
HoodieArrayWritableAvroUtils.getCacheForTable(tableName, writerSchema, jobConf);
+ this.recordKeyField = recordKeyField;
+ this.objectInspectorCache = objectInspectorCache;
this.columnTypeMap = objectInspectorCache.getColumnTypeMap();
}
- /**
- * If populate meta fields is false, then getRecordKeyFields()
- * should return exactly 1 recordkey field.
- */
- private static String getRecordKeyField(HoodieTableMetaClient metaClient) {
- if (metaClient.getTableConfig().populateMetaFields()) {
- return HoodieRecord.RECORD_KEY_METADATA_FIELD;
- }
-
- Option<String[]> recordKeyFieldsOpt =
metaClient.getTableConfig().getRecordKeyFields();
- ValidationUtils.checkArgument(recordKeyFieldsOpt.isPresent(), "No record
key field set in table config, but populateMetaFields is disabled");
- ValidationUtils.checkArgument(recordKeyFieldsOpt.get().length == 1, "More
than 1 record key set in table config, but populateMetaFields is disabled");
- return recordKeyFieldsOpt.get()[0];
- }
-
private void setSchemas(JobConf jobConf, Schema dataSchema, Schema
requiredSchema) {
List<String> dataColumnNameList = dataSchema.getFields().stream().map(f ->
f.name().toLowerCase(Locale.ROOT)).collect(Collectors.toList());
List<TypeInfo> dataColumnTypeList =
dataColumnNameList.stream().map(fieldName -> {
@@ -153,14 +123,23 @@ public class HiveHoodieReaderContext extends
HoodieReaderContext<ArrayWritable>
private ClosableIterator<ArrayWritable> getFileRecordIterator(StoragePath
filePath, String[] hosts, long start, long length, Schema dataSchema,
Schema
requiredSchema, HoodieStorage storage) throws IOException {
- JobConf jobConfCopy = new JobConf(jobConf);
+ JobConf jobConfCopy = new
JobConf(storage.getConf().unwrapAs(Configuration.class));
+ if (getNeedsBootstrapMerge()) {
+ // Hive PPD works at row-group level and only enabled when
hive.optimize.index.filter=true;
+ // The above config is disabled by default. But when enabled, would
cause misalignment between
+ // skeleton and bootstrap file. We will disable them specifically when
query needs bootstrap and skeleton
+ // file to be stitched.
+ // This disables row-group filtering
+ jobConfCopy.unset(TableScanDesc.FILTER_EXPR_CONF_STR);
+ jobConfCopy.unset(ConvertAstToSearchArg.SARG_PUSHDOWN);
+ }
//move the partition cols to the end, because in some cases it has issues
if we don't do that
Schema modifiedDataSchema =
HoodieAvroUtils.generateProjectionSchema(dataSchema,
Stream.concat(dataSchema.getFields().stream()
.map(f -> f.name().toLowerCase(Locale.ROOT)).filter(n ->
!partitionColSet.contains(n)),
partitionCols.stream().filter(c -> dataSchema.getField(c) !=
null)).collect(Collectors.toList()));
setSchemas(jobConfCopy, modifiedDataSchema, requiredSchema);
InputSplit inputSplit = new FileSplit(new Path(filePath.toString()),
start, length, hosts);
- RecordReader<NullWritable, ArrayWritable> recordReader =
readerCreator.getRecordReader(inputSplit, jobConfCopy, reporter);
+ RecordReader<NullWritable, ArrayWritable> recordReader =
readerCreator.getRecordReader(inputSplit, jobConfCopy);
if (firstRecordReader == null) {
firstRecordReader = recordReader;
}
@@ -179,10 +158,7 @@ public class HiveHoodieReaderContext extends
HoodieReaderContext<ArrayWritable>
@Override
public HoodieRecordMerger getRecordMerger(String mergerStrategy) {
- if (mergerStrategy.equals(DEFAULT_MERGER_STRATEGY_UUID)) {
- return new HoodieHiveRecordMerger();
- }
- throw new HoodieException(String.format("The merger strategy UUID is not
supported, Default: %s, Passed: %s", mergerStrategy,
DEFAULT_MERGER_STRATEGY_UUID));
+ return HoodieHiveRecordMerger.getRecordMerger(mergerStrategy);
}
@Override
diff --git
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieFileGroupReaderBasedRecordReader.java
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieFileGroupReaderBasedRecordReader.java
index f72b79b56e1..93b6ca8ad13 100644
---
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieFileGroupReaderBasedRecordReader.java
+++
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieFileGroupReaderBasedRecordReader.java
@@ -25,15 +25,19 @@ import org.apache.hudi.common.model.BaseFile;
import org.apache.hudi.common.model.FileSlice;
import org.apache.hudi.common.model.HoodieBaseFile;
import org.apache.hudi.common.model.HoodieFileGroupId;
+import org.apache.hudi.common.model.HoodieRecord;
import org.apache.hudi.common.table.HoodieTableMetaClient;
import org.apache.hudi.common.table.TableSchemaResolver;
import org.apache.hudi.common.table.read.HoodieFileGroupReader;
import org.apache.hudi.common.table.timeline.HoodieInstant;
import org.apache.hudi.common.util.Option;
import org.apache.hudi.common.util.StringUtils;
+import org.apache.hudi.common.util.ValidationUtils;
+import org.apache.hudi.common.util.VisibleForTesting;
import org.apache.hudi.hadoop.realtime.RealtimeSplit;
import org.apache.hudi.hadoop.utils.HoodieRealtimeInputFormatUtils;
import org.apache.hudi.hadoop.utils.HoodieRealtimeRecordReaderUtils;
+import org.apache.hudi.hadoop.utils.ObjectInspectorCache;
import org.apache.avro.Schema;
import org.apache.hadoop.fs.FileSystem;
@@ -47,7 +51,6 @@ import org.apache.hadoop.mapred.FileSplit;
import org.apache.hadoop.mapred.InputSplit;
import org.apache.hadoop.mapred.JobConf;
import org.apache.hadoop.mapred.RecordReader;
-import org.apache.hadoop.mapred.Reporter;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -55,6 +58,7 @@ import java.io.IOException;
import java.util.Arrays;
import java.util.Collections;
import java.util.HashSet;
+import java.util.List;
import java.util.Locale;
import java.util.Set;
import java.util.function.UnaryOperator;
@@ -84,8 +88,7 @@ public class HoodieFileGroupReaderBasedRecordReader
implements RecordReader<Null
public interface HiveReaderCreator {
org.apache.hadoop.mapred.RecordReader<NullWritable, ArrayWritable>
getRecordReader(
final org.apache.hadoop.mapred.InputSplit split,
- final org.apache.hadoop.mapred.JobConf job,
- final org.apache.hadoop.mapred.Reporter reporter
+ final org.apache.hadoop.mapred.JobConf job
) throws IOException;
}
@@ -99,8 +102,7 @@ public class HoodieFileGroupReaderBasedRecordReader
implements RecordReader<Null
public HoodieFileGroupReaderBasedRecordReader(HiveReaderCreator
readerCreator,
final InputSplit split,
- final JobConf jobConf,
- final Reporter reporter)
throws IOException {
+ final JobConf jobConf) throws
IOException {
this.jobConfCopy = new JobConf(jobConf);
HoodieRealtimeInputFormatUtils.cleanProjectionColumnIds(jobConfCopy);
Set<String> partitionColumns = new
HashSet<>(getPartitionFieldNames(jobConfCopy));
@@ -115,7 +117,9 @@ public class HoodieFileGroupReaderBasedRecordReader
implements RecordReader<Null
String latestCommitTime = getLatestCommitTime(split, metaClient);
Schema tableSchema = getLatestTableSchema(metaClient, jobConfCopy,
latestCommitTime);
Schema requestedSchema = createRequestedSchema(tableSchema, jobConfCopy);
- this.readerContext = new HiveHoodieReaderContext(readerCreator, split,
jobConfCopy, reporter, tableSchema, metaClient);
+ this.readerContext = new HiveHoodieReaderContext(readerCreator,
getRecordKeyField(metaClient),
+ getStoredPartitionFieldNames(jobConfCopy, tableSchema),
+ new ObjectInspectorCache(tableSchema, jobConfCopy));
this.arrayWritable = new ArrayWritable(Writable.class, new
Writable[requestedSchema.getFields().size()]);
TypedProperties props = metaClient.getTableConfig().getProps();
jobConf.forEach(e -> {
@@ -171,6 +175,30 @@ public class HoodieFileGroupReaderBasedRecordReader
implements RecordReader<Null
return readerContext.getProgress();
}
+ /**
+ * If populate meta fields is false, then getRecordKeyFields()
+ * should return exactly 1 recordkey field.
+ */
+ @VisibleForTesting
+ static String getRecordKeyField(HoodieTableMetaClient metaClient) {
+ if (metaClient.getTableConfig().populateMetaFields()) {
+ return HoodieRecord.RECORD_KEY_METADATA_FIELD;
+ }
+
+ Option<String[]> recordKeyFieldsOpt =
metaClient.getTableConfig().getRecordKeyFields();
+ ValidationUtils.checkArgument(recordKeyFieldsOpt.isPresent(), "No record
key field set in table config, but populateMetaFields is disabled");
+ ValidationUtils.checkArgument(recordKeyFieldsOpt.get().length == 1, "More
than 1 record key set in table config, but populateMetaFields is disabled");
+ return recordKeyFieldsOpt.get()[0];
+ }
+
+ /**
+ * List of partition fields that are actually written to the file
+ */
+ @VisibleForTesting
+ static List<String> getStoredPartitionFieldNames(JobConf jobConf, Schema
writerSchema) {
+ return getPartitionFieldNames(jobConf).stream().filter(n ->
writerSchema.getField(n) != null).collect(Collectors.toList());
+ }
+
public RealtimeSplit getSplit() {
return (RealtimeSplit) inputSplit;
}
@@ -203,7 +231,7 @@ public class HoodieFileGroupReaderBasedRecordReader
implements RecordReader<Null
}
/**
- * Convert FileSplit to FileSlice, but save the locations in 'hosts' because
that data is otherwise lost.
+ * Convert FileSplit to FileSlice
*/
private static FileSlice getFileSliceFromSplit(FileSplit split, FileSystem
fs, String tableBasePath) throws IOException {
BaseFile bootstrapBaseFile = createBootstrapBaseFile(split, fs);
diff --git
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieHiveRecordMerger.java
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieHiveRecordMerger.java
index 17a4738569e..e35bc68fada 100644
---
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieHiveRecordMerger.java
+++
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieHiveRecordMerger.java
@@ -19,53 +19,24 @@
package org.apache.hudi.hadoop;
-import org.apache.hudi.common.config.TypedProperties;
import org.apache.hudi.common.model.HoodieRecord;
import org.apache.hudi.common.model.HoodieRecordMerger;
-import org.apache.hudi.common.util.Option;
-import org.apache.hudi.common.util.ValidationUtils;
-import org.apache.hudi.common.util.collection.Pair;
-
-import org.apache.avro.Schema;
-
-import java.io.IOException;
-
-public class HoodieHiveRecordMerger implements HoodieRecordMerger {
- @Override
- public Option<Pair<HoodieRecord, Schema>> merge(HoodieRecord older, Schema
oldSchema, HoodieRecord newer, Schema newSchema, TypedProperties props) throws
IOException {
- ValidationUtils.checkArgument(older.getRecordType() ==
HoodieRecord.HoodieRecordType.HIVE);
- ValidationUtils.checkArgument(newer.getRecordType() ==
HoodieRecord.HoodieRecordType.HIVE);
- if (newer instanceof HoodieHiveRecord) {
- HoodieHiveRecord newHiveRecord = (HoodieHiveRecord) newer;
- if (newHiveRecord.isDeleted()) {
- return Option.empty();
- }
- } else if (newer.getData() == null) {
- return Option.empty();
- }
-
- if (older instanceof HoodieHiveRecord) {
- HoodieHiveRecord oldHiveRecord = (HoodieHiveRecord) older;
- if (oldHiveRecord.isDeleted()) {
- return Option.of(Pair.of(newer, newSchema));
- }
- } else if (older.getData() == null) {
- return Option.empty();
- }
- if (older.getOrderingValue(oldSchema,
props).compareTo(newer.getOrderingValue(newSchema, props)) > 0) {
- return Option.of(Pair.of(older, oldSchema));
- } else {
- return Option.of(Pair.of(newer, newSchema));
- }
- }
+import org.apache.hudi.exception.HoodieException;
+abstract class HoodieHiveRecordMerger implements HoodieRecordMerger {
@Override
public HoodieRecord.HoodieRecordType getRecordType() {
return HoodieRecord.HoodieRecordType.HIVE;
}
- @Override
- public String getMergingStrategy() {
- return HoodieRecordMerger.DEFAULT_MERGER_STRATEGY_UUID;
+ static HoodieRecordMerger getRecordMerger(String mergerStrategy) {
+ switch (mergerStrategy) {
+ case DEFAULT_MERGER_STRATEGY_UUID:
+ return new DefaultHiveRecordMerger();
+ case OVERWRITE_MERGER_STRATEGY_UUID:
+ return new OverwriteWithLatestHiveRecordMerger();
+ default:
+ throw new HoodieException("This merger strategy UUID is not supported:
" + mergerStrategy);
+ }
}
}
diff --git
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieParquetInputFormat.java
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieParquetInputFormat.java
index 18b9e221978..ed044ef7ddb 100644
---
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieParquetInputFormat.java
+++
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/HoodieParquetInputFormat.java
@@ -121,15 +121,15 @@ public class HoodieParquetInputFormat extends
HoodieParquetInputFormatBase {
return super.getRecordReader(split, job, reporter);
}
if (supportAvroRead &&
HoodieColumnProjectionUtils.supportTimestamp(job)) {
- return new HoodieFileGroupReaderBasedRecordReader((s, j, r) -> {
+ return new HoodieFileGroupReaderBasedRecordReader((s, j) -> {
try {
- return new ParquetRecordReaderWrapper(new
HoodieTimestampAwareParquetInputFormat(), s, j, r);
+ return new ParquetRecordReaderWrapper(new
HoodieTimestampAwareParquetInputFormat(), s, j, reporter);
} catch (InterruptedException e) {
throw new RuntimeException(e);
}
- }, split, job, reporter);
+ }, split, job);
} else {
- return new
HoodieFileGroupReaderBasedRecordReader(super::getRecordReader, split, job,
reporter);
+ return new HoodieFileGroupReaderBasedRecordReader((s, j) ->
super.getRecordReader(s, j, reporter), split, job);
}
} catch (final IOException e) {
throw new RuntimeException("Cannot create a RecordReaderWrapper", e);
diff --git
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/OverwriteWithLatestSparkMerger.java
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/OverwriteWithLatestHiveRecordMerger.java
similarity index 89%
rename from
hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/OverwriteWithLatestSparkMerger.java
rename to
hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/OverwriteWithLatestHiveRecordMerger.java
index 611f045f645..4a4acb18a36 100644
---
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/OverwriteWithLatestSparkMerger.java
+++
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/OverwriteWithLatestHiveRecordMerger.java
@@ -17,7 +17,7 @@
* under the License.
*/
-package org.apache.hudi;
+package org.apache.hudi.hadoop;
import org.apache.hudi.common.config.TypedProperties;
import org.apache.hudi.common.model.HoodieRecord;
@@ -29,10 +29,9 @@ import org.apache.avro.Schema;
import java.io.IOException;
/**
- * Spark merger that always chooses the newer record
+ * Hive merger that always chooses the newer record
*/
-public class OverwriteWithLatestSparkMerger extends HoodieSparkRecordMerger {
-
+public class OverwriteWithLatestHiveRecordMerger extends
HoodieHiveRecordMerger {
@Override
public String getMergingStrategy() {
return OVERWRITE_MERGER_STRATEGY_UUID;
@@ -42,5 +41,4 @@ public class OverwriteWithLatestSparkMerger extends
HoodieSparkRecordMerger {
public Option<Pair<HoodieRecord, Schema>> merge(HoodieRecord older, Schema
oldSchema, HoodieRecord newer, Schema newSchema, TypedProperties props) throws
IOException {
return Option.of(Pair.of(newer, newSchema));
}
-
}
diff --git
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/utils/HoodieArrayWritableAvroUtils.java
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/utils/HoodieArrayWritableAvroUtils.java
index a2da796c6f7..3657c6bec3e 100644
---
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/utils/HoodieArrayWritableAvroUtils.java
+++
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/utils/HoodieArrayWritableAvroUtils.java
@@ -27,24 +27,12 @@ import org.apache.avro.Schema;
import org.apache.hadoop.hive.ql.io.parquet.serde.ArrayWritableObjectInspector;
import org.apache.hadoop.io.ArrayWritable;
import org.apache.hadoop.io.Writable;
-import org.apache.hadoop.mapred.JobConf;
import java.util.List;
import java.util.function.UnaryOperator;
public class HoodieArrayWritableAvroUtils {
- private static final Cache<String, ObjectInspectorCache>
- OBJECT_INSPECTOR_TABLE_CACHE =
Caffeine.newBuilder().maximumSize(1000).build();
-
- public static ObjectInspectorCache getCacheForTable(String table, Schema
tableSchema, JobConf jobConf) {
- ObjectInspectorCache cache =
OBJECT_INSPECTOR_TABLE_CACHE.getIfPresent(table);
- if (cache == null) {
- cache = new ObjectInspectorCache(tableSchema, jobConf);
- }
- return cache;
- }
-
private static final Cache<Pair<Schema, Schema>, int[]>
PROJECTION_CACHE = Caffeine.newBuilder().maximumSize(1000).build();
diff --git
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/utils/HoodieRealtimeRecordReaderUtils.java
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/utils/HoodieRealtimeRecordReaderUtils.java
index 07863ac5b97..32dd8920742 100644
---
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/utils/HoodieRealtimeRecordReaderUtils.java
+++
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/utils/HoodieRealtimeRecordReaderUtils.java
@@ -216,7 +216,7 @@ public class HoodieRealtimeRecordReaderUtils {
}
return new ArrayWritable(Writable.class, recordValues);
case ENUM:
- return new Text(value.toString());
+ return new BytesWritable(value.toString().getBytes());
case ARRAY:
GenericArray arrayValue = (GenericArray) value;
Writable[] arrayValues = new Writable[arrayValue.size()];
diff --git
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/utils/ObjectInspectorCache.java
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/utils/ObjectInspectorCache.java
index ddcc28851df..c5bb8b4fa6f 100644
---
a/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/utils/ObjectInspectorCache.java
+++
b/hudi-hadoop-mr/src/main/java/org/apache/hudi/hadoop/utils/ObjectInspectorCache.java
@@ -19,8 +19,6 @@
package org.apache.hudi.hadoop.utils;
-import com.github.benmanes.caffeine.cache.Cache;
-import com.github.benmanes.caffeine.cache.Caffeine;
import org.apache.avro.Schema;
import org.apache.hadoop.hive.ql.io.parquet.serde.ArrayWritableObjectInspector;
import org.apache.hadoop.hive.serde.serdeConstants;
@@ -46,8 +44,7 @@ import java.util.stream.IntStream;
*/
public class ObjectInspectorCache {
private final Map<String, TypeInfo> columnTypeMap = new HashMap<>();
- private final Cache<Schema, ArrayWritableObjectInspector>
- objectInspectorCache = Caffeine.newBuilder().maximumSize(1000).build();
+ private final Map<Schema, ArrayWritableObjectInspector> objectInspectorCache
= new HashMap<>();
public Map<String, TypeInfo> getColumnTypeMap() {
return columnTypeMap;
@@ -82,7 +79,7 @@ public class ObjectInspectorCache {
}
public ArrayWritableObjectInspector getObjectInspector(Schema schema) {
- return objectInspectorCache.get(schema, s -> {
+ return objectInspectorCache.computeIfAbsent(schema, s -> {
List<String> columnNameList =
s.getFields().stream().map(Schema.Field::name).collect(Collectors.toList());
List<TypeInfo> columnTypeList =
columnNameList.stream().map(columnTypeMap::get).collect(Collectors.toList());
StructTypeInfo rowTypeInfo = (StructTypeInfo)
TypeInfoFactory.getStructTypeInfo(columnNameList, columnTypeList);
@@ -92,12 +89,7 @@ public class ObjectInspectorCache {
}
public Object getValue(ArrayWritable record, Schema schema, String
fieldName) {
- try {
- ArrayWritableObjectInspector objectInspector =
getObjectInspector(schema);
- return objectInspector.getStructFieldData(record,
objectInspector.getStructFieldRef(fieldName));
- } catch (Exception e) {
- throw e;
- }
-
+ ArrayWritableObjectInspector objectInspector = getObjectInspector(schema);
+ return objectInspector.getStructFieldData(record,
objectInspector.getStructFieldRef(fieldName));
}
}
diff --git
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/hudi/command/HoodieSparkValidateDuplicateKeyRecordMerger.scala
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/hudi/command/HoodieSparkValidateDuplicateKeyRecordMerger.scala
index 8a64007f053..c157f58f516 100644
---
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/hudi/command/HoodieSparkValidateDuplicateKeyRecordMerger.scala
+++
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/hudi/command/HoodieSparkValidateDuplicateKeyRecordMerger.scala
@@ -18,10 +18,10 @@
package org.apache.spark.sql.hudi.command
import org.apache.avro.Schema
-import org.apache.hudi.HoodieSparkRecordMerger
+import org.apache.hudi.{DefaultSparkRecordMerger, HoodieSparkRecordMerger}
import org.apache.hudi.common.config.TypedProperties
import org.apache.hudi.common.model.{HoodieRecord, HoodieRecordMerger,
OperationModeAwareness}
-import org.apache.hudi.common.util.{collection, HoodieRecordUtils, Option =>
HOption}
+import org.apache.hudi.common.util.{HoodieRecordUtils, collection, Option =>
HOption}
import org.apache.hudi.exception.HoodieDuplicateKeyException
/**
@@ -37,6 +37,11 @@ class HoodieSparkValidateDuplicateKeyRecordMerger extends
HoodieSparkRecordMerge
}
override def asPreCombiningMode(): HoodieRecordMerger = {
-
HoodieRecordUtils.loadRecordMerger(classOf[HoodieSparkRecordMerger].getName)
+
HoodieRecordUtils.loadRecordMerger(classOf[DefaultSparkRecordMerger].getName)
}
+
+ /**
+ * The kind of merging strategy this recordMerger belongs to. An UUID
represents merging strategy.
+ */
+ override def getMergingStrategy: String =
HoodieRecordMerger.DEFAULT_MERGER_STRATEGY_UUID
}
diff --git
a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/TestHoodieMergeHandleWithSparkMerger.java
b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/TestHoodieMergeHandleWithSparkMerger.java
index 620bf4beb13..b9309632743 100644
---
a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/TestHoodieMergeHandleWithSparkMerger.java
+++
b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/TestHoodieMergeHandleWithSparkMerger.java
@@ -160,7 +160,7 @@ public class TestHoodieMergeHandleWithSparkMerger extends
SparkClientFunctionalT
Properties extraProperties = new Properties();
extraProperties.setProperty(
RECORD_MERGER_IMPLS.key(),
- "org.apache.hudi.HoodieSparkRecordMerger");
+ "org.apache.hudi.DefaultSparkRecordMerger");
extraProperties.setProperty(
LOGFILE_DATA_BLOCK_FORMAT.key(),
"parquet");
@@ -229,7 +229,7 @@ public class TestHoodieMergeHandleWithSparkMerger extends
SparkClientFunctionalT
Map<String, String> properties = new HashMap<>();
properties.put(
RECORD_MERGER_IMPLS.key(),
- "org.apache.hudi.HoodieSparkRecordMerger");
+ "org.apache.hudi.DefaultSparkRecordMerger");
properties.put(
LOGFILE_DATA_BLOCK_FORMAT.key(),
"parquet");
@@ -361,21 +361,21 @@ public class TestHoodieMergeHandleWithSparkMerger extends
SparkClientFunctionalT
}
}
- public static class DefaultMerger extends HoodieSparkRecordMerger {
+ public static class DefaultMerger extends DefaultSparkRecordMerger {
@Override
public boolean shouldFlush(HoodieRecord record, Schema schema,
TypedProperties props) {
return true;
}
}
- public static class NoFlushMerger extends HoodieSparkRecordMerger {
+ public static class NoFlushMerger extends DefaultSparkRecordMerger {
@Override
public boolean shouldFlush(HoodieRecord record, Schema schema,
TypedProperties props) {
return false;
}
}
- public static class CustomMerger extends HoodieSparkRecordMerger {
+ public static class CustomMerger extends DefaultSparkRecordMerger {
@Override
public boolean shouldFlush(HoodieRecord record, Schema schema,
TypedProperties props) throws IOException {
return !((HoodieSparkRecord)
record).getData().getString(0).equals("001");
diff --git
a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/TestHoodiePositionBasedFileGroupRecordBuffer.java
b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/TestHoodiePositionBasedFileGroupRecordBuffer.java
index edcc1f6dad2..c80a2f28ded 100644
---
a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/TestHoodiePositionBasedFileGroupRecordBuffer.java
+++
b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/TestHoodiePositionBasedFileGroupRecordBuffer.java
@@ -62,7 +62,6 @@ import java.util.stream.Collectors;
import static
org.apache.hudi.common.engine.HoodieReaderContext.INTERNAL_META_RECORD_KEY;
import static org.apache.hudi.common.model.WriteOperationType.INSERT;
import static
org.apache.hudi.common.testutils.HoodieTestUtils.createMetaClient;
-import static
org.apache.hudi.common.testutils.RawTripTestPayload.recordsToStrings;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertNull;
@@ -90,7 +89,7 @@ public class TestHoodiePositionBasedFileGroupRecordBuffer
extends TestHoodieFile
writeConfigs.put("hoodie.compact.inline", "false");
writeConfigs.put(HoodieWriteConfig.WRITE_RECORD_POSITIONS.key(), "true");
writeConfigs.put(HoodieWriteConfig.WRITE_PAYLOAD_CLASS_NAME.key(),
getRecordPayloadForMergeMode(mergeMode));
- commitToTable(recordsToStrings(dataGen.generateInserts("001", 100)),
INSERT.value(), writeConfigs);
+ commitToTable(dataGen.generateInserts("001", 100), INSERT.value(),
writeConfigs);
String[] partitionPaths = dataGen.getPartitionPaths();
String[] partitionValues = new String[1];
@@ -115,11 +114,11 @@ public class TestHoodiePositionBasedFileGroupRecordBuffer
extends TestHoodieFile
ctx.setRecordMerger(new CustomMerger());
break;
case EVENT_TIME_ORDERING:
- ctx.setRecordMerger(new HoodieSparkRecordMerger());
+ ctx.setRecordMerger(new DefaultSparkRecordMerger());
break;
case OVERWRITE_WITH_LATEST:
default:
- ctx.setRecordMerger(new OverwriteWithLatestSparkMerger());
+ ctx.setRecordMerger(new OverwriteWithLatestSparkRecordMerger());
break;
}
ctx.setSchemaHandler(new HoodiePositionBasedSchemaHandler<>(ctx,
avroSchema, avroSchema,
diff --git
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/common/table/read/TestHoodieFileGroupReaderOnSpark.scala
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/common/table/read/TestHoodieFileGroupReaderOnSpark.scala
index 23013db5d33..85cf2fc1cb9 100644
---
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/common/table/read/TestHoodieFileGroupReaderOnSpark.scala
+++
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/common/table/read/TestHoodieFileGroupReaderOnSpark.scala
@@ -24,9 +24,9 @@ import org.apache.hudi.common.config.RecordMergeMode
import org.apache.hudi.common.engine.HoodieReaderContext
import org.apache.hudi.common.model.{DefaultHoodieRecordPayload, HoodieRecord,
OverwriteWithLatestAvroPayload, WriteOperationType}
import org.apache.hudi.common.table.HoodieTableMetaClient
-import org.apache.hudi.common.testutils.HoodieTestUtils
+import org.apache.hudi.common.testutils.{HoodieTestUtils, RawTripTestPayload}
import org.apache.hudi.storage.StorageConfiguration
-import org.apache.hudi.{HoodieSparkRecordMerger, SparkAdapterSupport,
SparkFileFormatInternalRowReaderContext}
+import org.apache.hudi.{DefaultSparkRecordMerger, SparkAdapterSupport,
SparkFileFormatInternalRowReaderContext}
import org.apache.avro.Schema
import org.apache.hadoop.conf.Configuration
@@ -89,12 +89,13 @@ class TestHoodieFileGroupReaderOnSpark extends
TestHoodieFileGroupReaderBase[Int
override def getHoodieReaderContext(tablePath: String, avroSchema: Schema,
storageConf: StorageConfiguration[_]): HoodieReaderContext[InternalRow] = {
val reader = sparkAdapter.createParquetFileReader(vectorized = false,
spark.sessionState.conf, Map.empty,
storageConf.unwrapAs(classOf[Configuration]))
val metaClient =
HoodieTableMetaClient.builder().setConf(storageConf).setBasePath(tablePath).build
- val recordKeyField = new
HoodieSparkRecordMerger().getMandatoryFieldsForMerging(metaClient.getTableConfig)(0)
+ val recordKeyField = new
DefaultSparkRecordMerger().getMandatoryFieldsForMerging(metaClient.getTableConfig)(0)
new SparkFileFormatInternalRowReaderContext(reader, recordKeyField,
Seq.empty, Seq.empty)
}
- override def commitToTable(recordList: util.List[String], operation: String,
options: util.Map[String, String]): Unit = {
- val inputDF: Dataset[Row] =
spark.read.json(spark.sparkContext.parallelize(recordList.asScala.toList, 2))
+ override def commitToTable(recordList: util.List[HoodieRecord[_]],
operation: String, options: util.Map[String, String]): Unit = {
+ val recs = RawTripTestPayload.recordsToStrings(recordList)
+ val inputDF: Dataset[Row] =
spark.read.json(spark.sparkContext.parallelize(recs.asScala.toList, 2))
inputDF.write.format("hudi")
.options(options)
diff --git
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/common/table/read/TestSpark35RecordPositionMetadataColumn.scala
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/common/table/read/TestSpark35RecordPositionMetadataColumn.scala
index 59dc89c5a21..87bdd4d3685 100644
---
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/common/table/read/TestSpark35RecordPositionMetadataColumn.scala
+++
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/common/table/read/TestSpark35RecordPositionMetadataColumn.scala
@@ -39,7 +39,7 @@ import org.junit.jupiter.api.{BeforeEach, Test}
class TestSpark35RecordPositionMetadataColumn extends
SparkClientFunctionalTestHarness {
private val PARQUET_FORMAT = "parquet"
- private val SPARK_MERGER = "org.apache.hudi.HoodieSparkRecordMerger"
+ private val SPARK_MERGER = "org.apache.hudi.DefaultSparkRecordMerger"
@BeforeEach
def setUp(): Unit = {
diff --git
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/CommonOptionUtils.scala
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/CommonOptionUtils.scala
index dffe06d4611..45248bbc9ae 100644
---
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/CommonOptionUtils.scala
+++
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/CommonOptionUtils.scala
@@ -23,7 +23,7 @@ import org.apache.hudi.common.config.HoodieMetadataConfig
import org.apache.hudi.common.model.HoodieRecord.HoodieRecordType
import org.apache.hudi.common.table.HoodieTableConfig
import org.apache.hudi.config.HoodieWriteConfig
-import org.apache.hudi.{DataSourceReadOptions, DataSourceWriteOptions,
HoodieSparkRecordMerger}
+import org.apache.hudi.{DataSourceReadOptions, DataSourceWriteOptions,
DefaultSparkRecordMerger}
object CommonOptionUtils {
@@ -39,7 +39,7 @@ object CommonOptionUtils {
HoodieWriteConfig.TBL_NAME.key -> "hoodie_test",
HoodieMetadataConfig.COMPACT_NUM_DELTA_COMMITS.key -> "1"
)
- val sparkOpts = Map(HoodieWriteConfig.RECORD_MERGER_IMPLS.key ->
classOf[HoodieSparkRecordMerger].getName)
+ val sparkOpts = Map(HoodieWriteConfig.RECORD_MERGER_IMPLS.key ->
classOf[DefaultSparkRecordMerger].getName)
def getWriterReaderOpts(recordType: HoodieRecordType,
opt: Map[String, String] = commonOpts,
diff --git
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestCOWDataSource.scala
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestCOWDataSource.scala
index 593034865a3..af98f85db6a 100644
---
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestCOWDataSource.scala
+++
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestCOWDataSource.scala
@@ -1762,7 +1762,7 @@ class TestCOWDataSource extends HoodieSparkClientTestBase
with ScalaAssertionSup
HoodieWriteConfig.KEYGENERATOR_CLASS_NAME.key() ->
"org.apache.hudi.keygen.ComplexKeyGenerator",
KeyGeneratorOptions.HIVE_STYLE_PARTITIONING_ENABLE.key() -> "true",
HiveSyncConfigHolder.HIVE_SYNC_ENABLED.key() -> "false",
- HoodieWriteConfig.RECORD_MERGER_IMPLS.key() ->
"org.apache.hudi.HoodieSparkRecordMerger"
+ HoodieWriteConfig.RECORD_MERGER_IMPLS.key() ->
"org.apache.hudi.DefaultSparkRecordMerger"
)
df1.write.format("hudi").options(hudiOptions).mode(SaveMode.Append).save(basePath)
diff --git
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestDataSourceForBootstrap.scala
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestDataSourceForBootstrap.scala
index fe3373f1a9f..5209cc3a80c 100644
---
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestDataSourceForBootstrap.scala
+++
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestDataSourceForBootstrap.scala
@@ -28,7 +28,7 @@ import
org.apache.hudi.functional.TestDataSourceForBootstrap.{dropMetaCols, sort
import org.apache.hudi.hadoop.fs.HadoopFSUtils
import org.apache.hudi.keygen.{NonpartitionedKeyGenerator, SimpleKeyGenerator}
import org.apache.hudi.testutils.HoodieClientTestUtils
-import org.apache.hudi.{DataSourceReadOptions, DataSourceWriteOptions,
HoodieDataSourceHelpers, HoodieSparkRecordMerger}
+import org.apache.hudi.{DataSourceReadOptions, DataSourceWriteOptions,
HoodieDataSourceHelpers, DefaultSparkRecordMerger}
import org.apache.hadoop.fs.{FileSystem, Path}
import org.apache.spark.api.java.JavaSparkContext
@@ -62,7 +62,7 @@ class TestDataSourceForBootstrap {
)
val sparkRecordTypeOpts = Map(
- HoodieWriteConfig.RECORD_MERGER_IMPLS.key ->
classOf[HoodieSparkRecordMerger].getName,
+ HoodieWriteConfig.RECORD_MERGER_IMPLS.key ->
classOf[DefaultSparkRecordMerger].getName,
HoodieStorageConfig.LOGFILE_DATA_BLOCK_FORMAT.key -> "parquet"
)
diff --git
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestHoodieMultipleBaseFileFormat.scala
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestHoodieMultipleBaseFileFormat.scala
index 5b4e7950708..f67ac29f8a7 100644
---
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestHoodieMultipleBaseFileFormat.scala
+++
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestHoodieMultipleBaseFileFormat.scala
@@ -26,7 +26,7 @@ import
org.apache.hudi.common.testutils.HoodieTestDataGenerator.{DEFAULT_FIRST_P
import org.apache.hudi.common.testutils.RawTripTestPayload.recordsToStrings
import org.apache.hudi.config.HoodieWriteConfig
import org.apache.hudi.testutils.HoodieSparkClientTestBase
-import org.apache.hudi.{DataSourceWriteOptions, HoodieSparkRecordMerger,
SparkDatasetMixin}
+import org.apache.hudi.{DataSourceWriteOptions, DefaultSparkRecordMerger,
SparkDatasetMixin}
import org.apache.spark.sql.{Dataset, Row, SaveMode, SparkSession}
import org.junit.jupiter.api.Assertions.assertEquals
@@ -52,7 +52,7 @@ class TestHoodieMultipleBaseFileFormat extends
HoodieSparkClientTestBase with Sp
HoodieWriteConfig.TBL_NAME.key -> "hoodie_test"
)
val sparkOpts = Map(
- HoodieWriteConfig.RECORD_MERGER_IMPLS.key ->
classOf[HoodieSparkRecordMerger].getName,
+ HoodieWriteConfig.RECORD_MERGER_IMPLS.key ->
classOf[DefaultSparkRecordMerger].getName,
HoodieStorageConfig.LOGFILE_DATA_BLOCK_FORMAT.key -> "parquet"
)
diff --git
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestMORDataSource.scala
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestMORDataSource.scala
index f9625f4aa8f..5ca998a5e3a 100644
---
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestMORDataSource.scala
+++
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestMORDataSource.scala
@@ -35,7 +35,7 @@ import org.apache.hudi.storage.StoragePath
import org.apache.hudi.table.action.compact.CompactionTriggerStrategy
import org.apache.hudi.testutils.{DataSourceTestUtils,
HoodieSparkClientTestBase}
import org.apache.hudi.util.JFunction
-import org.apache.hudi.{DataSourceReadOptions, DataSourceUtils,
DataSourceWriteOptions, HoodieDataSourceHelpers, HoodieSparkRecordMerger,
SparkDatasetMixin}
+import org.apache.hudi.{DataSourceReadOptions, DataSourceUtils,
DataSourceWriteOptions, HoodieDataSourceHelpers, DefaultSparkRecordMerger,
SparkDatasetMixin}
import org.apache.hadoop.fs.Path
import org.apache.spark.sql._
@@ -68,7 +68,7 @@ class TestMORDataSource extends HoodieSparkClientTestBase
with SparkDatasetMixin
HoodieWriteConfig.TBL_NAME.key -> "hoodie_test"
)
val sparkOpts = Map(
- HoodieWriteConfig.RECORD_MERGER_IMPLS.key ->
classOf[HoodieSparkRecordMerger].getName,
+ HoodieWriteConfig.RECORD_MERGER_IMPLS.key ->
classOf[DefaultSparkRecordMerger].getName,
HoodieStorageConfig.LOGFILE_DATA_BLOCK_FORMAT.key -> "parquet"
)
@@ -1019,7 +1019,7 @@ class TestMORDataSource extends HoodieSparkClientTestBase
with SparkDatasetMixin
"hoodie.datasource.write.row.writer.enable" -> "false"
)
if (recordType.equals(HoodieRecordType.SPARK)) {
- writeOpts = Map(HoodieWriteConfig.RECORD_MERGER_IMPLS.key ->
classOf[HoodieSparkRecordMerger].getName,
+ writeOpts = Map(HoodieWriteConfig.RECORD_MERGER_IMPLS.key ->
classOf[DefaultSparkRecordMerger].getName,
HoodieStorageConfig.LOGFILE_DATA_BLOCK_FORMAT.key -> "parquet") ++
writeOpts
}
val records1 = recordsToStrings(dataGen.generateInserts("001",
10)).asScala.toSeq
@@ -1061,7 +1061,7 @@ class TestMORDataSource extends HoodieSparkClientTestBase
with SparkDatasetMixin
"hoodie.datasource.write.row.writer.enable" -> "false"
)
if (recordType.equals(HoodieRecordType.SPARK)) {
- writeOpts = Map(HoodieWriteConfig.RECORD_MERGER_IMPLS.key ->
classOf[HoodieSparkRecordMerger].getName,
+ writeOpts = Map(HoodieWriteConfig.RECORD_MERGER_IMPLS.key ->
classOf[DefaultSparkRecordMerger].getName,
HoodieStorageConfig.LOGFILE_DATA_BLOCK_FORMAT.key -> "parquet") ++
writeOpts
}
val records1 = recordsToStrings(dataGen.generateInserts("001",
10)).asScala.toSeq
diff --git
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestSparkDataSourceDAGExecution.scala
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestSparkDataSourceDAGExecution.scala
index fff0046d56f..8c1591e6f78 100644
---
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestSparkDataSourceDAGExecution.scala
+++
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestSparkDataSourceDAGExecution.scala
@@ -25,7 +25,7 @@ import org.apache.hudi.common.util.Option
import org.apache.hudi.config.HoodieWriteConfig
import org.apache.hudi.testutils.HoodieSparkClientTestBase
import org.apache.hudi.util.JFunction
-import org.apache.hudi.{DataSourceWriteOptions, HoodieSparkRecordMerger,
ScalaAssertionSupport}
+import org.apache.hudi.{DataSourceWriteOptions, DefaultSparkRecordMerger,
ScalaAssertionSupport}
import org.apache.hadoop.fs.FileSystem
import org.apache.spark.scheduler.{SparkListener, SparkListenerStageCompleted}
@@ -58,7 +58,7 @@ class TestSparkDataSourceDAGExecution extends
HoodieSparkClientTestBase with Sca
HoodieWriteConfig.TBL_NAME.key -> "hoodie_test",
HoodieMetadataConfig.ENABLE.key -> "false"
)
- val sparkOpts = Map(HoodieWriteConfig.RECORD_MERGER_IMPLS.key ->
classOf[HoodieSparkRecordMerger].getName)
+ val sparkOpts = Map(HoodieWriteConfig.RECORD_MERGER_IMPLS.key ->
classOf[DefaultSparkRecordMerger].getName)
val verificationCol: String = "driver"
val updatedVerificationVal: String = "driver_update"
diff --git
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/execution/benchmark/ReadAndWriteWithoutAvroBenchmark.scala
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/execution/benchmark/ReadAndWriteWithoutAvroBenchmark.scala
index d7f7cf6a9f3..266740a4877 100644
---
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/execution/benchmark/ReadAndWriteWithoutAvroBenchmark.scala
+++
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/execution/benchmark/ReadAndWriteWithoutAvroBenchmark.scala
@@ -21,7 +21,7 @@ package org.apache.spark.sql.execution.benchmark
import org.apache.hudi.common.config.HoodieStorageConfig
import org.apache.hudi.common.model.HoodieAvroRecordMerger
import org.apache.hudi.config.{HoodieCompactionConfig, HoodieWriteConfig}
-import org.apache.hudi.{HoodieSparkRecordMerger, HoodieSparkUtils}
+import org.apache.hudi.{DefaultSparkRecordMerger, HoodieSparkUtils}
import org.apache.hadoop.fs.Path
import org.apache.spark.SparkConf
@@ -116,12 +116,12 @@ object ReadAndWriteWithoutAvroBenchmark extends
HoodieBenchmarkBase {
* pref insert overwrite: Best Time(ms) Avg
Time(ms) Stdev(ms) Rate(M/s) Per Row(ns) Relative
*
-----------------------------------------------------------------------------------------------------------------------------------
* org.apache.hudi.common.model.HoodieAvroRecordMerger 16714
17107 353 0.1 16714.5 1.0X
- * org.apache.hudi.HoodieSparkRecordMerger 12654
13924 1100 0.1 12653.8 1.3X
+ * org.apache.hudi.DefaultSparkRecordMerger 12654
13924 1100 0.1 12653.8 1.3X
*/
private def overwriteBenchmark(): Unit = {
val df = createComplexDataFrame(1000000)
val benchmark = new HoodieBenchmark("pref insert overwrite", 1000000, 3)
- Seq(classOf[HoodieAvroRecordMerger].getName,
classOf[HoodieSparkRecordMerger].getName).zip(Seq(avroTable,
sparkTable)).foreach {
+ Seq(classOf[HoodieAvroRecordMerger].getName,
classOf[DefaultSparkRecordMerger].getName).zip(Seq(avroTable,
sparkTable)).foreach {
case (merger, tableName) => benchmark.addCase(merger) { _ =>
withTempDir { f =>
prepareHoodieTable(tableName, new Path(f.getCanonicalPath,
tableName).toUri.toString, "mor", merger, df)
@@ -137,18 +137,18 @@ object ReadAndWriteWithoutAvroBenchmark extends
HoodieBenchmarkBase {
* pref upsert: Best Time(ms) Avg
Time(ms) Stdev(ms) Rate(M/s) Per Row(ns) Relative
*
-----------------------------------------------------------------------------------------------------------------------------------
* org.apache.hudi.common.model.HoodieAvroRecordMerger 6108
6383 257 0.0 610785.6 1.0X
- * org.apache.hudi.HoodieSparkRecordMerger 4833
5468 614 0.0 483300.0 1.3X
+ * org.apache.hudi.DefaultSparkRecordMerger 4833
5468 614 0.0 483300.0 1.3X
*
* Java HotSpot(TM) 64-Bit Server VM 1.8.0_211-b12 on Mac OS X 10.16
* Intel(R) Core(TM) i7-9750H CPU @ 2.60GHz
* pref read: Best Time(ms) Avg
Time(ms) Stdev(ms) Rate(M/s) Per Row(ns) Relative
*
-----------------------------------------------------------------------------------------------------------------------------------
* org.apache.hudi.common.model.HoodieAvroRecordMerger 813
818 8 0.0 81302.1 1.0X
- * org.apache.hudi.HoodieSparkRecordMerger 604
616 18 0.0 60430.1 1.3X
+ * org.apache.hudi.DefaultSparkRecordMerger 604
616 18 0.0 60430.1 1.3X
*/
private def upsertThenReadBenchmark(): Unit = {
val avroMergerImpl = classOf[HoodieAvroRecordMerger].getName
- val sparkMergerImpl = classOf[HoodieSparkRecordMerger].getName
+ val sparkMergerImpl = classOf[DefaultSparkRecordMerger].getName
val df = createComplexDataFrame(10000)
withTempDir { avroPath =>
withTempDir { sparkPath =>
diff --git
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/common/HoodieSparkSqlTestBase.scala
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/common/HoodieSparkSqlTestBase.scala
index 432d2ab767e..cb558b66b94 100644
---
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/common/HoodieSparkSqlTestBase.scala
+++
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/common/HoodieSparkSqlTestBase.scala
@@ -17,7 +17,7 @@
package org.apache.spark.sql.hudi.common
-import org.apache.hudi.HoodieSparkRecordMerger
+import org.apache.hudi.DefaultSparkRecordMerger
import org.apache.hudi.common.config.HoodieStorageConfig
import org.apache.hudi.common.model.HoodieAvroRecordMerger
import org.apache.hudi.common.model.HoodieRecord.HoodieRecordType
@@ -224,7 +224,7 @@ class HoodieSparkSqlTestBase extends FunSuite with
BeforeAndAfterAll {
// TODO HUDI-5264 Test parquet log with avro record in spark sql test
recordTypes.foreach { recordType =>
val (merger, format) = recordType match {
- case HoodieRecordType.SPARK =>
(classOf[HoodieSparkRecordMerger].getName, "parquet")
+ case HoodieRecordType.SPARK =>
(classOf[DefaultSparkRecordMerger].getName, "parquet")
case _ => (classOf[HoodieAvroRecordMerger].getName, "avro")
}
val config = Map(
@@ -240,7 +240,7 @@ class HoodieSparkSqlTestBase extends FunSuite with
BeforeAndAfterAll {
protected def getRecordType(): HoodieRecordType = {
val merger =
spark.sessionState.conf.getConfString(HoodieWriteConfig.RECORD_MERGER_IMPLS.key,
HoodieWriteConfig.RECORD_MERGER_IMPLS.defaultValue())
- if (merger.equals(classOf[HoodieSparkRecordMerger].getName)) {
+ if (merger.equals(classOf[DefaultSparkRecordMerger].getName)) {
HoodieRecordType.SPARK
} else {
HoodieRecordType.AVRO
diff --git
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/ddl/TestSpark3DDL.scala
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/ddl/TestSpark3DDL.scala
index fb39e8eaef7..c4b67799530 100644
---
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/ddl/TestSpark3DDL.scala
+++
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/ddl/TestSpark3DDL.scala
@@ -26,7 +26,7 @@ import org.apache.hudi.config.HoodieWriteConfig
import org.apache.hudi.index.inmemory.HoodieInMemoryHashIndex
import org.apache.hudi.testutils.DataSourceTestUtils
import org.apache.hudi.testutils.HoodieClientTestUtils.createMetaClient
-import org.apache.hudi.{DataSourceWriteOptions, HoodieSparkRecordMerger,
HoodieSparkUtils, QuickstartUtils}
+import org.apache.hudi.{DataSourceWriteOptions, DefaultSparkRecordMerger,
HoodieSparkUtils, QuickstartUtils}
import org.apache.hadoop.fs.Path
import org.apache.spark.sql.catalyst.TableIdentifier
@@ -698,7 +698,7 @@ class TestSpark3DDL extends HoodieSparkSqlTestBase {
}
val sparkOpts = Map(
- HoodieWriteConfig.RECORD_MERGER_IMPLS.key ->
classOf[HoodieSparkRecordMerger].getName,
+ HoodieWriteConfig.RECORD_MERGER_IMPLS.key ->
classOf[DefaultSparkRecordMerger].getName,
HoodieStorageConfig.LOGFILE_DATA_BLOCK_FORMAT.key -> "parquet"
)
diff --git
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/TestHoodieDeltaStreamer.java
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/TestHoodieDeltaStreamer.java
index e7f3acd6fc5..cd93bdbd9ce 100644
---
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/TestHoodieDeltaStreamer.java
+++
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/TestHoodieDeltaStreamer.java
@@ -21,7 +21,7 @@ package org.apache.hudi.utilities.deltastreamer;
import org.apache.hudi.DataSourceReadOptions;
import org.apache.hudi.DataSourceWriteOptions;
-import org.apache.hudi.HoodieSparkRecordMerger;
+import org.apache.hudi.DefaultSparkRecordMerger;
import org.apache.hudi.HoodieSparkUtils$;
import org.apache.hudi.client.SparkRDDWriteClient;
import org.apache.hudi.client.transaction.lock.InProcessLockProvider;
@@ -183,7 +183,7 @@ public class TestHoodieDeltaStreamer extends
HoodieDeltaStreamerTestBase {
private void addRecordMerger(HoodieRecordType type, List<String>
hoodieConfig) {
if (type == HoodieRecordType.SPARK) {
Map<String, String> opts = new HashMap<>();
- opts.put(HoodieWriteConfig.RECORD_MERGER_IMPLS.key(),
HoodieSparkRecordMerger.class.getName());
+ opts.put(HoodieWriteConfig.RECORD_MERGER_IMPLS.key(),
DefaultSparkRecordMerger.class.getName());
opts.put(HoodieStorageConfig.LOGFILE_DATA_BLOCK_FORMAT.key(), "parquet");
for (Map.Entry<String, String> entry : opts.entrySet()) {
hoodieConfig.add(String.format("%s=%s", entry.getKey(),
entry.getValue()));