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();
+ }
}