damccorm commented on code in PR #40251:
URL: https://github.com/apache/beam/pull/40251#discussion_r4094372040


##########
sdks/java/core/src/main/java/org/apache/beam/sdk/fn/splittabledofn/RestrictionTrackers.java:
##########
@@ -98,14 +99,27 @@ public RestrictionT currentRestriction() {
 
     @Override
     public SplitResult<RestrictionT> trySplit(double fractionOfRemainder) {
-      lock.lock();
+      return trySplit(fractionOfRemainder, SPLIT_TIMEOUT_SEC);
+    }
+
+    @VisibleForTesting
+    SplitResult<RestrictionT> trySplit(double fractionOfRemainder, int 
timeOutSec) {
       try {
-        SplitResult<RestrictionT> result = 
delegate.trySplit(fractionOfRemainder);
-        needsProgressUpdate = true;
-        return result;
-      } finally {
-        updateProgressAndUnlock();
+        // lock can be held long by 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;

Review Comment:
   That's a good catch. I think in that case, we have 2 options:
   
   1. Raise a runtime exception
   2. Time out (keep the current behavior).
   
   I lean towards keeping the current behavior for 2 reasons:
   
   1) It is less invasive
   2) It allows jobs to succeed if acquiring the split lock takes e.g. 8 
minutes (so no potential regression)
   
   I updated to do this, let me know what you think. It does mean we could fail 
a work item on a runner-initiated checkpoint, but at least the odds of doing 
that 4 times in a row should be low. I guess its also more motivation to move 
forward with your fix.
   
   Btw, I think that the null protection is mostly for this code path - 
https://github.com/apache/beam/blob/e3b3e7ca1f7a372d44ad563b1d4e4f4acb30b298/sdks/java/harness/src/main/java/org/apache/beam/fn/harness/FnApiDoFnRunner.java#L676
 - which would actually be ok since the lock would no longer be held after 
processElement returned. But I do think we should honor this to protect against 
data loss here nonetheless.



-- 
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