Abacn commented on code in PR #40205:
URL: https://github.com/apache/beam/pull/40205#discussion_r4064774188


##########
sdks/java/core/src/main/java/org/apache/beam/sdk/fn/splittabledofn/RestrictionTrackers.java:
##########
@@ -102,7 +113,7 @@ public void checkDone() throws IllegalStateException {
       try {
         delegate.checkDone();
       } finally {
-        lock.unlock();
+        unlock();

Review Comment:
   Do all 4 occurrences of `lock.unlock() -> updateProgressAndUnlock();` really 
need to update the progress cache inside the lock block?



##########
sdks/java/core/src/main/java/org/apache/beam/sdk/fn/splittabledofn/RestrictionTrackers.java:
##########
@@ -139,32 +151,29 @@ protected RestrictionTrackerObserverWithProgress(
 
     @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();

Review Comment:
   We now have mixed calls of `unlock()` and `lock.unlock()`. Consider name the 
new `unlock()` method `updateProgressAndUnlock()` ?



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to