This is an automated email from the ASF dual-hosted git repository.
codope pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/hudi.git
The following commit(s) were added to refs/heads/master by this push:
new 3c439a3a69f [HUDI-6873] fix clustering mor (#9774)
3c439a3a69f is described below
commit 3c439a3a69fb88dee551d94c8266e48b7c0d1e8f
Author: Jon Vexler <[email protected]>
AuthorDate: Wed Oct 11 23:04:55 2023 -0400
[HUDI-6873] fix clustering mor (#9774)
Currently during clustering of noncompacted mor filegroups with row writer
disabled
(currently the default for clustering), the records in the base file are
applied to the
log scanner after the log files have been scanned. If they have the same
precombine,
the base file records will be chosen over the log file records. This commit
mimics the
implementation in Iterators.scala to make the behavior consistent.
---------
Co-authored-by: Jonathan Vexler <=>
---
.../hudi/common/table/log/CachingIterator.java | 41 +++++++++++
.../common/table/log/HoodieFileSliceReader.java | 75 +++++++++++++------
.../hudi/common/table/log/LogFileIterator.java | 57 +++++++++++++++
.../run/strategy/JavaExecutionStrategy.java | 4 +-
.../MultipleSparkJobExecutionStrategy.java | 4 +-
.../hudi/sink/clustering/ClusteringOperator.java | 3 +-
.../TestHoodieSparkMergeOnReadTableClustering.java | 2 +-
.../apache/hudi/functional/TestMORDataSource.scala | 85 +++++++++++++++++++++-
8 files changed, 243 insertions(+), 28 deletions(-)
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/common/table/log/CachingIterator.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/common/table/log/CachingIterator.java
new file mode 100644
index 00000000000..d022b92ae22
--- /dev/null
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/common/table/log/CachingIterator.java
@@ -0,0 +1,41 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hudi.common.table.log;
+
+import java.util.Iterator;
+
+public abstract class CachingIterator<T> implements Iterator<T> {
+
+ protected T nextRecord;
+
+ protected abstract boolean doHasNext();
+
+ @Override
+ public final boolean hasNext() {
+ return nextRecord != null || doHasNext();
+ }
+
+ @Override
+ public final T next() {
+ T record = nextRecord;
+ nextRecord = null;
+ return record;
+ }
+
+}
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/common/table/log/HoodieFileSliceReader.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/common/table/log/HoodieFileSliceReader.java
index fc3ef4b8d92..1aa2f21fcb2 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/common/table/log/HoodieFileSliceReader.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/common/table/log/HoodieFileSliceReader.java
@@ -19,47 +19,80 @@
package org.apache.hudi.common.table.log;
+import org.apache.hudi.common.config.TypedProperties;
+import org.apache.hudi.common.model.HoodiePayloadProps;
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.collection.Pair;
+import org.apache.hudi.exception.HoodieClusteringException;
import org.apache.hudi.io.storage.HoodieFileReader;
import org.apache.avro.Schema;
import java.io.IOException;
import java.util.Iterator;
+import java.util.Map;
import java.util.Properties;
-/**
- * Reads records from base file and merges any updates from log files and
provides iterable over all records in the file slice.
- */
-public class HoodieFileSliceReader<T> implements Iterator<HoodieRecord<T>> {
+public class HoodieFileSliceReader<T> extends LogFileIterator<T> {
+ private Option<Iterator<HoodieRecord>> baseFileIterator;
+ private HoodieMergedLogRecordScanner scanner;
+ private Schema schema;
+ private Properties props;
- private final Iterator<HoodieRecord<T>> recordsIterator;
+ private TypedProperties payloadProps = new TypedProperties();
+ private Option<Pair<String, String>> simpleKeyGenFieldsOpt;
+ Map<String, HoodieRecord> records;
+ HoodieRecordMerger merger;
- public static HoodieFileSliceReader getFileSliceReader(
- Option<HoodieFileReader> baseFileReader, HoodieMergedLogRecordScanner
scanner, Schema schema, Properties props, Option<Pair<String, String>>
simpleKeyGenFieldsOpt) throws IOException {
+ public HoodieFileSliceReader(Option<HoodieFileReader> baseFileReader,
+ HoodieMergedLogRecordScanner scanner,
Schema schema, String preCombineField, HoodieRecordMerger merger,
+ Properties props, Option<Pair<String, String>>
simpleKeyGenFieldsOpt) throws IOException {
+ super(scanner);
if (baseFileReader.isPresent()) {
- Iterator<HoodieRecord> baseIterator =
baseFileReader.get().getRecordIterator(schema);
- while (baseIterator.hasNext()) {
-
scanner.processNextRecord(baseIterator.next().wrapIntoHoodieRecordPayloadWithParams(schema,
props,
- simpleKeyGenFieldsOpt, scanner.isWithOperationField(),
scanner.getPartitionNameOverride(), false, Option.empty()));
- }
+ this.baseFileIterator =
Option.of(baseFileReader.get().getRecordIterator(schema));
+ } else {
+ this.baseFileIterator = Option.empty();
}
- return new HoodieFileSliceReader(scanner.iterator());
+ this.scanner = scanner;
+ this.schema = schema;
+ this.merger = merger;
+ if (preCombineField != null) {
+
payloadProps.setProperty(HoodiePayloadProps.PAYLOAD_ORDERING_FIELD_PROP_KEY,
preCombineField);
+ }
+ this.props = props;
+ this.simpleKeyGenFieldsOpt = simpleKeyGenFieldsOpt;
+ this.records = scanner.getRecords();
}
- private HoodieFileSliceReader(Iterator<HoodieRecord<T>> recordsItr) {
- this.recordsIterator = recordsItr;
+ private boolean hasNextInternal() {
+ while (baseFileIterator.isPresent() && baseFileIterator.get().hasNext()) {
+ try {
+ HoodieRecord currentRecord =
baseFileIterator.get().next().wrapIntoHoodieRecordPayloadWithParams(schema,
props,
+ simpleKeyGenFieldsOpt, scanner.isWithOperationField(),
scanner.getPartitionNameOverride(), false, Option.empty());
+ Option<HoodieRecord> logRecord =
removeLogRecord(currentRecord.getRecordKey());
+ if (!logRecord.isPresent()) {
+ nextRecord = currentRecord;
+ return true;
+ }
+ Option<Pair<HoodieRecord, Schema>> mergedRecordOpt =
merger.merge(currentRecord, schema, logRecord.get(), schema, payloadProps);
+ if (mergedRecordOpt.isPresent()) {
+ HoodieRecord<T> mergedRecord = (HoodieRecord<T>)
mergedRecordOpt.get().getLeft();
+ nextRecord =
mergedRecord.wrapIntoHoodieRecordPayloadWithParams(schema, props,
simpleKeyGenFieldsOpt, scanner.isWithOperationField(),
+ scanner.getPartitionNameOverride(), false, Option.empty());
+ return true;
+ }
+ } catch (IOException e) {
+ throw new HoodieClusteringException("Failed to
wrapIntoHoodieRecordPayloadWithParams: " + e.getMessage());
+ }
+ }
+ return super.doHasNext();
}
@Override
- public boolean hasNext() {
- return recordsIterator.hasNext();
+ protected boolean doHasNext() {
+ return hasNextInternal();
}
- @Override
- public HoodieRecord<T> next() {
- return recordsIterator.next();
- }
}
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/common/table/log/LogFileIterator.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/common/table/log/LogFileIterator.java
new file mode 100644
index 00000000000..bf55a6ba06e
--- /dev/null
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/common/table/log/LogFileIterator.java
@@ -0,0 +1,57 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.hudi.common.table.log;
+
+import org.apache.hudi.common.model.HoodieRecord;
+import org.apache.hudi.common.util.Option;
+
+import java.util.Iterator;
+import java.util.Map;
+
+public class LogFileIterator<T> extends CachingIterator<HoodieRecord<T>> {
+ HoodieMergedLogRecordScanner scanner;
+ Map<String, HoodieRecord> records;
+ Iterator<HoodieRecord> iterator;
+
+ protected Option<HoodieRecord> removeLogRecord(String key) {
+ return Option.ofNullable(records.remove(key));
+ }
+
+ public LogFileIterator(HoodieMergedLogRecordScanner scanner) {
+ this.scanner = scanner;
+ this.records = scanner.getRecords();
+ }
+
+ private boolean hasNextInternal() {
+ if (iterator == null) {
+ iterator = records.values().iterator();
+ }
+ if (iterator.hasNext()) {
+ nextRecord = iterator.next();
+ return true;
+ }
+ return false;
+ }
+
+ @Override
+ protected boolean doHasNext() {
+ return hasNextInternal();
+ }
+}
diff --git
a/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/client/clustering/run/strategy/JavaExecutionStrategy.java
b/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/client/clustering/run/strategy/JavaExecutionStrategy.java
index dcd88b083fc..81786d88f8b 100644
---
a/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/client/clustering/run/strategy/JavaExecutionStrategy.java
+++
b/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/client/clustering/run/strategy/JavaExecutionStrategy.java
@@ -32,6 +32,7 @@ import org.apache.hudi.common.model.HoodieFileGroupId;
import org.apache.hudi.common.model.HoodieKey;
import org.apache.hudi.common.model.HoodieRecord;
import org.apache.hudi.common.table.HoodieTableConfig;
+import org.apache.hudi.common.table.log.HoodieFileSliceReader;
import org.apache.hudi.common.table.log.HoodieMergedLogRecordScanner;
import org.apache.hudi.common.util.Option;
import org.apache.hudi.common.util.StringUtils;
@@ -61,7 +62,6 @@ import java.util.Map;
import java.util.Properties;
import java.util.stream.Collectors;
-import static
org.apache.hudi.common.table.log.HoodieFileSliceReader.getFileSliceReader;
import static
org.apache.hudi.config.HoodieClusteringConfig.PLAN_STRATEGY_SORT_COLUMNS;
/**
@@ -195,7 +195,7 @@ public abstract class JavaExecutionStrategy<T>
? Option.empty()
:
Option.of(HoodieFileReaderFactory.getReaderFactory(recordType).getFileReader(table.getHadoopConf(),
new Path(clusteringOp.getDataFilePath())));
HoodieTableConfig tableConfig = table.getMetaClient().getTableConfig();
- Iterator<HoodieRecord<T>> fileSliceReader =
getFileSliceReader(baseFileReader, scanner, readerSchema,
+ Iterator<HoodieRecord<T>> fileSliceReader = new
HoodieFileSliceReader(baseFileReader, scanner, readerSchema,
tableConfig.getPreCombineField(), writeConfig.getRecordMerger(),
tableConfig.getProps(),
tableConfig.populateMetaFields() ? Option.empty() :
Option.of(Pair.of(tableConfig.getRecordKeyFieldProp(),
tableConfig.getPartitionFieldProp())));
diff --git
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/clustering/run/strategy/MultipleSparkJobExecutionStrategy.java
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/clustering/run/strategy/MultipleSparkJobExecutionStrategy.java
index e2f06de00ea..bf0511138d5 100644
---
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/clustering/run/strategy/MultipleSparkJobExecutionStrategy.java
+++
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/client/clustering/run/strategy/MultipleSparkJobExecutionStrategy.java
@@ -35,6 +35,7 @@ import org.apache.hudi.common.model.HoodieFileGroupId;
import org.apache.hudi.common.model.HoodieKey;
import org.apache.hudi.common.model.HoodieRecord;
import org.apache.hudi.common.table.HoodieTableConfig;
+import org.apache.hudi.common.table.log.HoodieFileSliceReader;
import org.apache.hudi.common.table.log.HoodieMergedLogRecordScanner;
import org.apache.hudi.common.util.CollectionUtils;
import org.apache.hudi.common.util.CustomizedThreadFactory;
@@ -91,7 +92,6 @@ import java.util.stream.Stream;
import static
org.apache.hudi.client.utils.SparkPartitionUtils.getPartitionFieldVals;
import static org.apache.hudi.common.config.HoodieCommonConfig.TIMESTAMP_AS_OF;
-import static
org.apache.hudi.common.table.log.HoodieFileSliceReader.getFileSliceReader;
import static
org.apache.hudi.config.HoodieClusteringConfig.PLAN_STRATEGY_SORT_COLUMNS;
/**
@@ -324,7 +324,7 @@ public abstract class MultipleSparkJobExecutionStrategy<T>
Option<HoodieFileReader> baseFileReader =
StringUtils.isNullOrEmpty(clusteringOp.getDataFilePath())
? Option.empty()
: Option.of(getBaseOrBootstrapFileReader(hadoopConf,
bootstrapBasePath, partitionFields, clusteringOp));
- recordIterators.add(getFileSliceReader(baseFileReader, scanner,
readerSchema,
+ recordIterators.add(new HoodieFileSliceReader(baseFileReader,
scanner, readerSchema, tableConfig.getPreCombineField(),
config.getRecordMerger(),
tableConfig.getProps(),
tableConfig.populateMetaFields() ? Option.empty() :
Option.of(Pair.of(tableConfig.getRecordKeyFieldProp(),
tableConfig.getPartitionFieldProp()))));
diff --git
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/clustering/ClusteringOperator.java
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/clustering/ClusteringOperator.java
index 144f2618017..75e63d69b5f 100644
---
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/clustering/ClusteringOperator.java
+++
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/clustering/ClusteringOperator.java
@@ -282,7 +282,8 @@ public class ClusteringOperator extends
TableStreamOperator<ClusteringCommitEven
.build();
HoodieTableConfig tableConfig = table.getMetaClient().getTableConfig();
- HoodieFileSliceReader<? extends IndexedRecord> hoodieFileSliceReader =
HoodieFileSliceReader.getFileSliceReader(baseFileReader, scanner, readerSchema,
+ HoodieFileSliceReader<? extends IndexedRecord> hoodieFileSliceReader =
new HoodieFileSliceReader(baseFileReader, scanner, readerSchema,
+ tableConfig.getPreCombineField(),writeConfig.getRecordMerger(),
tableConfig.getProps(),
tableConfig.populateMetaFields() ? Option.empty() :
Option.of(Pair.of(tableConfig.getRecordKeyFieldProp(),
tableConfig.getPartitionFieldProp())));
diff --git
a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/functional/TestHoodieSparkMergeOnReadTableClustering.java
b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/functional/TestHoodieSparkMergeOnReadTableClustering.java
index 46bd78fd775..415ad1fc8d2 100644
---
a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/functional/TestHoodieSparkMergeOnReadTableClustering.java
+++
b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/functional/TestHoodieSparkMergeOnReadTableClustering.java
@@ -61,7 +61,7 @@ class TestHoodieSparkMergeOnReadTableClustering extends
SparkClientFunctionalTes
private static Stream<Arguments> testClustering() {
// enableClusteringAsRow, doUpdates, populateMetaFields,
preserveCommitMetadata
return Stream.of(
- Arguments.of(true, true, true),
+ Arguments.of(false, true, true),
Arguments.of(true, true, false),
Arguments.of(true, false, true),
Arguments.of(true, false, false),
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 81915953252..4625c9c7bb6 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
@@ -27,7 +27,7 @@ import
org.apache.hudi.common.model.HoodieRecord.HoodieRecordType
import org.apache.hudi.common.model._
import org.apache.hudi.common.table.HoodieTableMetaClient
import org.apache.hudi.common.testutils.HoodieTestDataGenerator
-import org.apache.hudi.common.testutils.RawTripTestPayload.recordsToStrings
+import org.apache.hudi.common.testutils.RawTripTestPayload.{recordToString,
recordsToStrings}
import org.apache.hudi.common.util
import org.apache.hudi.config.{HoodieCompactionConfig, HoodieIndexConfig,
HoodieWriteConfig}
import org.apache.hudi.functional.TestCOWDataSource.convertColumnsToNullable
@@ -996,6 +996,89 @@ class TestMORDataSource extends HoodieSparkClientTestBase
with SparkDatasetMixin
.save(basePath)
}
+ @ParameterizedTest
+ @EnumSource(value = classOf[HoodieRecordType], names = Array("AVRO",
"SPARK"))
+ def testClusteringSamePrecombine(recordType: HoodieRecordType): Unit = {
+ var writeOpts = Map(
+ "hoodie.insert.shuffle.parallelism" -> "4",
+ "hoodie.upsert.shuffle.parallelism" -> "4",
+ DataSourceWriteOptions.RECORDKEY_FIELD.key -> "_row_key",
+ DataSourceWriteOptions.PARTITIONPATH_FIELD.key -> "partition",
+ DataSourceWriteOptions.PRECOMBINE_FIELD.key -> "timestamp",
+ HoodieWriteConfig.TBL_NAME.key -> "hoodie_test",
+ DataSourceWriteOptions.OPERATION.key() ->
DataSourceWriteOptions.UPSERT_OPERATION_OPT_VAL,
+ DataSourceWriteOptions.TABLE_TYPE.key()->
DataSourceWriteOptions.MOR_TABLE_TYPE_OPT_VAL,
+ "hoodie.clustering.inline"-> "true",
+ "hoodie.clustering.inline.max.commits" -> "2",
+ "hoodie.clustering.plan.strategy.sort.columns" -> "_row_key",
+ "hoodie.metadata.enable" -> "false",
+ "hoodie.datasource.write.row.writer.enable" -> "false"
+ )
+ if (recordType.equals(HoodieRecordType.SPARK)) {
+ writeOpts = Map(HoodieWriteConfig.RECORD_MERGER_IMPLS.key ->
classOf[HoodieSparkRecordMerger].getName,
+ HoodieStorageConfig.LOGFILE_DATA_BLOCK_FORMAT.key -> "parquet") ++
writeOpts
+ }
+ val records1 = recordsToStrings(dataGen.generateInserts("001", 10)).asScala
+ val inputDF1: Dataset[Row] =
spark.read.json(spark.sparkContext.parallelize(records1, 2))
+ inputDF1.write.format("org.apache.hudi")
+ .options(writeOpts)
+ .mode(SaveMode.Overwrite)
+ .save(basePath)
+
+ val records2 = recordsToStrings(dataGen.generateUniqueUpdates("002",
5)).asScala
+ val inputDF2: Dataset[Row] =
spark.read.json(spark.sparkContext.parallelize(records2, 2))
+ inputDF2.write.format("org.apache.hudi")
+ .options(writeOpts)
+ .mode(SaveMode.Append)
+ .save(basePath)
+
+ assertEquals(5,
+ spark.read.format("hudi").load(basePath)
+ .select("_row_key", "partition", "rider")
+ .except(inputDF2.select("_row_key", "partition", "rider")).count())
+ }
+
+ @ParameterizedTest
+ @EnumSource(value = classOf[HoodieRecordType], names = Array("AVRO",
"SPARK"))
+ def testClusteringSamePrecombineWithDelete(recordType: HoodieRecordType):
Unit = {
+ var writeOpts = Map(
+ "hoodie.insert.shuffle.parallelism" -> "4",
+ "hoodie.upsert.shuffle.parallelism" -> "4",
+ DataSourceWriteOptions.RECORDKEY_FIELD.key -> "_row_key",
+ DataSourceWriteOptions.PARTITIONPATH_FIELD.key -> "partition",
+ DataSourceWriteOptions.PRECOMBINE_FIELD.key -> "timestamp",
+ HoodieWriteConfig.TBL_NAME.key -> "hoodie_test",
+ DataSourceWriteOptions.OPERATION.key() ->
DataSourceWriteOptions.UPSERT_OPERATION_OPT_VAL,
+ DataSourceWriteOptions.TABLE_TYPE.key() ->
DataSourceWriteOptions.MOR_TABLE_TYPE_OPT_VAL,
+ "hoodie.clustering.inline" -> "true",
+ "hoodie.clustering.inline.max.commits" -> "2",
+ "hoodie.clustering.plan.strategy.sort.columns" -> "_row_key",
+ "hoodie.metadata.enable" -> "false",
+ "hoodie.datasource.write.row.writer.enable" -> "false"
+ )
+ if (recordType.equals(HoodieRecordType.SPARK)) {
+ writeOpts = Map(HoodieWriteConfig.RECORD_MERGER_IMPLS.key ->
classOf[HoodieSparkRecordMerger].getName,
+ HoodieStorageConfig.LOGFILE_DATA_BLOCK_FORMAT.key -> "parquet") ++
writeOpts
+ }
+ val records1 = recordsToStrings(dataGen.generateInserts("001", 10)).asScala
+ val inputDF1: Dataset[Row] =
spark.read.json(spark.sparkContext.parallelize(records1, 2))
+ inputDF1.write.format("org.apache.hudi")
+ .options(writeOpts)
+ .mode(SaveMode.Overwrite)
+ .save(basePath)
+
+ writeOpts = writeOpts + (DataSourceWriteOptions.OPERATION.key() ->
DataSourceWriteOptions.DELETE_OPERATION_OPT_VAL)
+ val records2 = recordsToStrings(dataGen.generateUniqueUpdates("002",
5)).asScala
+ val inputDF2: Dataset[Row] =
spark.read.json(spark.sparkContext.parallelize(records2, 2))
+ inputDF2.write.format("org.apache.hudi")
+ .options(writeOpts)
+ .mode(SaveMode.Append)
+ .save(basePath)
+
+ assertEquals(5,
+ spark.read.format("hudi").load(basePath).count())
+ }
+
@ParameterizedTest
@EnumSource(value = classOf[HoodieRecordType], names = Array("AVRO",
"SPARK"))
def testHoodieIsDeletedMOR(recordType: HoodieRecordType): Unit = {