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]