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

damccorm pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git


The following commit(s) were added to refs/heads/master by this push:
     new 404bde62095 Time out RestrictionTrackers.trySplit after 5 minutes on 
lock contention (#40251)
404bde62095 is described below

commit 404bde62095bd29bdec5cd1427cc3634d4a407cd
Author: Danny McCormick <[email protected]>
AuthorDate: Thu Sep 24 17:49:25 2026 +0000

    Time out RestrictionTrackers.trySplit after 5 minutes on lock contention 
(#40251)
    
    * Time out RestrictionTrackers.trySplit after 5 minutes on lock contention
    
    * Block on lock.lock() when fractionOfRemainder == 0 in trySplit
---
 .../sdk/fn/splittabledofn/RestrictionTrackers.java | 43 +++++++++++++---
 .../fn/splittabledofn/RestrictionTrackersTest.java | 57 +++++++++++++++++++++-
 2 files changed, 92 insertions(+), 8 deletions(-)

diff --git 
a/sdks/java/core/src/main/java/org/apache/beam/sdk/fn/splittabledofn/RestrictionTrackers.java
 
b/sdks/java/core/src/main/java/org/apache/beam/sdk/fn/splittabledofn/RestrictionTrackers.java
index 370871fd9ca..ded09b10649 100644
--- 
a/sdks/java/core/src/main/java/org/apache/beam/sdk/fn/splittabledofn/RestrictionTrackers.java
+++ 
b/sdks/java/core/src/main/java/org/apache/beam/sdk/fn/splittabledofn/RestrictionTrackers.java
@@ -45,8 +45,9 @@ public class RestrictionTrackers {
    * RestrictionTracker}.
    */
   @ThreadSafe
-  private static class RestrictionTrackerObserver<RestrictionT, PositionT>
+  static class RestrictionTrackerObserver<RestrictionT, PositionT>
       extends RestrictionTracker<RestrictionT, PositionT> {
+    private static final int SPLIT_TIMEOUT_SEC = 300;
     protected final RestrictionTracker<RestrictionT, PositionT> delegate;
     protected ReentrantLock lock = new ReentrantLock();
     protected volatile Progress lastProgress = Progress.NONE;
@@ -98,14 +99,42 @@ public class RestrictionTrackers {
 
     @Override
     public SplitResult<RestrictionT> trySplit(double fractionOfRemainder) {
-      lock.lock();
+      return trySplit(fractionOfRemainder, SPLIT_TIMEOUT_SEC);
+    }
+
+    @VisibleForTesting
+    SplitResult<RestrictionT> trySplit(double fractionOfRemainder, int 
timeOutSec) {
+      // When fractionOfRemainder == 0 (a checkpoint), returning null has a 
special meaning in the
+      // RestrictionTracker contract: it MUST imply that the restriction 
tracker is done and there
+      // is no more work left to do. Therefore, we cannot time out and return 
null when
+      // fractionOfRemainder == 0, and must block until the lock is acquired.
+      if (fractionOfRemainder == 0) {
+        lock.lock();
+        try {
+          SplitResult<RestrictionT> result = 
delegate.trySplit(fractionOfRemainder);
+          needsProgressUpdate = true;
+          return result;
+        } finally {
+          updateProgressAndUnlock();
+        }
+      }
       try {
-        SplitResult<RestrictionT> result = 
delegate.trySplit(fractionOfRemainder);
-        needsProgressUpdate = true;
-        return result;
-      } finally {
-        updateProgressAndUnlock();
+        // For dynamic splits (fractionOfRemainder > 0), lock can be held long 
by a long-running
+        // tryClaim. We tolerate this scenario by returning null (declining to 
split) when lock
+        // timeout occurs.
+        if (lock.tryLock(timeOutSec, TimeUnit.SECONDS)) {
+          try {
+            SplitResult<RestrictionT> result = 
delegate.trySplit(fractionOfRemainder);
+            needsProgressUpdate = true;
+            return result;
+          } finally {
+            updateProgressAndUnlock();
+          }
+        }
+      } catch (InterruptedException e) {
+        Thread.currentThread().interrupt();
       }
+      return null;
     }
 
     @Override
diff --git 
a/sdks/java/core/src/test/java/org/apache/beam/sdk/fn/splittabledofn/RestrictionTrackersTest.java
 
b/sdks/java/core/src/test/java/org/apache/beam/sdk/fn/splittabledofn/RestrictionTrackersTest.java
index e2403322478..3708f71587d 100644
--- 
a/sdks/java/core/src/test/java/org/apache/beam/sdk/fn/splittabledofn/RestrictionTrackersTest.java
+++ 
b/sdks/java/core/src/test/java/org/apache/beam/sdk/fn/splittabledofn/RestrictionTrackersTest.java
@@ -156,7 +156,7 @@ public class RestrictionTrackersTest {
         }
       }
       isBlocked = false;
-      return null;
+      return SplitResult.of("primary", "residual");
     }
 
     @Override
@@ -262,4 +262,59 @@ public class RestrictionTrackersTest {
     progress = ((HasProgress) tracker).getProgress();
     assertEquals(RestrictionTrackerWithProgress.UPDATED_PROGRESS, progress);
   }
+
+  @Test
+  public void testClaimObserversTrySplitNonBlockingOnTryClaim() throws 
InterruptedException {
+    RestrictionTrackerWithProgress withProgress = new 
RestrictionTrackerWithProgress(true, false);
+    RestrictionTracker<Object, Object> tracker =
+        RestrictionTrackers.observe(withProgress, new 
RestrictionTrackers.NoopClaimObserver<>());
+    Thread blocking = new Thread(() -> tracker.tryClaim(new Object()));
+    blocking.start();
+    withProgress.waitUntilBlocking(true);
+
+    // While tryClaim holds the lock, trySplit times out and returns null 
instead of blocking
+    // indefinitely.
+    SplitResult<Object> splitResult =
+        ((RestrictionTrackers.RestrictionTrackerObserver<Object, Object>) 
tracker).trySplit(0.5, 1);
+    assertEquals(null, splitResult);
+
+    withProgress.releaseLock();
+    withProgress.waitUntilBlocking(false);
+    blocking.join();
+
+    // Once tryClaim releases the lock, trySplit succeeds.
+    splitResult =
+        ((RestrictionTrackers.RestrictionTrackerObserver<Object, Object>) 
tracker).trySplit(0.5, 1);
+    assertEquals(SplitResult.of("primary", "residual"), splitResult);
+  }
+
+  @Test
+  public void testClaimObserversTrySplitZeroFractionBlocksOnTryClaim() throws 
InterruptedException {
+    RestrictionTrackerWithProgress withProgress = new 
RestrictionTrackerWithProgress(true, false);
+    RestrictionTracker<Object, Object> tracker =
+        RestrictionTrackers.observe(withProgress, new 
RestrictionTrackers.NoopClaimObserver<>());
+    Thread blocking = new Thread(() -> tracker.tryClaim(new Object()));
+    blocking.start();
+    withProgress.waitUntilBlocking(true);
+
+    // When fractionOfRemainder == 0, trySplit must block until the lock is 
released rather than
+    // timing out and returning null (since null implies the restriction 
tracker is done).
+    Thread releaseLater =
+        new Thread(
+            () -> {
+              try {
+                Thread.sleep(1500);
+              } catch (InterruptedException e) {
+                Thread.currentThread().interrupt();
+              }
+              withProgress.releaseLock();
+            });
+    releaseLater.start();
+
+    SplitResult<Object> splitResult =
+        ((RestrictionTrackers.RestrictionTrackerObserver<Object, Object>) 
tracker).trySplit(0.0, 1);
+    assertEquals(SplitResult.of("primary", "residual"), splitResult);
+    blocking.join();
+    releaseLater.join();
+  }
 }

Reply via email to