From 6879a4939c162e931bc19610f0ff9329e0d21191 Mon Sep 17 00:00:00 2001 From: Sivabalan Narayanan Date: Wed, 10 Nov 2021 18:47:51 -0500 Subject: [PATCH 1/4] Enabling timeline server based marker as default --- .../src/main/java/org/apache/hudi/config/HoodieWriteConfig.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 7988e93075220..a0ebcb3dda6a0 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 @@ -248,7 +248,7 @@ public class HoodieWriteConfig extends HoodieConfig { public static final ConfigProperty MARKERS_TYPE = ConfigProperty .key("hoodie.write.markers.type") - .defaultValue(MarkerType.DIRECT.toString()) + .defaultValue(MarkerType.TIMELINE_SERVER_BASED.toString()) .sinceVersion("0.9.0") .withDocumentation("Marker type to use. Two modes are supported: " + "- DIRECT: individual marker file corresponding to each data file is directly " From 6e6aa9e80ac80a7e47fe1c9e71f5513f2e44cbf8 Mon Sep 17 00:00:00 2001 From: Sivabalan Narayanan Date: Thu, 11 Nov 2021 18:56:31 -0500 Subject: [PATCH 2/4] Fixing tests --- .../table/marker/TimelineServerBasedWriteMarkers.java | 6 +++++- .../hudi/table/upgrade/OneToZeroDowngradeHandler.java | 1 + .../hudi/client/TestHoodieClientMultiWriter.java | 2 ++ .../test/java/org/apache/hudi/client/TestMultiFS.java | 2 ++ .../apache/hudi/client/TestUpdateSchemaEvolution.java | 5 ++++- .../apache/hudi/io/TestHoodieTimelineArchiveLog.java | 6 ++++++ .../action/commit/TestCopyOnWriteActionExecutor.java | 4 +++- .../hudi/table/marker/TestWriteMarkersBase.java | 10 ++++++++++ .../hudi/table/upgrade/TestUpgradeDowngrade.java | 11 ++++++++--- .../hudi/functional/TestStructuredStreaming.scala | 4 +++- 10 files changed, 44 insertions(+), 7 deletions(-) diff --git a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/marker/TimelineServerBasedWriteMarkers.java b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/marker/TimelineServerBasedWriteMarkers.java index 7ba6cb6bef1d5..f280f0e886e95 100644 --- a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/marker/TimelineServerBasedWriteMarkers.java +++ b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/marker/TimelineServerBasedWriteMarkers.java @@ -143,7 +143,11 @@ protected Option create(String partitionPath, String dataFileName, IOType LOG.info("[timeline-server-based] Created marker file " + partitionPath + "/" + markerFileName + " in " + timer.endTimer() + " ms"); if (success) { - return Option.of(new Path(new Path(markerDirPath, partitionPath), markerFileName)); + if (partitionPath.isEmpty()) { + return Option.of(new Path(markerDirPath, markerFileName)); + } else { + return Option.of(new Path(new Path(markerDirPath, partitionPath), markerFileName)); + } } else { return Option.empty(); } diff --git a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/upgrade/OneToZeroDowngradeHandler.java b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/upgrade/OneToZeroDowngradeHandler.java index e6051cf321b50..131c380982826 100644 --- a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/upgrade/OneToZeroDowngradeHandler.java +++ b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/upgrade/OneToZeroDowngradeHandler.java @@ -47,6 +47,7 @@ public Map downgrade( List commits = inflightTimeline.getReverseOrderedInstants().collect(Collectors.toList()); for (HoodieInstant inflightInstant : commits) { // delete existing markers + String markerType = config.getMarkersType().name(); WriteMarkers writeMarkers = WriteMarkersFactory.get(config.getMarkersType(), table, inflightInstant.getTimestamp()); writeMarkers.quietDeleteMarkerDir(context, config.getMarkersDeleteParallelism()); } diff --git a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/TestHoodieClientMultiWriter.java b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/TestHoodieClientMultiWriter.java index cef6641f1743e..d93d0e0914a70 100644 --- a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/TestHoodieClientMultiWriter.java +++ b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/TestHoodieClientMultiWriter.java @@ -31,6 +31,7 @@ import org.apache.hudi.common.model.HoodieTableType; import org.apache.hudi.common.model.TableServiceType; import org.apache.hudi.common.model.WriteConcurrencyMode; +import org.apache.hudi.common.table.marker.MarkerType; import org.apache.hudi.common.table.timeline.HoodieInstant; import org.apache.hudi.common.table.view.FileSystemViewStorageConfig; import org.apache.hudi.common.table.view.FileSystemViewStorageType; @@ -212,6 +213,7 @@ private void testMultiWriterWithAsyncTableServicesWithConflict(HoodieTableType t .withFailedWritesCleaningPolicy(HoodieFailedWritesCleaningPolicy.LAZY) .withMaxNumDeltaCommitsBeforeCompaction(2).build()) .withEmbeddedTimelineServerEnabled(false) + .withMarkersType(MarkerType.DIRECT.name()) .withFileSystemViewConfig(FileSystemViewStorageConfig.newBuilder().withStorageType( FileSystemViewStorageType.MEMORY).build()) .withClusteringConfig(HoodieClusteringConfig.newBuilder().withInlineClusteringNumCommits(1).build()) diff --git a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/TestMultiFS.java b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/TestMultiFS.java index 457b8b526aa04..894fcd2d4ca5d 100644 --- a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/TestMultiFS.java +++ b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/TestMultiFS.java @@ -24,6 +24,7 @@ import org.apache.hudi.common.model.HoodieRecord; import org.apache.hudi.common.model.HoodieTableType; import org.apache.hudi.common.table.HoodieTableMetaClient; +import org.apache.hudi.common.table.marker.MarkerType; import org.apache.hudi.common.table.timeline.HoodieActiveTimeline; import org.apache.hudi.common.table.timeline.HoodieTimeline; import org.apache.hudi.common.testutils.HoodieTestDataGenerator; @@ -77,6 +78,7 @@ protected HoodieWriteConfig getHoodieWriteConfig(String basePath, boolean enable .withSchema(HoodieTestDataGenerator.TRIP_EXAMPLE_SCHEMA).withParallelism(2, 2).forTable(tableName) .withIndexConfig(HoodieIndexConfig.newBuilder().withIndexType(HoodieIndex.IndexType.BLOOM).build()) .withMetadataConfig(HoodieMetadataConfig.newBuilder().enable(enableMetadata).build()) + .withMarkersType(MarkerType.DIRECT.name()) .build(); } diff --git a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/TestUpdateSchemaEvolution.java b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/TestUpdateSchemaEvolution.java index 96782f49428cb..7995ecac02279 100644 --- a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/TestUpdateSchemaEvolution.java +++ b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/TestUpdateSchemaEvolution.java @@ -22,6 +22,7 @@ import org.apache.hudi.common.model.HoodieKey; import org.apache.hudi.common.model.HoodieRecord; import org.apache.hudi.common.model.HoodieRecordLocation; +import org.apache.hudi.common.table.marker.MarkerType; import org.apache.hudi.common.table.timeline.HoodieTimeline; import org.apache.hudi.common.testutils.HoodieTestUtils; import org.apache.hudi.common.testutils.RawTripTestPayload; @@ -228,6 +229,8 @@ public void testSchemaEvolutionOnUpdateMisMatchWithChangeColumnType() throws Exc private HoodieWriteConfig makeHoodieClientConfig(String name) { Schema schema = getSchemaFromResource(getClass(), name); - return HoodieWriteConfig.newBuilder().withPath(basePath).withSchema(schema.toString()).build(); + return HoodieWriteConfig.newBuilder().withPath(basePath) + .withMarkersType(MarkerType.DIRECT.name()) + .withSchema(schema.toString()).build(); } } diff --git a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/io/TestHoodieTimelineArchiveLog.java b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/io/TestHoodieTimelineArchiveLog.java index 7cb9740a8c6cc..de8acaaaf4f08 100644 --- a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/io/TestHoodieTimelineArchiveLog.java +++ b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/io/TestHoodieTimelineArchiveLog.java @@ -31,6 +31,7 @@ import org.apache.hudi.common.model.HoodieTableType; import org.apache.hudi.common.model.WriteOperationType; import org.apache.hudi.common.table.HoodieTableMetaClient; +import org.apache.hudi.common.table.marker.MarkerType; import org.apache.hudi.common.table.timeline.HoodieArchivedTimeline; import org.apache.hudi.common.table.timeline.HoodieInstant; import org.apache.hudi.common.table.timeline.HoodieInstant.State; @@ -124,6 +125,7 @@ private HoodieWriteConfig initTestTableAndGetWriteConfig(boolean enableMetadata, HoodieTableType tableType) throws Exception { init(tableType); HoodieWriteConfig writeConfig = HoodieWriteConfig.newBuilder().withPath(basePath) + .withMarkersType(MarkerType.DIRECT.name()) .withSchema(HoodieTestDataGenerator.TRIP_EXAMPLE_SCHEMA).withParallelism(2, 2) .withCompactionConfig(HoodieCompactionConfig.newBuilder().retainCommits(1).archiveCommitsWith(minArchivalCommits, maxArchivalCommits).build()) .withMetadataConfig(HoodieMetadataConfig.newBuilder().enable(enableMetadata) @@ -211,6 +213,7 @@ public void testArchiveCommitSavepointNoHole() throws Exception { .withSchema(HoodieTestDataGenerator.TRIP_EXAMPLE_SCHEMA).withParallelism(2, 2).forTable("test-trip-table") .withCompactionConfig(HoodieCompactionConfig.newBuilder().retainCommits(1).archiveCommitsWith(2, 5).build()) .withMetadataConfig(HoodieMetadataConfig.newBuilder().enable(false).build()) + .withMarkersType(MarkerType.DIRECT.name()) .build(); HoodieTestDataGenerator.createCommitFile(basePath, "100", wrapperFs.getConf()); @@ -329,6 +332,7 @@ public void testArchiveCommitTimeline() throws Exception { .withParallelism(2, 2).forTable("test-trip-table") .withCompactionConfig(HoodieCompactionConfig.newBuilder().retainCommits(1).archiveCommitsWith(2, 3).build()) .withMetadataConfig(HoodieMetadataConfig.newBuilder().enable(false).build()) + .withMarkersType(MarkerType.DIRECT.name()) .build(); metaClient = HoodieTableMetaClient.reload(metaClient); @@ -485,6 +489,7 @@ public void testArchiveCompletedRollbackAndClean() throws Exception { .withParallelism(2, 2).forTable("test-trip-table") .withCompactionConfig(HoodieCompactionConfig.newBuilder().retainCommits(1).archiveCommitsWith(minInstantsToKeep, maxInstantsToKeep).build()) .withMetadataConfig(HoodieMetadataConfig.newBuilder().enable(false).build()) + .withMarkersType(MarkerType.DIRECT.name()) .build(); metaClient = HoodieTableMetaClient.reload(metaClient); @@ -520,6 +525,7 @@ public void testArchiveInflightClean() throws Exception { .withParallelism(2, 2).forTable("test-trip-table") .withCompactionConfig(HoodieCompactionConfig.newBuilder().retainCommits(1).archiveCommitsWith(2, 3).build()) .withMetadataConfig(HoodieMetadataConfig.newBuilder().enable(false).build()) + .withMarkersType(MarkerType.DIRECT.name()) .build(); metaClient = HoodieTableMetaClient.reload(metaClient); diff --git a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/table/action/commit/TestCopyOnWriteActionExecutor.java b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/table/action/commit/TestCopyOnWriteActionExecutor.java index 40df1af898ea3..ee58a2de3c8a6 100644 --- a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/table/action/commit/TestCopyOnWriteActionExecutor.java +++ b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/table/action/commit/TestCopyOnWriteActionExecutor.java @@ -26,6 +26,7 @@ import org.apache.hudi.common.model.HoodieRecord; import org.apache.hudi.common.model.HoodieTableType; import org.apache.hudi.common.table.HoodieTableMetaClient; +import org.apache.hudi.common.table.marker.MarkerType; import org.apache.hudi.common.testutils.HoodieTestUtils; import org.apache.hudi.common.testutils.RawTripTestPayload; import org.apache.hudi.common.testutils.Transformations; @@ -112,7 +113,7 @@ private HoodieWriteConfig makeHoodieClientConfig() { private HoodieWriteConfig.Builder makeHoodieClientConfigBuilder() { // Prepare the AvroParquetIO - return HoodieWriteConfig.newBuilder().withPath(basePath).withSchema(SCHEMA.toString()); + return HoodieWriteConfig.newBuilder().withPath(basePath).withSchema(SCHEMA.toString()).withMarkersType(MarkerType.DIRECT.name()); } // TODO (weiy): Add testcases for crossing file writing. @@ -405,6 +406,7 @@ public void testFileSizeUpsertRecords() throws Exception { public void testInsertUpsertWithHoodieAvroPayload() throws Exception { Schema schema = getSchemaFromResource(TestCopyOnWriteActionExecutor.class, "/testDataGeneratorSchema.txt"); HoodieWriteConfig config = HoodieWriteConfig.newBuilder().withPath(basePath).withSchema(schema.toString()) + .withMarkersType(MarkerType.DIRECT.name()) .withStorageConfig(HoodieStorageConfig.newBuilder() .parquetMaxFileSize(1000 * 1024).hfileMaxFileSize(1000 * 1024).build()).build(); metaClient = HoodieTableMetaClient.reload(metaClient); diff --git a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/table/marker/TestWriteMarkersBase.java b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/table/marker/TestWriteMarkersBase.java index 0298ed37a6383..1e0558177399b 100644 --- a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/table/marker/TestWriteMarkersBase.java +++ b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/table/marker/TestWriteMarkersBase.java @@ -76,6 +76,16 @@ public void testCreation() throws Exception { verifyMarkersInFileSystem(); } + @Test + public void testNonPartitionedDataset() throws IOException { + writeMarkers.create("", "file1", IOType.CREATE); + writeMarkers.create("", "file2", IOType.MERGE); + + assertTrue(writeMarkers.doesMarkerDirExist()); + assertTrue(writeMarkers.deleteMarkerDir(context, 2)); + assertFalse(writeMarkers.doesMarkerDirExist()); + } + @Test public void testDeletionWhenMarkerDirExists() throws IOException { //when diff --git a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/table/upgrade/TestUpgradeDowngrade.java b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/table/upgrade/TestUpgradeDowngrade.java index 19ec4e6d0654c..d921be4860f9f 100644 --- a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/table/upgrade/TestUpgradeDowngrade.java +++ b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/table/upgrade/TestUpgradeDowngrade.java @@ -327,7 +327,7 @@ public void testDowngrade(boolean deletePartialMarkerFiles, HoodieTableType tabl .run(toVersion, null); // assert marker files - assertMarkerFilesForDowngrade(table, commitInstant, toVersion == HoodieTableVersion.ONE); + assertMarkerFilesForDowngrade(table, commitInstant, toVersion == HoodieTableVersion.ONE, markerType); // verify hoodie.table.version got downgraded metaClient = HoodieTableMetaClient.builder().setConf(context.getHadoopConf().get()).setBasePath(cfg.getBasePath()) @@ -344,9 +344,14 @@ public void testDowngrade(boolean deletePartialMarkerFiles, HoodieTableType tabl */ } - private void assertMarkerFilesForDowngrade(HoodieTable table, HoodieInstant commitInstant, boolean assertExists) throws IOException { + private void assertMarkerFilesForDowngrade(HoodieTable table, HoodieInstant commitInstant, boolean assertExists, MarkerType markerType) throws IOException { // Verify recreated marker files are as expected - WriteMarkers writeMarkers = WriteMarkersFactory.get(getConfig().getMarkersType(), table, commitInstant.getTimestamp()); + WriteMarkers writeMarkers = null; + if (!assertExists) { + writeMarkers = WriteMarkersFactory.get(markerType, table, commitInstant.getTimestamp()); + } else { + writeMarkers = WriteMarkersFactory.get(getConfig().getMarkersType(), table, commitInstant.getTimestamp()); + } if (assertExists) { assertTrue(writeMarkers.doesMarkerDirExist()); assertEquals(0, getTimelineServerBasedMarkerFileCount(table.getMetaClient().getMarkerFolderPath(commitInstant.getTimestamp()), diff --git a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestStructuredStreaming.scala b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestStructuredStreaming.scala index 68b630be5d5e9..2e704e17e59c3 100644 --- a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestStructuredStreaming.scala +++ b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestStructuredStreaming.scala @@ -20,6 +20,7 @@ package org.apache.hudi.functional import org.apache.hadoop.fs.{FileSystem, Path} import org.apache.hudi.common.model.FileSlice import org.apache.hudi.common.table.HoodieTableMetaClient +import org.apache.hudi.common.table.marker.MarkerType import org.apache.hudi.common.testutils.RawTripTestPayload.recordsToStrings import org.apache.hudi.common.testutils.{HoodieTestDataGenerator, HoodieTestTable} import org.apache.hudi.config.{HoodieClusteringConfig, HoodieStorageConfig, HoodieWriteConfig} @@ -50,7 +51,8 @@ class TestStructuredStreaming extends HoodieClientTestBase { DataSourceWriteOptions.RECORDKEY_FIELD.key -> "_row_key", DataSourceWriteOptions.PARTITIONPATH_FIELD.key -> "partition", DataSourceWriteOptions.PRECOMBINE_FIELD.key -> "timestamp", - HoodieWriteConfig.TBL_NAME.key -> "hoodie_test" + HoodieWriteConfig.TBL_NAME.key -> "hoodie_test", + HoodieWriteConfig.MARKERS_TYPE.key() -> MarkerType.DIRECT.name() ) @BeforeEach override def setUp() { From 0ab4fdbdc71f727e7b2b6824f1d3718b16c79fca Mon Sep 17 00:00:00 2001 From: Sivabalan Narayanan Date: Mon, 15 Nov 2021 14:22:42 -0500 Subject: [PATCH 3/4] Fixing structured streaming to use DIRECT style markers --- .../TestJavaCopyOnWriteActionExecutor.java | 5 +++- .../TestHoodieClientOnCopyOnWriteStorage.java | 23 ++++++++++++++----- .../org/apache/hudi/util/StreamerUtil.java | 2 ++ .../org/apache/hudi/HoodieStreamingSink.scala | 8 +++++-- .../functional/TestStructuredStreaming.scala | 3 +-- .../functional/TestHoodieDeltaStreamer.java | 3 --- 6 files changed, 30 insertions(+), 14 deletions(-) diff --git a/hudi-client/hudi-java-client/src/test/java/org/apache/hudi/table/action/commit/TestJavaCopyOnWriteActionExecutor.java b/hudi-client/hudi-java-client/src/test/java/org/apache/hudi/table/action/commit/TestJavaCopyOnWriteActionExecutor.java index 796d7b74a83c5..6045b65cf618d 100644 --- a/hudi-client/hudi-java-client/src/test/java/org/apache/hudi/table/action/commit/TestJavaCopyOnWriteActionExecutor.java +++ b/hudi-client/hudi-java-client/src/test/java/org/apache/hudi/table/action/commit/TestJavaCopyOnWriteActionExecutor.java @@ -27,6 +27,7 @@ import org.apache.hudi.common.model.HoodieRecord; import org.apache.hudi.common.model.HoodieTableType; import org.apache.hudi.common.table.HoodieTableMetaClient; +import org.apache.hudi.common.table.marker.MarkerType; import org.apache.hudi.common.testutils.HoodieTestDataGenerator; import org.apache.hudi.common.testutils.HoodieTestUtils; import org.apache.hudi.common.testutils.RawTripTestPayload; @@ -114,7 +115,8 @@ private HoodieWriteConfig.Builder makeHoodieClientConfigBuilder() { return HoodieWriteConfig.newBuilder() .withEngineType(EngineType.JAVA) .withPath(basePath) - .withSchema(SCHEMA.toString()); + .withSchema(SCHEMA.toString()) + .withMarkersType(MarkerType.DIRECT.name()); } @Test @@ -413,6 +415,7 @@ public void testInsertUpsertWithHoodieAvroPayload() throws Exception { .withEngineType(EngineType.JAVA) .withPath(basePath) .withSchema(schema.toString()) + .withMarkersType(MarkerType.DIRECT.name()) .withStorageConfig(HoodieStorageConfig.newBuilder() .parquetMaxFileSize(1000 * 1024).hfileMaxFileSize(1000 * 1024).build()) .build(); diff --git a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/functional/TestHoodieClientOnCopyOnWriteStorage.java b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/functional/TestHoodieClientOnCopyOnWriteStorage.java index 86d18fe28b006..6befb574868c8 100644 --- a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/functional/TestHoodieClientOnCopyOnWriteStorage.java +++ b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/functional/TestHoodieClientOnCopyOnWriteStorage.java @@ -47,6 +47,7 @@ import org.apache.hudi.common.model.IOType; import org.apache.hudi.common.model.WriteOperationType; import org.apache.hudi.common.table.HoodieTableMetaClient; +import org.apache.hudi.common.table.marker.MarkerType; import org.apache.hudi.common.table.timeline.HoodieActiveTimeline; import org.apache.hudi.common.table.timeline.HoodieInstant; import org.apache.hudi.common.table.timeline.HoodieTimeline; @@ -495,7 +496,7 @@ void assertNodupesInPartition(List records) { @ParameterizedTest @MethodSource("populateMetaFieldsParams") public void testUpserts(boolean populateMetaFields) throws Exception { - HoodieWriteConfig.Builder cfgBuilder = getConfigBuilder().withRollbackUsingMarkers(true); + HoodieWriteConfig.Builder cfgBuilder = getConfigBuilder().withRollbackUsingMarkers(true).withMarkersType(MarkerType.DIRECT.name()); addConfigsForPopulateMetaFields(cfgBuilder, populateMetaFields); testUpsertsInternal(cfgBuilder.build(), SparkRDDWriteClient::upsert, false); } @@ -506,7 +507,7 @@ public void testUpserts(boolean populateMetaFields) throws Exception { @ParameterizedTest @MethodSource("populateMetaFieldsParams") public void testUpsertsPrepped(boolean populateMetaFields) throws Exception { - HoodieWriteConfig.Builder cfgBuilder = getConfigBuilder().withRollbackUsingMarkers(true); + HoodieWriteConfig.Builder cfgBuilder = getConfigBuilder().withRollbackUsingMarkers(true).withMarkersType(MarkerType.DIRECT.name()); addConfigsForPopulateMetaFields(cfgBuilder, populateMetaFields); testUpsertsInternal(cfgBuilder.build(), SparkRDDWriteClient::upsertPreppedRecords, true); } @@ -524,8 +525,11 @@ private void testUpsertsInternal(HoodieWriteConfig config, // Force using older timeline layout HoodieWriteConfig hoodieWriteConfig = getConfigBuilder(HoodieFailedWritesCleaningPolicy.LAZY) .withRollbackUsingMarkers(true) + .withMarkersType(MarkerType.DIRECT.name()) // in this test, we directly operate with HoodieMergeHandle and so timeline server is not available for such + // operations .withProps(config.getProps()).withTimelineLayoutVersion( VERSION_0).build(); + LOG.warn("Enabled " + hoodieWriteConfig.isEmbeddedTimelineServerEnabled() + ", " + config.isEmbeddedTimelineServerEnabled()); HoodieTableMetaClient.withPropertyBuilder() .fromMetaClient(metaClient) @@ -561,8 +565,10 @@ private void testUpsertsInternal(HoodieWriteConfig config, 0, 150, config.populateMetaFields()); // Now simulate an upgrade and perform a restore operation - HoodieWriteConfig newConfig = getConfigBuilder().withProps(config.getProps()).withTimelineLayoutVersion( - TimelineLayoutVersion.CURR_VERSION).build(); + HoodieWriteConfig newConfig = getConfigBuilder().withProps(config.getProps()) + .withTimelineLayoutVersion(TimelineLayoutVersion.CURR_VERSION) + .withMarkersType(MarkerType.DIRECT.name()).build(); + LOG.warn("timeline server " + newConfig.isEmbeddedTimelineServerEnabled()); client = getHoodieWriteClient(newConfig); client.restoreToInstant("004"); @@ -2014,7 +2020,7 @@ public void testConsistencyCheckDuringFinalize(boolean enableOptimisticConsisten HoodieTableMetaClient metaClient = HoodieTableMetaClient.builder().setConf(hadoopConf).setBasePath(basePath).build(); String instantTime = "000"; HoodieWriteConfig cfg = getConfigBuilder().withAutoCommit(false).withConsistencyGuardConfig(ConsistencyGuardConfig.newBuilder() - .withEnableOptimisticConsistencyGuard(enableOptimisticConsistencyGuard).build()).build(); + .withEnableOptimisticConsistencyGuard(enableOptimisticConsistencyGuard).build()).withMarkersType(MarkerType.DIRECT.name()).build(); SparkRDDWriteClient client = getHoodieWriteClient(cfg); Pair> result = testConsistencyCheck(metaClient, instantTime, enableOptimisticConsistencyGuard); @@ -2045,10 +2051,12 @@ private void testRollbackAfterConsistencyCheckFailureUsingFileList(boolean rollb properties = getPropertiesForKeyGen(); } - HoodieWriteConfig cfg = !enableOptimisticConsistencyGuard ? getConfigBuilder().withRollbackUsingMarkers(rollbackUsingMarkers).withAutoCommit(false) + HoodieWriteConfig cfg = !enableOptimisticConsistencyGuard ? getConfigBuilder().withRollbackUsingMarkers(rollbackUsingMarkers) + .withAutoCommit(false).withMarkersType(MarkerType.DIRECT.name()) .withConsistencyGuardConfig(ConsistencyGuardConfig.newBuilder().withConsistencyCheckEnabled(true) .withMaxConsistencyCheckIntervalMs(1).withInitialConsistencyCheckIntervalMs(1).withEnableOptimisticConsistencyGuard(enableOptimisticConsistencyGuard).build()).build() : getConfigBuilder().withRollbackUsingMarkers(rollbackUsingMarkers).withAutoCommit(false) + .withMarkersType(MarkerType.DIRECT.name()) .withConsistencyGuardConfig(ConsistencyGuardConfig.newBuilder() .withConsistencyCheckEnabled(true) .withEnableOptimisticConsistencyGuard(enableOptimisticConsistencyGuard) @@ -2280,10 +2288,12 @@ private Pair> testConsistencyCheck(HoodieTableMetaCli HoodieWriteConfig cfg = !enableOptimisticConsistencyGuard ? (getConfigBuilder().withAutoCommit(false) .withConsistencyGuardConfig(ConsistencyGuardConfig.newBuilder().withConsistencyCheckEnabled(true) .withMaxConsistencyCheckIntervalMs(1).withInitialConsistencyCheckIntervalMs(1).withEnableOptimisticConsistencyGuard(enableOptimisticConsistencyGuard).build()) + .withMarkersType(MarkerType.DIRECT.name()) .build()) : (getConfigBuilder().withAutoCommit(false) .withConsistencyGuardConfig(ConsistencyGuardConfig.newBuilder().withConsistencyCheckEnabled(true) .withEnableOptimisticConsistencyGuard(enableOptimisticConsistencyGuard) .withOptimisticConsistencyGuardSleepTimeMs(1).build()) + .withMarkersType(MarkerType.DIRECT.name()) .build()); SparkRDDWriteClient client = getHoodieWriteClient(cfg); @@ -2448,6 +2458,7 @@ private HoodieWriteConfig getParallelWritingWriteConfig(HoodieFailedWritesCleani .withTimelineLayoutVersion(1) .withHeartbeatIntervalInMs(3 * 1000) .withAutoCommit(false) + .withMarkersType(MarkerType.DIRECT.name()) .withProperties(populateMetaFields ? new Properties() : getPropertiesForKeyGen()).build(); } diff --git a/hudi-flink/src/main/java/org/apache/hudi/util/StreamerUtil.java b/hudi-flink/src/main/java/org/apache/hudi/util/StreamerUtil.java index 4b9a516106426..e6b9cc46116da 100644 --- a/hudi-flink/src/main/java/org/apache/hudi/util/StreamerUtil.java +++ b/hudi-flink/src/main/java/org/apache/hudi/util/StreamerUtil.java @@ -31,6 +31,7 @@ import org.apache.hudi.common.table.HoodieTableConfig; import org.apache.hudi.common.table.HoodieTableMetaClient; import org.apache.hudi.common.table.log.HoodieLogFormat; +import org.apache.hudi.common.table.marker.MarkerType; import org.apache.hudi.common.table.timeline.HoodieActiveTimeline; import org.apache.hudi.common.table.timeline.HoodieInstant; import org.apache.hudi.common.table.timeline.HoodieTimeline; @@ -167,6 +168,7 @@ public static HoodieWriteConfig getHoodieClientConfig( .withPath(conf.getString(FlinkOptions.PATH)) .combineInput(conf.getBoolean(FlinkOptions.PRE_COMBINE), true) .withMergeAllowDuplicateOnInserts(OptionsResolver.insertClustering(conf)) + .withMarkersType(MarkerType.DIRECT.name()) .withCompactionConfig( HoodieCompactionConfig.newBuilder() .withPayloadClass(conf.getString(FlinkOptions.PAYLOAD_CLASS_NAME)) diff --git a/hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/hudi/HoodieStreamingSink.scala b/hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/hudi/HoodieStreamingSink.scala index 6e736d225a523..3070ca7c514ea 100644 --- a/hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/hudi/HoodieStreamingSink.scala +++ b/hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/hudi/HoodieStreamingSink.scala @@ -18,16 +18,17 @@ package org.apache.hudi import java.lang import java.util.function.Function - import org.apache.hudi.async.{AsyncClusteringService, AsyncCompactService, SparkStreamingAsyncClusteringService, SparkStreamingAsyncCompactService} import org.apache.hudi.client.SparkRDDWriteClient import org.apache.hudi.client.common.HoodieSparkEngineContext import org.apache.hudi.common.model.HoodieRecordPayload +import org.apache.hudi.common.table.marker.MarkerType import org.apache.hudi.common.table.{HoodieTableConfig, HoodieTableMetaClient} import org.apache.hudi.common.table.timeline.HoodieInstant.State import org.apache.hudi.common.table.timeline.{HoodieInstant, HoodieTimeline} import org.apache.hudi.common.util.CompactionUtils import org.apache.hudi.common.util.ClusteringUtils +import org.apache.hudi.config.HoodieWriteConfig import org.apache.hudi.exception.HoodieCorruptedDataException import org.apache.log4j.LogManager import org.apache.spark.api.java.JavaSparkContext @@ -78,11 +79,14 @@ class HoodieStreamingSink(sqlContext: SQLContext, log.error("Async clustering service shutdown unexpectedly") throw new IllegalStateException("Async clustering service shutdown unexpectedly") } + // override to use DIRECT style markers. In Structured streaming, timeline server is closed after first micro-batch + // and subsequent micro-batches does not have timeline server running. hence we can't use TIMELINE server based markers. + val updatedOptions = options.updated(HoodieWriteConfig.MARKERS_TYPE.key(), MarkerType.DIRECT.name()) retry(retryCnt, retryIntervalMs)( Try( HoodieSparkSqlWriter.write( - sqlContext, mode, options, data, hoodieTableConfig, writeClient, Some(triggerAsyncCompactor), Some(triggerAsyncClustering)) + sqlContext, mode, updatedOptions, data, hoodieTableConfig, writeClient, Some(triggerAsyncCompactor), Some(triggerAsyncClustering)) ) match { case Success((true, commitOps, compactionInstantOps, clusteringInstant, client, tableConfig)) => log.info(s"Micro batch id=$batchId succeeded" diff --git a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestStructuredStreaming.scala b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestStructuredStreaming.scala index 2e704e17e59c3..1bcb365540734 100644 --- a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestStructuredStreaming.scala +++ b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestStructuredStreaming.scala @@ -51,8 +51,7 @@ class TestStructuredStreaming extends HoodieClientTestBase { DataSourceWriteOptions.RECORDKEY_FIELD.key -> "_row_key", DataSourceWriteOptions.PARTITIONPATH_FIELD.key -> "partition", DataSourceWriteOptions.PRECOMBINE_FIELD.key -> "timestamp", - HoodieWriteConfig.TBL_NAME.key -> "hoodie_test", - HoodieWriteConfig.MARKERS_TYPE.key() -> MarkerType.DIRECT.name() + HoodieWriteConfig.TBL_NAME.key -> "hoodie_test" ) @BeforeEach override def setUp() { diff --git a/hudi-utilities/src/test/java/org/apache/hudi/utilities/functional/TestHoodieDeltaStreamer.java b/hudi-utilities/src/test/java/org/apache/hudi/utilities/functional/TestHoodieDeltaStreamer.java index 9cf040ceaa5b0..1ca3c4e0368a7 100644 --- a/hudi-utilities/src/test/java/org/apache/hudi/utilities/functional/TestHoodieDeltaStreamer.java +++ b/hudi-utilities/src/test/java/org/apache/hudi/utilities/functional/TestHoodieDeltaStreamer.java @@ -1222,7 +1222,6 @@ private void prepareParquetDFSSource(boolean useSchemaProvider, boolean hasTrans } parquetProps.setProperty("include", "base.properties"); - parquetProps.setProperty("hoodie.embed.timeline.server", "false"); parquetProps.setProperty("hoodie.datasource.write.recordkey.field", "_row_key"); parquetProps.setProperty("hoodie.datasource.write.partitionpath.field", "not_there"); if (useSchemaProvider) { @@ -1253,7 +1252,6 @@ private void testORCDFSSource(boolean useSchemaProvider, List transforme // Properties used for testing delta-streamer with orc source orcProps.setProperty("include", "base.properties"); - orcProps.setProperty("hoodie.embed.timeline.server","false"); orcProps.setProperty("hoodie.datasource.write.recordkey.field", "_row_key"); orcProps.setProperty("hoodie.datasource.write.partitionpath.field", "not_there"); if (useSchemaProvider) { @@ -1280,7 +1278,6 @@ private void prepareJsonKafkaDFSSource(String propsFileName, String autoResetVal TypedProperties props = new TypedProperties(); populateAllCommonProps(props, dfsBasePath, testUtils.brokerAddress()); props.setProperty("include", "base.properties"); - props.setProperty("hoodie.embed.timeline.server", "false"); props.setProperty("hoodie.datasource.write.recordkey.field", "_row_key"); props.setProperty("hoodie.datasource.write.partitionpath.field", "not_there"); props.setProperty("hoodie.deltastreamer.source.dfs.root", JSON_KAFKA_SOURCE_ROOT); From 6f83276f81a57d93fdb0bd2a85bc16ee4706e32f Mon Sep 17 00:00:00 2001 From: Sivabalan Narayanan Date: Thu, 18 Nov 2021 17:41:10 -0500 Subject: [PATCH 4/4] fixing more tests --- .../client/functional/TestHoodieBackedMetadata.java | 2 ++ .../functional/TestHoodieMetadataBootstrap.java | 2 ++ .../io/storage/row/TestHoodieRowCreateHandle.java | 6 ++++-- .../test/java/org/apache/hudi/table/TestCleaner.java | 3 ++- .../TestHoodieSparkMergeOnReadTableRollback.java | 12 +++++++++--- .../functional/TestHoodieDeltaStreamer.java | 6 ++++++ 6 files changed, 25 insertions(+), 6 deletions(-) diff --git a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/functional/TestHoodieBackedMetadata.java b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/functional/TestHoodieBackedMetadata.java index 59e12a1515ad5..f9949ccfc135b 100644 --- a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/functional/TestHoodieBackedMetadata.java +++ b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/functional/TestHoodieBackedMetadata.java @@ -41,6 +41,7 @@ import org.apache.hudi.common.table.HoodieTableConfig; import org.apache.hudi.common.table.HoodieTableMetaClient; import org.apache.hudi.common.table.HoodieTableVersion; +import org.apache.hudi.common.table.marker.MarkerType; import org.apache.hudi.common.table.timeline.HoodieActiveTimeline; import org.apache.hudi.common.table.timeline.HoodieInstant; import org.apache.hudi.common.table.timeline.HoodieTimeline; @@ -1265,6 +1266,7 @@ private HoodieWriteConfig getSmallInsertWriteConfig(int insertSplitSize, String .hfileMaxFileSize(dataGen.getEstimatedFileSizeInBytes(200)) .parquetMaxFileSize(dataGen.getEstimatedFileSizeInBytes(200)).build()) .withMergeAllowDuplicateOnInserts(mergeAllowDuplicateInserts) + .withMarkersType(MarkerType.DIRECT.name()) .build(); } diff --git a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/functional/TestHoodieMetadataBootstrap.java b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/functional/TestHoodieMetadataBootstrap.java index 12c8410c35e02..593e79d5543f7 100644 --- a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/functional/TestHoodieMetadataBootstrap.java +++ b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/client/functional/TestHoodieMetadataBootstrap.java @@ -21,6 +21,7 @@ import org.apache.hudi.common.config.HoodieMetadataConfig; import org.apache.hudi.common.model.HoodieCommitMetadata; import org.apache.hudi.common.model.HoodieTableType; +import org.apache.hudi.common.table.marker.MarkerType; import org.apache.hudi.common.testutils.HoodieMetadataTestTable; import org.apache.hudi.common.testutils.HoodieTestDataGenerator; import org.apache.hudi.common.testutils.HoodieTestTable; @@ -242,6 +243,7 @@ private HoodieWriteConfig getWriteConfig(int minArchivalCommits, int maxArchival .withSchema(HoodieTestDataGenerator.TRIP_EXAMPLE_SCHEMA).withParallelism(2, 2) .withCompactionConfig(HoodieCompactionConfig.newBuilder().retainCommits(1).archiveCommitsWith(minArchivalCommits, maxArchivalCommits).build()) .withMetadataConfig(HoodieMetadataConfig.newBuilder().enable(false).build()) + .withMarkersType(MarkerType.DIRECT.name()) .forTable("test-trip-table").build(); } diff --git a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/io/storage/row/TestHoodieRowCreateHandle.java b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/io/storage/row/TestHoodieRowCreateHandle.java index 76a91ef124bb7..2e2a2cae52400 100644 --- a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/io/storage/row/TestHoodieRowCreateHandle.java +++ b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/io/storage/row/TestHoodieRowCreateHandle.java @@ -22,6 +22,7 @@ import org.apache.hudi.common.config.HoodieMetadataConfig; import org.apache.hudi.common.model.HoodieRecord; import org.apache.hudi.common.model.HoodieWriteStat; +import org.apache.hudi.common.table.marker.MarkerType; import org.apache.hudi.common.testutils.HoodieTestDataGenerator; import org.apache.hudi.config.HoodieWriteConfig; import org.apache.hudi.exception.HoodieInsertException; @@ -75,7 +76,7 @@ public void tearDown() throws Exception { @Test public void testRowCreateHandle() throws Exception { // init config and table - HoodieWriteConfig cfg = SparkDatasetTestUtils.getConfigBuilder(basePath).build(); + HoodieWriteConfig cfg = SparkDatasetTestUtils.getConfigBuilder(basePath).withMarkersType(MarkerType.DIRECT.name()).build(); HoodieTable table = HoodieSparkTable.create(cfg, context, metaClient); List fileNames = new ArrayList<>(); List fileAbsPaths = new ArrayList<>(); @@ -116,7 +117,7 @@ public void testRowCreateHandle() throws Exception { @Test public void testGlobalFailure() throws Exception { // init config and table - HoodieWriteConfig cfg = SparkDatasetTestUtils.getConfigBuilder(basePath).build(); + HoodieWriteConfig cfg = SparkDatasetTestUtils.getConfigBuilder(basePath).withMarkersType(MarkerType.DIRECT.name()).build(); HoodieTable table = HoodieSparkTable.create(cfg, context, metaClient); String partitionPath = HoodieTestDataGenerator.DEFAULT_PARTITION_PATHS[0]; @@ -170,6 +171,7 @@ public void testGlobalFailure() throws Exception { public void testInstantiationFailure() throws IOException { // init config and table HoodieWriteConfig cfg = SparkDatasetTestUtils.getConfigBuilder(basePath).withMetadataConfig(HoodieMetadataConfig.newBuilder().enable(false).build()) + .withMarkersType(MarkerType.DIRECT.name()) .withPath("/dummypath/abc/").build(); HoodieTable table = HoodieSparkTable.create(cfg, context, metaClient); diff --git a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/table/TestCleaner.java b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/table/TestCleaner.java index 2305d7bdeb4d1..23c39a319fa14 100644 --- a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/table/TestCleaner.java +++ b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/table/TestCleaner.java @@ -54,6 +54,7 @@ import org.apache.hudi.common.model.IOType; import org.apache.hudi.common.model.WriteOperationType; import org.apache.hudi.common.table.HoodieTableMetaClient; +import org.apache.hudi.common.table.marker.MarkerType; import org.apache.hudi.common.table.timeline.HoodieActiveTimeline; import org.apache.hudi.common.table.timeline.HoodieInstant; import org.apache.hudi.common.table.timeline.HoodieInstant.State; @@ -1300,7 +1301,7 @@ public void testCleanMarkerDataFilesOnRollback() throws Exception { final int numTempFilesBefore = testTable.listAllFilesInTempFolder().length; assertEquals(10, numTempFilesBefore, "Some marker files are created."); - HoodieWriteConfig config = HoodieWriteConfig.newBuilder().withPath(basePath).build(); + HoodieWriteConfig config = HoodieWriteConfig.newBuilder().withPath(basePath).withMarkersType(MarkerType.DIRECT.name()).build(); metaClient = HoodieTableMetaClient.reload(metaClient); HoodieTable table = HoodieSparkTable.create(config, context, metaClient); table.getActiveTimeline().transitionRequestedToInflight( diff --git a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/table/functional/TestHoodieSparkMergeOnReadTableRollback.java b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/table/functional/TestHoodieSparkMergeOnReadTableRollback.java index 6bbb0f655bb8e..1775916eab1dd 100644 --- a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/table/functional/TestHoodieSparkMergeOnReadTableRollback.java +++ b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/table/functional/TestHoodieSparkMergeOnReadTableRollback.java @@ -29,6 +29,7 @@ import org.apache.hudi.common.model.HoodieTableType; import org.apache.hudi.common.table.HoodieTableConfig; import org.apache.hudi.common.table.HoodieTableMetaClient; +import org.apache.hudi.common.table.marker.MarkerType; import org.apache.hudi.common.table.timeline.HoodieInstant; import org.apache.hudi.common.table.timeline.HoodieTimeline; import org.apache.hudi.common.table.view.HoodieTableFileSystemView; @@ -294,7 +295,9 @@ void testRollbackWithDeltaAndCompactionCommit(boolean rollbackUsingMarkers) thro @Test void testMultiRollbackWithDeltaAndCompactionCommit() throws Exception { boolean populateMetaFields = true; - HoodieWriteConfig.Builder cfgBuilder = getConfigBuilder(false).withMetadataConfig(HoodieMetadataConfig.newBuilder().enable(false).build()); + HoodieWriteConfig.Builder cfgBuilder = getConfigBuilder(false) + .withMetadataConfig(HoodieMetadataConfig.newBuilder().enable(false).build()) + .withMarkersType(MarkerType.DIRECT.name()); addConfigsForPopulateMetaFields(cfgBuilder, populateMetaFields); HoodieWriteConfig cfg = cfgBuilder.build(); @@ -344,7 +347,9 @@ void testMultiRollbackWithDeltaAndCompactionCommit() throws Exception { newCommitTime = "002"; // WriteClient with custom config (disable small file handling) HoodieWriteConfig smallFileWriteConfig = getHoodieWriteConfigWithSmallFileHandlingOffBuilder(populateMetaFields) - .withMetadataConfig(HoodieMetadataConfig.newBuilder().enable(false).build()).build(); + .withMetadataConfig(HoodieMetadataConfig.newBuilder().enable(false).build()) + .withMarkersType(MarkerType.DIRECT.name()) + .build(); try (SparkRDDWriteClient nClient = getHoodieWriteClient(smallFileWriteConfig)) { nClient.startCommitWithTime(newCommitTime); @@ -483,7 +488,8 @@ void testInsertsGeneratedIntoLogFilesRollback(boolean rollbackUsingMarkers) thro HoodieTestDataGenerator dataGen = new HoodieTestDataGenerator(); // insert 100 records // Setting IndexType to be InMemory to simulate Global Index nature - HoodieWriteConfig config = getConfigBuilder(false, rollbackUsingMarkers, HoodieIndex.IndexType.INMEMORY).build(); + HoodieWriteConfig config = getConfigBuilder(false, rollbackUsingMarkers, HoodieIndex.IndexType.INMEMORY) + .withMarkersType(MarkerType.DIRECT.name()).build(); try (SparkRDDWriteClient writeClient = getHoodieWriteClient(config)) { String newCommitTime = "100"; diff --git a/hudi-utilities/src/test/java/org/apache/hudi/utilities/functional/TestHoodieDeltaStreamer.java b/hudi-utilities/src/test/java/org/apache/hudi/utilities/functional/TestHoodieDeltaStreamer.java index 1ca3c4e0368a7..1313d45149e09 100644 --- a/hudi-utilities/src/test/java/org/apache/hudi/utilities/functional/TestHoodieDeltaStreamer.java +++ b/hudi-utilities/src/test/java/org/apache/hudi/utilities/functional/TestHoodieDeltaStreamer.java @@ -34,6 +34,7 @@ import org.apache.hudi.common.table.HoodieTableConfig; import org.apache.hudi.common.table.HoodieTableMetaClient; import org.apache.hudi.common.table.TableSchemaResolver; +import org.apache.hudi.common.table.marker.MarkerType; import org.apache.hudi.common.table.timeline.HoodieInstant; import org.apache.hudi.common.table.timeline.HoodieTimeline; import org.apache.hudi.common.testutils.HoodieTestDataGenerator; @@ -41,6 +42,7 @@ import org.apache.hudi.common.util.StringUtils; import org.apache.hudi.config.HoodieClusteringConfig; import org.apache.hudi.config.HoodieCompactionConfig; +import org.apache.hudi.config.HoodieWriteConfig; import org.apache.hudi.exception.HoodieException; import org.apache.hudi.exception.TableNotFoundException; import org.apache.hudi.hive.HiveSyncConfig; @@ -1224,6 +1226,8 @@ private void prepareParquetDFSSource(boolean useSchemaProvider, boolean hasTrans parquetProps.setProperty("include", "base.properties"); parquetProps.setProperty("hoodie.datasource.write.recordkey.field", "_row_key"); parquetProps.setProperty("hoodie.datasource.write.partitionpath.field", "not_there"); + parquetProps.setProperty("hoodie.embed.timeline.server", "false"); + parquetProps.setProperty(HoodieWriteConfig.MARKERS_TYPE.key(), MarkerType.DIRECT.name()); if (useSchemaProvider) { parquetProps.setProperty("hoodie.deltastreamer.schemaprovider.source.schema.file", dfsBasePath + "/" + sourceSchemaFile); if (hasTransformer) { @@ -1280,6 +1284,8 @@ private void prepareJsonKafkaDFSSource(String propsFileName, String autoResetVal props.setProperty("include", "base.properties"); props.setProperty("hoodie.datasource.write.recordkey.field", "_row_key"); props.setProperty("hoodie.datasource.write.partitionpath.field", "not_there"); + props.setProperty("hoodie.embed.timeline.server", "false"); + props.setProperty(HoodieWriteConfig.MARKERS_TYPE.key(), MarkerType.DIRECT.name()); props.setProperty("hoodie.deltastreamer.source.dfs.root", JSON_KAFKA_SOURCE_ROOT); props.setProperty("hoodie.deltastreamer.source.kafka.topic", topicName); props.setProperty("hoodie.deltastreamer.source.kafka.checkpoint.type", kafkaCheckpointType);