This is an automated email from the ASF dual-hosted git repository.

voonhous pushed a commit to branch release-1.2.1
in repository https://gitbox.apache.org/repos/asf/hudi.git

commit 3e671a234fa3314191b73dc1b7d87ca77447b74e
Author: Sivabalan Narayanan <[email protected]>
AuthorDate: Thu Jun 25 19:24:00 2026 -0700

    fix(test): close TestingServer + write clients on all exit paths in 
TestHoodieClientMultiWriter (#19062)
    
    Test flakiness fix in 
TestHoodieClientMultiWriter.testHoodieClientBasicMultiWriterWithEarlyConflictDetection.
    Test-only change. No production code is touched.
    
    Co-authored-by: Claude Opus 4.7 (1M context) <[email protected]>
    (cherry picked from commit 014e8a281319059f900311f97bf28fdf453bd3c9)
---
 .../hudi/client/TestHoodieClientMultiWriter.java   | 121 ++++++++++++---------
 1 file changed, 67 insertions(+), 54 deletions(-)

diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/client/TestHoodieClientMultiWriter.java
 
b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/client/TestHoodieClientMultiWriter.java
index bebdf9d88e60..a53f167c2efd 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/client/TestHoodieClientMultiWriter.java
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/java/org/apache/hudi/client/TestHoodieClientMultiWriter.java
@@ -85,6 +85,7 @@ import org.junit.jupiter.params.provider.Arguments;
 import org.junit.jupiter.params.provider.EnumSource;
 import org.junit.jupiter.params.provider.MethodSource;
 
+import java.io.Closeable;
 import java.io.File;
 import java.io.IOException;
 import java.util.ArrayList;
@@ -500,71 +501,83 @@ public class TestHoodieClientMultiWriter extends 
HoodieClientTestBase {
       writeConfig = buildWriteConfigForEarlyConflictDetect(markerType, 
properties, InProcessLockProvider.class, earlyConflictDetectionStrategy);
     }
 
-    final SparkRDDWriteClient client1 = getHoodieWriteClient(writeConfig);
+    SparkRDDWriteClient client1 = null;
+    SparkRDDWriteClient client2 = null;
+    SparkRDDWriteClient client3 = null;
+    SparkRDDWriteClient client4 = null;
+    try {
+      client1 = getHoodieWriteClient(writeConfig);
 
-    // Create the first commit
-    final String nextCommitTime1 = "001";
-    createCommitWithInserts(writeConfig, client1, "000", nextCommitTime1, 200);
+      // Create the first commit
+      final String nextCommitTime1 = "001";
+      createCommitWithInserts(writeConfig, client1, "000", nextCommitTime1, 
200);
 
-    final SparkRDDWriteClient client2 = getHoodieWriteClient(writeConfig);
-    final SparkRDDWriteClient client3 = getHoodieWriteClient(writeConfig);
+      client2 = getHoodieWriteClient(writeConfig);
+      client3 = getHoodieWriteClient(writeConfig);
 
-    final String nextCommitTime2 = "002";
+      final String nextCommitTime2 = "002";
 
-    // start to write commit 002
-    final JavaRDD<WriteStatus> writeStatusList2 = 
startCommitForUpdate(writeConfig, client2, nextCommitTime2, 100);
+      // start to write commit 002
+      final SparkRDDWriteClient finalClient2 = client2;
+      final JavaRDD<WriteStatus> writeStatusList2 = 
startCommitForUpdate(writeConfig, client2, nextCommitTime2, 100);
 
-    // start to write commit 003
-    // this commit 003 will fail quickly because early conflict detection 
before create marker.
-    final String nextCommitTime3 = "003";
-    assertThrows(SparkException.class, () -> {
-      final JavaRDD<WriteStatus> writeStatusList3 =
-          startCommitForUpdate(writeConfig, client3, nextCommitTime3, 100);
-      client3.commit(nextCommitTime3, writeStatusList3);
-    }, "Early conflict detected but cannot resolve conflicts for overlapping 
writes");
+      // start to write commit 003
+      // this commit 003 will fail quickly because early conflict detection 
before create marker.
+      final String nextCommitTime3 = "003";
+      final SparkRDDWriteClient finalClient3 = client3;
+      assertThrows(SparkException.class, () -> {
+        final JavaRDD<WriteStatus> writeStatusList3 =
+            startCommitForUpdate(writeConfig, finalClient3, nextCommitTime3, 
100);
+        finalClient3.commit(nextCommitTime3, writeStatusList3);
+      }, "Early conflict detected but cannot resolve conflicts for overlapping 
writes");
 
-    // start to commit 002 and success
-    assertDoesNotThrow(() -> {
-      client2.commit(nextCommitTime2, writeStatusList2);
-    });
-
-    HoodieWriteConfig config4 =
-        HoodieWriteConfig.newBuilder().withProperties(writeConfig.getProps())
-            .withHeartbeatIntervalInMs(heartBeatIntervalForCommit4).build();
-    final SparkRDDWriteClient client4 = getHoodieWriteClient(config4);
-
-    StoragePath heartbeatFilePath = new StoragePath(
-        HoodieTableMetaClient.getHeartbeatFolderPath(basePath) + 
StoragePath.SEPARATOR + nextCommitTime3);
-    storage.create(heartbeatFilePath, true);
+      // start to commit 002 and success
+      assertDoesNotThrow(() -> {
+        finalClient2.commit(nextCommitTime2, writeStatusList2);
+      });
 
-    // Wait for heart beat expired for failed commitTime3 "003"
-    // Otherwise commit4 still can see conflict between failed write 003. The 
early-conflict
-    // check treats 003 as alive until its heartbeat is older than
-    // (tolerable misses * heartbeat interval); tolerable misses is pinned in
-    // buildWriteConfigForEarlyConflictDetect, so wait one interval past that 
window.
-    Thread.sleep(heartBeatIntervalForCommit4 * 
(EARLY_CONFLICT_HEARTBEAT_TOLERABLE_MISSES + 1));
+      HoodieWriteConfig config4 =
+          HoodieWriteConfig.newBuilder().withProperties(writeConfig.getProps())
+              .withHeartbeatIntervalInMs(heartBeatIntervalForCommit4).build();
+      client4 = getHoodieWriteClient(config4);
 
-    final String nextCommitTime4 = "004";
-    assertDoesNotThrow(() -> {
-      final JavaRDD<WriteStatus> writeStatusList4 =
-          startCommitForUpdate(writeConfig, client4, nextCommitTime4, 100);
-      client4.commit(nextCommitTime4, writeStatusList4);
-    });
+      StoragePath heartbeatFilePath = new StoragePath(
+          HoodieTableMetaClient.getHeartbeatFolderPath(basePath) + 
StoragePath.SEPARATOR + nextCommitTime3);
+      storage.create(heartbeatFilePath, true);
 
-    List<String> completedInstant = 
metaClient.reloadActiveTimeline().getCommitsTimeline()
-        .filterCompletedInstants().getInstants().stream()
-        .map(HoodieInstant::requestedTime).collect(Collectors.toList());
+      // Wait for heart beat expired for failed commitTime3 "003"
+      // Otherwise commit4 still can see conflict between failed write 003. 
The early-conflict
+      // check treats 003 as alive until its heartbeat is older than
+      // (tolerable misses * heartbeat interval); tolerable misses is pinned in
+      // buildWriteConfigForEarlyConflictDetect, so wait one interval past 
that window.
+      Thread.sleep(heartBeatIntervalForCommit4 * 
(EARLY_CONFLICT_HEARTBEAT_TOLERABLE_MISSES + 1));
 
-    assertEquals(3, completedInstant.size());
-    assertTrue(completedInstant.contains(nextCommitTime1));
-    assertTrue(completedInstant.contains(nextCommitTime2));
-    assertTrue(completedInstant.contains(nextCommitTime4));
+      final String nextCommitTime4 = "004";
+      final SparkRDDWriteClient finalClient4 = client4;
+      assertDoesNotThrow(() -> {
+        final JavaRDD<WriteStatus> writeStatusList4 =
+            startCommitForUpdate(writeConfig, finalClient4, nextCommitTime4, 
100);
+        finalClient4.commit(nextCommitTime4, writeStatusList4);
+      });
 
-    FileIOUtils.deleteDirectory(new File(basePath));
-    client1.close();
-    client2.close();
-    client3.close();
-    client4.close();
+      List<String> completedInstant = 
metaClient.reloadActiveTimeline().getCommitsTimeline()
+          .filterCompletedInstants().getInstants().stream()
+          .map(HoodieInstant::requestedTime).collect(Collectors.toList());
+
+      assertEquals(3, completedInstant.size());
+      assertTrue(completedInstant.contains(nextCommitTime1));
+      assertTrue(completedInstant.contains(nextCommitTime2));
+      assertTrue(completedInstant.contains(nextCommitTime4));
+
+      FileIOUtils.deleteDirectory(new File(basePath));
+    } finally {
+      // Close the write clients on every exit path, including failed 
assertions.
+      // The TestingServer itself is closed by @AfterEach (see zkTestingServer 
above).
+      FileIOUtils.closeQuietly(client1 == null ? null : (Closeable) 
client1::close);
+      FileIOUtils.closeQuietly(client2 == null ? null : (Closeable) 
client2::close);
+      FileIOUtils.closeQuietly(client3 == null ? null : (Closeable) 
client3::close);
+      FileIOUtils.closeQuietly(client4 == null ? null : (Closeable) 
client4::close);
+    }
   }
 
   @Test

Reply via email to