This is an automated email from the ASF dual-hosted git repository. voonhous pushed a commit to branch release-1.2.1 in repository https://gitbox.apache.org/repos/asf/hudi.git
commit 96bb7b2ccc93d13b773d0397d3f08df8f50a1e2f Author: Joy <[email protected]> AuthorDate: Sat Jul 11 07:27:59 2026 +0800 fix(flink): prevent data loss on global failover for streaming writes (#19237) * fix(flink): prevent data loss on global failover for streaming writes --------- Co-authored-by: jiangyu84 <[email protected]> Co-authored-by: danny0405 <[email protected]> (cherry picked from commit be8a54c4ebdaf287993d24fefba60a58952a9c9e) --- .../apache/hudi/client/BaseHoodieWriteClient.java | 23 +++++ .../metadata/HoodieBackedTableMetadataWriter.java | 9 ++ .../TestHoodieBackedTableMetadataWriter.java | 27 ++++++ .../client/FlinkStreamingMetadataWriteHandler.java | 30 ++++++- .../apache/hudi/client/HoodieFlinkWriteClient.java | 25 ++++++ .../apache/hudi/client/TestFlinkWriteClient.java | 98 +++++++++++++++++++++- .../hudi/sink/StreamWriteOperatorCoordinator.java | 28 ++++--- .../org/apache/hudi/sink/utils/EventBuffers.java | 10 +-- .../sink/TestStreamWriteOperatorCoordinator.java | 94 ++++++++++++++------- .../org/apache/hudi/sink/TestWriteMergeOnRead.java | 35 ++++++++ 10 files changed, 327 insertions(+), 52 deletions(-) diff --git a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/BaseHoodieWriteClient.java b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/BaseHoodieWriteClient.java index f5d050d88098..73f5de7b94ef 100644 --- a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/BaseHoodieWriteClient.java +++ b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/client/BaseHoodieWriteClient.java @@ -679,6 +679,29 @@ public abstract class BaseHoodieWriteClient<T, I, K, O> extends BaseHoodieClient } } + /** + * Performs post-commit cleanup when the instant is already completed and commit metadata is not + * available to invoke the regular post-commit hook. This can happen while recovering a streaming + * metadata-table write after failover. The table is recreated from the write configuration so its + * marker directory can still be removed, and the heartbeat is always stopped even if marker cleanup + * fails. + * + * @param instantTime the completed instant to clean up + */ + public void postCommit(String instantTime) { + try { + HoodieTable table = createTable(config); + context.setJobStatus(this.getClass().getSimpleName(), "Cleaning up marker directories for commit " + instantTime + " in table " + + config.getTableName()); + // Delete the marker directory for the instant. + WriteMarkersFactory.get(config.getMarkersType(), table, instantTime) + .quietDeleteMarkerDir(context, config.getMarkersDeleteParallelism()); + metrics.updateTableServiceInstantMetrics(table.getActiveTimeline()); + } finally { + this.heartbeatClient.stop(instantTime); + } + } + /** * Triggers cleaning and archival for the table of interest. This method is called outside of locks. So, internal callers should ensure they acquire lock whereever applicable. * @param table instance of {@link HoodieTable} of interest. diff --git a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/HoodieBackedTableMetadataWriter.java b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/HoodieBackedTableMetadataWriter.java index df61ed996b1f..2185418fe8eb 100644 --- a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/HoodieBackedTableMetadataWriter.java +++ b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/HoodieBackedTableMetadataWriter.java @@ -1469,6 +1469,12 @@ public abstract class HoodieBackedTableMetadataWriter<I, O> implements HoodieTab @Override public void completeStreamingCommit(String instantTime, HoodieEngineContext context, List<HoodieWriteStat> partialWriteStats, HoodieCommitMetadata metadata) { + if (metadataMetaClient.getActiveTimeline().filterCompletedInstants().containsInstant(instantTime)) { + LOG.info("Skipping streaming metadata commit completion for already completed instant {}", instantTime); + getWriteClient().postCommit(instantTime); + return; + } + List<HoodieWriteStat> allWriteStats = new ArrayList<>(partialWriteStats); // update metadata for left over partitions which does not have streaming writes support. allWriteStats.addAll(prepareAndWriteToNonStreamingPartitions(metadata, instantTime).map(WriteStatus::getStat).collectAsList()); @@ -1907,6 +1913,9 @@ public abstract class HoodieBackedTableMetadataWriter<I, O> implements HoodieTab public void close() throws Exception { if (metadata != null) { metadata.close(); + // Keep the closed reader reference: guarded update paths use its presence to proceed and + // mayBeReinitMetadataReader() detects the closed file-system view and reopens the reader. + // Nullifying it here would silently skip subsequent metadata updates, such as rollbacks. } if (writeClient != null) { writeClient.close(); diff --git a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/metadata/TestHoodieBackedTableMetadataWriter.java b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/metadata/TestHoodieBackedTableMetadataWriter.java index 7fcc2ce9a5d8..449578f3968d 100644 --- a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/metadata/TestHoodieBackedTableMetadataWriter.java +++ b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/metadata/TestHoodieBackedTableMetadataWriter.java @@ -23,6 +23,7 @@ import org.apache.hudi.common.config.HoodieMetadataConfig; import org.apache.hudi.common.config.HoodieTableServiceManagerConfig; import org.apache.hudi.common.data.HoodieData; import org.apache.hudi.common.engine.HoodieEngineContext; +import org.apache.hudi.common.model.HoodieCommitMetadata; import org.apache.hudi.common.model.HoodieFailedWritesCleaningPolicy; import org.apache.hudi.common.model.HoodieRecord; import org.apache.hudi.common.schema.HoodieSchema; @@ -47,6 +48,7 @@ import org.junit.jupiter.params.provider.MethodSource; import org.mockito.MockedStatic; import java.util.ArrayList; +import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -63,6 +65,7 @@ import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyInt; import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.CALLS_REAL_METHODS; import static org.mockito.Mockito.RETURNS_DEEP_STUBS; import static org.mockito.Mockito.doCallRealMethod; import static org.mockito.Mockito.doThrow; @@ -71,6 +74,7 @@ import static org.mockito.Mockito.mockConstruction; import static org.mockito.Mockito.mockStatic; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoMoreInteractions; import static org.mockito.Mockito.when; class TestHoodieBackedTableMetadataWriter { @@ -89,6 +93,29 @@ class TestHoodieBackedTableMetadataWriter { when(metadataConfig.getMaxReaderBufferSize()).thenReturn(1024); } + @Test + void completeStreamingCommitSkipsAlreadyCompletedMetadataInstant() { + String instantTime = "20260709120000000"; + HoodieBackedTableMetadataWriter<List<HoodieRecord>, List<?>> metadataWriter = + mock(HoodieBackedTableMetadataWriter.class, CALLS_REAL_METHODS); + HoodieEngineContext engineContext = mock(HoodieEngineContext.class); + HoodieTableMetaClient metadataMetaClient = mock(HoodieTableMetaClient.class); + HoodieActiveTimeline activeTimeline = mock(HoodieActiveTimeline.class); + HoodieTimeline completedTimeline = mock(HoodieTimeline.class); + BaseHoodieWriteClient writeClient = mock(BaseHoodieWriteClient.class); + + metadataWriter.metadataMetaClient = metadataMetaClient; + when(metadataMetaClient.getActiveTimeline()).thenReturn(activeTimeline); + when(activeTimeline.filterCompletedInstants()).thenReturn(completedTimeline); + when(completedTimeline.containsInstant(instantTime)).thenReturn(true); + when(metadataWriter.initializeWriteClient()).thenReturn(writeClient); + + metadataWriter.completeStreamingCommit(instantTime, engineContext, Collections.emptyList(), mock(HoodieCommitMetadata.class)); + + verify(writeClient).postCommit(instantTime); + verifyNoMoreInteractions(writeClient); + } + @ParameterizedTest @CsvSource(value = { "true,true,false,true", diff --git a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/FlinkStreamingMetadataWriteHandler.java b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/FlinkStreamingMetadataWriteHandler.java index f500f99e0635..33c34b8c00a1 100644 --- a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/FlinkStreamingMetadataWriteHandler.java +++ b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/FlinkStreamingMetadataWriteHandler.java @@ -23,6 +23,7 @@ import org.apache.hudi.common.model.HoodieRecord; import org.apache.hudi.common.util.Option; import org.apache.hudi.common.util.ValidationUtils; import org.apache.hudi.exception.HoodieException; +import org.apache.hudi.metadata.FlinkHoodieBackedTableMetadataWriter; import org.apache.hudi.metadata.HoodieTableMetadataWriter; import org.apache.hudi.table.HoodieTable; @@ -83,6 +84,27 @@ public class FlinkStreamingMetadataWriteHandler extends StreamingMetadataWriteHa metadataWriterOpt.get().startCommit(instantTime); } + /** + * Start only the metadata table heartbeat for an existing streaming write instant. + * + * <p>This is used by coordinator recommit where the metadata table instant and + * the streaming index files may already exist and must not be rolled back. + * + * @param instantTime The instant time + * @param table The hoodie table + */ + public void startHeartbeat(String instantTime, HoodieTable table) { + Option<HoodieTableMetadataWriter> metadataWriterOpt = getMetadataWriter(instantTime, table); + ValidationUtils.checkState(metadataWriterOpt.isPresent(), + "Should not be reachable. Metadata Writer should have been instantiated by now"); + ValidationUtils.checkState(metadataWriterOpt.get() instanceof FlinkHoodieBackedTableMetadataWriter, + "Flink streaming metadata writes expect a Flink metadata writer"); + FlinkHoodieBackedTableMetadataWriter metadataWriter = (FlinkHoodieBackedTableMetadataWriter) metadataWriterOpt.get(); + if (metadataWriter.getWriteClient().getConfig().getFailedWritesCleanPolicy().isLazy()) { + metadataWriter.getWriteClient().getHeartbeatClient().start(instantTime); + } + } + /** * Clean resources after streaming write to the metadata table in index write function or stop * heartbeat for instant in the coordinator. This method removes the metadata writer associated @@ -93,11 +115,13 @@ public class FlinkStreamingMetadataWriteHandler extends StreamingMetadataWriteHa public void cleanResources(String instantTime) { Option<HoodieTableMetadataWriter> metadataWriterOpt = this.metadataWriterMap.remove(instantTime); if (metadataWriterOpt == null || metadataWriterOpt.isEmpty()) { - log.warn("Metadata writer for {} has not been initialized, no need to stop heartbeat.", instantTime); + log.debug("Metadata writer for {} has already been closed, skip closing.", instantTime); return; } - try { - metadataWriterOpt.get().close(); + try (HoodieTableMetadataWriter metadataWriter = metadataWriterOpt.get()) { + if (metadataWriter instanceof FlinkHoodieBackedTableMetadataWriter) { + ((FlinkHoodieBackedTableMetadataWriter) metadataWriter).getWriteClient().getHeartbeatClient().stop(instantTime); + } } catch (Exception e) { throw new HoodieException("Failed to close the metadata writer for instant: " + instantTime, e); } diff --git a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/HoodieFlinkWriteClient.java b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/HoodieFlinkWriteClient.java index 228209cf1df6..51ef11b39013 100644 --- a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/HoodieFlinkWriteClient.java +++ b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/client/HoodieFlinkWriteClient.java @@ -138,6 +138,20 @@ public class HoodieFlinkWriteClient<T> } } + /** + * Restart the heartbeat for a recommitted instant. + * + * @param instantTime The instant time + */ + public void restartHeartbeat(String instantTime) { + if (getConfig().getFailedWritesCleanPolicy().isLazy()) { + getHeartbeatClient().start(instantTime); + } + if (isStreamingWriteMetadataTable) { + this.streamingMetadataWriteHandler.startHeartbeat(instantTime, getHoodieTable()); + } + } + /** * Performs streaming write operations to metadata partitions. * This method retrieves the metadata writer for the given instant time and table, @@ -629,4 +643,15 @@ public class HoodieFlinkWriteClient<T> ((MiniBatchHandle) writeHandle).closeGracefully(); } } + + /** + * Flink keeps the heartbeat active when a commit attempt fails because the coordinator may need to + * recommit the instant after failover. Successful commits stop the heartbeat from {@code postCommit}, + * while {@link #cleanResources(String)} cleans up data-table and metadata-table resources for + * instants discarded or resolved by the coordinator. Therefore this generic hook is intentionally + * a no-op. + */ + @Override + public void releaseResources(String instantTime) { + } } diff --git a/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/client/TestFlinkWriteClient.java b/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/client/TestFlinkWriteClient.java index 33ae4238962a..d0ae6a30a54a 100644 --- a/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/client/TestFlinkWriteClient.java +++ b/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/client/TestFlinkWriteClient.java @@ -19,27 +19,47 @@ package org.apache.hudi.client; +import org.apache.hudi.client.heartbeat.HoodieHeartbeatClient; +import org.apache.hudi.common.config.HoodieMetadataConfig; +import org.apache.hudi.common.engine.EngineType; +import org.apache.hudi.common.model.HoodieFailedWritesCleaningPolicy; +import org.apache.hudi.common.table.HoodieTableMetaClient; +import org.apache.hudi.config.HoodieCleanConfig; +import org.apache.hudi.config.HoodieIndexConfig; import org.apache.hudi.config.HoodieWriteConfig; +import org.apache.hudi.exception.HoodieException; +import org.apache.hudi.index.HoodieIndex; +import org.apache.hudi.table.HoodieTable; import org.apache.hudi.testutils.HoodieFlinkClientTestHarness; +import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.ValueSource; import java.io.IOException; +import java.util.concurrent.atomic.AtomicBoolean; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; public class TestFlinkWriteClient extends HoodieFlinkClientTestHarness { @BeforeEach - private void setup() throws IOException { + void setup() throws IOException { initPath(); initFileSystem(); initMetaClient(); } + @AfterEach + void teardown() throws IOException { + cleanupResources(); + } + @ParameterizedTest @ValueSource(booleans = {true, false}) public void testWriteClientAndTableServiceClientWithTimelineServer( @@ -61,4 +81,80 @@ public class TestFlinkWriteClient extends HoodieFlinkClientTestHarness { writeClient.close(); } + + @Test + public void testReleaseAndPostCommitResourcesForLazyFailedWrites() throws IOException { + HoodieWriteConfig writeConfig = HoodieWriteConfig.newBuilder() + .withPath(metaClient.getBasePath()) + .withEngineType(EngineType.FLINK) + .withCleanConfig(HoodieCleanConfig.newBuilder() + .withFailedWritesCleaningPolicy(HoodieFailedWritesCleaningPolicy.LAZY) + .build()) + .build(); + + AtomicBoolean failTableCreation = new AtomicBoolean(false); + writeClient = new HoodieFlinkWriteClient(context, writeConfig) { + @Override + protected HoodieTable createTable(HoodieWriteConfig config) { + if (failTableCreation.get()) { + throw new HoodieException("Expected table creation failure"); + } + return super.createTable(config); + } + }; + String instantTime = "20260709120000000"; + writeClient.restartHeartbeat(instantTime); + + assertTrue(HoodieHeartbeatClient.heartbeatExists( + metaClient.getStorage(), metaClient.getBasePath().toString(), instantTime)); + + writeClient.releaseResources(instantTime); + assertTrue(HoodieHeartbeatClient.heartbeatExists( + metaClient.getStorage(), metaClient.getBasePath().toString(), instantTime)); + + writeClient.postCommit(instantTime); + assertFalse(HoodieHeartbeatClient.heartbeatExists( + metaClient.getStorage(), metaClient.getBasePath().toString(), instantTime)); + + String failedPostCommitInstantTime = "20260709120000001"; + writeClient.restartHeartbeat(failedPostCommitInstantTime); + failTableCreation.set(true); + assertThrows(HoodieException.class, () -> writeClient.postCommit(failedPostCommitInstantTime)); + assertFalse(HoodieHeartbeatClient.heartbeatExists( + metaClient.getStorage(), metaClient.getBasePath().toString(), failedPostCommitInstantTime)); + } + + @Test + public void testCleanResourcesCleansMetadataTableHeartbeatForStreamingMetadataWrites() throws IOException { + HoodieWriteConfig writeConfig = HoodieWriteConfig.newBuilder() + .withPath(metaClient.getBasePath()) + .withEngineType(EngineType.FLINK) + .withIndexConfig(HoodieIndexConfig.newBuilder() + .withIndexType(HoodieIndex.IndexType.GLOBAL_RECORD_LEVEL_INDEX) + .build()) + .withMetadataConfig(HoodieMetadataConfig.newBuilder() + .enable(true) + .withStreamingWriteEnabled(true) + .withEnableGlobalRecordLevelIndex(true) + .build()) + .withCleanConfig(HoodieCleanConfig.newBuilder() + .withFailedWritesCleaningPolicy(HoodieFailedWritesCleaningPolicy.LAZY) + .build()) + .build(); + + writeClient = new HoodieFlinkWriteClient(context, writeConfig, true); + String instantTime = "20260709120000000"; + writeClient.restartHeartbeat(instantTime); + + String metadataTableBasePath = metaClient.getBasePath() + + "/" + HoodieTableMetaClient.METADATA_TABLE_FOLDER_PATH; + assertTrue(HoodieHeartbeatClient.heartbeatExists( + metaClient.getStorage(), metadataTableBasePath, instantTime)); + + writeClient.cleanResources(instantTime); + assertFalse(HoodieHeartbeatClient.heartbeatExists( + metaClient.getStorage(), metaClient.getBasePath().toString(), instantTime)); + assertFalse(HoodieHeartbeatClient.heartbeatExists( + metaClient.getStorage(), metadataTableBasePath, instantTime)); + } } diff --git a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteOperatorCoordinator.java b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteOperatorCoordinator.java index 7e9d8cfa4cf5..686aaaa9509c 100644 --- a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteOperatorCoordinator.java +++ b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteOperatorCoordinator.java @@ -252,7 +252,7 @@ public class StreamWriteOperatorCoordinator if (OptionsResolver.isMultiWriter(conf)) { initClientIds(conf); } - restoreEvents(); + restoreEvents(Long.MAX_VALUE); } catch (Throwable throwable) { log.error("Failed to start operator coordinator.", throwable); context.failJob(throwable); @@ -325,14 +325,16 @@ public class StreamWriteOperatorCoordinator @Override public void resetToCheckpoint(long checkpointID, byte[] checkpointData) { if (checkpointData != null) { - initEventBufferIfNecessary(); - this.eventBuffers.addEventsToBuffer(SerializationUtils.deserialize(checkpointData)); // resetToCheckpoint() is called in two cases: // 1. The job is restarted from state, start() will be called later. // 2. The job is recovered from global failover. The coordinator is already started, and start() will not be called again. - if (executor != null && tableState.isRecordLevelIndex) { - // use sync execution here to make sure the recommitting finishes before RLI bootstrapping - this.executor.executeSync(this::restoreEvents, "Recommit pending instants on resetting to checkpoint: %s.", checkpointID); + if (this.eventBuffers == null) { + // case1: the events restore is moved to start() since it requires the write/meta client instantiation. + initEventBuffer(); + this.eventBuffers.addEventsToBuffer(SerializationUtils.deserialize(checkpointData)); + } else { + // case2: restore the events directly. + restoreEvents(checkpointID); } } } @@ -433,10 +435,11 @@ public class StreamWriteOperatorCoordinator // Utilities // ------------------------------------------------------------------------- - private void restoreEvents() { + private void restoreEvents(long checkpointId) { if (this.eventBuffers.nonEmpty()) { - final HoodieTimeline completedTimeline = this.metaClient.getActiveTimeline().filterCompletedInstants(); + final HoodieTimeline completedTimeline = this.metaClient.reloadActiveTimeline().filterCompletedInstants(); this.eventBuffers.getEventBufferStream() + .filter(entry -> entry.getKey() < checkpointId) .forEach(entry -> recommitInstant(completedTimeline, entry.getKey(), entry.getValue().getLeft(), entry.getValue().getRight())); this.metaClient.reloadActiveTimeline(); } @@ -503,6 +506,10 @@ public class StreamWriteOperatorCoordinator if (this.eventBuffers != null) { return; } + initEventBuffer(); + } + + private void initEventBuffer() { // initialize event buffer this.eventBuffers = EventBuffers.getInstance(conf, this.parallelism); } @@ -544,13 +551,12 @@ public class StreamWriteOperatorCoordinator log.info("Recommit instant {}", instant); // Recommit should start heartbeat for lazy failed writes clean policy to avoid aborting for heartbeat expired; // The following up checkpoints would recommit the instant. - if (writeClient.getConfig().getFailedWritesCleanPolicy().isLazy()) { - writeClient.getHeartbeatClient().start(instant); - } + writeClient.restartHeartbeat(instant); return commitInstant(checkpointId, instant, bootstrapBuffer); } else { // clean the corresponding event buffer if the instant is already committed. eventBuffers.reset(checkpointId); + writeClient.cleanResources(instant); return false; } } diff --git a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/EventBuffers.java b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/EventBuffers.java index 87e1010e2227..10f89637e342 100644 --- a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/EventBuffers.java +++ b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/EventBuffers.java @@ -171,7 +171,7 @@ public class EventBuffers implements Serializable { */ public Map<Long, Pair<String, EventBuffer>> getAllCompletedEvents() { return this.eventBuffers.entrySet().stream() - .filter(entry -> entry.getValue().getRight().allEventsCompleted()) + .filter(entry -> entry.getValue().getRight().allEventsReceived()) .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); } @@ -241,14 +241,6 @@ public class EventBuffers implements Serializable { } } - /** - * Return true if there is no event sent by eager flushing from writers. - */ - public boolean allEventsCompleted() { - return Stream.concat(Arrays.stream(dataWriteEventBuffer), Arrays.stream(indexWriteEventBuffer)) - .allMatch(event -> event == null || event.isLastBatch()); - } - /** * Return true if all the events in the data write buffer are null. */ diff --git a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestStreamWriteOperatorCoordinator.java b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestStreamWriteOperatorCoordinator.java index d5903ef946dc..bb57dc45314c 100644 --- a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestStreamWriteOperatorCoordinator.java +++ b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestStreamWriteOperatorCoordinator.java @@ -20,6 +20,7 @@ package org.apache.hudi.sink; import org.apache.hudi.client.WriteStatus; import org.apache.hudi.client.heartbeat.HoodieHeartbeatClient; +import org.apache.hudi.common.fs.FSUtils; import org.apache.hudi.common.model.HoodieFailedWritesCleaningPolicy; import org.apache.hudi.common.model.HoodieTableType; import org.apache.hudi.common.model.HoodieWriteStat; @@ -106,7 +107,7 @@ public class TestStreamWriteOperatorCoordinator { @BeforeEach public void before() throws Exception { - coordinator = createCoordinator(TestConfigurations.getDefaultConf(tempFile.getAbsolutePath()), 2); + coordinator = startCoordinator(TestConfigurations.getDefaultConf(tempFile.getAbsolutePath()), 2); } @AfterEach @@ -142,6 +143,12 @@ public class TestStreamWriteOperatorCoordinator { } } + /** + * Verifies both coordinator restore paths. In case 1, a newly constructed coordinator receives + * checkpoint data before {@link StreamWriteOperatorCoordinator#start()}, so it must deserialize + * the event buffers for restoration during start. In case 2, global failover resets an already + * started coordinator, so it recommits its live event buffers directly. + */ @ParameterizedTest @ValueSource(booleans = {true, false}) public void testCheckpointAndRestore(boolean isStreamingIndexWriteEnabled) throws Exception { @@ -151,7 +158,7 @@ public class TestStreamWriteOperatorCoordinator { conf.set(FlinkOptions.INDEX_TYPE, GLOBAL_RECORD_LEVEL_INDEX.name()); conf.set(FlinkOptions.INDEX_WRITE_TASKS, 2); } - coordinator = createCoordinator(conf, 2); + coordinator = startCoordinator(conf, 2); requestInstantTime(-1); String instant = coordinator.getInstant(); @@ -171,16 +178,32 @@ public class TestStreamWriteOperatorCoordinator { CompletableFuture<byte[]> future = new CompletableFuture<>(); coordinator.checkpointCoordinator(1, future); - coordinator.notifyCheckpointComplete(1); + + // Case 1: job restart restores checkpoint data before the coordinator starts. + try (StreamWriteOperatorCoordinator restoredCoordinator = createCoordinator(conf, 2)) { + restoredCoordinator.resetToCheckpoint(1, future.get()); + + EventBuffers.EventBuffer eventBuffer = restoredCoordinator.getEventBuffer(-1); + assertEquals(2, eventBuffer.getDataWriteEventBuffer().length); + assertEquals(isStreamingIndexWriteEnabled ? 2 : 0, eventBuffer.getIndexWriteEventBuffer().length); + } + + // Case 2: global failover recommits the live buffers of the already started coordinator. coordinator.resetToCheckpoint(1, future.get()); - EventBuffers.EventBuffer eventBuffer = coordinator.getEventBuffer(); - assertEquals(2, eventBuffer.getDataWriteEventBuffer().length); - assertEquals(isStreamingIndexWriteEnabled ? 2 : 0, eventBuffer.getIndexWriteEventBuffer().length); + assertNull(coordinator.getEventBuffer()); + assertTrue(StreamerUtil.createMetaClient(conf).reloadActiveTimeline() + .filterCompletedInstants().containsInstant(instant)); } + /** + * Verifies legacy checkpoint compatibility for both restore paths. Case 1 deserializes the legacy + * checkpoint into a newly constructed coordinator. Case 2 intentionally does not deserialize the + * checkpoint because a coordinator surviving global failover recommits its live buffers directly. + */ @Test public void testRestoreFromLegacyState() throws Exception { + Configuration conf = TestConfigurations.getDefaultConf(tempFile.getAbsolutePath()); requestInstantTime(-1); String instant = coordinator.getInstant(); assertNotEquals("", instant); @@ -192,7 +215,6 @@ public class TestStreamWriteOperatorCoordinator { CompletableFuture<byte[]> future = new CompletableFuture<>(); coordinator.checkpointCoordinator(1, future); - coordinator.notifyCheckpointComplete(1); Map<Long, Pair<String, EventBuffers.EventBuffer>> eventBuffers = SerializationUtils.deserialize(future.get()); // convert to legacy event buffers @@ -200,11 +222,22 @@ public class TestStreamWriteOperatorCoordinator { eventBuffers.forEach((ckpId, eventBuffer) -> { legacyEventBuffers.put(ckpId, Pair.of(eventBuffer.getLeft(), eventBuffer.getRight().getDataWriteEventBuffer())); }); - // simulate recovering from legacy state + + // Case 1: job restart restores legacy checkpoint data before the coordinator starts. + try (StreamWriteOperatorCoordinator restoredCoordinator = createCoordinator(conf, 2)) { + restoredCoordinator.resetToCheckpoint(1, SerializationUtils.serialize(legacyEventBuffers)); + + EventBuffers.EventBuffer eventBuffer = restoredCoordinator.getEventBuffer(-1); + assertEquals(2, eventBuffer.getDataWriteEventBuffer().length); + assertEquals(0, eventBuffer.getIndexWriteEventBuffer().length); + } + + // Case 2: global failover ignores checkpoint bytes and recommits the live buffers. coordinator.resetToCheckpoint(1, SerializationUtils.serialize(legacyEventBuffers)); - EventBuffers.EventBuffer eventBuffer = coordinator.getEventBuffer(); - assertEquals(2, eventBuffer.getDataWriteEventBuffer().length); - assertEquals(0, eventBuffer.getIndexWriteEventBuffer().length); + + assertNull(coordinator.getEventBuffer()); + assertTrue(StreamerUtil.createMetaClient(conf).reloadActiveTimeline() + .filterCompletedInstants().containsInstant(instant)); } @Test @@ -226,7 +259,7 @@ public class TestStreamWriteOperatorCoordinator { public void testEventReset() throws Exception { Configuration conf = TestConfigurations.getDefaultConf(tempFile.getAbsolutePath()); conf.set(FlinkOptions.TABLE_TYPE, HoodieTableType.MERGE_ON_READ.name()); - coordinator = createCoordinator(conf, 2); + coordinator = startCoordinator(conf, 2); CompletableFuture<byte[]> future = new CompletableFuture<>(); coordinator.checkpointCoordinator(1, future); String instant = requestInstantTime(0); @@ -269,7 +302,7 @@ public class TestStreamWriteOperatorCoordinator { conf.setString(HoodieWriteConfig.ALLOW_EMPTY_COMMIT.key(), "false"); OperatorCoordinator.Context context = new MockOperatorCoordinatorContext(new OperatorID(), 2); - coordinator = createCoordinator(conf, 2); + coordinator = startCoordinator(conf, 2); final CompletableFuture<byte[]> future = new CompletableFuture<>(); coordinator.checkpointCoordinator(1, future); @@ -315,7 +348,7 @@ public class TestStreamWriteOperatorCoordinator { conf.set(FlinkOptions.INDEX_TYPE, indexType); conf.set(FlinkOptions.INDEX_WRITE_TASKS, 1); - try (StreamWriteOperatorCoordinator coordinator = createCoordinator(conf, 1)) { + try (StreamWriteOperatorCoordinator coordinator = startCoordinator(conf, 1)) { assertTrue(getRecordLevelIndexFlag(coordinator)); } } @@ -327,7 +360,7 @@ public class TestStreamWriteOperatorCoordinator { // override the default configuration Configuration conf = TestConfigurations.getDefaultConf(tempFile.getAbsolutePath()); conf.setString(HoodieCleanConfig.FAILED_WRITES_CLEANER_POLICY.key(), HoodieFailedWritesCleaningPolicy.LAZY.name()); - coordinator = createCoordinator(conf, 1); + coordinator = startCoordinator(conf, 1); assertTrue(coordinator.getWriteClient().getConfig().getFailedWritesCleanPolicy().isLazy()); @@ -373,7 +406,7 @@ public class TestStreamWriteOperatorCoordinator { // override the default configuration Configuration conf = TestConfigurations.getDefaultConf(tempFile.getAbsolutePath()); conf.set(FlinkOptions.HIVE_SYNC_ENABLED, true); - coordinator = createCoordinator(conf, 1); + coordinator = startCoordinator(conf, 1); String instant = mockWriteWithMetadata(0); assertNotEquals("", instant); @@ -391,7 +424,7 @@ public class TestStreamWriteOperatorCoordinator { int metadataCompactionDeltaCommits = 5; conf.set(FlinkOptions.METADATA_ENABLED, true); conf.set(FlinkOptions.METADATA_COMPACTION_DELTA_COMMITS, metadataCompactionDeltaCommits); - coordinator = createCoordinator(conf, 1); + coordinator = startCoordinator(conf, 1); String instant = coordinator.getInstant(); assertEquals("", instant); @@ -468,7 +501,7 @@ public class TestStreamWriteOperatorCoordinator { conf.set(FlinkOptions.METADATA_ENABLED, true); conf.set(FlinkOptions.METADATA_COMPACTION_DELTA_COMMITS, 20); conf.setString("hoodie.metadata.log.compaction.enable", "true"); - coordinator = createCoordinator(conf, 1); + coordinator = startCoordinator(conf, 1); String instant = coordinator.getInstant(); assertEquals("", instant); @@ -512,7 +545,7 @@ public class TestStreamWriteOperatorCoordinator { // override the default configuration Configuration conf = TestConfigurations.getDefaultConf(tempFile.getAbsolutePath()); conf.set(FlinkOptions.METADATA_ENABLED, true); - coordinator = createCoordinator(conf, 1); + coordinator = startCoordinator(conf, 1); String instant = coordinator.getInstant(); assertEquals("", instant); @@ -536,7 +569,7 @@ public class TestStreamWriteOperatorCoordinator { metadataTableMetaClient.getActiveTimeline().transitionRequestedToInflight(HoodieActiveTimeline.DELTA_COMMIT_ACTION, instant); metadataTableMetaClient.reloadActiveTimeline(); // reset the coordinator to mimic the job failover. - coordinator = createCoordinator(conf, 1); + coordinator = startCoordinator(conf, 1); // write another commit with new instant on the metadata timeline instant = mockWriteWithMetadata(ckp); @@ -555,7 +588,7 @@ public class TestStreamWriteOperatorCoordinator { Logger logger = Mockito.mock(Logger.class); // avoid too many logs by executor NonThrownExecutor executor = NonThrownExecutor.builder(logger).waitForTasksFinish(true).build(); - try (StreamWriteOperatorCoordinator coordinator = createCoordinator(conf, 1)) { + try (StreamWriteOperatorCoordinator coordinator = startCoordinator(conf, 1)) { coordinator.start(); coordinator.setExecutor(executor); TimeUnit.SECONDS.sleep(5); // wait for handled bootstrap event @@ -593,7 +626,7 @@ public class TestStreamWriteOperatorCoordinator { conf.setString(HoodieWriteConfig.WRITE_CONCURRENCY_MODE.key(), WriteConcurrencyMode.OPTIMISTIC_CONCURRENCY_CONTROL.name()); conf.setString("hoodie.write.lock.client.num_retries", "1"); - coordinator = createCoordinator(conf, 1); + coordinator = startCoordinator(conf, 1); String instant = coordinator.getInstant(); assertEquals("", instant); @@ -618,7 +651,7 @@ public class TestStreamWriteOperatorCoordinator { public void testCommitOnEmptyBatch() throws Exception { Configuration conf = TestConfigurations.getDefaultConf(tempFile.getAbsolutePath()); conf.setString(HoodieWriteConfig.ALLOW_EMPTY_COMMIT.key(), "true"); - try (StreamWriteOperatorCoordinator coordinator = createCoordinator(conf, 2)) { + try (StreamWriteOperatorCoordinator coordinator = startCoordinator(conf, 2)) { // Coordinator start the instant String instant = requestInstantTime(coordinator, -1); @@ -649,7 +682,7 @@ public class TestStreamWriteOperatorCoordinator { void testHandleInFlightInstantsRequest() throws Exception { Configuration conf = TestConfigurations.getDefaultConf(tempFile.getAbsolutePath()); conf.set(FlinkOptions.TABLE_TYPE, HoodieTableType.MERGE_ON_READ.name()); - coordinator = createCoordinator(conf, 2); + coordinator = startCoordinator(conf, 2); // Request an instant time to create an initial instant String instant1 = requestInstantTime(1); @@ -701,15 +734,20 @@ public class TestStreamWriteOperatorCoordinator { } } - private static StreamWriteOperatorCoordinator createCoordinator(Configuration conf, int subTasks) throws Exception { - MockOperatorCoordinatorContext coordinatorContext = new MockOperatorCoordinatorContext(new OperatorID(), subTasks); - StreamWriteOperatorCoordinator coordinator = new StreamWriteOperatorCoordinator(conf, coordinatorContext); + private static StreamWriteOperatorCoordinator startCoordinator(Configuration conf, int subTasks) throws Exception { + StreamWriteOperatorCoordinator coordinator = createCoordinator(conf, subTasks); coordinator.start(); + MockOperatorCoordinatorContext coordinatorContext = (MockOperatorCoordinatorContext) coordinator.getContext(); coordinator.setExecutor(new MockCoordinatorExecutor(coordinatorContext)); coordinator.setInstantRequestExecutor(new MockCoordinatorExecutor(coordinatorContext)); return coordinator; } + private static StreamWriteOperatorCoordinator createCoordinator(Configuration conf, int subTasks) { + MockOperatorCoordinatorContext coordinatorContext = new MockOperatorCoordinatorContext(new OperatorID(), subTasks); + return new StreamWriteOperatorCoordinator(conf, coordinatorContext); + } + private String mockWriteWithMetadata(long checkpointId) { String instant = requestInstantTime(checkpointId); OperatorEvent event = createOperatorEvent(0, checkpointId, instant, "par1", false, true, 0.1); @@ -759,7 +797,7 @@ public class TestStreamWriteOperatorCoordinator { HoodieWriteStat writeStat = new HoodieWriteStat(); writeStat.setPartitionPath(partitionPath); writeStat.setFileId("fileId123"); - writeStat.setPath("path123"); + writeStat.setPath(partitionPath + "/" + FSUtils.makeBaseFileName(instant, "1-0-1", "fileId123", ".parquet")); writeStat.setFileSizeInBytes(123); writeStat.setTotalWriteBytes(123); writeStat.setNumWrites(1); diff --git a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestWriteMergeOnRead.java b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestWriteMergeOnRead.java index 5e4deee54c96..a0ddf33eb703 100644 --- a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestWriteMergeOnRead.java +++ b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestWriteMergeOnRead.java @@ -317,6 +317,41 @@ public class TestWriteMergeOnRead extends TestWriteCopyOnWrite { .end(); } + @Test + public void testRecommitOnGlobalFailoverWithStreamingIndex() throws Exception { + conf.set(FlinkOptions.INDEX_TYPE, HoodieIndex.IndexType.GLOBAL_RECORD_LEVEL_INDEX.name()); + conf.setString(HoodieMetadataConfig.GLOBAL_RECORD_LEVEL_INDEX_ENABLE_PROP.key(), "true"); + conf.setString(HoodieMetadataConfig.STREAMING_WRITE_ENABLED.key(), "true"); + conf.set(FlinkOptions.WRITE_COMMIT_ACK_TIMEOUT, 10_000L); + conf.set(FlinkOptions.INDEX_BOOTSTRAP_ENABLED, true); + + Map<String, String> expected = new HashMap<>(); + expected.put("par1", "[id1,par1,id1,Danny,23,1,par1]"); + expected.put("par2", "[id3,par2,id3,Julian,53,3,par2]"); + + preparePipeline(conf) + .consume(TestData.DATA_SET_PART1) + .emptyEventBuffer() + .checkpoint(1) + .assertNextEvent(1, "par1") + .consume(TestData.DATA_SET_PART3) + .checkpoint(2) + // both ckp-1 and ckp-2 are not committing + .assertNextEvent(1, "par2") + .checkCompletedInstantCount(0) + // Simulate the global failover path. The coordinator resets to ckp-2, recommits ckp-1 + // from the coordinator state, and recommits ckp-2 from the restored writer state. + .jobFailover() + .checkCompletedInstantCount(2) + // This is already the global failover path, so recommit does not need to trigger another failover. + .assertGlobalFailure(false) + .checkIndexLoaded( + new HoodieKey("id1", "par1"), + new HoodieKey("id3", "par2")) + .checkWrittenData(expected, 2) + .end(); + } + @Test public void testInsertDuplicateRecordsWithCDCMode() throws Exception { conf.set(FlinkOptions.WRITE_COMMIT_ACK_TIMEOUT, 10_000L);
