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

Reply via email to