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 4d11a60552e4 fix(flink): deduplicate delete keys and cover write paths 
(#19354)
4d11a60552e4 is described below

commit 4d11a60552e43aae88e5dfcfbbdb31d287505273
Author: Danny Chan <[email protected]>
AuthorDate: Thu Jul 23 18:12:29 2026 +0800

    fix(flink): deduplicate delete keys and cover write paths (#19354)
    
    * fix(flink): deduplicate delete keys and cover write paths
---
 .../table/action/commit/FlinkDeleteHelper.java     |   7 +-
 .../apache/hudi/client/TestFlinkWriteClient.java   |  36 ++
 .../client/TestFlinkWriteClientFunctional.java     | 491 +++++++++++++++++++++
 .../client/TestHoodieFlinkTableServiceClient.java  | 185 ++++++++
 .../hudi/io/TestFlinkWriteHandleFactory.java       | 153 ++++++-
 .../table/TestHoodieFlinkTableActionRouting.java   | 264 +++++++++++
 .../table/action/commit/TestFlinkDeleteHelper.java | 174 ++++++++
 .../hudi/table/action/commit/JavaDeleteHelper.java |  16 +-
 .../table/action/commit/TestJavaDeleteHelper.java  |  72 +++
 9 files changed, 1383 insertions(+), 15 deletions(-)

diff --git 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/action/commit/FlinkDeleteHelper.java
 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/action/commit/FlinkDeleteHelper.java
index 7dfc8a336c32..733b79fa9987 100644
--- 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/action/commit/FlinkDeleteHelper.java
+++ 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/action/commit/FlinkDeleteHelper.java
@@ -38,6 +38,7 @@ import java.time.Instant;
 import java.util.Collections;
 import java.util.HashMap;
 import java.util.HashSet;
+import java.util.LinkedHashSet;
 import java.util.LinkedList;
 import java.util.List;
 import java.util.stream.Collectors;
@@ -65,16 +66,16 @@ public class FlinkDeleteHelper<R> extends
   public List<HoodieKey> deduplicateKeys(List<HoodieKey> keys, 
HoodieTable<EmptyHoodieRecordPayload, 
List<HoodieRecord<EmptyHoodieRecordPayload>>, List<HoodieKey>, 
List<WriteStatus>> table, int parallelism) {
     boolean isIndexingGlobal = table.getIndex().isGlobal();
     if (isIndexingGlobal) {
-      HashSet<String> recordKeys = 
keys.stream().map(HoodieKey::getRecordKey).collect(Collectors.toCollection(HashSet::new));
+      HashSet<String> recordKeys = new HashSet<>();
       List<HoodieKey> deduplicatedKeys = new LinkedList<>();
       keys.forEach(x -> {
-        if (recordKeys.contains(x.getRecordKey())) {
+        if (recordKeys.add(x.getRecordKey())) {
           deduplicatedKeys.add(x);
         }
       });
       return deduplicatedKeys;
     } else {
-      HashSet<HoodieKey> set = new HashSet<>(keys);
+      LinkedHashSet<HoodieKey> set = new LinkedHashSet<>(keys);
       keys.clear();
       keys.addAll(set);
       return keys;
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 d0ae6a30a54a..bf41c13426bb 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
@@ -23,11 +23,15 @@ 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.model.HoodieKey;
+import org.apache.hudi.common.model.TableServiceType;
 import org.apache.hudi.common.table.HoodieTableMetaClient;
+import org.apache.hudi.common.util.Option;
 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.exception.HoodieNotSupportedException;
 import org.apache.hudi.index.HoodieIndex;
 import org.apache.hudi.table.HoodieTable;
 import org.apache.hudi.testutils.HoodieFlinkClientTestHarness;
@@ -39,10 +43,12 @@ import org.junit.jupiter.params.ParameterizedTest;
 import org.junit.jupiter.params.provider.ValueSource;
 
 import java.io.IOException;
+import java.util.Collections;
 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.assertNotNull;
 import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 
@@ -157,4 +163,34 @@ public class TestFlinkWriteClient extends 
HoodieFlinkClientTestHarness {
     assertFalse(HoodieHeartbeatClient.heartbeatExists(
         metaClient.getStorage(), metadataTableBasePath, instantTime));
   }
+
+  @Test
+  void testUnsupportedWriteEntryPointsAndInvalidTableServiceFailFast() {
+    HoodieWriteConfig writeConfig = HoodieWriteConfig.newBuilder()
+        .withPath(metaClient.getBasePath())
+        .withEngineType(EngineType.FLINK)
+        .withEmbeddedTimelineServerEnabled(false)
+        .build();
+    writeClient = new HoodieFlinkWriteClient(context, writeConfig);
+
+    assertThrows(HoodieNotSupportedException.class, () -> 
writeClient.bootstrap(Option.empty()));
+    assertThrows(HoodieNotSupportedException.class,
+        () -> writeClient.insertPreppedRecords(Collections.emptyList(), 
"001"));
+    assertThrows(HoodieNotSupportedException.class,
+        () -> writeClient.bulkInsert(Collections.emptyList(), "001"));
+    assertThrows(HoodieNotSupportedException.class,
+        () -> writeClient.bulkInsert(Collections.emptyList(), "001", 
Option.empty()));
+    assertThrows(HoodieNotSupportedException.class,
+        () -> writeClient.cluster("001", false));
+    assertThrows(HoodieException.class,
+        () -> writeClient.delete(Collections.singletonList(new HoodieKey("id", 
"partition")), "001"));
+    assertThrows(HoodieException.class,
+        () -> writeClient.deletePrepped(Collections.emptyList(), "001"));
+    assertThrows(IllegalArgumentException.class,
+        () -> writeClient.completeTableService(TableServiceType.CLEAN, null, 
null, "001"));
+    assertFalse(writeClient.loadActiveTimelineOnTableInit());
+    writeClient.waitForCleaningFinish();
+    writeClient.cleanHandles();
+    assertNotNull(writeClient.getHoodieTable(false));
+  }
 }
diff --git 
a/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/client/TestFlinkWriteClientFunctional.java
 
b/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/client/TestFlinkWriteClientFunctional.java
new file mode 100644
index 000000000000..4d5e2670db14
--- /dev/null
+++ 
b/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/client/TestFlinkWriteClientFunctional.java
@@ -0,0 +1,491 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.hudi.client;
+
+import org.apache.hudi.avro.model.HoodieClusteringPlan;
+import 
org.apache.hudi.client.clustering.plan.strategy.FlinkSizeBasedClusteringPlanStrategyRecently;
+import org.apache.hudi.client.common.HoodieFlinkEngineContext;
+import org.apache.hudi.client.model.HoodieFlinkRecord;
+import org.apache.hudi.common.config.HoodieMetadataConfig;
+import org.apache.hudi.common.engine.EngineType;
+import org.apache.hudi.common.fs.FSUtils;
+import org.apache.hudi.common.model.HoodieKey;
+import org.apache.hudi.common.model.HoodieOperation;
+import org.apache.hudi.common.model.HoodieRecord;
+import org.apache.hudi.common.model.HoodieRecordLocation;
+import org.apache.hudi.common.model.HoodieTableType;
+import org.apache.hudi.common.model.WriteOperationType;
+import org.apache.hudi.common.table.HoodieTableConfig;
+import org.apache.hudi.common.table.HoodieTableMetaClient;
+import org.apache.hudi.common.table.cdc.HoodieCDCSupplementalLoggingMode;
+import org.apache.hudi.common.table.timeline.HoodieInstant;
+import org.apache.hudi.common.util.ClusteringUtils;
+import org.apache.hudi.common.util.Option;
+import org.apache.hudi.config.HoodieClusteringConfig;
+import org.apache.hudi.config.HoodieWriteConfig;
+import org.apache.hudi.io.FlinkCreateHandle;
+import org.apache.hudi.io.FlinkMergeHandle;
+import org.apache.hudi.io.FlinkWriteHandleFactory;
+import org.apache.hudi.io.HoodieWriteMergeHandle;
+import org.apache.hudi.storage.StoragePath;
+import org.apache.hudi.table.HoodieFlinkTable;
+import org.apache.hudi.table.action.HoodieWriteMetadata;
+import org.apache.hudi.table.action.commit.BucketInfo;
+import org.apache.hudi.table.action.commit.BucketType;
+import org.apache.hudi.testutils.HoodieFlinkClientTestHarness;
+
+import org.apache.flink.table.data.GenericRowData;
+import org.apache.flink.table.data.StringData;
+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.Arguments;
+import org.junit.jupiter.params.provider.MethodSource;
+import org.mockito.MockedStatic;
+import org.mockito.Mockito;
+
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Properties;
+import java.util.function.Supplier;
+import java.util.stream.Stream;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNotEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyBoolean;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * Functional coverage for the Flink client write boundary.
+ *
+ * <p>The datasource bucket assigner hands this client records that already 
carry a target file group.
+ * These tests construct that same input directly so client and handle 
behavior can be exercised without
+ * depending on the datasource module.
+ */
+class TestFlinkWriteClientFunctional extends HoodieFlinkClientTestHarness {
+
+  private static final String PARTITION_PATH = "2026/07/23";
+  private static final String FILE_ID = "f0";
+  private static final String SCHEMA = "{"
+      + "\"type\":\"record\","
+      + "\"name\":\"flink_write_test\","
+      + "\"fields\":["
+      + "{\"name\":\"id\",\"type\":\"string\"},"
+      + "{\"name\":\"name\",\"type\":\"string\"},"
+      + "{\"name\":\"ts\",\"type\":\"long\"}"
+      + "]}";
+
+  private HoodieWriteConfig writeConfig;
+
+  @BeforeEach
+  void setUp() {
+    initPath();
+    initFileSystem();
+  }
+
+  @AfterEach
+  void tearDown() throws IOException {
+    cleanupResources();
+  }
+
+  static Stream<Arguments> tableTypesAndCdc() {
+    return Stream.of(
+        Arguments.of(HoodieTableType.COPY_ON_WRITE, false),
+        Arguments.of(HoodieTableType.COPY_ON_WRITE, true),
+        Arguments.of(HoodieTableType.MERGE_ON_READ, false),
+        Arguments.of(HoodieTableType.MERGE_ON_READ, true));
+  }
+
+  @ParameterizedTest
+  @MethodSource("tableTypesAndCdc")
+  void testInsertAndUpsertWriteFilesAndCommitMetadata(HoodieTableType 
tableType, boolean cdcEnabled)
+      throws IOException {
+    initWriteClient(tableType, cdcEnabled, false);
+
+    String insertInstant = writeClient.startCommit();
+    transitionToInflight(insertInstant);
+    List<HoodieRecord> firstBatch = new ArrayList<>(Arrays.asList(
+        insertRecord("id1", "one", 1L),
+        insertRecord("id2", "two", 2L)));
+    if (tableType == HoodieTableType.MERGE_ON_READ) {
+      firstBatch.add(insertRecord("id3", "three", 3L));
+    }
+    List<WriteStatus> firstInsertStatuses = writeClient.insert(firstBatch, 
insertInstant);
+    assertWriteStatuses(firstInsertStatuses, firstBatch.size());
+
+    // A second COW mini-batch for the same bucket exercises the 
incremental/replace handle.
+    List<WriteStatus> insertStatuses = tableType == 
HoodieTableType.COPY_ON_WRITE
+        ? writeClient.insert(Arrays.asList(insertRecord("id3", "three", 3L)), 
insertInstant)
+        : firstInsertStatuses;
+    if (tableType == HoodieTableType.COPY_ON_WRITE && cdcEnabled) {
+      insertStatuses = writeClient.upsert(
+          Arrays.asList(updateRecord("id1", "one-mini-batch-update", 4L, 
insertInstant)),
+          insertInstant);
+    }
+    assertWriteStatuses(insertStatuses, 3);
+    assertTrue(writeClient.commit(insertInstant, insertStatuses));
+    assertCommitMetadata(insertInstant, tableType, 3);
+
+    writeClient.cleanHandles();
+    String updateInstant = writeClient.startCommit();
+    transitionToInflight(updateInstant);
+    List<HoodieRecord> updates = Arrays.asList(
+        updateRecord("id1", "one-updated", 11L, insertInstant),
+        deleteRecord("id2", 12L, insertInstant));
+    List<WriteStatus> updateStatuses = writeClient.upsert(updates, 
updateInstant);
+    long expectedWrites = tableType == HoodieTableType.COPY_ON_WRITE ? 2 : 1;
+    assertWriteStatuses(updateStatuses, expectedWrites);
+    assertEquals(1,
+        updateStatuses.stream().map(WriteStatus::getStat).mapToLong(stat -> 
stat.getNumDeletes()).sum());
+    assertTrue(writeClient.commit(updateInstant, updateStatuses));
+    assertCommitMetadata(updateInstant, tableType, expectedWrites);
+
+    if (tableType == HoodieTableType.COPY_ON_WRITE && cdcEnabled) {
+      assertTrue(updateStatuses.stream()
+          .map(WriteStatus::getStat)
+          .anyMatch(stat -> stat.getCdcStats() != null && 
!stat.getCdcStats().isEmpty()));
+    }
+  }
+
+  @Test
+  void testCopyOnWriteCleansRetryFiles() throws IOException {
+    context = new HoodieFlinkEngineContext(
+        new HoodieFlinkEngineContext.DefaultTaskContextSupplier() {
+          @Override
+          public Supplier<Long> getAttemptIdSupplier() {
+            return () -> 1L;
+          }
+        });
+    initWriteClient(HoodieTableType.COPY_ON_WRITE, false, false);
+
+    String insertInstant = writeClient.startCommit();
+    transitionToInflight(insertInstant);
+    StoragePath staleInsertPath = createInvalidRetryFile(insertInstant);
+    List<WriteStatus> insertStatuses = writeClient.insert(
+        Collections.singletonList(insertRecord("id1", "one", 1L)), 
insertInstant);
+    assertWriteStatuses(insertStatuses, 1);
+    assertFalse(metaClient.getStorage().exists(staleInsertPath));
+    assertTrue(writeClient.commit(insertInstant, insertStatuses));
+
+    writeClient.cleanHandles();
+    String updateInstant = writeClient.startCommit();
+    transitionToInflight(updateInstant);
+    StoragePath staleUpdatePath = createInvalidRetryFile(updateInstant);
+    List<WriteStatus> updateStatuses = writeClient.upsert(
+        Collections.singletonList(updateRecord("id1", "one-updated", 2L, 
insertInstant)),
+        updateInstant);
+    assertWriteStatuses(updateStatuses, 1);
+    assertFalse(metaClient.getStorage().exists(staleUpdatePath));
+    assertTrue(writeClient.commit(updateInstant, updateStatuses));
+  }
+
+  @Test
+  void testCopyOnWriteHandleRolloverAndGracefulCloseCleanup() throws 
IOException {
+    initWriteClient(HoodieTableType.COPY_ON_WRITE, false, false);
+
+    String insertInstant = writeClient.startCommit();
+    transitionToInflight(insertInstant);
+    writeClient.insert(Arrays.asList(
+        insertRecord("id1", "one", 1L),
+        insertRecord("id2", "two", 2L)), insertInstant);
+    List<WriteStatus> insertStatuses = writeClient.insert(
+        Collections.singletonList(insertRecord("id3", "three", 3L)), 
insertInstant);
+    assertWriteStatuses(insertStatuses, 3);
+
+    HoodieFlinkTable table = writeClient.getHoodieTable();
+    FlinkCreateHandle rolloverHandle = new FlinkCreateHandle(
+        writeConfig, insertInstant, table, PARTITION_PATH, FILE_ID, 
table.getTaskContextSupplier());
+    assertTrue(rolloverHandle.canWrite(insertRecord("id4", "four", 4L)));
+    assertNotEquals(insertStatuses.get(0).getStat().getPath(), 
rolloverHandle.getWritePath().toString());
+    rolloverHandle.closeGracefully();
+    // Closing gracefully is intentionally idempotent.
+    rolloverHandle.closeGracefully();
+
+    FlinkCreateHandle failingHandle = new FlinkCreateHandle(
+        writeConfig, insertInstant, table, PARTITION_PATH, FILE_ID, 
table.getTaskContextSupplier()) {
+      @Override
+      public List<WriteStatus> close() {
+        super.close();
+        throw new IllegalStateException("expected close failure");
+      }
+    };
+    StoragePath failedCreatePath = failingHandle.getWritePath();
+    failingHandle.closeGracefully();
+    assertFalse(metaClient.getStorage().exists(failedCreatePath));
+    assertTrue(writeClient.commit(insertInstant, insertStatuses));
+
+    String mergeInstant = writeClient.startCommit();
+    transitionToInflight(mergeInstant);
+    table = writeClient.getHoodieTable();
+    FlinkMergeHandle failingMergeHandle = new FlinkMergeHandle(
+        writeConfig,
+        mergeInstant,
+        table,
+        Collections.<HoodieRecord>emptyList().iterator(),
+        PARTITION_PATH,
+        FILE_ID,
+        table.getTaskContextSupplier()) {
+      @Override
+      public List<WriteStatus> close() {
+        super.close();
+        throw new IllegalStateException("expected close failure");
+      }
+    };
+    StoragePath failedMergePath = failingMergeHandle.getWritePath();
+    failingMergeHandle.closeGracefully();
+    assertFalse(metaClient.getStorage().exists(failedMergePath));
+  }
+
+  @Test
+  void testScheduleClusteringFromRecentlyWrittenPartition() throws IOException 
{
+    initWriteClient(HoodieTableType.COPY_ON_WRITE, false, true);
+
+    String insertInstant = writeClient.startCommit();
+    transitionToInflight(insertInstant);
+    List<WriteStatus> statuses = writeClient.insert(Arrays.asList(
+        insertRecord("id1", "one", 1L),
+        insertRecord("id2", "two", 2L)), insertInstant);
+    assertTrue(writeClient.commit(insertInstant, statuses));
+
+    Option<String> clusteringInstant = 
writeClient.scheduleClustering(Option.empty());
+    assertTrue(clusteringInstant.isPresent());
+    metaClient = HoodieTableMetaClient.reload(metaClient);
+    HoodieInstant requestedInstant = metaClient.getInstantGenerator()
+        .getClusteringCommitRequestedInstant(clusteringInstant.get());
+    HoodieClusteringPlan clusteringPlan = ClusteringUtils
+        .getClusteringPlan(metaClient, requestedInstant).get().getRight();
+    assertEquals(1, clusteringPlan.getInputGroups().size());
+    assertEquals(PARTITION_PATH,
+        
clusteringPlan.getInputGroups().get(0).getSlices().get(0).getPartitionPath());
+  }
+
+  @Test
+  void testPreppedWriteEntryPointsCommitMetadata() throws IOException {
+    initWriteClient(HoodieTableType.COPY_ON_WRITE, false, false);
+
+    String insertInstant = writeClient.startCommit();
+    transitionToInflight(insertInstant);
+    List<WriteStatus> insertStatuses = writeClient.insert(Arrays.asList(
+        insertRecord("id1", "one", 1L),
+        insertRecord("id2", "two", 2L)), insertInstant);
+    assertTrue(writeClient.commit(insertInstant, insertStatuses));
+
+    writeClient.cleanHandles();
+    String upsertInstant = writeClient.startCommit();
+    transitionToInflight(upsertInstant);
+    List<WriteStatus> upsertStatuses = writeClient.upsertPreppedRecords(
+        Collections.singletonList(updateRecord("id1", "one-prepped", 3L, 
insertInstant)),
+        upsertInstant);
+    assertWriteStatuses(upsertStatuses, 2);
+    assertTrue(writeClient.commit(upsertInstant, upsertStatuses));
+    assertCommitMetadata(upsertInstant, HoodieTableType.COPY_ON_WRITE, 2);
+
+    writeClient.cleanHandles();
+    String bulkInsertInstant = writeClient.startCommit();
+    transitionToInflight(bulkInsertInstant);
+    List<WriteStatus> bulkInsertStatuses = 
writeClient.bulkInsertPreppedRecords(
+        Collections.singletonList(insertRecord("id3", "three", 4L)),
+        bulkInsertInstant,
+        Option.empty());
+    assertWriteStatuses(bulkInsertStatuses, 1);
+    assertTrue(writeClient.commit(bulkInsertInstant, bulkInsertStatuses));
+    assertCommitMetadata(bulkInsertInstant, HoodieTableType.COPY_ON_WRITE, 1);
+  }
+
+  @Test
+  @SuppressWarnings({"rawtypes", "unchecked"})
+  void testClientRoutesOverwriteAndDeleteActionsAfterPriorCommit() throws 
IOException {
+    initWriteClient(HoodieTableType.COPY_ON_WRITE, false, false);
+    String insertInstant = writeClient.startCommit();
+    transitionToInflight(insertInstant);
+    List<WriteStatus> insertStatuses = writeClient.insert(
+        Collections.singletonList(insertRecord("id1", "one", 1L)), 
insertInstant);
+    assertTrue(writeClient.commit(insertInstant, insertStatuses));
+    Map<String, List<String>> replacedFileIds =
+        
writeClient.getPartitionToReplacedFileIds(WriteOperationType.INSERT_OVERWRITE, 
insertStatuses);
+    assertEquals(Collections.singleton(FILE_ID),
+        new HashSet<>(replacedFileIds.get(PARTITION_PATH)));
+
+    HoodieFlinkTable table = mock(HoodieFlinkTable.class);
+    when(table.getMetaClient()).thenReturn(metaClient);
+    HoodieWriteMetadata<List<WriteStatus>> metadata = new 
HoodieWriteMetadata<>();
+    metadata.setWriteStatuses(Collections.emptyList());
+    when(table.insertOverwrite(any(), any(), any(), anyString(), 
any())).thenReturn(metadata);
+    when(table.insertOverwriteTable(any(), any(), any(), anyString(), 
any())).thenReturn(metadata);
+    when(table.delete(any(), anyString(), any())).thenReturn(metadata);
+    when(table.deletePrepped(any(), anyString(), any())).thenReturn(metadata);
+    when(table.deletePartitions(any(), anyString(), 
any())).thenReturn(metadata);
+
+    HoodieFlinkWriteClient routingClient = new HoodieFlinkWriteClient(context, 
writeConfig) {
+      @Override
+      protected org.apache.hudi.table.HoodieTable createTable(
+          HoodieWriteConfig config, HoodieTableMetaClient ignoredMetaClient) {
+        return table;
+      }
+    };
+    FlinkWriteHandleFactory.Factory handleFactory = 
mock(FlinkWriteHandleFactory.Factory.class);
+    FlinkCreateHandle writeHandle = mock(FlinkCreateHandle.class);
+    when(handleFactory.create(any(), any(), any(), anyString(), any(), 
any())).thenReturn(writeHandle);
+    BucketInfo bucketInfo = new BucketInfo(BucketType.INSERT, FILE_ID, 
PARTITION_PATH);
+    List<HoodieKey> deleteKeys = Collections.singletonList(new 
HoodieKey("id1", PARTITION_PATH));
+    List<HoodieRecord> preppedDeletes = Collections.emptyList();
+    List<String> partitions = Collections.singletonList(PARTITION_PATH);
+
+    try (MockedStatic<FlinkWriteHandleFactory> factory = 
Mockito.mockStatic(FlinkWriteHandleFactory.class)) {
+      factory.when(() -> FlinkWriteHandleFactory.getFactory(any(), any(), 
anyBoolean()))
+          .thenReturn(handleFactory);
+      assertTrue(routingClient.insertOverwrite(
+          Collections.<HoodieRecord>emptyList().iterator(), bucketInfo, 
"001").isEmpty());
+      assertTrue(routingClient.insertOverwriteTable(
+          Collections.<HoodieRecord>emptyList().iterator(), bucketInfo, 
"002").isEmpty());
+      assertTrue(routingClient.delete(deleteKeys, "003").isEmpty());
+      assertTrue(routingClient.deletePrepped(preppedDeletes, "004").isEmpty());
+      assertTrue(routingClient.deletePartitions(partitions, "005").isEmpty());
+
+      verify(table).insertOverwrite(eq(context), eq(writeHandle), 
eq(bucketInfo), eq("001"), any());
+      verify(table).insertOverwriteTable(eq(context), eq(writeHandle), 
eq(bucketInfo), eq("002"), any());
+      verify(table).delete(eq(context), eq("003"), eq(deleteKeys));
+      verify(table).deletePrepped(eq(context), eq("004"), eq(preppedDeletes));
+      verify(table).deletePartitions(eq(context), eq("005"), eq(partitions));
+    } finally {
+      routingClient.close();
+    }
+  }
+
+  private HoodieRecord insertRecord(String key, String name, long ts) {
+    return record(key, name, ts, HoodieOperation.INSERT, "I");
+  }
+
+  private HoodieRecord updateRecord(String key, String name, long ts, String 
instantTime) {
+    return record(key, name, ts, HoodieOperation.UPDATE_AFTER, instantTime);
+  }
+
+  private HoodieRecord deleteRecord(String key, long ts, String instantTime) {
+    return record(key, "deleted", ts, HoodieOperation.DELETE, instantTime);
+  }
+
+  private HoodieRecord record(
+      String key, String name, long ts, HoodieOperation operation, String 
locationInstant) {
+    GenericRowData row = GenericRowData.of(
+        StringData.fromString(key), StringData.fromString(name), ts);
+    HoodieFlinkRecord record = new HoodieFlinkRecord(
+        new HoodieKey(key, PARTITION_PATH), operation, ts, row);
+    record.setCurrentLocation(new HoodieRecordLocation(locationInstant, 
FILE_ID));
+    return record;
+  }
+
+  private void initWriteClient(
+      HoodieTableType tableType, boolean cdcEnabled, boolean 
useRecentClusteringStrategy)
+      throws IOException {
+    Properties tableProperties = new Properties();
+    tableProperties.setProperty(HoodieTableConfig.CDC_ENABLED.key(), 
Boolean.toString(cdcEnabled));
+    tableProperties.setProperty(
+        HoodieTableConfig.CDC_SUPPLEMENTAL_LOGGING_MODE.key(),
+        HoodieCDCSupplementalLoggingMode.DATA_BEFORE_AFTER.name());
+    tableProperties.setProperty(HoodieTableConfig.RECORDKEY_FIELDS.key(), 
"id");
+    tableProperties.setProperty(HoodieTableConfig.PARTITION_FIELDS.key(), 
"partition_path");
+    tableProperties.setProperty(HoodieTableConfig.ORDERING_FIELDS.key(), "ts");
+    tableProperties.setProperty(
+        HoodieWriteConfig.MERGE_ALLOW_DUPLICATE_ON_INSERTS_ENABLE.key(), 
"false");
+    metaClient = HoodieTableMetaClient.newTableBuilder()
+        .setTableName("flink_write_client_test")
+        .setTableType(tableType)
+        .fromProperties(tableProperties)
+        .initTable(storageConf, basePath);
+
+    HoodieWriteConfig.Builder builder = HoodieWriteConfig.newBuilder()
+        .withPath(basePath)
+        .withEngineType(EngineType.FLINK)
+        .withSchema(SCHEMA)
+        .withProperties(tableProperties)
+        
.withMetadataConfig(HoodieMetadataConfig.newBuilder().enable(false).build())
+        .withMergeHandleClassName(HoodieWriteMergeHandle.class.getName());
+    if (useRecentClusteringStrategy) {
+      builder.withClusteringConfig(HoodieClusteringConfig.newBuilder()
+          .withEngineType(EngineType.FLINK)
+          
.withClusteringPlanStrategyClass(FlinkSizeBasedClusteringPlanStrategyRecently.class.getName())
+          .withClusteringPlanSmallFileLimit(Long.MAX_VALUE)
+          .withClusteringMaxNumGroups(10)
+          .withClusteringSortColumns("id")
+          .build());
+    }
+    writeConfig = builder.build();
+    writeClient = new HoodieFlinkWriteClient<>(context, writeConfig);
+  }
+
+  private void assertWriteStatuses(List<WriteStatus> statuses, long 
expectedRecords) {
+    assertFalse(statuses.isEmpty());
+    assertTrue(statuses.stream().noneMatch(WriteStatus::hasErrors));
+    assertEquals(expectedRecords,
+        statuses.stream().map(WriteStatus::getStat).mapToLong(stat -> 
stat.getNumWrites()).sum());
+    statuses.forEach(status -> {
+      assertEquals(PARTITION_PATH, status.getStat().getPartitionPath());
+      assertNotNull(status.getStat().getPath());
+    });
+  }
+
+  private void transitionToInflight(String instantTime) {
+    metaClient.reloadActiveTimeline();
+    metaClient.getActiveTimeline().transitionRequestedToInflight(
+        metaClient.getCommitActionType(), instantTime);
+  }
+
+  private StoragePath createInvalidRetryFile(String instantTime) throws 
IOException {
+    StoragePath partitionPath = new StoragePath(metaClient.getBasePath(), 
PARTITION_PATH);
+    metaClient.getStorage().createDirectory(partitionPath);
+    String fileName = FSUtils.makeBaseFileName(
+        instantTime,
+        FSUtils.makeWriteToken(0, 1, 0),
+        FILE_ID,
+        metaClient.getTableConfig().getBaseFileFormat().getFileExtension());
+    StoragePath retryPath = new StoragePath(partitionPath, fileName);
+    metaClient.getStorage().create(retryPath).close();
+    return retryPath;
+  }
+
+  private void assertCommitMetadata(String instantTime, HoodieTableType 
tableType, long expectedRecords)
+      throws IOException {
+    metaClient = HoodieTableMetaClient.reload(metaClient);
+    String action = metaClient.getCommitActionType();
+    HoodieInstant instant = metaClient.getActiveTimeline()
+        .getTimelineOfActions(Collections.singleton(action))
+        .filterCompletedInstants()
+        .getInstantsAsStream()
+        .filter(candidate -> candidate.requestedTime().equals(instantTime))
+        .findFirst()
+        .orElseThrow(() -> new AssertionError("Missing completed instant " + 
instantTime));
+    assertEquals(expectedRecords,
+        
metaClient.getActiveTimeline().readCommitMetadata(instant).fetchTotalRecordsWritten());
+  }
+}
diff --git 
a/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/client/TestHoodieFlinkTableServiceClient.java
 
b/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/client/TestHoodieFlinkTableServiceClient.java
index e120470cfeee..25d12795e465 100644
--- 
a/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/client/TestHoodieFlinkTableServiceClient.java
+++ 
b/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/client/TestHoodieFlinkTableServiceClient.java
@@ -21,26 +21,43 @@ package org.apache.hudi.client;
 import org.apache.hudi.client.common.HoodieFlinkEngineContext;
 import org.apache.hudi.client.embedded.EmbeddedTimelineService;
 import org.apache.hudi.common.config.HoodieMetadataConfig;
+import org.apache.hudi.common.model.HoodieCommitMetadata;
+import org.apache.hudi.common.model.HoodieReplaceCommitMetadata;
+import org.apache.hudi.common.model.HoodieWriteStat;
+import org.apache.hudi.common.model.TableServiceType;
 import org.apache.hudi.common.table.timeline.HoodieActiveTimeline;
+import org.apache.hudi.common.table.timeline.HoodieInstant;
 import org.apache.hudi.common.table.timeline.HoodieTimeline;
+import org.apache.hudi.common.util.ClusteringUtils;
 import org.apache.hudi.common.util.Option;
 import org.apache.hudi.config.HoodieLockConfig;
 import org.apache.hudi.config.HoodieWriteConfig;
 import org.apache.hudi.core.transaction.lock.InProcessLockProvider;
+import org.apache.hudi.exception.HoodieException;
 import org.apache.hudi.metadata.FlinkHoodieBackedTableMetadataWriter;
 import org.apache.hudi.storage.StorageConfiguration;
 import org.apache.hudi.table.HoodieFlinkTable;
 import org.apache.hudi.table.HoodieTable;
+import org.apache.hudi.table.action.HoodieWriteMetadata;
+import org.apache.hudi.table.action.compact.CompactHelpers;
+import org.apache.hudi.table.marker.WriteMarkers;
+import org.apache.hudi.table.marker.WriteMarkersFactory;
 import org.apache.hudi.testutils.HoodieFlinkClientTestHarness;
 
 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 org.mockito.MockedStatic;
 import org.mockito.Mockito;
 
 import java.io.IOException;
+import java.util.Collections;
 
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.never;
@@ -75,11 +92,13 @@ public class TestHoodieFlinkTableServiceClient extends 
HoodieFlinkClientTestHarn
     HoodieActiveTimeline activeTimeline = mock(HoodieActiveTimeline.class);
     HoodieTimeline inflightAndRequestedTimeline = mock(HoodieTimeline.class);
     when(table.getActiveTimeline()).thenReturn(activeTimeline);
+    when(table.getMetaClient()).thenReturn(metaClient);
     
when(activeTimeline.filterInflightsAndRequested()).thenReturn(inflightAndRequestedTimeline);
     
when(inflightAndRequestedTimeline.lastInstant()).thenReturn(Option.empty());
 
     FlinkHoodieBackedTableMetadataWriter metadataWriter = 
mock(FlinkHoodieBackedTableMetadataWriter.class);
     when(metadataWriter.isInitialized()).thenReturn(true);
+    when(metadataWriter.hasPartitionsStateChanged()).thenReturn(true);
 
     TestableHoodieFlinkTableServiceClient tableServiceClient =
         new TestableHoodieFlinkTableServiceClient(context, writeConfig, 
Option.empty(), table);
@@ -98,6 +117,135 @@ public class TestHoodieFlinkTableServiceClient extends 
HoodieFlinkClientTestHarn
     verify(table, never()).maybeDeleteMetadataTable();
   }
 
+  @Test
+  void testInitMetadataTableWrapsMetadataWriterFailure() {
+    HoodieWriteConfig writeConfig = HoodieWriteConfig.newBuilder()
+        .withPath(metaClient.getBasePath())
+        .withLockConfig(HoodieLockConfig.newBuilder()
+            .withLockProvider(InProcessLockProvider.class)
+            .build())
+        
.withMetadataConfig(HoodieMetadataConfig.newBuilder().enable(true).build())
+        .build();
+    HoodieFlinkTable<?> table = mock(HoodieFlinkTable.class);
+    HoodieActiveTimeline activeTimeline = mock(HoodieActiveTimeline.class);
+    HoodieTimeline pendingTimeline = mock(HoodieTimeline.class);
+    when(table.getActiveTimeline()).thenReturn(activeTimeline);
+    
when(activeTimeline.filterInflightsAndRequested()).thenReturn(pendingTimeline);
+    when(pendingTimeline.lastInstant()).thenReturn(Option.empty());
+
+    TestableHoodieFlinkTableServiceClient client =
+        new TestableHoodieFlinkTableServiceClient(context, writeConfig, 
Option.empty(), table);
+    try (MockedStatic<FlinkHoodieBackedTableMetadataWriter> writerFactory =
+             Mockito.mockStatic(FlinkHoodieBackedTableMetadataWriter.class)) {
+      writerFactory.when(() -> 
FlinkHoodieBackedTableMetadataWriter.create(any(), any(), any(), any()))
+          .thenThrow(new IllegalStateException("expected metadata writer 
failure"));
+      assertThrows(HoodieException.class, client::initMetadataTable);
+    } finally {
+      client.close();
+    }
+  }
+
+  @Test
+  void testMetadataDisabledDeletesStaleMetadataTable() {
+    HoodieWriteConfig writeConfig = HoodieWriteConfig.newBuilder()
+        .withPath(metaClient.getBasePath())
+        
.withMetadataConfig(HoodieMetadataConfig.newBuilder().enable(false).build())
+        .build();
+    HoodieFlinkTable<?> table = mock(HoodieFlinkTable.class);
+    TestableHoodieFlinkTableServiceClient client =
+        new TestableHoodieFlinkTableServiceClient(context, writeConfig, 
Option.empty(), table);
+    try {
+      client.initMetadataTable();
+    } finally {
+      client.close();
+    }
+
+    verify(table).maybeDeleteMetadataTable();
+    verify(table, never()).deleteMetadataIndexIfNecessary();
+  }
+
+  @Test
+  void testOutputConversionAndNoOpHooks() {
+    HoodieWriteConfig writeConfig = HoodieWriteConfig.newBuilder()
+        .withPath(metaClient.getBasePath())
+        
.withMetadataConfig(HoodieMetadataConfig.newBuilder().enable(false).build())
+        .build();
+    TestableHoodieFlinkTableServiceClient client =
+        new TestableHoodieFlinkTableServiceClient(context, writeConfig, 
Option.empty(), mock(HoodieTable.class));
+    try {
+      WriteStatus status = new WriteStatus(false, 0.0);
+      status.setStat(new HoodieWriteStat());
+      HoodieWriteMetadata<java.util.List<WriteStatus>> metadata = new 
HoodieWriteMetadata<>();
+      metadata.setWriteStatuses(Collections.singletonList(status));
+
+      client.callTriggerWritesAndFetchWriteStats(metadata);
+      assertSame(metadata, client.callConvertToOutputMetadata(metadata));
+      
client.callHandleWriteErrors(Collections.singletonList(status.getStat()));
+      // cluster() is intentionally unimplemented for Flink.
+      assertNull(client.cluster("001", false));
+      assertNotNull(client.createRealTable());
+    } finally {
+      client.close();
+    }
+  }
+
+  @Test
+  void testCompleteCompactionCommitsAndCleansMarkers() {
+    HoodieWriteConfig writeConfig = HoodieWriteConfig.newBuilder()
+        .withPath(metaClient.getBasePath())
+        
.withMetadataConfig(HoodieMetadataConfig.newBuilder().enable(false).build())
+        .build();
+    HoodieTable table = mock(HoodieTable.class);
+    
when(table.getInstantGenerator()).thenReturn(metaClient.getInstantGenerator());
+    HoodieCommitMetadata metadata = new HoodieCommitMetadata();
+    TestableHoodieFlinkTableServiceClient client =
+        new TestableHoodieFlinkTableServiceClient(context, writeConfig, 
Option.empty(), table);
+    CompactHelpers compactHelpers = mock(CompactHelpers.class);
+    WriteMarkers writeMarkers = mock(WriteMarkers.class);
+
+    try (MockedStatic<CompactHelpers> helpersFactory = 
Mockito.mockStatic(CompactHelpers.class);
+         MockedStatic<WriteMarkersFactory> markersFactory = 
Mockito.mockStatic(WriteMarkersFactory.class)) {
+      
helpersFactory.when(CompactHelpers::getInstance).thenReturn(compactHelpers);
+      markersFactory.when(() -> WriteMarkersFactory.get(any(), any(), 
any())).thenReturn(writeMarkers);
+      client.callCompleteCompaction(metadata, table, "20260723120000000");
+    } finally {
+      client.close();
+    }
+
+    verify(compactHelpers).completeInflightCompaction(table, 
"20260723120000000", metadata);
+    verify(writeMarkers).quietDeleteMarkerDir(any(), any(Integer.class));
+  }
+
+  @Test
+  void testCompleteClusteringCommitsAndCleansMarkers() {
+    HoodieWriteConfig writeConfig = HoodieWriteConfig.newBuilder()
+        .withPath(metaClient.getBasePath())
+        
.withMetadataConfig(HoodieMetadataConfig.newBuilder().enable(false).build())
+        .build();
+    HoodieFlinkTable table = mock(HoodieFlinkTable.class);
+    HoodieActiveTimeline activeTimeline = mock(HoodieActiveTimeline.class);
+    when(table.getActiveTimeline()).thenReturn(activeTimeline);
+    
when(table.getInstantGenerator()).thenReturn(metaClient.getInstantGenerator());
+    HoodieInstant clusteringInstant = mock(HoodieInstant.class);
+    HoodieReplaceCommitMetadata metadata = new HoodieReplaceCommitMetadata();
+    TestableHoodieFlinkTableServiceClient client =
+        new TestableHoodieFlinkTableServiceClient(context, writeConfig, 
Option.empty(), table);
+    WriteMarkers writeMarkers = mock(WriteMarkers.class);
+
+    try (MockedStatic<ClusteringUtils> clusteringUtils = 
Mockito.mockStatic(ClusteringUtils.class);
+         MockedStatic<WriteMarkersFactory> markersFactory = 
Mockito.mockStatic(WriteMarkersFactory.class)) {
+      clusteringUtils.when(() -> ClusteringUtils.getInflightClusteringInstant(
+          "20260723120000001", activeTimeline, 
metaClient.getInstantGenerator()))
+          .thenReturn(Option.of(clusteringInstant));
+      markersFactory.when(() -> WriteMarkersFactory.get(any(), any(), 
any())).thenReturn(writeMarkers);
+      client.callCompleteClustering(metadata, table, "20260723120000001");
+    } finally {
+      client.close();
+    }
+
+    verify(writeMarkers).quietDeleteMarkerDir(any(), any(Integer.class));
+  }
+
   private static class TestableHoodieFlinkTableServiceClient extends 
HoodieFlinkTableServiceClient<Object> {
     private final HoodieTable mockedTable;
 
@@ -113,5 +261,42 @@ public class TestHoodieFlinkTableServiceClient extends 
HoodieFlinkClientTestHarn
     protected HoodieTable createTable(HoodieWriteConfig config, 
StorageConfiguration<?> storageConf, boolean skipValidation) {
       return mockedTable;
     }
+
+    private void 
callTriggerWritesAndFetchWriteStats(HoodieWriteMetadata<java.util.List<WriteStatus>>
 metadata) {
+      triggerWritesAndFetchWriteStats(metadata);
+    }
+
+    private HoodieWriteMetadata<java.util.List<WriteStatus>> 
callConvertToOutputMetadata(
+        HoodieWriteMetadata<java.util.List<WriteStatus>> metadata) {
+      return convertToOutputMetadata(metadata);
+    }
+
+    private void callHandleWriteErrors(java.util.List<HoodieWriteStat> 
writeStats) {
+      handleWriteErrors(writeStats, TableServiceType.COMPACT);
+    }
+
+    private HoodieTable createRealTable() {
+      return super.createTable(config, storageConf, false);
+    }
+
+    private void callCompleteCompaction(
+        HoodieCommitMetadata metadata, HoodieTable table, String instantTime) {
+      completeCompaction(metadata, table, instantTime, 
Collections.emptyList());
+    }
+
+    private void callCompleteClustering(
+        HoodieReplaceCommitMetadata metadata, HoodieTable table, String 
instantTime) {
+      completeClustering(metadata, table, instantTime);
+    }
+
+    @Override
+    protected void finalizeWrite(HoodieTable table, String instantTime, 
java.util.List<HoodieWriteStat> stats) {
+      // The test covers Flink orchestration; base finalize-write behavior is 
covered by client-common tests.
+    }
+
+    @Override
+    protected void writeTableMetadata(HoodieTable table, String instantTime, 
HoodieCommitMetadata metadata) {
+      // The test covers Flink orchestration; metadata writer behavior is 
covered independently.
+    }
   }
 }
diff --git 
a/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/io/TestFlinkWriteHandleFactory.java
 
b/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/io/TestFlinkWriteHandleFactory.java
index 9962537beefc..de42467d6348 100644
--- 
a/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/io/TestFlinkWriteHandleFactory.java
+++ 
b/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/io/TestFlinkWriteHandleFactory.java
@@ -19,24 +19,169 @@
 
 package org.apache.hudi.io;
 
+import org.apache.hudi.common.engine.TaskContextSupplier;
+import org.apache.hudi.common.model.HoodieRecord;
+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.HoodieTableVersion;
+import org.apache.hudi.config.HoodieWriteConfig;
+import org.apache.hudi.io.v2.RowDataInlineLogWriteHandle;
+import org.apache.hudi.storage.StoragePath;
 import org.apache.hudi.table.HoodieTable;
+import org.apache.hudi.table.action.commit.BucketInfo;
+import org.apache.hudi.table.action.commit.BucketType;
 
+import org.apache.hadoop.fs.Path;
+import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
+import org.mockito.MockedConstruction;
+import org.mockito.Mockito;
 
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.Iterator;
+import java.util.Map;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertSame;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.when;
 
+/** Verifies operation-specific Flink write handle selection. */
+@SuppressWarnings({"rawtypes", "unchecked"})
 class TestFlinkWriteHandleFactory {
 
+  private HoodieTableConfig tableConfig;
+  private HoodieWriteConfig writeConfig;
+  private HoodieTable table;
+  private Iterator<HoodieRecord<Object>> records;
+
+  @BeforeEach
+  void setUp() {
+    tableConfig = mock(HoodieTableConfig.class);
+    writeConfig = mock(HoodieWriteConfig.class);
+    table = mock(HoodieTable.class);
+    records = Collections.<HoodieRecord<Object>>emptyList().iterator();
+    
when(table.getTaskContextSupplier()).thenReturn(mock(TaskContextSupplier.class));
+    HoodieTableMetaClient metaClient = mock(HoodieTableMetaClient.class);
+    when(table.getMetaClient()).thenReturn(metaClient);
+    when(metaClient.getTableConfig()).thenReturn(tableConfig);
+    when(tableConfig.isLSMTreeStorageLayout()).thenReturn(false);
+    
when(writeConfig.getBasePath()).thenReturn("/tmp/flink-handle-factory-test");
+    when(writeConfig.getWriteVersion()).thenReturn(HoodieTableVersion.NINE);
+  }
+
+  @Test
+  void testCommitFactoryCreatesBaseFileHandle() {
+    when(tableConfig.getTableType()).thenReturn(HoodieTableType.COPY_ON_WRITE);
+    when(tableConfig.isCDCEnabled()).thenReturn(false);
+    BucketInfo bucketInfo = new BucketInfo(BucketType.INSERT, "file-1", 
"partition");
+    Map<String, Path> handles = new HashMap<>();
+    StoragePath writePath = new 
StoragePath("/tmp/flink-handle-factory-test/file-1.parquet");
+
+    try (MockedConstruction<FlinkCreateHandle> mocked = 
Mockito.mockConstruction(
+        FlinkCreateHandle.class,
+        (handle, context) -> 
when(handle.getWritePath()).thenReturn(writePath))) {
+      HoodieWriteHandle handle = factory(false).create(
+          handles, bucketInfo, writeConfig, "001", table, records);
+      assertSame(mocked.constructed().get(0), handle);
+      assertEquals(writePath.toUri(), handles.get("file-1").toUri());
+    }
+  }
+
+  @Test
+  void testClusterFactoryCreatesConcatAndIncrementalConcatHandles() {
+    when(writeConfig.allowDuplicateInserts()).thenReturn(true);
+    BucketInfo bucketInfo = new BucketInfo(BucketType.UPDATE, "file-1", 
"partition");
+    Map<String, Path> handles = new HashMap<>();
+    StoragePath firstPath = new 
StoragePath("/tmp/flink-handle-factory-test/first.parquet");
+
+    try (MockedConstruction<FlinkConcatHandle> mocked = 
Mockito.mockConstruction(
+        FlinkConcatHandle.class,
+        (handle, context) -> 
when(handle.getWritePath()).thenReturn(firstPath))) {
+      HoodieWriteHandle handle = factory(false).create(
+          handles, bucketInfo, writeConfig, "001", table, records);
+      assertSame(mocked.constructed().get(0), handle);
+    }
+
+    StoragePath replacementPath = new 
StoragePath("/tmp/flink-handle-factory-test/replacement.parquet");
+    try (MockedConstruction<FlinkIncrementalConcatHandle> mocked = 
Mockito.mockConstruction(
+        FlinkIncrementalConcatHandle.class,
+        (handle, context) -> 
when(handle.getWritePath()).thenReturn(replacementPath))) {
+      HoodieWriteHandle handle = factory(false).create(
+          handles, bucketInfo, writeConfig, "001", table, records);
+      assertSame(mocked.constructed().get(0), handle);
+      assertEquals(replacementPath.toUri(), handles.get("file-1").toUri());
+    }
+  }
+
+  @Test
+  void testCdcFactoryCreatesIncrementalChangeLogHandle() {
+    when(tableConfig.getTableType()).thenReturn(HoodieTableType.COPY_ON_WRITE);
+    when(tableConfig.isCDCEnabled()).thenReturn(true);
+    BucketInfo bucketInfo = new BucketInfo(BucketType.UPDATE, "file-1", 
"partition");
+    Map<String, Path> handles = new HashMap<>();
+    handles.put("file-1", new 
Path("/tmp/flink-handle-factory-test/original.parquet"));
+    StoragePath replacementPath = new 
StoragePath("/tmp/flink-handle-factory-test/cdc.parquet");
+
+    try (MockedConstruction<FlinkIncrementalMergeHandleWithChangeLog> mocked = 
Mockito.mockConstruction(
+        FlinkIncrementalMergeHandleWithChangeLog.class,
+        (handle, context) -> 
when(handle.getWritePath()).thenReturn(replacementPath))) {
+      HoodieWriteHandle handle = factory(false).create(
+          handles, bucketInfo, writeConfig, "001", table, records);
+      assertSame(mocked.constructed().get(0), handle);
+      assertEquals(replacementPath.toUri(), handles.get("file-1").toUri());
+    }
+  }
+
+  @Test
+  void testCommitFactoryCreatesFileGroupReaderMergeHandle() {
+    when(tableConfig.getTableType()).thenReturn(HoodieTableType.COPY_ON_WRITE);
+    
when(writeConfig.getMergeHandleClassName()).thenReturn(FileGroupReaderBasedMergeHandle.class.getName());
+    HoodieTableMetaClient metaClient = mock(HoodieTableMetaClient.class);
+    when(table.getMetaClient()).thenReturn(metaClient);
+    when(metaClient.getTableConfig()).thenReturn(tableConfig);
+    when(tableConfig.isLSMTreeStorageLayout()).thenReturn(false);
+    BucketInfo bucketInfo = new BucketInfo(BucketType.UPDATE, "file-1", 
"partition");
+    StoragePath writePath = new 
StoragePath("/tmp/flink-handle-factory-test/reader.parquet");
+
+    try (MockedConstruction<FlinkFileGroupReaderBasedMergeHandle> mocked = 
Mockito.mockConstruction(
+        FlinkFileGroupReaderBasedMergeHandle.class,
+        (handle, context) -> 
when(handle.getWritePath()).thenReturn(writePath))) {
+      HoodieWriteHandle handle = factory(false).create(
+          new HashMap<>(), bucketInfo, writeConfig, "001", table, records);
+      assertSame(mocked.constructed().get(0), handle);
+    }
+  }
+
+  @Test
+  void testMergeOnReadSelectsMetadataAndRowDataLogHandles() {
+    when(tableConfig.getTableType()).thenReturn(HoodieTableType.MERGE_ON_READ);
+    BucketInfo bucketInfo = new BucketInfo(BucketType.UPDATE, "file-1", 
"partition");
+
+    when(writeConfig.getBasePath()).thenReturn("/tmp/table/.hoodie/metadata");
+    try (MockedConstruction<FlinkInlineLogAppendHandle> mocked =
+             Mockito.mockConstruction(FlinkInlineLogAppendHandle.class)) {
+      HoodieWriteHandle handle = factory(false).create(
+          new HashMap<>(), bucketInfo, writeConfig, "001", table, records);
+      assertSame(mocked.constructed().get(0), handle);
+    }
+
+    when(writeConfig.getBasePath()).thenReturn("/tmp/table");
+    try (MockedConstruction<RowDataInlineLogWriteHandle> mocked =
+             Mockito.mockConstruction(RowDataInlineLogWriteHandle.class)) {
+      HoodieWriteHandle handle = factory(false).create(
+          new HashMap<>(), bucketInfo, writeConfig, "002", table, records);
+      assertSame(mocked.constructed().get(0), handle);
+    }
+  }
+
   @Test
   void testLsmTreeStorageLayoutDetection() {
-    HoodieTable table = mock(HoodieTable.class);
     HoodieTableMetaClient metaClient = mock(HoodieTableMetaClient.class);
-    HoodieTableConfig tableConfig = mock(HoodieTableConfig.class);
     when(table.getMetaClient()).thenReturn(metaClient);
     when(metaClient.getTableConfig()).thenReturn(tableConfig);
     when(tableConfig.isLSMTreeStorageLayout()).thenReturn(false);
@@ -45,4 +190,8 @@ class TestFlinkWriteHandleFactory {
     when(tableConfig.isLSMTreeStorageLayout()).thenReturn(true);
     assertTrue(FlinkWriteHandleFactory.isLsmTreeStorageLayout(table));
   }
+
+  private FlinkWriteHandleFactory.Factory factory(boolean overwrite) {
+    return FlinkWriteHandleFactory.getFactory(tableConfig, writeConfig, 
overwrite);
+  }
 }
diff --git 
a/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/table/TestHoodieFlinkTableActionRouting.java
 
b/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/table/TestHoodieFlinkTableActionRouting.java
new file mode 100644
index 000000000000..07cc09897673
--- /dev/null
+++ 
b/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/table/TestHoodieFlinkTableActionRouting.java
@@ -0,0 +1,264 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.hudi.table;
+
+import org.apache.hudi.avro.model.HoodieRollbackPlan;
+import org.apache.hudi.client.WriteStatus;
+import org.apache.hudi.common.data.HoodieListData;
+import org.apache.hudi.common.engine.EngineType;
+import org.apache.hudi.common.model.HoodieKey;
+import org.apache.hudi.common.model.HoodieRecord;
+import org.apache.hudi.common.model.HoodieTableType;
+import org.apache.hudi.common.table.HoodieTableVersion;
+import org.apache.hudi.common.table.timeline.HoodieInstant;
+import org.apache.hudi.common.util.Option;
+import org.apache.hudi.config.HoodieWriteConfig;
+import org.apache.hudi.exception.HoodieNotSupportedException;
+import org.apache.hudi.io.HoodieAppendHandle;
+import org.apache.hudi.io.HoodieCreateHandle;
+import org.apache.hudi.io.HoodieInlineLogAppendHandle;
+import org.apache.hudi.io.HoodieWriteHandle;
+import org.apache.hudi.table.action.BaseActionExecutor;
+import org.apache.hudi.table.action.HoodieWriteMetadata;
+import org.apache.hudi.table.action.clean.CleanPlanActionExecutor;
+import org.apache.hudi.table.action.cluster.ClusteringPlanActionExecutor;
+import org.apache.hudi.table.action.commit.BucketInfo;
+import org.apache.hudi.table.action.commit.BucketType;
+import 
org.apache.hudi.table.action.commit.FlinkDeletePreppedCommitActionExecutor;
+import org.apache.hudi.table.action.commit.FlinkInsertCommitActionExecutor;
+import 
org.apache.hudi.table.action.commit.FlinkInsertOverwriteCommitActionExecutor;
+import 
org.apache.hudi.table.action.commit.FlinkInsertOverwriteTableCommitActionExecutor;
+import 
org.apache.hudi.table.action.commit.FlinkInsertPreppedCommitActionExecutor;
+import org.apache.hudi.table.action.commit.FlinkPartitionTTLActionExecutor;
+import org.apache.hudi.table.action.commit.FlinkUpsertCommitActionExecutor;
+import 
org.apache.hudi.table.action.commit.FlinkUpsertPreppedCommitActionExecutor;
+import 
org.apache.hudi.table.action.commit.delta.FlinkUpsertDeltaCommitActionExecutor;
+import 
org.apache.hudi.table.action.commit.delta.FlinkUpsertPreppedDeltaCommitActionExecutor;
+import org.apache.hudi.table.action.compact.RunCompactionActionExecutor;
+import org.apache.hudi.table.action.compact.ScheduleCompactionActionExecutor;
+import org.apache.hudi.table.action.rollback.BaseRollbackPlanActionExecutor;
+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.mockito.MockedConstruction;
+import org.mockito.Mockito;
+
+import java.io.IOException;
+import java.util.Collections;
+import java.util.List;
+import java.util.function.Supplier;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+/** Tests Flink table action routing and explicitly unsupported engine APIs. */
+@SuppressWarnings({"rawtypes", "unchecked"})
+class TestHoodieFlinkTableActionRouting extends HoodieFlinkClientTestHarness {
+
+  @BeforeEach
+  void setUp() {
+    initPath();
+    initFileSystem();
+  }
+
+  @AfterEach
+  void tearDown() throws IOException {
+    cleanupResources();
+  }
+
+  @Test
+  void testCopyOnWriteUnsupportedActionsFailFast() throws IOException {
+    initMetaClient(HoodieTableType.COPY_ON_WRITE);
+    HoodieFlinkCopyOnWriteTable table = new 
HoodieFlinkCopyOnWriteTable(config(), context, metaClient);
+
+    assertUnsupported(() -> table.upsert(context, "001", 
Collections.emptyList()));
+    assertUnsupported(() -> table.insert(context, "001", 
Collections.emptyList()));
+    assertUnsupported(() -> table.bulkInsert(context, "001", 
Collections.emptyList(), Option.empty()));
+    assertUnsupported(() -> table.delete(context, "001", 
Collections.<HoodieKey>emptyList()));
+    assertUnsupported(() -> table.deletePrepped(context, "001", 
Collections.<HoodieRecord>emptyList()));
+    assertUnsupported(() -> table.upsertPrepped(context, "001", 
Collections.<HoodieRecord>emptyList()));
+    assertUnsupported(() -> table.insertPrepped(context, "001", 
Collections.<HoodieRecord>emptyList()));
+    assertUnsupported(() -> table.bulkInsertPrepped(
+        context, "001", Collections.<HoodieRecord>emptyList(), 
Option.empty()));
+    assertUnsupported(() -> table.insertOverwrite(context, "001", 
Collections.<HoodieRecord>emptyList()));
+    assertUnsupported(() -> table.insertOverwriteTable(context, "001", 
Collections.<HoodieRecord>emptyList()));
+    assertUnsupported(() -> table.scheduleCompaction(context, "001", 
Option.empty()));
+    assertUnsupported(() -> table.compact(context, "001"));
+    assertUnsupported(() -> table.cluster(context, "001"));
+    assertUnsupported(() -> table.bootstrap(context, Option.empty()));
+    assertUnsupported(() -> table.rollbackBootstrap(context, "001"));
+    assertUnsupported(() -> table.scheduleIndexing(context, "001", 
Collections.emptyList(), Collections.emptyList()));
+    assertUnsupported(() -> table.index(context, "001"));
+    assertUnsupported(() -> table.savepoint(context, "001", "user", 
"comment"));
+    assertUnsupported(() -> table.scheduleRestore(context, "002", "001"));
+    assertUnsupported(() -> table.restore(context, "002", "001"));
+  }
+
+  @Test
+  void testCopyOnWriteRoutesSupportedPlanningActions() throws IOException {
+    initMetaClient(HoodieTableType.COPY_ON_WRITE);
+    HoodieFlinkCopyOnWriteTable table = new 
HoodieFlinkCopyOnWriteTable(config(), context, metaClient);
+
+    try (MockedConstruction<ClusteringPlanActionExecutor> ignored = 
Mockito.mockConstruction(
+        ClusteringPlanActionExecutor.class,
+        (executor, constructionContext) -> 
when(executor.execute()).thenReturn(Option.empty()))) {
+      assertFalse(table.scheduleClustering(context, "001", 
Option.empty()).isPresent());
+    }
+    try (MockedConstruction<CleanPlanActionExecutor> ignored = 
Mockito.mockConstruction(
+        CleanPlanActionExecutor.class,
+        (executor, constructionContext) -> 
when(executor.execute()).thenReturn(Option.empty()))) {
+      assertFalse(table.createCleanerPlan(context, 
Option.empty()).isPresent());
+    }
+    try (MockedConstruction<FlinkPartitionTTLActionExecutor> ignored = 
Mockito.mockConstruction(
+        FlinkPartitionTTLActionExecutor.class,
+        (executor, constructionContext) -> {
+          HoodieWriteMetadata<List<WriteStatus>> metadata = new 
HoodieWriteMetadata<>();
+          metadata.setWriteStatuses(Collections.emptyList());
+          when(executor.execute()).thenReturn(metadata);
+        })) {
+      assertEquals(Collections.emptyList(), table.managePartitionTTL(context, 
"002").getWriteStatuses());
+    }
+    Option<HoodieRollbackPlan> rollbackPlan = 
Option.of(mock(HoodieRollbackPlan.class));
+    assertResultPropagated(BaseRollbackPlanActionExecutor.class, rollbackPlan,
+        () -> table.scheduleRollback(
+            context, "003", mock(HoodieInstant.class), false, false, false));
+  }
+
+  @Test
+  void testCopyOnWriteRoutesWriteActionsAndCompactionInsert() throws 
IOException {
+    initMetaClient(HoodieTableType.COPY_ON_WRITE);
+    HoodieFlinkCopyOnWriteTable table = new 
HoodieFlinkCopyOnWriteTable(config(), context, metaClient);
+    HoodieWriteHandle writeHandle = mock(HoodieWriteHandle.class);
+    BucketInfo bucketInfo = new BucketInfo(BucketType.INSERT, "file-1", 
"partition");
+
+    assertWriteMetadataPropagated(FlinkUpsertCommitActionExecutor.class,
+        () -> table.upsert(context, writeHandle, bucketInfo, "001", 
Collections.emptyIterator()));
+    assertWriteMetadataPropagated(FlinkInsertCommitActionExecutor.class,
+        () -> table.insert(context, writeHandle, bucketInfo, "001", 
Collections.emptyIterator()));
+    assertWriteMetadataPropagated(FlinkDeletePreppedCommitActionExecutor.class,
+        () -> table.deletePrepped(context, writeHandle, bucketInfo, "001", 
Collections.emptyList()));
+    assertWriteMetadataPropagated(FlinkUpsertPreppedCommitActionExecutor.class,
+        () -> table.upsertPrepped(context, writeHandle, bucketInfo, "001", 
Collections.emptyList()));
+    assertWriteMetadataPropagated(FlinkInsertPreppedCommitActionExecutor.class,
+        () -> table.insertPrepped(context, writeHandle, bucketInfo, "001", 
Collections.emptyList()));
+    
assertWriteMetadataPropagated(FlinkInsertOverwriteCommitActionExecutor.class,
+        () -> table.insertOverwrite(context, writeHandle, bucketInfo, "001", 
Collections.emptyIterator()));
+    
assertWriteMetadataPropagated(FlinkInsertOverwriteTableCommitActionExecutor.class,
+        () -> table.insertOverwriteTable(context, writeHandle, bucketInfo, 
"001", Collections.emptyIterator()));
+
+    try (MockedConstruction<HoodieCreateHandle> ignored = 
Mockito.mockConstruction(HoodieCreateHandle.class)) {
+      assertEquals(Collections.emptyList(),
+          table.handleInsert("001", "partition", "file-1", 
Collections.emptyMap()).next());
+    }
+  }
+
+  @Test
+  void testMergeOnReadValidatesHandlesAndRoutesScheduling() throws IOException 
{
+    initMetaClient(HoodieTableType.MERGE_ON_READ);
+    HoodieFlinkMergeOnReadTable table = new 
HoodieFlinkMergeOnReadTable(config(), context, metaClient);
+    HoodieWriteHandle writeHandle = mock(HoodieWriteHandle.class);
+    BucketInfo bucketInfo = new BucketInfo(BucketType.UPDATE, "file-1", 
"partition");
+
+    assertThrows(IllegalArgumentException.class,
+        () -> table.upsert(context, writeHandle, bucketInfo, "001", 
Collections.emptyIterator()));
+    assertThrows(IllegalArgumentException.class,
+        () -> table.upsertPrepped(context, writeHandle, bucketInfo, "001", 
Collections.emptyList()));
+
+    try (MockedConstruction<ScheduleCompactionActionExecutor> mocked = 
Mockito.mockConstruction(
+        ScheduleCompactionActionExecutor.class,
+        (executor, constructionContext) -> 
when(executor.execute()).thenReturn(Option.empty()))) {
+      assertFalse(table.scheduleCompaction(context, "002", 
Option.empty()).isPresent());
+      assertFalse(table.scheduleLogCompaction(context, "003", 
Option.empty()).isPresent());
+      assertEquals(2, mocked.constructed().size());
+    }
+  }
+
+  @Test
+  void testMergeOnReadRoutesAppendAndCompactionActions() throws IOException {
+    initMetaClient(HoodieTableType.MERGE_ON_READ);
+    HoodieFlinkMergeOnReadTable table = new 
HoodieFlinkMergeOnReadTable(config(), context, metaClient);
+    HoodieAppendHandle appendHandle = mock(HoodieAppendHandle.class);
+    BucketInfo bucketInfo = new BucketInfo(BucketType.UPDATE, "file-1", 
"partition");
+
+    assertWriteMetadataPropagated(FlinkUpsertDeltaCommitActionExecutor.class,
+        () -> table.upsert(context, appendHandle, bucketInfo, "001", 
Collections.emptyIterator()));
+    
assertWriteMetadataPropagated(FlinkUpsertPreppedDeltaCommitActionExecutor.class,
+        () -> table.upsertPrepped(context, appendHandle, bucketInfo, "001", 
Collections.emptyList()));
+    assertWriteMetadataPropagated(FlinkUpsertDeltaCommitActionExecutor.class,
+        () -> table.insert(context, appendHandle, bucketInfo, "001", 
Collections.emptyIterator()));
+
+    HoodieWriteMetadata compactionMetadata = new HoodieWriteMetadata();
+    
compactionMetadata.setWriteStatuses(HoodieListData.eager(Collections.emptyList()));
+    try (MockedConstruction<RunCompactionActionExecutor> ignored = 
Mockito.mockConstruction(
+        RunCompactionActionExecutor.class,
+        (executor, constructionContext) -> 
when(executor.execute()).thenReturn(compactionMetadata))) {
+      assertEquals(Collections.emptyList(), table.compact(context, 
"002").getWriteStatuses());
+      assertEquals(Collections.emptyList(), table.logCompact(context, 
"003").getWriteStatuses());
+    }
+
+    try (MockedConstruction<HoodieInlineLogAppendHandle> ignored = 
Mockito.mockConstruction(
+        HoodieInlineLogAppendHandle.class,
+        (handle, constructionContext) -> 
when(handle.close()).thenReturn(Collections.emptyList()))) {
+      assertEquals(Collections.emptyList(),
+          table.handleInsertsForLogCompaction(
+              "004", "partition", "file-1", Collections.emptyMap(), 
Collections.emptyMap()).next());
+    }
+  }
+
+  private <E extends BaseActionExecutor> void assertWriteMetadataPropagated(
+      Class<E> executorClass, Supplier<Object> invocation) {
+    HoodieWriteMetadata<List<WriteStatus>> metadata = new 
HoodieWriteMetadata<>();
+    metadata.setWriteStatuses(Collections.emptyList());
+    assertResultPropagated(executorClass, metadata, invocation);
+  }
+
+  private <E extends BaseActionExecutor> void assertResultPropagated(
+      Class<E> executorClass, Object expected, Supplier<Object> invocation) {
+    try (MockedConstruction<E> mocked = Mockito.mockConstruction(
+        executorClass,
+        (executor, constructionContext) -> 
when(executor.execute()).thenReturn(expected))) {
+      assertSame(expected, invocation.get());
+      assertEquals(1, mocked.constructed().size());
+    }
+  }
+
+  private HoodieWriteConfig config() {
+    return HoodieWriteConfig.newBuilder()
+        .withPath(basePath)
+        .withEngineType(EngineType.FLINK)
+        .withWriteTableVersion(HoodieTableVersion.NINE.versionCode())
+        .build();
+  }
+
+  private void assertUnsupported(ThrowingRunnable runnable) {
+    assertThrows(HoodieNotSupportedException.class, runnable::run);
+  }
+
+  @FunctionalInterface
+  private interface ThrowingRunnable {
+    void run() throws Exception;
+  }
+}
diff --git 
a/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/table/action/commit/TestFlinkDeleteHelper.java
 
b/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/table/action/commit/TestFlinkDeleteHelper.java
new file mode 100644
index 000000000000..3b9b6e8bc423
--- /dev/null
+++ 
b/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/table/action/commit/TestFlinkDeleteHelper.java
@@ -0,0 +1,174 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.hudi.table.action.commit;
+
+import org.apache.hudi.client.WriteStatus;
+import org.apache.hudi.common.data.HoodieData;
+import org.apache.hudi.common.data.HoodieListData;
+import org.apache.hudi.common.engine.HoodieEngineContext;
+import org.apache.hudi.common.model.EmptyHoodieRecordPayload;
+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.config.HoodieWriteConfig;
+import org.apache.hudi.exception.HoodieUpsertException;
+import org.apache.hudi.index.HoodieIndex;
+import org.apache.hudi.table.HoodieTable;
+import org.apache.hudi.table.action.HoodieWriteMetadata;
+
+import org.junit.jupiter.api.Test;
+
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.List;
+import java.util.stream.Collectors;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotSame;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.reset;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/** Tests for {@link FlinkDeleteHelper}. */
+@SuppressWarnings({"rawtypes", "unchecked"})
+class TestFlinkDeleteHelper {
+
+  @Test
+  void testDeduplicateKeysForGlobalAndPartitionedIndexes() {
+    FlinkDeleteHelper helper = FlinkDeleteHelper.newInstance();
+    assertSame(helper, FlinkDeleteHelper.newInstance());
+
+    HoodieTable table = mock(HoodieTable.class);
+    HoodieIndex index = mock(HoodieIndex.class);
+    when(table.getIndex()).thenReturn(index);
+    List<HoodieKey> keys = new ArrayList<>(Arrays.asList(
+        new HoodieKey("id1", "p1"),
+        new HoodieKey("id1", "p2"),
+        new HoodieKey("id2", "p1"),
+        new HoodieKey("id2", "p1")));
+
+    when(index.isGlobal()).thenReturn(true);
+    List<HoodieKey> globalResult = helper.deduplicateKeys(keys, table, 1);
+    assertEquals(Arrays.asList("id1", "id2"), globalResult.stream()
+        .map(HoodieKey::getRecordKey).collect(Collectors.toList()));
+    assertEquals(Arrays.asList("p1", "p1"), globalResult.stream()
+        .map(HoodieKey::getPartitionPath).collect(Collectors.toList()));
+    assertNotSame(keys, globalResult);
+
+    when(index.isGlobal()).thenReturn(false);
+    List<HoodieKey> partitionedResult = helper.deduplicateKeys(keys, table, 1);
+    assertSame(keys, partitionedResult);
+    assertEquals(Arrays.asList(
+        new HoodieKey("id1", "p1"),
+        new HoodieKey("id1", "p2"),
+        new HoodieKey("id2", "p1")), partitionedResult);
+  }
+
+  @Test
+  void testExecuteTagsExistingRecordsAndDelegatesDelete() {
+    HoodieTable table = mock(HoodieTable.class);
+    HoodieIndex index = mock(HoodieIndex.class);
+    HoodieEngineContext context = mock(HoodieEngineContext.class);
+    BaseCommitActionExecutor executor = mock(BaseCommitActionExecutor.class);
+    when(table.getIndex()).thenReturn(index);
+    when(index.isGlobal()).thenReturn(false);
+
+    when(index.tagLocation(any(HoodieData.class), eq(context), 
eq(table))).thenAnswer(invocation -> {
+      HoodieData<HoodieRecord<EmptyHoodieRecordPayload>> records = 
invocation.getArgument(0);
+      List<HoodieRecord<EmptyHoodieRecordPayload>> tagged = 
records.collectAsList();
+      tagged.get(0).setCurrentLocation(new HoodieRecordLocation("001", 
"file-1"));
+      return HoodieListData.eager(tagged);
+    });
+
+    HoodieWriteMetadata<List<WriteStatus>> expected = new 
HoodieWriteMetadata<>();
+    expected.setWriteStatuses(Collections.singletonList(new WriteStatus(false, 
0.0)));
+    when(executor.execute(any(List.class))).thenReturn(expected);
+
+    HoodieWriteConfig config = HoodieWriteConfig.newBuilder()
+        .withPath("/tmp/flink-delete-helper")
+        .combineDeleteInput(true)
+        .build();
+    List<HoodieKey> keys = new ArrayList<>(Arrays.asList(
+        new HoodieKey("id1", "p1"), new HoodieKey("id1", "p1"), new 
HoodieKey("missing", "p1")));
+
+    HoodieWriteMetadata<List<WriteStatus>> result = 
FlinkDeleteHelper.newInstance()
+        .execute("002", keys, context, config, table, executor);
+
+    assertSame(expected, result);
+    assertTrue(result.getIndexLookupDuration().isPresent());
+    verify(executor).execute(any(List.class));
+    verify(executor, never()).saveWorkloadProfileMetadataToInflight(any(), 
any());
+  }
+
+  @Test
+  void testExecuteWithNoExistingRecordsCreatesEmptyMetadata() {
+    HoodieTable table = mock(HoodieTable.class);
+    HoodieIndex index = mock(HoodieIndex.class);
+    HoodieEngineContext context = mock(HoodieEngineContext.class);
+    BaseCommitActionExecutor executor = mock(BaseCommitActionExecutor.class);
+    when(table.getIndex()).thenReturn(index);
+    when(index.tagLocation(any(HoodieData.class), eq(context), eq(table)))
+        .thenAnswer(invocation -> invocation.getArgument(0));
+
+    HoodieWriteConfig config = HoodieWriteConfig.newBuilder()
+        .withPath("/tmp/flink-delete-helper")
+        .combineDeleteInput(false)
+        .build();
+    HoodieWriteMetadata<List<WriteStatus>> result = 
FlinkDeleteHelper.newInstance().execute(
+        "003", Collections.singletonList(new HoodieKey("missing", "p1")),
+        context, config, table, executor);
+
+    assertTrue(result.getWriteStatuses().isEmpty());
+    verify(executor).saveWorkloadProfileMetadataToInflight(any(), eq("003"));
+    verify(executor).runPrecommitValidators(result);
+    verify(executor, never()).execute(any(List.class));
+  }
+
+  @Test
+  void testExecutePreservesOrWrapsFailures() {
+    HoodieTable table = mock(HoodieTable.class);
+    HoodieEngineContext context = mock(HoodieEngineContext.class);
+    BaseCommitActionExecutor executor = mock(BaseCommitActionExecutor.class);
+    HoodieWriteConfig config = HoodieWriteConfig.newBuilder()
+        .withPath("/tmp/flink-delete-helper")
+        .build();
+    List<HoodieKey> keys = Collections.singletonList(new HoodieKey("id1", 
"p1"));
+
+    HoodieUpsertException upsertException = new 
HoodieUpsertException("expected");
+    when(table.getIndex()).thenThrow(upsertException);
+    HoodieUpsertException firstFailure = 
assertThrows(HoodieUpsertException.class,
+        () -> FlinkDeleteHelper.newInstance().execute("004", keys, context, 
config, table, executor));
+    assertSame(upsertException, firstFailure);
+
+    reset(table);
+    when(table.getIndex()).thenThrow(new IllegalStateException("boom"));
+    HoodieUpsertException wrapped = assertThrows(HoodieUpsertException.class,
+        () -> FlinkDeleteHelper.newInstance().execute("005", keys, context, 
config, table, executor));
+    assertTrue(wrapped.getCause() instanceof IllegalStateException);
+  }
+}
diff --git 
a/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/table/action/commit/JavaDeleteHelper.java
 
b/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/table/action/commit/JavaDeleteHelper.java
index bc077c12c00c..45c2637259dd 100644
--- 
a/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/table/action/commit/JavaDeleteHelper.java
+++ 
b/hudi-client/hudi-java-client/src/main/java/org/apache/hudi/table/action/commit/JavaDeleteHelper.java
@@ -38,7 +38,7 @@ import java.time.Instant;
 import java.util.Collections;
 import java.util.HashMap;
 import java.util.HashSet;
-import java.util.LinkedList;
+import java.util.LinkedHashSet;
 import java.util.List;
 import java.util.stream.Collectors;
 
@@ -64,16 +64,12 @@ public class JavaDeleteHelper<R> extends
                                          int parallelism) {
     boolean isIndexingGlobal = table.getIndex().isGlobal();
     if (isIndexingGlobal) {
-      HashSet<String> recordKeys = 
keys.stream().map(HoodieKey::getRecordKey).collect(Collectors.toCollection(HashSet::new));
-      List<HoodieKey> deduplicatedKeys = new LinkedList<>();
-      keys.forEach(x -> {
-        if (recordKeys.contains(x.getRecordKey())) {
-          deduplicatedKeys.add(x);
-        }
-      });
-      return deduplicatedKeys;
+      HashSet<String> recordKeys = new HashSet<>();
+      return keys.stream()
+          .filter(key -> recordKeys.add(key.getRecordKey()))
+          .collect(Collectors.toList());
     } else {
-      HashSet<HoodieKey> set = new HashSet<>(keys);
+      LinkedHashSet<HoodieKey> set = new LinkedHashSet<>(keys);
       keys.clear();
       keys.addAll(set);
       return keys;
diff --git 
a/hudi-client/hudi-java-client/src/test/java/org/apache/hudi/table/action/commit/TestJavaDeleteHelper.java
 
b/hudi-client/hudi-java-client/src/test/java/org/apache/hudi/table/action/commit/TestJavaDeleteHelper.java
new file mode 100644
index 000000000000..fab7d0423163
--- /dev/null
+++ 
b/hudi-client/hudi-java-client/src/test/java/org/apache/hudi/table/action/commit/TestJavaDeleteHelper.java
@@ -0,0 +1,72 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *      http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hudi.table.action.commit;
+
+import org.apache.hudi.common.model.HoodieKey;
+import org.apache.hudi.index.HoodieIndex;
+import org.apache.hudi.table.HoodieTable;
+
+import org.junit.jupiter.api.Test;
+
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.List;
+import java.util.stream.Collectors;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotSame;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+/** Tests for {@link JavaDeleteHelper}. */
+@SuppressWarnings({"rawtypes", "unchecked"})
+class TestJavaDeleteHelper {
+
+  @Test
+  void testDeduplicateKeysForGlobalAndPartitionedIndexes() {
+    JavaDeleteHelper helper = JavaDeleteHelper.newInstance();
+    assertSame(helper, JavaDeleteHelper.newInstance());
+
+    HoodieTable table = mock(HoodieTable.class);
+    HoodieIndex index = mock(HoodieIndex.class);
+    when(table.getIndex()).thenReturn(index);
+    List<HoodieKey> keys = new ArrayList<>(Arrays.asList(
+        new HoodieKey("id1", "p1"),
+        new HoodieKey("id1", "p2"),
+        new HoodieKey("id2", "p1"),
+        new HoodieKey("id2", "p1")));
+
+    when(index.isGlobal()).thenReturn(true);
+    List<HoodieKey> globalResult = helper.deduplicateKeys(keys, table, 1);
+    assertEquals(Arrays.asList("id1", "id2"), globalResult.stream()
+        .map(HoodieKey::getRecordKey).collect(Collectors.toList()));
+    assertEquals(Arrays.asList("p1", "p1"), globalResult.stream()
+        .map(HoodieKey::getPartitionPath).collect(Collectors.toList()));
+    assertNotSame(keys, globalResult);
+
+    when(index.isGlobal()).thenReturn(false);
+    List<HoodieKey> partitionedResult = helper.deduplicateKeys(keys, table, 1);
+    assertSame(keys, partitionedResult);
+    assertEquals(Arrays.asList(
+        new HoodieKey("id1", "p1"),
+        new HoodieKey("id1", "p2"),
+        new HoodieKey("id2", "p1")), partitionedResult);
+  }
+}

Reply via email to