This is an automated email from the ASF dual-hosted git repository.
danny0405 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 3a864ec6359 [HUDI-6089] Handle default insert behaviour to ingest
duplicates (#10728)
3a864ec6359 is described below
commit 3a864ec63598c2919c06ed03422cf54416b31b43
Author: wombatu-kun <[email protected]>
AuthorDate: Sun Mar 3 07:44:26 2024 +0700
[HUDI-6089] Handle default insert behaviour to ingest duplicates (#10728)
Co-authored-by: Vova Kolmakov <[email protected]>
---
.../src/main/java/org/apache/hudi/config/HoodieWriteConfig.java | 2 +-
.../main/java/org/apache/hudi/metadata/HoodieMetadataWriteUtils.java | 1 +
.../src/test/java/org/apache/hudi/config/TestHoodieWriteConfig.java | 1 +
.../org/apache/spark/sql/hudi/TestHoodieTableValuedFunction.scala | 1 +
.../src/test/scala/org/apache/spark/sql/hudi/TestInsertTable.scala | 4 +++-
.../test/scala/org/apache/spark/sql/hudi/TestMergeIntoTable2.scala | 2 ++
.../apache/hudi/utilities/deltastreamer/TestHoodieDeltaStreamer.java | 2 ++
7 files changed, 11 insertions(+), 2 deletions(-)
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/config/HoodieWriteConfig.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/config/HoodieWriteConfig.java
index f4cb386d271..9447069a995 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/config/HoodieWriteConfig.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/config/HoodieWriteConfig.java
@@ -562,7 +562,7 @@ public class HoodieWriteConfig extends HoodieConfig {
public static final ConfigProperty<String>
MERGE_ALLOW_DUPLICATE_ON_INSERTS_ENABLE = ConfigProperty
.key("hoodie.merge.allow.duplicate.on.inserts")
- .defaultValue("false")
+ .defaultValue("true")
.markAdvanced()
.withDocumentation("When enabled, we allow duplicate keys even if
inserts are routed to merge with an existing file (for ensuring file sizing)."
+ " This is only relevant for insert operation, since upsert, delete
operations will ensure unique key constraints are maintained.");
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/HoodieMetadataWriteUtils.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/HoodieMetadataWriteUtils.java
index 7c42ccf5016..243b74b9199 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/HoodieMetadataWriteUtils.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/HoodieMetadataWriteUtils.java
@@ -86,6 +86,7 @@ public class HoodieMetadataWriteUtils {
HoodieWriteConfig.Builder builder = HoodieWriteConfig.newBuilder()
.withEngineType(writeConfig.getEngineType())
.withTimelineLayoutVersion(TimelineLayoutVersion.CURR_VERSION)
+ .withMergeAllowDuplicateOnInserts(false)
.withConsistencyGuardConfig(ConsistencyGuardConfig.newBuilder()
.withConsistencyCheckEnabled(writeConfig.getConsistencyGuardConfig().isConsistencyCheckEnabled())
.withInitialConsistencyCheckIntervalMs(writeConfig.getConsistencyGuardConfig().getInitialConsistencyCheckIntervalMs())
diff --git
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/config/TestHoodieWriteConfig.java
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/config/TestHoodieWriteConfig.java
index 5c93f924ece..90fcfd4fd7a 100644
---
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/config/TestHoodieWriteConfig.java
+++
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/config/TestHoodieWriteConfig.java
@@ -89,6 +89,7 @@ public class TestHoodieWriteConfig {
assertEquals(5, config.getMaxCommitsToKeep());
assertEquals(2, config.getMinCommitsToKeep());
assertTrue(config.shouldUseExternalSchemaTransformation());
+ assertTrue(config.allowDuplicateInserts());
}
@Test
diff --git
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/TestHoodieTableValuedFunction.scala
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/TestHoodieTableValuedFunction.scala
index bdf512d3451..aa6ff39431f 100644
---
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/TestHoodieTableValuedFunction.scala
+++
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/TestHoodieTableValuedFunction.scala
@@ -450,6 +450,7 @@ class TestHoodieTableValuedFunction extends
HoodieSparkSqlTestBase {
|""".stripMargin
)
+ spark.sql("set hoodie.merge.allow.duplicate.on.inserts = false")
spark.sql(
s"""
| insert into $tableName
diff --git
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/TestInsertTable.scala
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/TestInsertTable.scala
index 7ee3626e34b..693b2039043 100644
---
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/TestInsertTable.scala
+++
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/TestInsertTable.scala
@@ -392,6 +392,7 @@ class TestInsertTable extends HoodieSparkSqlTestBase {
Seq(2, "a2", 12.0)
)
+ spark.sql("set hoodie.merge.allow.duplicate.on.inserts = false")
assertThrows[HoodieDuplicateKeyException] {
try {
spark.sql(s"insert into $tableName select 1, 'a1', 10")
@@ -1186,7 +1187,7 @@ class TestInsertTable extends HoodieSparkSqlTestBase {
}
test("Test combine before insert") {
- withSQLConf("hoodie.sql.bulk.insert.enable" -> "false") {
+ withSQLConf("hoodie.sql.bulk.insert.enable" -> "false",
"hoodie.merge.allow.duplicate.on.inserts" -> "false") {
withRecordType()(withTempDir{tmp =>
val tableName = generateTableName
spark.sql(
@@ -1500,6 +1501,7 @@ class TestInsertTable extends HoodieSparkSqlTestBase {
Seq(3, "a3", 30.0, 3000, "2021-01-07")
)
+ spark.sql("set hoodie.merge.allow.duplicate.on.inserts = false")
spark.sql(
s"""
| insert into $tableName values
diff --git
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/TestMergeIntoTable2.scala
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/TestMergeIntoTable2.scala
index b8f315575dd..b76b8d16130 100644
---
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/TestMergeIntoTable2.scala
+++
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/TestMergeIntoTable2.scala
@@ -923,6 +923,7 @@ class TestMergeIntoTable2 extends HoodieSparkSqlTestBase {
| partitioned by(dt)
| location '${tmp.getCanonicalPath}'
""".stripMargin)
+ spark.sql("set hoodie.merge.allow.duplicate.on.inserts = false")
spark.sql(
s"""
@@ -971,6 +972,7 @@ class TestMergeIntoTable2 extends HoodieSparkSqlTestBase {
| partitioned by(dt)
| location '${path1}'
""".stripMargin)
+ spark.sql("set hoodie.merge.allow.duplicate.on.inserts = false")
spark.sql(s"insert into $sourceTable values(1, 'a1', cast(3.01 as
double), 11, '2022-09-26'),(2, 'a2', cast(3.02 as double), 12,
'2022-09-27'),(3, 'a3', cast(3.03 as double), 13, '2022-09-28'),(4, 'a4',
cast(3.04 as double), 14, '2022-09-29')")
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 5294ae1b4c4..2a3e99b29a3 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
@@ -1137,6 +1137,7 @@ public class TestHoodieDeltaStreamer extends
HoodieDeltaStreamerTestBase {
cfg.tableType = HoodieTableType.COPY_ON_WRITE.name();
cfg.configs.addAll(getTableServicesConfigs(totalRecords, "false", "", "",
"true", "3"));
cfg.configs.add(String.format("%s=%s",
"hoodie.datasource.write.row.writer.enable", "false"));
+ cfg.configs.add(String.format("%s=%s",
"hoodie.merge.allow.duplicate.on.inserts", "false"));
HoodieDeltaStreamer ds = new HoodieDeltaStreamer(cfg, jsc);
deltaStreamerTestRunner(ds, cfg, (r) -> {
TestHelpers.assertAtLeastNReplaceCommits(1, tableBasePath, fs);
@@ -1200,6 +1201,7 @@ public class TestHoodieDeltaStreamer extends
HoodieDeltaStreamerTestBase {
cfg.continuousMode = true;
cfg.tableType = HoodieTableType.MERGE_ON_READ.name();
cfg.configs.addAll(getTableServicesConfigs(totalRecords, "false", "", "",
"true", "3"));
+ cfg.configs.add(String.format("%s=%s",
"hoodie.merge.allow.duplicate.on.inserts", "false"));
HoodieDeltaStreamer ds = new HoodieDeltaStreamer(cfg, jsc);
deltaStreamerTestRunner(ds, cfg, (r) -> {
TestHelpers.assertAtleastNCompactionCommits(2, tableBasePath, fs);