danny0405 commented on code in PR #19354:
URL: https://github.com/apache/hudi/pull/19354#discussion_r3636430397


##########
hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/client/TestFlinkWriteClient.java:
##########
@@ -157,4 +163,41 @@ public void 
testCleanResourcesCleansMetadataTableHeartbeatForStreamingMetadataWr
     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(IndexOutOfBoundsException.class,
+        () -> writeClient.insert(Collections.emptyList(), "001"));
+    assertThrows(IndexOutOfBoundsException.class,
+        () -> writeClient.upsert(Collections.emptyList(), "001"));
+    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"));
+    assertThrows(AssertionError.class,
+        () -> 
writeClient.getPartitionToReplacedFileIds(WriteOperationType.UPSERT, 
Collections.emptyList()));
+
+    assertFalse(writeClient.loadActiveTimelineOnTableInit());
+    writeClient.waitForCleaningFinish();
+    writeClient.cleanHandles();
+    assertTrue(writeClient.getHoodieTable(false) != null);

Review Comment:
   Updated to assertNotNull(writeClient.getHoodieTable(false)) in b5a5310f8923.



##########
hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/client/TestFlinkWriteClient.java:
##########
@@ -157,4 +163,41 @@ public void 
testCleanResourcesCleansMetadataTableHeartbeatForStreamingMetadataWr
     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(IndexOutOfBoundsException.class,

Review Comment:
   Agreed. I removed the empty-input insert and upsert exception assertions in 
b5a5310f8923 so this test no longer pins the incidental records.get(0) behavior.



##########
hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/table/TestHoodieFlinkTableActionRouting.java:
##########
@@ -0,0 +1,250 @@
+/*
+ * 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.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.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.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.assertNull;
+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<java.util.List<WriteStatus>> metadata = new 
HoodieWriteMetadata<>();
+          metadata.setWriteStatuses(Collections.emptyList());
+          when(executor.execute()).thenReturn(metadata);
+        })) {
+      assertEquals(Collections.emptyList(), table.managePartitionTTL(context, 
"002").getWriteStatuses());
+    }
+    assertConstructed(BaseRollbackPlanActionExecutor.class,
+        () -> 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");
+
+    assertConstructed(FlinkUpsertCommitActionExecutor.class,
+        () -> table.upsert(context, writeHandle, bucketInfo, "001", 
Collections.emptyIterator()));
+    assertConstructed(FlinkInsertCommitActionExecutor.class,
+        () -> table.insert(context, writeHandle, bucketInfo, "001", 
Collections.emptyIterator()));
+    assertConstructed(FlinkDeletePreppedCommitActionExecutor.class,
+        () -> table.deletePrepped(context, writeHandle, bucketInfo, "001", 
Collections.emptyList()));
+    assertConstructed(FlinkUpsertPreppedCommitActionExecutor.class,
+        () -> table.upsertPrepped(context, writeHandle, bucketInfo, "001", 
Collections.emptyList()));
+    assertConstructed(FlinkInsertPreppedCommitActionExecutor.class,
+        () -> table.insertPrepped(context, writeHandle, bucketInfo, "001", 
Collections.emptyList()));
+    assertConstructed(FlinkInsertOverwriteCommitActionExecutor.class,
+        () -> table.insertOverwrite(context, writeHandle, bucketInfo, "001", 
Collections.emptyIterator()));
+    assertConstructed(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");
+
+    assertConstructed(FlinkUpsertDeltaCommitActionExecutor.class,
+        () -> table.upsert(context, appendHandle, bucketInfo, "001", 
Collections.emptyIterator()));
+    assertConstructed(FlinkUpsertPreppedDeltaCommitActionExecutor.class,
+        () -> table.upsertPrepped(context, appendHandle, bucketInfo, "001", 
Collections.emptyList()));
+    assertConstructed(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> void assertConstructed(Class<E> executorClass, Supplier<Object> 
invocation) {

Review Comment:
   Addressed in b5a5310f8923. The helper now stubs each executor execute() call 
with sentinel metadata and asserts the table method returns that exact 
instance; rollback routing likewise uses a sentinel Option<HoodieRollbackPlan>.



##########
hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/client/TestFlinkWriteClientFunctional.java:
##########
@@ -0,0 +1,446 @@
+/*
+ * 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.List;
+import java.util.Properties;
+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.Mockito.mock;
+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 {
+    if (tableType == HoodieTableType.COPY_ON_WRITE && !cdcEnabled) {
+      context = new HoodieFlinkEngineContext(
+          new HoodieFlinkEngineContext.DefaultTaskContextSupplier() {
+            @Override
+            public java.util.function.Supplier<Long> getAttemptIdSupplier() {
+              return () -> 1L;
+            }
+          });
+    }
+    initWriteClient(tableType, cdcEnabled, false);
+
+    String insertInstant = writeClient.startCommit();
+    transitionToInflight(insertInstant);
+    if (tableType == HoodieTableType.COPY_ON_WRITE && !cdcEnabled) {
+      createInvalidRetryFile(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);
+    if (tableType == HoodieTableType.COPY_ON_WRITE && !cdcEnabled) {
+      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();
+      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 failedWritePath = failingHandle.getWritePath();
+      failingHandle.closeGracefully();
+      assertFalse(metaClient.getStorage().exists(failedWritePath));
+    }
+    assertTrue(writeClient.commit(insertInstant, insertStatuses));
+    assertCommitMetadata(insertInstant, tableType, 3);
+
+    writeClient.cleanHandles();
+    String updateInstant = writeClient.startCommit();
+    transitionToInflight(updateInstant);
+    if (tableType == HoodieTableType.COPY_ON_WRITE && !cdcEnabled) {
+      createInvalidRetryFile(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) {
+      String failedMergeInstant = writeClient.startCommit();
+      transitionToInflight(failedMergeInstant);
+      HoodieFlinkTable table = writeClient.getHoodieTable();
+      FlinkMergeHandle failingHandle = new FlinkMergeHandle(
+          writeConfig,
+          failedMergeInstant,
+          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 failedWritePath = failingHandle.getWritePath();
+      failingHandle.closeGracefully();
+      assertFalse(metaClient.getStorage().exists(failedWritePath));
+    }
+
+    if (tableType == HoodieTableType.COPY_ON_WRITE && cdcEnabled) {
+      assertTrue(updateStatuses.stream()
+          .map(WriteStatus::getStat)
+          .anyMatch(stat -> stat.getCdcStats() != null && 
!stat.getCdcStats().isEmpty()));
+    }
+  }
+
+  @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));
+    java.util.Map<String, List<String>> replacedFileIds =
+        
writeClient.getPartitionToReplacedFileIds(WriteOperationType.INSERT_OVERWRITE, 
insertStatuses);
+    assertEquals(Collections.singleton(FILE_ID),
+        new java.util.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);
+
+    try (MockedStatic<FlinkWriteHandleFactory> factory = 
Mockito.mockStatic(FlinkWriteHandleFactory.class)) {
+      factory.when(() -> FlinkWriteHandleFactory.getFactory(any(), any(), 
anyBoolean()))
+          .thenReturn(handleFactory);
+      assertTrue(routingClient.insertOverwrite(

Review Comment:
   Addressed in b5a5310f8923. The test now verifies the exact insertOverwrite, 
insertOverwriteTable, delete, deletePrepped, and deletePartitions calls, 
including instant and inputs, so a routing regression will fail even when every 
stub returns an empty result.



##########
hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/action/commit/FlinkDeleteHelper.java:
##########
@@ -65,16 +66,16 @@ public static FlinkDeleteHelper newInstance() {
   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())) {

Review Comment:
   Good catch. I mirrored the fix in JavaDeleteHelper: global indexes now 
retain the first key via Set.add, and partitioned indexes use LinkedHashSet for 
deterministic order. I also added TestJavaDeleteHelper coverage for both paths 
in b5a5310f8923.



##########
hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/client/TestFlinkWriteClientFunctional.java:
##########
@@ -0,0 +1,446 @@
+/*
+ * 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.List;
+import java.util.Properties;
+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.Mockito.mock;
+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 {
+    if (tableType == HoodieTableType.COPY_ON_WRITE && !cdcEnabled) {
+      context = new HoodieFlinkEngineContext(
+          new HoodieFlinkEngineContext.DefaultTaskContextSupplier() {
+            @Override
+            public java.util.function.Supplier<Long> getAttemptIdSupplier() {
+              return () -> 1L;
+            }
+          });
+    }
+    initWriteClient(tableType, cdcEnabled, false);
+
+    String insertInstant = writeClient.startCommit();
+    transitionToInflight(insertInstant);
+    if (tableType == HoodieTableType.COPY_ON_WRITE && !cdcEnabled) {
+      createInvalidRetryFile(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);
+    if (tableType == HoodieTableType.COPY_ON_WRITE && !cdcEnabled) {
+      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();
+      rolloverHandle.closeGracefully();

Review Comment:
   Yes, this is an idempotency check. I moved it into the dedicated 
graceful-close test and added an explicit comment in b5a5310f8923.



##########
hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/client/TestFlinkWriteClientFunctional.java:
##########
@@ -0,0 +1,446 @@
+/*
+ * 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.List;
+import java.util.Properties;
+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.Mockito.mock;
+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)

Review Comment:
   Addressed in b5a5310f8923. The parameterized test now contains only the 
shared COW/MOR and CDC insert/upsert/commit flow. Retry cleanup and 
rollover/failing create/merge handle behavior are isolated in 
testCopyOnWriteCleansRetryFiles and 
testCopyOnWriteHandleRolloverAndGracefulCloseCleanup.



##########
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 testCommitFactoryCreatesAndReusesBaseFileHandles() {

Review Comment:
   Renamed to testCommitFactoryCreatesBaseFileHandle in b5a5310f8923 so the 
name matches the single-create behavior being tested.



##########
hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/client/TestFlinkWriteClient.java:
##########
@@ -157,4 +163,41 @@ public void 
testCleanResourcesCleansMetadataTableHeartbeatForStreamingMetadataWr
     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(IndexOutOfBoundsException.class,
+        () -> writeClient.insert(Collections.emptyList(), "001"));
+    assertThrows(IndexOutOfBoundsException.class,
+        () -> writeClient.upsert(Collections.emptyList(), "001"));
+    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"));
+    assertThrows(AssertionError.class,
+        () -> 
writeClient.getPartitionToReplacedFileIds(WriteOperationType.UPSERT, 
Collections.emptyList()));

Review Comment:
   Agreed. I dropped the two incidental IndexOutOfBoundsException assertions 
and the internal AssertionError assertion in b5a5310f8923. The test now keeps 
only stable, intentional API behavior.



##########
hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/client/TestFlinkWriteClientFunctional.java:
##########
@@ -0,0 +1,446 @@
+/*
+ * 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.List;
+import java.util.Properties;
+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.Mockito.mock;
+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 {
+    if (tableType == HoodieTableType.COPY_ON_WRITE && !cdcEnabled) {
+      context = new HoodieFlinkEngineContext(
+          new HoodieFlinkEngineContext.DefaultTaskContextSupplier() {
+            @Override
+            public java.util.function.Supplier<Long> getAttemptIdSupplier() {
+              return () -> 1L;
+            }
+          });
+    }
+    initWriteClient(tableType, cdcEnabled, false);
+
+    String insertInstant = writeClient.startCommit();
+    transitionToInflight(insertInstant);
+    if (tableType == HoodieTableType.COPY_ON_WRITE && !cdcEnabled) {
+      createInvalidRetryFile(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);
+    if (tableType == HoodieTableType.COPY_ON_WRITE && !cdcEnabled) {
+      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();
+      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 failedWritePath = failingHandle.getWritePath();
+      failingHandle.closeGracefully();
+      assertFalse(metaClient.getStorage().exists(failedWritePath));
+    }
+    assertTrue(writeClient.commit(insertInstant, insertStatuses));
+    assertCommitMetadata(insertInstant, tableType, 3);
+
+    writeClient.cleanHandles();
+    String updateInstant = writeClient.startCommit();
+    transitionToInflight(updateInstant);
+    if (tableType == HoodieTableType.COPY_ON_WRITE && !cdcEnabled) {
+      createInvalidRetryFile(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) {
+      String failedMergeInstant = writeClient.startCommit();
+      transitionToInflight(failedMergeInstant);
+      HoodieFlinkTable table = writeClient.getHoodieTable();
+      FlinkMergeHandle failingHandle = new FlinkMergeHandle(
+          writeConfig,
+          failedMergeInstant,
+          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 failedWritePath = failingHandle.getWritePath();
+      failingHandle.closeGracefully();
+      assertFalse(metaClient.getStorage().exists(failedWritePath));
+    }
+
+    if (tableType == HoodieTableType.COPY_ON_WRITE && cdcEnabled) {
+      assertTrue(updateStatuses.stream()
+          .map(WriteStatus::getStat)
+          .anyMatch(stat -> stat.getCdcStats() != null && 
!stat.getCdcStats().isEmpty()));
+    }
+  }
+
+  @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));
+    java.util.Map<String, List<String>> replacedFileIds =
+        
writeClient.getPartitionToReplacedFileIds(WriteOperationType.INSERT_OVERWRITE, 
insertStatuses);
+    assertEquals(Collections.singleton(FILE_ID),
+        new java.util.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);
+
+    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(
+          Collections.singletonList(new HoodieKey("id1", PARTITION_PATH)), 
"003").isEmpty());
+      assertTrue(routingClient.deletePrepped(Collections.emptyList(), 
"004").isEmpty());
+      
assertTrue(routingClient.deletePartitions(Collections.singletonList(PARTITION_PATH),
 "005").isEmpty());
+    } 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 void 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());
+    metaClient.getStorage().create(new StoragePath(partitionPath, 
fileName)).close();
+  }
+
+  private void assertCommitMetadata(String instantTime, HoodieTableType 
tableType, long expectedRecords)
+      throws IOException {
+    metaClient = HoodieTableMetaClient.reload(metaClient);
+    String action = metaClient.getCommitActionType();
+    HoodieInstant instant = metaClient.getActiveTimeline()
+        .getTimelineOfActions(java.util.Collections.singleton(action))

Review Comment:
   Cleaned up in b5a5310f8923: Map, HashSet, Supplier, and Collections now use 
imports consistently.



##########
hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/client/TestHoodieFlinkTableServiceClient.java:
##########
@@ -98,6 +117,134 @@ void 
testInitMetadataTableRespectsStreamingWriteFlag(boolean metadataStreamingWr
     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()));
+      assertNull(client.cluster("001", false));

Review Comment:
   Added the explanatory comment that cluster() is intentionally unimplemented 
for Flink in b5a5310f8923.



##########
hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/action/commit/FlinkDeleteHelper.java:
##########
@@ -65,16 +66,16 @@ public static FlinkDeleteHelper newInstance() {
   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())) {

Review Comment:
   Applied in this PR in b5a5310f8923. JavaDeleteHelper now has the same 
first-occurrence global dedup and deterministic LinkedHashSet partitioned 
behavior, with a focused regression test covering both cases.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to