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 a40d8cd9bd9 Return last evaluated progress on lock timeout in
RestrictionTrackers.getProgress (#40205)
a40d8cd9bd9 is described below
commit a40d8cd9bd9ae0bb157ba33a45ae07575b53b411
Author: Danny McCormick <[email protected]>
AuthorDate: Tue Sep 22 22:31:00 2026 +0000
Return last evaluated progress on lock timeout in
RestrictionTrackers.getProgress (#40205)
* Return last evaluated progress on lock timeout in
RestrictionTrackers.getProgress
* Update progress before unlocking when needsProgressUpdate is set
* Unconditionally update progress after trySplit and rename unlock to
updateProgressAndUnlock
---
.../sdk/fn/splittabledofn/RestrictionTrackers.java | 66 +++++++++++++---------
.../fn/splittabledofn/RestrictionTrackersTest.java | 55 ++++++++++++++++--
2 files changed, 90 insertions(+), 31 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 6fefc6b184a..370871fd9ca 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
@@ -49,7 +49,8 @@ public class RestrictionTrackers {
extends RestrictionTracker<RestrictionT, PositionT> {
protected final RestrictionTracker<RestrictionT, PositionT> delegate;
protected ReentrantLock lock = new ReentrantLock();
- protected volatile boolean hasInitialProgress = false;
+ protected volatile Progress lastProgress = Progress.NONE;
+ protected volatile boolean needsProgressUpdate = false;
private final ClaimObserver<PositionT> claimObserver;
protected RestrictionTrackerObserver(
@@ -59,6 +60,16 @@ public class RestrictionTrackers {
this.claimObserver = claimObserver;
}
+ protected void updateProgressAndUnlock() {
+ try {
+ if (needsProgressUpdate) {
+ updateProgressBlocking();
+ }
+ } finally {
+ lock.unlock();
+ }
+ }
+
@Override
public boolean tryClaim(PositionT position) {
lock.lock();
@@ -71,7 +82,7 @@ public class RestrictionTrackers {
return false;
}
} finally {
- lock.unlock();
+ updateProgressAndUnlock();
}
}
@@ -81,7 +92,7 @@ public class RestrictionTrackers {
try {
return delegate.currentRestriction();
} finally {
- lock.unlock();
+ updateProgressAndUnlock();
}
}
@@ -90,9 +101,10 @@ public class RestrictionTrackers {
lock.lock();
try {
SplitResult<RestrictionT> result =
delegate.trySplit(fractionOfRemainder);
+ needsProgressUpdate = true;
return result;
} finally {
- lock.unlock();
+ updateProgressAndUnlock();
}
}
@@ -102,7 +114,7 @@ public class RestrictionTrackers {
try {
delegate.checkDone();
} finally {
- lock.unlock();
+ updateProgressAndUnlock();
}
}
@@ -112,10 +124,13 @@ public class RestrictionTrackers {
}
/** Evaluate progress if requested. */
- protected Progress getProgressBlocking() {
+ protected void updateProgressBlocking() {
lock.lock();
try {
- return ((HasProgress) delegate).getProgress();
+ needsProgressUpdate = false;
+ if (delegate instanceof HasProgress) {
+ lastProgress = ((HasProgress) delegate).getProgress();
+ }
} finally {
lock.unlock();
}
@@ -129,7 +144,7 @@ public class RestrictionTrackers {
@ThreadSafe
static class RestrictionTrackerObserverWithProgress<RestrictionT, PositionT>
extends RestrictionTrackerObserver<RestrictionT, PositionT> implements
HasProgress {
- private static final int FIRST_PROGRESS_TIMEOUT_SEC = 60;
+ private static final int PROGRESS_TIMEOUT_SEC = 60;
protected RestrictionTrackerObserverWithProgress(
RestrictionTracker<RestrictionT, PositionT> delegate,
@@ -139,32 +154,29 @@ public class RestrictionTrackers {
@Override
public Progress getProgress() {
- return getProgress(FIRST_PROGRESS_TIMEOUT_SEC);
+ return getProgress(PROGRESS_TIMEOUT_SEC);
}
@VisibleForTesting
Progress getProgress(int timeOutSec) {
- if (!hasInitialProgress) {
- Progress progress = Progress.NONE;
- try {
- // lock can be held long by long-running tryClaim/trySplit. We
tolerate this scenario
- // by returning zero progress when initial progress never evaluated
before due to lock
- // timeout.
- if (lock.tryLock(timeOutSec, TimeUnit.SECONDS)) {
- try {
- progress = getProgressBlocking();
- hasInitialProgress = true;
- } finally {
- lock.unlock();
- }
+ try {
+ // lock can be held long by long-running tryClaim/trySplit. We
tolerate this scenario
+ // by returning the last evaluated progress (or zero progress if never
evaluated before)
+ // when lock timeout occurs.
+ if (lock.tryLock(timeOutSec, TimeUnit.SECONDS)) {
+ try {
+ updateProgressBlocking();
+ } finally {
+ lock.unlock();
}
- } catch (InterruptedException e) {
- Thread.currentThread().interrupt();
+ } else {
+ needsProgressUpdate = true;
}
- return progress;
- } else {
- return getProgressBlocking();
+ } catch (InterruptedException e) {
+ needsProgressUpdate = true;
+ Thread.currentThread().interrupt();
}
+ return lastProgress;
}
}
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 7b6f3d47c27..e2403322478 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
@@ -103,7 +103,9 @@ public class RestrictionTrackersTest {
private boolean blockTryClaim;
private boolean blockTrySplit;
private volatile boolean isBlocked;
+ private volatile Progress currentProgress = REPORT_PROGRESS;
public static final Progress REPORT_PROGRESS = Progress.from(2.0, 3.0);
+ public static final Progress UPDATED_PROGRESS = Progress.from(4.0, 1.0);
public RestrictionTrackerWithProgress() {
this(false, false);
@@ -117,7 +119,11 @@ public class RestrictionTrackersTest {
@Override
public Progress getProgress() {
- return REPORT_PROGRESS;
+ return currentProgress;
+ }
+
+ public void setProgress(Progress progress) {
+ this.currentProgress = progress;
}
@Override
@@ -161,6 +167,14 @@ public class RestrictionTrackersTest {
return IsBounded.BOUNDED;
}
+ public synchronized void setBlockTryClaim(boolean blockTryClaim) {
+ this.blockTryClaim = blockTryClaim;
+ }
+
+ public synchronized void setBlockTrySplit(boolean blockTrySplit) {
+ this.blockTrySplit = blockTrySplit;
+ }
+
public synchronized void releaseLock() {
blockTrySplit = false;
blockTryClaim = false;
@@ -190,13 +204,32 @@ public class RestrictionTrackersTest {
Thread blocking = new Thread(() -> tracker.tryClaim(new Object()));
blocking.start();
withProgress.waitUntilBlocking(true);
+ // Times out while first tryClaim holds lock; returns NONE and sets
needsProgressUpdate = true
RestrictionTracker.Progress progress =
((RestrictionTrackers.RestrictionTrackerObserverWithProgress)
tracker).getProgress(1);
assertEquals(RestrictionTracker.Progress.NONE, progress);
+ // When first tryClaim finishes, updateProgressAndUnlock() sees
needsProgressUpdate == true and
+ // evaluates REPORT_PROGRESS before releasing the lock.
withProgress.releaseLock();
withProgress.waitUntilBlocking(false);
- progress = ((HasProgress) tracker).getProgress();
+ blocking.join();
+
+ // Even if a second blocking tryClaim immediately grabs the lock before
getProgress is called
+ // again, getProgress(1) returns REPORT_PROGRESS (updated during first
tryClaim's
+ // updateProgressAndUnlock).
+ withProgress.setProgress(RestrictionTrackerWithProgress.UPDATED_PROGRESS);
+ withProgress.setBlockTryClaim(true);
+ Thread secondBlocking = new Thread(() -> tracker.tryClaim(new Object()));
+ secondBlocking.start();
+ withProgress.waitUntilBlocking(true);
+ progress =
+ ((RestrictionTrackers.RestrictionTrackerObserverWithProgress)
tracker).getProgress(1);
assertEquals(RestrictionTrackerWithProgress.REPORT_PROGRESS, progress);
+ withProgress.releaseLock();
+ withProgress.waitUntilBlocking(false);
+ secondBlocking.join();
+ progress = ((HasProgress) tracker).getProgress();
+ assertEquals(RestrictionTrackerWithProgress.UPDATED_PROGRESS, progress);
}
@Test
@@ -204,15 +237,29 @@ public class RestrictionTrackersTest {
RestrictionTrackerWithProgress withProgress = new
RestrictionTrackerWithProgress(false, true);
RestrictionTracker<Object, Object> tracker =
RestrictionTrackers.observe(withProgress, new
RestrictionTrackers.NoopClaimObserver<>());
+ // trySplit unconditionally refreshes lastProgress via
updateProgressAndUnlock()
Thread blocking = new Thread(() -> tracker.trySplit(0.5));
blocking.start();
withProgress.waitUntilBlocking(true);
+ withProgress.releaseLock();
+ withProgress.waitUntilBlocking(false);
+ blocking.join();
+
+ // Even though getProgress was never called before, trySplit
unconditionally updated
+ // lastProgress to REPORT_PROGRESS. If a subsequent tryClaim blocks,
getProgress(1) returns
+ // REPORT_PROGRESS rather than NONE or a stale pre-split progress.
+ withProgress.setProgress(RestrictionTrackerWithProgress.UPDATED_PROGRESS);
+ withProgress.setBlockTryClaim(true);
+ Thread secondBlocking = new Thread(() -> tracker.tryClaim(new Object()));
+ secondBlocking.start();
+ withProgress.waitUntilBlocking(true);
RestrictionTracker.Progress progress =
((RestrictionTrackers.RestrictionTrackerObserverWithProgress)
tracker).getProgress(1);
- assertEquals(RestrictionTracker.Progress.NONE, progress);
+ assertEquals(RestrictionTrackerWithProgress.REPORT_PROGRESS, progress);
withProgress.releaseLock();
withProgress.waitUntilBlocking(false);
+ secondBlocking.join();
progress = ((HasProgress) tracker).getProgress();
- assertEquals(RestrictionTrackerWithProgress.REPORT_PROGRESS, progress);
+ assertEquals(RestrictionTrackerWithProgress.UPDATED_PROGRESS, progress);
}
}