scwhittle commented on code in PR #40251:
URL: https://github.com/apache/beam/pull/40251#discussion_r4092050410
##########
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:
trySplit has this comment:
* @return a {@link SplitResult} if a split was possible, otherwise
returns {@code null}. If the
* {@code fractionOfRemainder == 0}, a {@code null} result MUST imply
that the restriction
* tracker is done and there is no more work left to do.
So it seems like data loss to return null if fractionOfRemainder == 0.
--
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]