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 a2bdd5b981b4 test(flink): add sink v2 operator coverage (#19543)
a2bdd5b981b4 is described below

commit a2bdd5b981b438384e2ff613b68168311610599b
Author: Danny Chan <[email protected]>
AuthorDate: Fri Aug 7 17:15:47 2026 +0800

    test(flink): add sink v2 operator coverage (#19543)
    
    * test(flink): add sink v2 operator coverage
---
 .../apache/hudi/sink/v2/TestCleanFunctionV2.java   | 142 ++++++++++++
 .../v2/clustering/TestClusteringCommitSinkV2.java  | 253 +++++++++++++++++++++
 .../v2/compact/TestCompactionCommitSinkV2.java     | 232 +++++++++++++++++++
 .../apache/hudi/sink/v2/utils/TestPipelinesV2.java | 235 +++++++++++++++++++
 4 files changed, 862 insertions(+)

diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/v2/TestCleanFunctionV2.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/v2/TestCleanFunctionV2.java
new file mode 100644
index 000000000000..ca2d75cd3ff0
--- /dev/null
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/v2/TestCleanFunctionV2.java
@@ -0,0 +1,142 @@
+/*
+ * 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.sink.v2;
+
+import org.apache.hudi.client.HoodieFlinkWriteClient;
+import org.apache.hudi.configuration.FlinkOptions;
+import org.apache.hudi.util.FlinkWriteClients;
+
+import org.apache.flink.configuration.Configuration;
+import org.apache.flink.streaming.api.operators.ProcessOperator;
+import org.apache.flink.streaming.runtime.streamrecord.StreamRecord;
+import org.apache.flink.streaming.util.OneInputStreamOperatorTestHarness;
+import org.apache.flink.table.data.RowData;
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedStatic;
+
+import java.util.concurrent.Callable;
+import java.util.concurrent.TimeUnit;
+
+import static org.junit.jupiter.api.Assertions.assertFalse;
+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.doThrow;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.timeout;
+import static org.mockito.Mockito.verify;
+
+/** Tests checkpoint-driven cleaning in {@link CleanFunctionV2}. */
+class TestCleanFunctionV2 {
+
+  @Test
+  void testAsyncCleaningLifecycle() throws Exception {
+    Configuration conf = new Configuration();
+    conf.set(FlinkOptions.CLEAN_ASYNC_ENABLED, true);
+    HoodieFlinkWriteClient writeClient = mock(HoodieFlinkWriteClient.class);
+    CleanFunctionV2<String> function = new CleanFunctionV2<>(conf);
+
+    try (OneInputStreamOperatorTestHarness<String, RowData> harness =
+             openHarness(conf, function, writeClient)) {
+      // Opening waits for any cleaning left by a previous job attempt.
+      verify(writeClient, timeout(5000)).clean();
+      assertTrue(waitUntil(() -> !function.isCleaning));
+
+      harness.processElement(new StreamRecord<>("ignored"));
+      assertTrue(harness.getOutput().isEmpty());
+
+      harness.snapshot(1, 1);
+      verify(writeClient).startAsyncCleaning();
+      assertTrue(function.isCleaning);
+
+      harness.notifyOfCompletedCheckpoint(1);
+      verify(writeClient, timeout(5000)).waitForCleaningFinish();
+      assertTrue(waitUntil(() -> !function.isCleaning));
+    }
+
+    verify(writeClient).close();
+  }
+
+  @Test
+  void testSnapshotDoesNotPropagateCleaningFailure() throws Exception {
+    Configuration conf = new Configuration();
+    conf.set(FlinkOptions.CLEAN_ASYNC_ENABLED, true);
+    HoodieFlinkWriteClient writeClient = mock(HoodieFlinkWriteClient.class);
+    CleanFunctionV2<String> function = new CleanFunctionV2<>(conf);
+    doThrow(new 
RuntimeException("expected")).when(writeClient).startAsyncCleaning();
+
+    try (OneInputStreamOperatorTestHarness<String, RowData> harness =
+             openHarness(conf, function, writeClient)) {
+      verify(writeClient, timeout(5000)).clean();
+      assertTrue(waitUntil(() -> !function.isCleaning));
+
+      harness.snapshot(1, 1);
+
+      assertFalse(function.isCleaning);
+      verify(writeClient).startAsyncCleaning();
+    }
+  }
+
+  @Test
+  void testCleaningDisabled() throws Exception {
+    Configuration conf = new Configuration();
+    conf.set(FlinkOptions.CLEAN_ASYNC_ENABLED, false);
+    HoodieFlinkWriteClient writeClient = mock(HoodieFlinkWriteClient.class);
+    CleanFunctionV2<String> function = new CleanFunctionV2<>(conf);
+
+    try (OneInputStreamOperatorTestHarness<String, RowData> harness =
+             openHarness(conf, function, writeClient)) {
+      harness.snapshot(1, 1);
+      harness.notifyOfCompletedCheckpoint(1);
+      harness.processElement(new StreamRecord<>("ignored"));
+
+      verify(writeClient, never()).clean();
+      verify(writeClient, never()).startAsyncCleaning();
+      verify(writeClient, never()).waitForCleaningFinish();
+      assertTrue(harness.getOutput().isEmpty());
+    }
+  }
+
+  private OneInputStreamOperatorTestHarness<String, RowData> openHarness(
+      Configuration conf,
+      CleanFunctionV2<String> function,
+      HoodieFlinkWriteClient writeClient) throws Exception {
+    OneInputStreamOperatorTestHarness<String, RowData> harness =
+        new OneInputStreamOperatorTestHarness<>(new 
ProcessOperator<>(function), 1, 1, 0);
+    try (MockedStatic<FlinkWriteClients> writeClients = 
mockStatic(FlinkWriteClients.class)) {
+      writeClients.when(() -> FlinkWriteClients.createWriteClient(
+          eq(conf), any())).thenReturn(writeClient);
+      harness.open();
+    }
+    return harness;
+  }
+
+  private boolean waitUntil(Callable<Boolean> condition) throws Exception {
+    long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(5);
+    while (System.nanoTime() < deadline) {
+      if (condition.call()) {
+        return true;
+      }
+      Thread.sleep(10);
+    }
+    return condition.call();
+  }
+}
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/v2/clustering/TestClusteringCommitSinkV2.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/v2/clustering/TestClusteringCommitSinkV2.java
new file mode 100644
index 000000000000..1614dfdd172f
--- /dev/null
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/v2/clustering/TestClusteringCommitSinkV2.java
@@ -0,0 +1,253 @@
+/*
+ * 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.sink.v2.clustering;
+
+import org.apache.hudi.avro.model.HoodieClusteringGroup;
+import org.apache.hudi.avro.model.HoodieClusteringPlan;
+import org.apache.hudi.client.HoodieFlinkWriteClient;
+import org.apache.hudi.client.WriteStatus;
+import org.apache.hudi.common.model.HoodieCommitMetadata;
+import org.apache.hudi.common.model.HoodieFileGroupId;
+import org.apache.hudi.common.model.HoodieWriteStat;
+import org.apache.hudi.common.model.TableServiceType;
+import org.apache.hudi.common.table.HoodieTableMetaClient;
+import org.apache.hudi.common.table.timeline.HoodieActiveTimeline;
+import org.apache.hudi.common.table.timeline.HoodieInstant;
+import org.apache.hudi.common.table.timeline.InstantGenerator;
+import org.apache.hudi.common.util.ClusteringUtils;
+import org.apache.hudi.common.util.Option;
+import org.apache.hudi.common.util.collection.Pair;
+import org.apache.hudi.config.HoodieWriteConfig;
+import org.apache.hudi.configuration.FlinkOptions;
+import org.apache.hudi.sink.clustering.ClusteringCommitEvent;
+import org.apache.hudi.table.HoodieFlinkTable;
+import org.apache.hudi.util.ClusteringUtil;
+import org.apache.hudi.util.FlinkWriteClients;
+
+import org.apache.flink.configuration.Configuration;
+import org.apache.flink.streaming.api.operators.ProcessOperator;
+import org.apache.flink.streaming.runtime.streamrecord.StreamRecord;
+import org.apache.flink.streaming.util.OneInputStreamOperatorTestHarness;
+import org.apache.flink.table.data.RowData;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedStatic;
+
+import java.util.Arrays;
+import java.util.List;
+import java.util.stream.Stream;
+
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.ArgumentMatchers.same;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/** Tests {@link ClusteringCommitSinkV2} with Flink's operator test harness. */
+class TestClusteringCommitSinkV2 {
+
+  private static final String INSTANT = "20240101000000000";
+
+  private Configuration conf;
+  private HoodieFlinkWriteClient writeClient;
+  private HoodieFlinkTable table;
+  private HoodieTableMetaClient metaClient;
+  private HoodieActiveTimeline activeTimeline;
+  private InstantGenerator instantGenerator;
+  private HoodieInstant inflightInstant;
+
+  @BeforeEach
+  @SuppressWarnings("unchecked")
+  void setUp() {
+    conf = new Configuration();
+    conf.set(FlinkOptions.CLEAN_ASYNC_ENABLED, false);
+    writeClient = mock(HoodieFlinkWriteClient.class);
+    table = mock(HoodieFlinkTable.class);
+    metaClient = mock(HoodieTableMetaClient.class);
+    activeTimeline = mock(HoodieActiveTimeline.class);
+    instantGenerator = mock(InstantGenerator.class);
+    inflightInstant = mock(HoodieInstant.class);
+
+    when(writeClient.getHoodieTable()).thenReturn(table);
+    when(table.getMetaClient()).thenReturn(metaClient);
+    when(table.getInstantGenerator()).thenReturn(instantGenerator);
+    when(metaClient.getActiveTimeline()).thenReturn(activeTimeline);
+  }
+
+  @Test
+  void testIgnoreEventsWhoseClusteringPlanNoLongerExists() throws Exception {
+    try (OneInputStreamOperatorTestHarness<ClusteringCommitEvent, RowData> 
harness = openHarness();
+         MockedStatic<ClusteringUtils> clusteringUtils = 
mockStatic(ClusteringUtils.class);
+         MockedStatic<ClusteringUtil> clusteringUtil = 
mockStatic(ClusteringUtil.class)) {
+      clusteringUtils.when(() -> ClusteringUtils.getInflightClusteringInstant(
+          INSTANT, activeTimeline, 
instantGenerator)).thenReturn(Option.empty());
+
+      process(harness, new ClusteringCommitEvent(INSTANT, "file-1", 0));
+
+      clusteringUtil.verifyNoInteractions();
+      verify(writeClient, never()).completeTableService(any(), any(), any(), 
any());
+    }
+  }
+
+  @Test
+  void testWaitForEveryGroupThenRollbackFailedClustering() throws Exception {
+    try (OneInputStreamOperatorTestHarness<ClusteringCommitEvent, RowData> 
harness = openHarness();
+         MockedStatic<ClusteringUtils> clusteringUtils = 
mockStatic(ClusteringUtils.class);
+         MockedStatic<ClusteringUtil> clusteringUtil = 
mockStatic(ClusteringUtil.class)) {
+      stubClusteringPlan(clusteringUtils, planWithGroups(2));
+
+      process(harness, successEvent("file-1", new WriteStatus()));
+      clusteringUtil.verifyNoInteractions();
+
+      process(harness, new ClusteringCommitEvent(INSTANT, "file-2", 1));
+      clusteringUtil.verify(
+          () -> ClusteringUtil.rollbackClustering(table, writeClient, 
INSTANT), times(1));
+
+      // Resetting after rollback makes the next event reload the plan.
+      process(harness, successEvent("file-1", new WriteStatus()));
+      clusteringUtils.verify(
+          () -> ClusteringUtils.getClusteringPlan(metaClient, 
inflightInstant), times(2));
+    }
+  }
+
+  @Test
+  void testRollbackWriteStatusesWithErrorsUnlessConfiguredToIgnore() throws 
Exception {
+    WriteStatus failedStatus = new WriteStatus();
+    failedStatus.setTotalErrorRecords(3);
+
+    try (OneInputStreamOperatorTestHarness<ClusteringCommitEvent, RowData> 
harness = openHarness();
+         MockedStatic<ClusteringUtils> clusteringUtils = 
mockStatic(ClusteringUtils.class);
+         MockedStatic<ClusteringUtil> clusteringUtil = 
mockStatic(ClusteringUtil.class)) {
+      stubClusteringPlan(clusteringUtils, planWithGroups(1));
+
+      process(harness, successEvent("file-1", failedStatus));
+
+      clusteringUtil.verify(
+          () -> ClusteringUtil.rollbackClustering(table, writeClient, 
INSTANT));
+      verify(writeClient, never()).completeTableService(any(), any(), any(), 
any());
+    }
+  }
+
+  @Test
+  void testCommitWriteStatusesWithErrorsWhenConfiguredToIgnore() throws 
Exception {
+    conf.set(FlinkOptions.IGNORE_FAILED, true);
+    HoodieClusteringPlan plan = planWithGroups(1);
+    WriteStatus failedStatus = writeStatus("partition", "new-file", 3);
+    stubWriteConfig();
+
+    try (OneInputStreamOperatorTestHarness<ClusteringCommitEvent, RowData> 
harness = openHarness();
+         MockedStatic<ClusteringUtils> clusteringUtils = 
mockStatic(ClusteringUtils.class);
+         MockedStatic<ClusteringUtil> clusteringUtil = 
mockStatic(ClusteringUtil.class)) {
+      stubClusteringPlan(clusteringUtils, plan);
+      clusteringUtils.when(() -> 
ClusteringUtils.getFileGroupsFromClusteringPlan(plan))
+          .thenReturn(Stream.of(new HoodieFileGroupId("partition", 
"old-file")));
+
+      process(harness, successEvent("file-1", failedStatus));
+
+      clusteringUtil.verifyNoInteractions();
+      verify(writeClient).completeTableService(
+          eq(TableServiceType.CLUSTER), any(HoodieCommitMetadata.class), 
same(table), eq(INSTANT));
+    }
+  }
+
+  @Test
+  void testCommitSuccessfulClusteringAndCleanInline() throws Exception {
+    HoodieClusteringPlan plan = planWithGroups(1);
+    WriteStatus status = writeStatus("partition", "new-file", 0);
+    stubWriteConfig();
+
+    try (OneInputStreamOperatorTestHarness<ClusteringCommitEvent, RowData> 
harness = openHarness();
+         MockedStatic<ClusteringUtils> clusteringUtils = 
mockStatic(ClusteringUtils.class)) {
+      stubClusteringPlan(clusteringUtils, plan);
+      clusteringUtils.when(() -> 
ClusteringUtils.getFileGroupsFromClusteringPlan(plan))
+          .thenReturn(Stream.of(
+              new HoodieFileGroupId("partition", "old-file"),
+              new HoodieFileGroupId("partition", "new-file")));
+
+      process(harness, successEvent("file-1", status));
+    }
+
+    verify(metaClient).reloadActiveTimeline();
+    verify(writeClient).completeTableService(
+        eq(TableServiceType.CLUSTER), any(HoodieCommitMetadata.class), 
same(table), eq(INSTANT));
+    verify(writeClient).clean();
+  }
+
+  private OneInputStreamOperatorTestHarness<ClusteringCommitEvent, RowData> 
openHarness() throws Exception {
+    ProcessOperator<ClusteringCommitEvent, RowData> operator =
+        new ProcessOperator<>(new ClusteringCommitSinkV2(conf));
+    OneInputStreamOperatorTestHarness<ClusteringCommitEvent, RowData> harness =
+        new OneInputStreamOperatorTestHarness<>(operator, 1, 1, 0);
+    try (MockedStatic<FlinkWriteClients> writeClients = 
mockStatic(FlinkWriteClients.class)) {
+      writeClients.when(() -> FlinkWriteClients.createWriteClient(
+          eq(conf), any())).thenReturn(writeClient);
+      harness.open();
+    }
+    return harness;
+  }
+
+  private void process(
+      OneInputStreamOperatorTestHarness<ClusteringCommitEvent, RowData> 
harness,
+      ClusteringCommitEvent event) throws Exception {
+    harness.processElement(new StreamRecord<>(event));
+  }
+
+  private void stubClusteringPlan(
+      MockedStatic<ClusteringUtils> clusteringUtils,
+      HoodieClusteringPlan plan) {
+    clusteringUtils.when(() -> ClusteringUtils.getInflightClusteringInstant(
+        INSTANT, activeTimeline, 
instantGenerator)).thenReturn(Option.of(inflightInstant));
+    clusteringUtils.when(() -> ClusteringUtils.getClusteringPlan(metaClient, 
inflightInstant))
+        .thenReturn(Option.of(Pair.of(inflightInstant, plan)));
+  }
+
+  private HoodieClusteringPlan planWithGroups(int numGroups) {
+    HoodieClusteringPlan plan = mock(HoodieClusteringPlan.class);
+    HoodieClusteringGroup[] groups = new HoodieClusteringGroup[numGroups];
+    Arrays.setAll(groups, ignored -> mock(HoodieClusteringGroup.class));
+    when(plan.getInputGroups()).thenReturn(Arrays.asList(groups));
+    return plan;
+  }
+
+  private WriteStatus writeStatus(String partition, String fileId, long 
errors) {
+    HoodieWriteStat writeStat = new HoodieWriteStat();
+    writeStat.setPartitionPath(partition);
+    writeStat.setFileId(fileId);
+    writeStat.setNumWrites(1);
+    WriteStatus status = new WriteStatus();
+    status.setStat(writeStat);
+    status.setTotalErrorRecords(errors);
+    return status;
+  }
+
+  private void stubWriteConfig() {
+    HoodieWriteConfig writeConfig = mock(HoodieWriteConfig.class);
+    when(writeConfig.getSchema()).thenReturn("{}");
+    when(writeClient.getConfig()).thenReturn(writeConfig);
+  }
+
+  private ClusteringCommitEvent successEvent(String fileId, WriteStatus... 
statuses) {
+    List<WriteStatus> statusList = Arrays.asList(statuses);
+    return new ClusteringCommitEvent(INSTANT, fileId, statusList, 0);
+  }
+}
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/v2/compact/TestCompactionCommitSinkV2.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/v2/compact/TestCompactionCommitSinkV2.java
new file mode 100644
index 000000000000..925e5bb2d479
--- /dev/null
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/v2/compact/TestCompactionCommitSinkV2.java
@@ -0,0 +1,232 @@
+/*
+ * 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.sink.v2.compact;
+
+import org.apache.hudi.avro.model.HoodieCompactionOperation;
+import org.apache.hudi.avro.model.HoodieCompactionPlan;
+import org.apache.hudi.client.HoodieFlinkWriteClient;
+import org.apache.hudi.client.WriteStatus;
+import org.apache.hudi.client.transaction.TransactionManager;
+import org.apache.hudi.common.model.HoodieCommitMetadata;
+import org.apache.hudi.common.model.HoodieWriteStat;
+import org.apache.hudi.common.table.HoodieTableMetaClient;
+import org.apache.hudi.common.util.CompactionUtils;
+import org.apache.hudi.config.HoodieWriteConfig;
+import org.apache.hudi.configuration.FlinkOptions;
+import org.apache.hudi.sink.compact.CompactionCommitEvent;
+import org.apache.hudi.table.HoodieFlinkTable;
+import org.apache.hudi.table.action.compact.CompactHelpers;
+import org.apache.hudi.util.CompactionUtil;
+import org.apache.hudi.util.FlinkWriteClients;
+
+import org.apache.flink.configuration.Configuration;
+import org.apache.flink.streaming.api.operators.ProcessOperator;
+import org.apache.flink.streaming.runtime.streamrecord.StreamRecord;
+import org.apache.flink.streaming.util.OneInputStreamOperatorTestHarness;
+import org.apache.flink.table.data.RowData;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedStatic;
+
+import java.util.Arrays;
+import java.util.List;
+
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.ArgumentMatchers.same;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/** Tests {@link CompactionCommitSinkV2} with Flink's operator test harness. */
+class TestCompactionCommitSinkV2 {
+
+  private static final String INSTANT = "20240101000000000";
+
+  private Configuration conf;
+  private HoodieFlinkWriteClient writeClient;
+  private HoodieFlinkTable table;
+  private HoodieTableMetaClient metaClient;
+  private TransactionManager transactionManager;
+
+  @BeforeEach
+  @SuppressWarnings("unchecked")
+  void setUp() {
+    conf = new Configuration();
+    conf.set(FlinkOptions.CLEAN_ASYNC_ENABLED, false);
+    writeClient = mock(HoodieFlinkWriteClient.class);
+    table = mock(HoodieFlinkTable.class);
+    metaClient = mock(HoodieTableMetaClient.class);
+    transactionManager = mock(TransactionManager.class);
+
+    when(writeClient.getHoodieTable()).thenReturn(table);
+    when(writeClient.getTransactionManager()).thenReturn(transactionManager);
+    when(table.getMetaClient()).thenReturn(metaClient);
+  }
+
+  @Test
+  void testWaitForEveryOperationThenRollbackFailedCompaction() throws 
Exception {
+    HoodieCompactionPlan plan = planWithOperations(2);
+    try (OneInputStreamOperatorTestHarness<CompactionCommitEvent, RowData> 
harness = openHarness();
+         MockedStatic<CompactionUtils> compactionUtils = 
mockStatic(CompactionUtils.class);
+         MockedStatic<CompactionUtil> compactionUtil = 
mockStatic(CompactionUtil.class)) {
+      compactionUtils.when(() -> CompactionUtils.getCompactionPlan(metaClient, 
INSTANT))
+          .thenReturn(plan);
+
+      process(harness, successEvent("file-1", new WriteStatus()));
+      compactionUtil.verifyNoInteractions();
+
+      process(harness, failedEvent("file-2"));
+      compactionUtil.verify(
+          () -> CompactionUtil.rollbackCompaction(table, INSTANT, 
transactionManager), times(1));
+
+      // Resetting after rollback makes the next event reload the plan.
+      process(harness, successEvent("file-1", new WriteStatus()));
+      compactionUtils.verify(
+          () -> CompactionUtils.getCompactionPlan(metaClient, INSTANT), 
times(2));
+    }
+  }
+
+  @Test
+  void testRollbackWriteStatusesWithErrorsUnlessConfiguredToIgnore() throws 
Exception {
+    WriteStatus failedStatus = new WriteStatus();
+    failedStatus.setTotalErrorRecords(3);
+    HoodieCompactionPlan plan = planWithOperations(1);
+
+    try (OneInputStreamOperatorTestHarness<CompactionCommitEvent, RowData> 
harness = openHarness();
+         MockedStatic<CompactionUtils> compactionUtils = 
mockStatic(CompactionUtils.class);
+         MockedStatic<CompactionUtil> compactionUtil = 
mockStatic(CompactionUtil.class)) {
+      compactionUtils.when(() -> CompactionUtils.getCompactionPlan(metaClient, 
INSTANT))
+          .thenReturn(plan);
+
+      process(harness, successEvent("file-1", failedStatus));
+
+      compactionUtil.verify(
+          () -> CompactionUtil.rollbackCompaction(table, INSTANT, 
transactionManager));
+      verify(writeClient, never()).completeCompaction(any(), any(), any());
+    }
+  }
+
+  @Test
+  @SuppressWarnings("unchecked")
+  void testCommitWriteStatusesWithErrorsWhenConfiguredToIgnore() throws 
Exception {
+    conf.set(FlinkOptions.IGNORE_FAILED, true);
+    WriteStatus failedStatus = writeStatus("partition", "file-1", 3);
+    HoodieCommitMetadata metadata = new HoodieCommitMetadata();
+    CompactHelpers compactHelpers = mock(CompactHelpers.class);
+    HoodieCompactionPlan plan = planWithOperations(1);
+    stubWriteConfig();
+
+    try (OneInputStreamOperatorTestHarness<CompactionCommitEvent, RowData> 
harness = openHarness();
+         MockedStatic<CompactionUtils> compactionUtils = 
mockStatic(CompactionUtils.class);
+         MockedStatic<CompactionUtil> compactionUtil = 
mockStatic(CompactionUtil.class);
+         MockedStatic<CompactHelpers> helpersFactory = 
mockStatic(CompactHelpers.class)) {
+      compactionUtils.when(() -> CompactionUtils.getCompactionPlan(metaClient, 
INSTANT))
+          .thenReturn(plan);
+      
helpersFactory.when(CompactHelpers::getInstance).thenReturn(compactHelpers);
+      when(compactHelpers.createCompactionMetadata(
+          same(table), eq(INSTANT), any(), eq("{}"))).thenReturn(metadata);
+
+      process(harness, successEvent("file-1", failedStatus));
+
+      compactionUtil.verifyNoInteractions();
+      verify(writeClient).completeCompaction(same(metadata), same(table), 
eq(INSTANT));
+    }
+  }
+
+  @Test
+  @SuppressWarnings("unchecked")
+  void testCommitSuccessfulCompactionAndCleanInline() throws Exception {
+    WriteStatus status = writeStatus("partition", "file-1", 0);
+    HoodieCommitMetadata metadata = new HoodieCommitMetadata();
+    CompactHelpers compactHelpers = mock(CompactHelpers.class);
+    HoodieCompactionPlan plan = planWithOperations(1);
+    stubWriteConfig();
+
+    try (OneInputStreamOperatorTestHarness<CompactionCommitEvent, RowData> 
harness = openHarness();
+         MockedStatic<CompactionUtils> compactionUtils = 
mockStatic(CompactionUtils.class);
+         MockedStatic<CompactHelpers> helpersFactory = 
mockStatic(CompactHelpers.class)) {
+      compactionUtils.when(() -> CompactionUtils.getCompactionPlan(metaClient, 
INSTANT))
+          .thenReturn(plan);
+      
helpersFactory.when(CompactHelpers::getInstance).thenReturn(compactHelpers);
+      when(compactHelpers.createCompactionMetadata(
+          same(table), eq(INSTANT), any(), eq("{}"))).thenReturn(metadata);
+
+      process(harness, successEvent("file-1", status));
+    }
+
+    verify(writeClient).completeCompaction(same(metadata), same(table), 
eq(INSTANT));
+    verify(writeClient).clean();
+  }
+
+  private OneInputStreamOperatorTestHarness<CompactionCommitEvent, RowData> 
openHarness() throws Exception {
+    ProcessOperator<CompactionCommitEvent, RowData> operator =
+        new ProcessOperator<>(new CompactionCommitSinkV2(conf));
+    OneInputStreamOperatorTestHarness<CompactionCommitEvent, RowData> harness =
+        new OneInputStreamOperatorTestHarness<>(operator, 1, 1, 0);
+    try (MockedStatic<FlinkWriteClients> writeClients = 
mockStatic(FlinkWriteClients.class)) {
+      writeClients.when(() -> FlinkWriteClients.createWriteClient(
+          eq(conf), any())).thenReturn(writeClient);
+      harness.open();
+    }
+    return harness;
+  }
+
+  private void process(
+      OneInputStreamOperatorTestHarness<CompactionCommitEvent, RowData> 
harness,
+      CompactionCommitEvent event) throws Exception {
+    harness.processElement(new StreamRecord<>(event));
+  }
+
+  private HoodieCompactionPlan planWithOperations(int numOperations) {
+    HoodieCompactionPlan plan = mock(HoodieCompactionPlan.class);
+    HoodieCompactionOperation[] operations = new 
HoodieCompactionOperation[numOperations];
+    Arrays.setAll(operations, ignored -> 
mock(HoodieCompactionOperation.class));
+    when(plan.getOperations()).thenReturn(Arrays.asList(operations));
+    return plan;
+  }
+
+  private WriteStatus writeStatus(String partition, String fileId, long 
errors) {
+    HoodieWriteStat writeStat = new HoodieWriteStat();
+    writeStat.setPartitionPath(partition);
+    writeStat.setFileId(fileId);
+    WriteStatus status = new WriteStatus();
+    status.setStat(writeStat);
+    status.setTotalErrorRecords(errors);
+    return status;
+  }
+
+  private void stubWriteConfig() {
+    HoodieWriteConfig writeConfig = mock(HoodieWriteConfig.class);
+    when(writeConfig.getSchema()).thenReturn("{}");
+    when(writeClient.getConfig()).thenReturn(writeConfig);
+  }
+
+  private CompactionCommitEvent successEvent(String fileId, WriteStatus... 
statuses) {
+    List<WriteStatus> statusList = Arrays.asList(statuses);
+    return new CompactionCommitEvent(INSTANT, fileId, statusList, 0, false, 
false);
+  }
+
+  private CompactionCommitEvent failedEvent(String fileId) {
+    return new CompactionCommitEvent(INSTANT, fileId, 0, false, false);
+  }
+}
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/v2/utils/TestPipelinesV2.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/v2/utils/TestPipelinesV2.java
new file mode 100644
index 000000000000..5e20e264094b
--- /dev/null
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/v2/utils/TestPipelinesV2.java
@@ -0,0 +1,235 @@
+/*
+ * 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.sink.v2.utils;
+
+import org.apache.hudi.client.model.HoodieFlinkInternalRow;
+import org.apache.hudi.configuration.FlinkOptions;
+import org.apache.hudi.configuration.OptionsResolver;
+import org.apache.hudi.exception.HoodieException;
+import org.apache.hudi.sink.utils.Pipelines;
+import org.apache.hudi.utils.TestConfigurations;
+
+import org.apache.flink.api.common.typeinfo.TypeInformation;
+import org.apache.flink.api.dag.Transformation;
+import org.apache.flink.configuration.Configuration;
+import org.apache.flink.streaming.api.datastream.DataStream;
+import org.apache.flink.streaming.api.datastream.DataStreamSink;
+import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
+import org.apache.flink.streaming.api.transformations.PartitionTransformation;
+import org.apache.flink.streaming.runtime.partitioner.CustomPartitionerWrapper;
+import org.apache.flink.streaming.runtime.partitioner.StreamPartitioner;
+import org.apache.flink.table.data.RowData;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedStatic;
+
+import java.lang.reflect.Field;
+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.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.mockStatic;
+
+/** Tests topology construction and routing in {@link PipelinesV2}. */
+class TestPipelinesV2 {
+
+  private Configuration conf;
+  private DataStream<RowData> input;
+
+  @BeforeEach
+  void setUp() {
+    conf = new Configuration();
+    StreamExecutionEnvironment env = 
StreamExecutionEnvironment.getExecutionEnvironment();
+    env.setParallelism(7);
+    input = env.fromCollection(
+        Collections.<RowData>emptyList(),
+        TypeInformation.of(RowData.class));
+  }
+
+  @Test
+  void testSinkUsesModeSpecificParallelismAndStableIdentity() {
+    conf.set(FlinkOptions.OPERATION, "bulk_insert");
+    conf.set(FlinkOptions.TABLE_NAME, "sink_v2_test");
+    conf.set(FlinkOptions.WRITE_TASKS, 4);
+
+    DataStreamSink<RowData> sink = PipelinesV2.sink(
+        input, conf, TestConfigurations.ROW_TYPE, false, true);
+
+    assertEquals(4, sink.getTransformation().getParallelism());
+    assertEquals("sink_v2", sink.getTransformation().getName());
+    assertEquals("uid_sink_v2_sink_v2_test", 
sink.getTransformation().getUid());
+  }
+
+  @Test
+  void testServiceTopologiesUseExpectedSingletonOperatorsAndPartitioners() 
throws Exception {
+    conf.set(FlinkOptions.CLUSTERING_TASKS, 3);
+    conf.set(FlinkOptions.COMPACTION_TASKS, 2);
+
+    DataStream<RowData> clean = PipelinesV2.cleanV2(conf, input);
+    DataStream<RowData> cluster = PipelinesV2.clusterV2(
+        conf, TestConfigurations.ROW_TYPE, input);
+    DataStream<RowData> compact = PipelinesV2.compactV2(conf, input);
+
+    assertOperator(clean, "clean_commits", 1, 1);
+    assertOperator(cluster, "cluster_plan_generate", 1, 1);
+    assertOperator(cluster, "clustering_commit", 1, 1);
+    assertOperator(compact, "compact_plan_generate", 1, 1);
+    assertOperator(compact, "compact_commit", 1, 1);
+    assertEquals(3, transformation(cluster, 
"clustering_task").getParallelism());
+    assertEquals(2, transformation(compact, "compact_task").getParallelism());
+    assertEquals(1, countCustomPartitions(cluster, 
Pipelines.IndexPartitioner.class));
+    assertEquals(1, countCustomPartitions(compact, 
Pipelines.IndexPartitioner.class));
+  }
+
+  @Test
+  @SuppressWarnings("unchecked")
+  void testComposePipelineRoutesEveryWriteMode() {
+    DataStream<HoodieFlinkInternalRow> bootstrapped = mock(DataStream.class);
+    Configuration bulk = new Configuration();
+    Configuration append = new Configuration();
+    Configuration appendCluster = new Configuration();
+    Configuration appendClean = new Configuration();
+    Configuration compact = new Configuration();
+    Configuration clean = new Configuration();
+    Configuration plain = new Configuration();
+    bulk.set(FlinkOptions.PATH, "bulk");
+    append.set(FlinkOptions.PATH, "append");
+    appendCluster.set(FlinkOptions.PATH, "append-cluster");
+    appendClean.set(FlinkOptions.PATH, "append-clean");
+    compact.set(FlinkOptions.PATH, "compact");
+    clean.set(FlinkOptions.PATH, "clean");
+    plain.set(FlinkOptions.PATH, "plain");
+    appendCluster.set(FlinkOptions.CLUSTERING_TASKS, 2);
+    compact.set(FlinkOptions.COMPACTION_TASKS, 2);
+
+    try (MockedStatic<OptionsResolver> options = 
mockStatic(OptionsResolver.class);
+         MockedStatic<Pipelines> pipelines = mockStatic(Pipelines.class)) {
+      options.when(() -> 
OptionsResolver.isBulkInsertOperation(bulk)).thenReturn(true);
+      pipelines.when(() -> Pipelines.bulkInsert(bulk, 
TestConfigurations.ROW_TYPE, input))
+          .thenReturn(input);
+      assertSame(input, PipelinesV2.composePipeline(
+          input, bulk, TestConfigurations.ROW_TYPE, false, true));
+      assertThrows(HoodieException.class, () -> PipelinesV2.composePipeline(
+          input, bulk, TestConfigurations.ROW_TYPE, false, false));
+
+      stubAppend(options, pipelines, append);
+      assertSame(input, PipelinesV2.composePipeline(
+          input, append, TestConfigurations.ROW_TYPE, false, false));
+      assertFalse(append.get(FlinkOptions.COMPACTION_SCHEDULE_ENABLED));
+
+      stubAppend(options, pipelines, appendCluster);
+      options.when(() -> 
OptionsResolver.needsAsyncClustering(appendCluster)).thenReturn(true);
+      assertOperator(PipelinesV2.composePipeline(
+          input, appendCluster, TestConfigurations.ROW_TYPE, false, false),
+          "clustering_commit", 1, 1);
+
+      stubAppend(options, pipelines, appendClean);
+      options.when(() -> 
OptionsResolver.isLazyFailedWritesCleaning(appendClean)).thenReturn(true);
+      assertOperator(PipelinesV2.composePipeline(
+          input, appendClean, TestConfigurations.ROW_TYPE, false, false),
+          "clean_commits", 1, 1);
+
+      stubRegular(pipelines, compact, bootstrapped, true);
+      options.when(() -> 
OptionsResolver.needsAsyncCompaction(compact)).thenReturn(true);
+      assertOperator(PipelinesV2.composePipeline(
+          input, compact, TestConfigurations.ROW_TYPE, false, true),
+          "compact_commit", 1, 1);
+      
assertFalse(compact.get(FlinkOptions.COMPACTION_OPERATION_EXECUTE_ASYNC_ENABLED));
+
+      stubRegular(pipelines, clean, bootstrapped, false);
+      options.when(() -> 
OptionsResolver.needsAsyncCleaning(clean)).thenReturn(true);
+      assertOperator(PipelinesV2.composePipeline(
+          input, clean, TestConfigurations.ROW_TYPE, false, false),
+          "clean_commits", 1, 1);
+
+      stubRegular(pipelines, plain, bootstrapped, false);
+      assertSame(input, PipelinesV2.composePipeline(
+          input, plain, TestConfigurations.ROW_TYPE, false, false));
+    }
+  }
+
+  private void stubAppend(
+      MockedStatic<OptionsResolver> options,
+      MockedStatic<Pipelines> pipelines,
+      Configuration configuration) {
+    options.when(() -> 
OptionsResolver.isAppendMode(configuration)).thenReturn(true);
+    pipelines.when(() -> Pipelines.append(
+        configuration, TestConfigurations.ROW_TYPE, input)).thenReturn(input);
+  }
+
+  private void stubRegular(
+      MockedStatic<Pipelines> pipelines,
+      Configuration configuration,
+      DataStream<HoodieFlinkInternalRow> bootstrapped,
+      boolean isBounded) {
+    pipelines.when(() -> Pipelines.bootstrap(
+        configuration, TestConfigurations.ROW_TYPE, input, isBounded, false))
+        .thenReturn(bootstrapped);
+    pipelines.when(() -> Pipelines.hoodieStreamWrite(
+        configuration, TestConfigurations.ROW_TYPE, 
bootstrapped)).thenReturn(input);
+  }
+
+  private void assertOperator(
+      DataStream<?> stream,
+      String name,
+      int parallelism,
+      int maxParallelism) {
+    Transformation<?> transformation = transformation(stream, name);
+    assertEquals(parallelism, transformation.getParallelism());
+    assertEquals(maxParallelism, transformation.getMaxParallelism());
+  }
+
+  private Transformation<?> transformation(DataStream<?> stream, String name) {
+    List<Transformation<?>> matches = 
stream.getTransformation().getTransitivePredecessors()
+        .stream()
+        .filter(candidate -> name.equals(candidate.getName()))
+        .collect(Collectors.toList());
+    assertEquals(1, matches.size(), "Expected one transformation named " + 
name);
+    return matches.get(0);
+  }
+
+  private long countCustomPartitions(DataStream<?> stream, Class<?> 
partitionerClass)
+      throws Exception {
+    long count = 0;
+    for (Transformation<?> transformation : 
stream.getTransformation().getTransitivePredecessors()) {
+      if (transformation instanceof PartitionTransformation) {
+        StreamPartitioner<?> partitioner =
+            ((PartitionTransformation<?>) transformation).getPartitioner();
+        if (partitioner instanceof CustomPartitionerWrapper
+            && partitionerClass.isInstance(
+                getCustomPartitioner((CustomPartitionerWrapper<?, ?>) 
partitioner))) {
+          count++;
+        }
+      }
+    }
+    return count;
+  }
+
+  private Object getCustomPartitioner(CustomPartitionerWrapper<?, ?> 
partitionerWrapper)
+      throws Exception {
+    Field partitionerField = 
CustomPartitionerWrapper.class.getDeclaredField("partitioner");
+    partitionerField.setAccessible(true);
+    return partitionerField.get(partitionerWrapper);
+  }
+}

Reply via email to